Runtime: batch, replay y stream¶
Qué vas a aprender¶
- Cómo el worker batch convierte una petición en un resultado publicado: outbox, leases, admisión idempotente y recuperación de huérfanos.
- Qué es el replay: un log normalizado de llegadas, reproducible, con un fault schedule que inyecta fallos controlados.
- Cómo funciona el stream sobre Redpanda: watermarks, backpressure, diario (journal) de offsets,
fencing por lease PG con
fence_tokeny handoff entre owners. - Qué se demostró en el drill de caos de 2 owners del 2026-09-22.
- Cómo se comprueba la paridad de los tres modos y qué son los checkpoints CAS.
Fuentes
packages/finaz_runtime/{worker.py,core/,recovery/,stream/}, packages/finaz_control/,
packages/finaz_replay/, docs/v2/{MOTORES,OPERACION,GATES}.md, capítulo 06 del pack,
deployment/v2/release-manifest.s2.json y la evidencia del servidor
qa_reports/v2/runs/2026-09-22/g3/ (drill 2 owners) y 2026-09-20/ (caos kill -9).
1. La idea común: al menos una vez + publicación idempotente¶
Antes de entrar en cada modo, fija la garantía que persigue todo el runtime (capítulo 06 del pack):
Entrega al menos una vez, publicación lógica idempotente y corte recuperable.
No se promete «exactamente una vez» entre PostgreSQL, ClickHouse y Redpanda (no hay transacción distribuida). Lo que se promete es que, aunque un trabajo se procese dos veces, el resultado visible es uno solo, porque la publicación en PostgreSQL es publish-if-absent y verifica el dueño.
Analogía: el cartero y el buzón con candado
El cartero puede intentar entregar la misma carta dos veces (al menos una vez). El buzón tiene un candado que sólo abre la primera carta con ese número de envío y que, además, comprueba que el cartero sigue teniendo la ruta asignada (fence). La segunda entrega rebota sin romper nada.
Las piezas de control viven en finaz_control (AG-004):
| Módulo | Tabla PG | Qué garantiza |
|---|---|---|
admission |
admission |
claim-or-reject por idempotency_key |
leases |
leases |
un único dueño por run, con fence_token monótono |
outbox |
outbox |
cola de eventos ordenada; reclamo por lote con CAS; marca de entrega |
publication |
publications, checkpoints, outbox |
publish-if-absent en un solo commit |
reconciler |
outbox, publications |
repara «reclamado sin entregar» y verifica hashes |
2. Worker batch¶
2.1 El recorrido de un run¶
Cuando la API admite un run (capítulo 4), escribe en una sola transacción la fila en api_runs
(estado QUEUED) y un evento run.requested en el outbox. El worker-batch es un bucle que
sondea ese outbox.
sequenceDiagram
autonumber
participant API as api
participant PG as PostgreSQL
participant W as worker-batch
participant CAS as CAS /artifacts
participant CH as ClickHouse
API->>PG: INSERT api_runs (QUEUED) + INSERT outbox run.requested
loop cada FINAZ_WORKER_POLL_SEC
W->>PG: outbox.reclamar_lote(consumidor, limite=50)
end
W->>PG: UPDATE api_runs SET estado='RUNNING' WHERE estado='QUEUED'
alt cancel_requested_at no nulo
W->>PG: estado = CANCELLED (sin ejecutar)
else
W->>W: BatchRunner: snapshot → clock → EMA → persist
W->>CAS: FeatureBatch (content-addressed)
W->>CH: staging de filas
W->>PG: publicar ResultRef (publish-if-absent)
W->>PG: estado = SUCCEEDED, result_refs, run_revision+1
end
W->>PG: outbox.marcar_entregado
Detalles que conviene retener (de docs/v2/M2_M3_INTERFACES.md y worker.py):
- CAS de estado:
QUEUED → RUNNINGse hace conUPDATE … WHERE run_id=? AND estado='QUEUED'. Si actualiza 0 filas, otro worker ya lo cogió (o no existe) y el evento se libera como ignorado. - Cancelación en frontera segura: si alguien pidió cancelar, se marca
CANCELLEDantes de ejecutar. No se interrumpe un cálculo a medias. - Fallo honesto: si la ejecución falla,
FAILEDcon el mensaje enerrores, evento entregado. No se reintenta en bucle infinito (este fue precisamente un defecto encontrado en la auditoría del 2026-09-19: con el fixture ausente, el worker giraba en bucle; hoy libera una vez y avisa). - Modo
--once: para humos; carga el fixture, ejecuta el vertical G1, publica e imprime elResultRefpor stdout.
2.2 Admisión idempotente¶
La admisión usa una clave con alcance principal \0 operación \0 clave. El run_id es
determinista: run_ + sha256 de (principal, operación, clave, hash del cuerpo). Consecuencias:
| Caso | Resultado |
|---|---|
| Misma clave + mismo cuerpo | Mismo run_id, misma respuesta (no se crea otro run) |
| Misma clave + otro cuerpo | 409 «clave de idempotencia ya usada con otra solicitud» |
| Otra clave | Run nuevo |
Lo verás en vivo en el capítulo 4.
2.3 Leases y fence tokens¶
Un lease es un préstamo con caducidad: «este run es tuyo hasta tal hora». La función
leases.adquirir tiene cuatro ramas, y merece la pena leerlas despacio:
| Situación | Qué pasa | fence_token |
|---|---|---|
| No hay lease | Se crea | 1 |
| Lease vigente de otro owner | Error LeaseOcupado (LEASE_HELD) |
sin cambio |
| Lease vigente mío | Se renueva | sin cambio |
| Lease expirado | El nuevo owner lo reclama | +1 |
El fence_token sólo sube. Cualquier efecto externo (publicar, escribir el diario) comprueba que el
fence que tú tienes es el vigente; si alguien te lo ha quitado, tu efecto se rechaza con
STALE_FENCE / FenceObsoleto.
Analogía: la llave de la habitación del hotel
Cada vez que la recepción reasigna una habitación, reprograma la cerradura (fence +1). Tu llave vieja sigue existiendo, puedes caminar hasta la puerta (gastar CPU), pero no abre. Renovar tu estancia no cambia la cerradura; perderla, sí.
Incidente real: comparar fechas como texto
En docs/v2/OPERACION.md hay un caso instructivo: STALE_FENCE con un lease recién renovado.
La causa era comparar timestamps como strings (" " < "T" en ASCII). La corrección:
comparar datetime. Moraleja: el fencing es tan fiable como su comparación de tiempo.
2.4 Recuperación de huérfanos¶
¿Qué pasa si matas el worker con kill -9 en mitad de un lote? Al arrancar, el bucle ejecuta
_recuperar_huerfanos:
- (a) run en
RUNNINGsin lease vigente → el worker anterior murió: vuelve aQUEUEDy su evento a pendiente. - (b) run en
QUEUEDcon evento ya entregado → el evento vuelve a pendiente. - Con lease vigente no se toca: hay otro worker vivo y el fencing lo protege.
Esto no salió de la nada: el ensayo de caos del hito 3 (G9) encontró que tras un SIGKILL a mitad de
lote un run podía quedar RUNNING con un reclamo huérfano. Se corrigió re-encolando eventos
reclamados sin entregar, con test propio. La evidencia de G9 (qa_reports/v2/runs/2026-09-20/caos-worker-batch-kill9.json): kill -9 de
worker-batch con un lote de 20 corridas en vuelo → 20/20 SUCCEEDED, fallos: [], sin pérdida ni
duplicados (la tabla de runs pasa de 412 a 432). docs/v2/GATES.md habla de «60 corridas → 60/60»;
el fichero de evidencia registra 20.
3. Replay: el pasado, reproducible y con fallos a la carta¶
El replay (finaz_replay, AG-015) toma un snapshot y produce un log normalizado de llegadas:
la secuencia exacta de eventos tal y como «habrían llegado», con un ritmo (speed_mode,
speed_multiplier) y una semilla. Ese log es la base de la paridad: el mismo log alimenta al stream.
El fault schedule¶
Un programa de fallos congelado inyecta problemas realistas. El del fixture
(tests_v2/fixtures_shared/fault_schedule.json, semilla 7) tiene 7 fallos:
| Evento | Tipo | Efecto |
|---|---|---|
fixture.event.0.2 |
DUPLICATE |
el evento llega dos veces |
fixture.event.0.3 |
DELAY |
retraso de 2 s en dominio |
fixture.event.1.1 |
GAP |
falta un evento |
fixture.event.2.2 |
REVISION |
llega una revisión 120 s después |
fixture.event.0.4 |
PAUSE |
pausa de 100 ms de pared |
fixture.event.2.4 |
BURST |
ráfaga de 3 |
fixture.event.1.4 |
SOURCE_FAILURE |
caída de la fuente 300 ms |
El resultado vivo (MOTORES, G2): 23 llegadas, 7 fallos aplicados y auditados, con un log_hash
idéntico en cada ejecución (sha256:b0238dc3… en el manifest S5).
# Replay como job one-shot (perfil replay-job; termina solo, exit 0)
docker compose -p finaz-trading-engine -f deployment/v2/compose.yaml \
--profile replay-job run --rm replay
Replay no es un daemon
replay tiene restart: "no" y no tiene healthcheck: es un trabajo que publica el log en un topic
de Redpanda (finaz-te-sim.fixture.replay.faults.normalized.v1) y termina.
4. Stream sobre Redpanda¶
El worker-stream consume ese topic y aplica el mismo cálculo que el batch: reutiliza la
proyección de reloj del batch (project_batch, no una réplica) y una EMA incremental (EmaStateful)
con paridad bit a bit con ema_sma_seed_reject_gap.
# Worker stream acotado (humo), grupo fresco
docker compose -p finaz-trading-engine -f deployment/v2/compose.yaml --profile stream \
run --rm --no-deps worker-stream python -m finaz_runtime.stream.worker \
--topic finaz-te-sim.fixture.replay.faults.normalized.v1 \
--group s5-hito2-N --flow /flows/01_A_feature.json \
--once-n <total-topic> --emit /artifacts/emit-hito2.json
4.1 Qué hace con cada registro¶
| Registro | Tratamiento |
|---|---|
| Instrumento no seleccionado por el flow | filtered: se journala (avanza el high-water) pero no toca reloj ni EMA |
Mismo event_id visto otra vez (copia DUPLICATE) |
duplicates_skipped: se salta por identidad, nunca por timestamp |
Cierre null o no finito |
alimenta None a la EMA → GAP (reject_gap) |
| Valor ilegible o tombstone | omitted (veneno): no tumba el stream; se reintenta en la reanudación sin doble conteo |
Desviación documentada
El batch falla cerrado ante un valor ilegible; el stream no puede abortar por un mensaje
venenoso, así que lo cuenta como omitted. La cuarentena formal queda como trabajo futuro.
4.2 Watermarks¶
Un watermark responde a «¿hasta qué instante de evento puedo dar por cerrado el pasado?».
- Por partición:
W_p = max_event_seen − allowed_lateness, y sólo si no hay huecos abiertos. - En un fan-in: el mínimo de las particiones activas. Una partición callada no hace avanzar el watermark: la falta de tráfico no prueba ociosidad (Q031).
- Sólo eventos explícitos
IDLE_CONFIRMEDoEND_OF_BOUNDED_SOURCEsacan a una partición del mínimo. - Un evento anterior al watermark es tardío: se registra o cuarentena y sólo se adopta en decisiones futuras; nunca reescribe una decisión pasada.
4.3 Backpressure¶
backpressure.py decide pausas con umbrales explícitos (lag, bytes en vuelo, espacio libre):
stateDiagram-v2
[*] --> FLOWING
FLOWING --> SOURCE_PAUSED: lag/bytes/espacio sobre umbral
FLOWING --> PROCESS_PAUSED: buffer de proceso lleno
SOURCE_PAUSED --> FLOWING: por debajo del umbral
PROCESS_PAUSED --> FLOWING: buffer drenado
FLOWING --> RECOVERY_BLOCKED: la retención borraría inputs de un checkpoint
RECOVERY_BLOCKED es la parte honesta: si el broker va a borrar offsets que un checkpoint necesita,
no se «promete» un replay imposible.
4.4 Diario de offsets y fencing por lease PG¶
El worker guarda su progreso en un diario (journal) de offsets:
<artifacts>/stream-journal/<group>__<topic>.json, con el estado de la EMA y el fence. El offset
del consumer group de Redpanda es una optimización; la fuente de verdad del estado FINAZ es el
diario/checkpoint.
Desde el 2026-09-22, con --run-id, el worker usa finaz_runtime.stream.fencing.GuardiaFence:
- Adquiere el lease del run en PG antes de consumir.
- Cada efecto externo (commit del diario y escritura del emit) se ejecuta con la fila del lease
bloqueada (
SELECT … FOR UPDATE), tras verificar owner + fence + vigencia. - Un latido renueva el lease cada
ttl/3. - Un owner obsoleto falla con
STREAM_FENCE_OBSOLETOsin ejecutar el efecto; un intruso con el lease ajeno vigente falla conSTREAM_LEASE_OCUPADO.
4.5 Handoff entre owners¶
El manifiesto de handoff (handoff.py) formaliza el traspaso: el nuevo owner reclama con fence
superior y reanuda desde el último commit del anterior. Un worker obsoleto puede gastar CPU o
escribir staging, pero no publicar. El mismo módulo describe el handoff histórico → stream: se captura
el high-watermark por partición y el stream abre desde los next_offsets exactos, con dedupe por
identidad de evento.
5. El drill de caos de 2 owners (2026-09-22)¶
Es la prueba viva de G3. Script deployment/v2/smoke/caos_2owners_g3.py, imagen
finaz-te-core:s4, TTL de lease 6 s, topic con 299 registros (particiones 0 y 3 con high-water
78 y 221). Evidencia: qa_reports/v2/runs/2026-09-22/g3/caos-2owners.json.
sequenceDiagram
autonumber
participant A as owner-A
participant B0 as intruso B0
participant B as owner-B
participant PG as lease PG
A->>PG: adquirir (fence 1)
A->>A: consume y avanza el diario bajo fence
B0->>PG: adquirir
PG-->>B0: STREAM_LEASE_OCUPADO (exit 1, sin emit)
Note over A: docker pause (stall real)
B->>PG: lease de A expirado → reclamar (fence 2)
B->>B: reanuda desde el último commit de A
Note over A: docker unpause
A->>PG: commit del diario con fence 1
PG-->>A: STREAM_FENCE_OBSOLETO (sin diario ni emit)
B->>B: termina OK: 299/299, values_hash = golden
Una segunda fase repite la idea con SIGKILL: B2 avanza, recibe kill -9; C espera
(LEASE_HELD, no roba el lease fresco de B2), reclama con fence superior cuando expira y reanuda
exactamente desde el diario de B2.
Resultado: 36/36 checks PASS. Los más importantes:
| Check | Qué demuestra |
|---|---|
| Z3 | El intruso con A sano falla STREAM_LEASE_OCUPADO y no escribe emit |
| Z7 | El lease pasa a B con fence_token 2 (antes 1) |
| Z10–Z12 | A, al despertar, es rechazado en su primer intento de diario y no publica nada |
| Z14 | B reanuda desde el último commit de A |
| Z15 / K8 | values_hash = golden S5, eventos = referencia, event_id únicos, diario final = high-water |
| K4 | C respeta un lease fresco tras el kill |
Hallazgo del drill: el daemon en producción aún no tiene esto
El mismo día se comprobó el worker-stream que corre como servicio: sigue en
finaz-te-core:s3, sin --run-id (sin fencing), y su emit vivo no es golden
(values_hash sha256:fe5812e2…) por un defecto de reanudación (high-water sobre aceptados sin
drenar y event_id no persistidos) que ya está corregido en :s4. Promoverlo exige imagen :s4,
FINAZ_STREAM_RUN_ID en el compose, grupo nuevo y verify_release. Por eso S5 sigue
INCOMPLETO (capítulo 6).
6. Paridad de los tres modos¶
La promesa del primer vertical: mismo FlowSpec, mismo resultado lógico en BATCH, REPLAY y STREAM.
flowchart LR
SN[(fixture.snapshot.01)] --> B[BATCH<br/>BatchRunner]
SN --> R[REPLAY<br/>log normalizado<br/>+ 7 fallos]
R --> T[[topic Redpanda]]
T --> S[STREAM<br/>worker-stream]
B --> G{{golden<br/>null,null,2,3,4,5}}
S --> G
La referencia del manifest S5 (deployment/v2/release-manifest.s2.json):
"stream_parity": {
"values": [null, null, 2, 3, 4, 5],
"values_hash": "sha256:9e906a09a1c594964c16291a77cb38b0ffb021e69c230d71d1dd77ebbd6a7877",
"consumed": 6, "filtered": 15, "omitted": 1, "duplicates_skipped": 1,
"verdict": "EQUAL_TO_BATCH_GOLDEN"
}
De las 23 llegadas del replay: 6 consumidas para el instrumento del flow, 15 filtradas (otros
instrumentos), 1 omitida y 1 duplicado saltado. Y la EMA sale igual que en batch, a pesar de los
fallos inyectados. Si comparas el emit no-golden del daemon :s3 (valores 3.046875, 2.5234375, …)
verás por qué el hash es el juez: basta un duplicado aceptado como nuevo para que la EMA se desvíe.
7. Checkpoints CAS¶
Un checkpoint (finaz_runtime/recovery, AG-014) es el corte consistente de un epoch:
- estado por nodo (reloj, EMA…),
next_offsetspor partición,- refs de side inputs,
- dedupe, RNG, watermark y manifiesto del sink.
Los blobs de estado se guardan en CAS (almacenamiento direccionado por contenido, JSON canónico
FINAZ_JSON_V1) y el puntero al checkpoint se registra en PG (checkpoints) mediante el coordinador de
commit con fence. El manifest_hash cubre todo el cuerpo salvo su propio campo (regla anti
self-hash).
Restore nunca salta a «latest»
El restore comprueba compatibilidad (plan, schema de estado, hashes) y restaura el checkpoint correcto, no el más reciente que encuentre. En vivo: commit + restore con estado idéntico byte a byte (MOTORES, AG-014).
Analogía: guardar partida
Un checkpoint es un «guardar partida» que incluye no sólo tu posición, sino qué enemigos ya derrotaste (dedupe) y la semilla del azar (RNG). Si sólo guardaras la posición, al cargar volverías a derrotar a los mismos enemigos: duplicados.
Resumen¶
- Garantía global: al menos una vez + publicación idempotente + corte recuperable; PostgreSQL es la única puerta de visibilidad.
- Batch: outbox → CAS
QUEUED→RUNNING→ pipeline → publish-if-absent →SUCCEEDED; cancelación en frontera segura; recuperación de huérfanos al arrancar (20/20 tras kill -9 en G9). - Leases: fence +1 sólo al reclamar uno expirado; renovar no cambia el fence; efectos con fence obsoleto se rechazan.
- Replay: log normalizado reproducible (23 llegadas, 7 fallos,
log_hashestable). - Stream: watermarks conservadores, backpressure explícito, diario de offsets, fencing por lease PG
con
SELECT … FOR UPDATE; drill de 2 owners 36/36 PASS con:s4; el daemon:s3aún pendiente. - Paridad:
values_hash sha256:9e906a09…= golden[null, null, 2, 3, 4, 5].
Para practicar¶
- Dibuja la tabla de
leasestras esta secuencia: A adquiere; pasan 7 s (TTL 6); B adquiere; A intenta renovar. ¿Quéfence_tokentiene cada uno y qué error recibe A? - En el fault schedule, ¿qué fallo explica
duplicates_skipped: 1en la paridad? - ¿Por qué el watermark de un fan-in no avanza si una partición está callada? Pon un ejemplo de mercado donde avanzar sería un error.
- Lee
_recuperar_huerfanosenpackages/finaz_runtime/worker.pyy explica por qué no toca runs con lease vigente. - Abre
qa_reports/v2/runs/2026-09-22/g3/caos-2owners.jsonen el servidor y localiza el check Z10: ¿qué acción de A fue la primera rechazada y a qué hora UTC?