Saltar a contenido

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_token y 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 → RUNNING se hace con UPDATE … 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 CANCELLED antes de ejecutar. No se interrumpe un cálculo a medias.
  • Fallo honesto: si la ejecución falla, FAILED con el mensaje en errores, 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 el ResultRef por 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 RUNNING sin lease vigente → el worker anterior murió: vuelve a QUEUED y su evento a pendiente.
  • (b) run en QUEUED con 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_CONFIRMED o END_OF_BOUNDED_SOURCE sacan 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:

  1. Adquiere el lease del run en PG antes de consumir.
  2. 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.
  3. Un latido renueva el lease cada ttl/3.
  4. Un owner obsoleto falla con STREAM_FENCE_OBSOLETO sin ejecutar el efecto; un intruso con el lease ajeno vigente falla con STREAM_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_offsets por 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_hash estable).
  • 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 :s3 aún pendiente.
  • Paridad: values_hash sha256:9e906a09… = golden [null, null, 2, 3, 4, 5].

Para practicar

  1. Dibuja la tabla de leases tras esta secuencia: A adquiere; pasan 7 s (TTL 6); B adquiere; A intenta renovar. ¿Qué fence_token tiene cada uno y qué error recibe A?
  2. En el fault schedule, ¿qué fallo explica duplicates_skipped: 1 en la paridad?
  3. ¿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.
  4. Lee _recuperar_huerfanos en packages/finaz_runtime/worker.py y explica por qué no toca runs con lease vigente.
  5. Abre qa_reports/v2/runs/2026-09-22/g3/caos-2owners.json en el servidor y localiza el check Z10: ¿qué acción de A fue la primera rechazada y a qué hora UTC?