FlowSpec, compilador y ClockEngine¶
Qué vas a aprender¶
- Qué es un FlowSpec y por qué es un DAG declarativo sin código arbitrario.
- Cómo funciona el registry de operadores (la allowlist de 22 operadores) y qué significa que
un operador esté
IMPLEMENTEDoPROPOSED. - Qué hace el compilador (
finaz_flow): de un FlowSpec a un CompiledPlan por modo (BATCH, REPLAY, STREAM), y qué rechaza. - Qué es el ClockEngine (
finaz_clock): proyección a rejillas de trading/sesión y relojes de volumen y de valor. - El Flow A del primer vertical, paso a paso, hasta el famoso golden
[null, null, 2, 3, 4, 5].
Fuentes
packages/finaz_flow/compile.py, registry/operators/operator_registry.json,
packages/finaz_clock/{clocks,project}.py, packages/finaz_features/operators/ema.py,
tests_v2/fixtures_shared/ema_period3_golden.json, el ejemplo del pack
contracts/examples/01_A_feature.json y los capítulos 06, 07 y 16 del pack. Los bloques de salida
que ves en este capítulo se obtuvieron ejecutando el código del repositorio (sólo lectura).
1. FlowSpec: la receta, no la cocina¶
Un FlowSpec describe qué se quiere calcular, no cómo ni dónde. Es un grafo dirigido acíclico (DAG) de nodos (cada uno es un operador con parámetros) unidos por aristas tipadas (de un puerto de salida a un puerto de entrada).
Analogía: la partitura
Un FlowSpec es una partitura: dice qué notas tocar y en qué orden. No dice si la toca una orquesta en directo (STREAM), un ensayo grabado que se reproduce (REPLAY) o alguien leyendo la partitura entera de una vez (BATCH). Esa decisión la toma el compilador, y el resultado —la partitura anotada para un intérprete concreto— es el CompiledPlan.
Anatomía de un FlowSpec¶
Estos son los campos del Flow A (ejemplo 01_A_feature.json del pack), resumidos:
{
"schema_version": "1.0.0",
"kind": "FlowSpec",
"flow_id": "fixture.flow.01_A_feature",
"flow_version": 1,
"nodes": [
{"node_id": "source", "operator_id": "source.snapshot", "operator_version": "1.0.0", "parameters": {"snapshot": {…}, "universe": {…}, "replay": {…}}},
{"node_id": "clock", "operator_id": "clock.project", "operator_version": "1.0.0", "parameters": {"clock": {…}}},
{"node_id": "ema", "operator_id": "feature.ema.sma_seed.reject_gap", "operator_version": "1.0.0", "parameters": {"period": 3, "field": "close", "clock": {…}}},
{"node_id": "persist", "operator_id": "result.persist", "operator_version": "1.0.0", "parameters": {"dataset_id": "fixture.result.ema"}}
],
"edges": [
{"from_node": "source", "from_port": "data", "to_node": "clock", "to_port": "data"},
{"from_node": "clock", "from_port": "clocked", "to_node": "ema", "to_port": "clocked"},
{"from_node": "ema", "from_port": "features", "to_node": "persist", "to_port": "value"}
],
"requested_modes": ["BATCH", "REPLAY", "STREAM"],
"time_domain": {"kind": "TimeDomainRef", "time_domain_id": "fixture.sim.01", "type": "SIMULATED", …},
"pit_policy": "SYNTHETIC_PIT_FIXTURE",
"failure_policy": "FAIL_CLOSED",
"outputs": ["persist.result"]
}
| Campo | Para qué sirve |
|---|---|
nodes[].operator_id + operator_version |
Operador pineado: nunca «la última versión». |
nodes[].parameters |
Parámetros validados contra el schema del operador. Las referencias a datos, universo o reloj van con hash (manifest_hash, definition_hash, calendar_hash…). |
edges |
Qué puerto alimenta a qué puerto. El compilador comprueba tipos, unidades, reloj y universo. |
requested_modes |
En qué modos se quiere poder ejecutar. Si un operador no soporta uno, el compilador lo rechaza, no lo degrada. |
time_domain |
El dominio temporal (aquí SIMULATED, con un simulation_start en nanosegundos como string). |
pit_policy |
Política point-in-time; en el fixture es sintética y declarada como tal. |
failure_policy |
FAIL_CLOSED: ante duda, se para; no se inventa. |
Sin código arbitrario
Un FlowSpec no puede contener SQL ni Python. La API (finaz_api/flujos.py) busca además
claves prohibidas en el cuerpo antes de aceptarlo. Todo cálculo pasa por un operador del registry.
2. El registry de operadores: una allowlist de 22¶
El fichero registry/operators/operator_registry.json (versión 1.0.0, compilador
finaz_flow/1.0.0 (AG-012)) contiene 22 operadores. Cada entrada tiene un descriptor
(modos soportados, semántica temporal, determinismo, si tiene estado, warmup…) y un
finaz_binding que dice qué función Python lo implementa y en qué estado está.
| Operador | Modos | Estado del binding | Implementación |
|---|---|---|---|
clock.project |
BATCH, REPLAY, STREAM | IMPLEMENTED | finaz_clock.project.project_batch (AG-010) |
feature.ema.sma_seed.reject_gap |
BATCH, REPLAY, STREAM | IMPLEMENTED | finaz_features.operators.ema.ema_sma_seed_reject_gap (AG-011) |
source.snapshot |
BATCH, REPLAY, STREAM | PROPOSED | (AG-008/AG-009) |
result.persist |
BATCH, REPLAY, STREAM | PROPOSED | (AG-009/AG-013) |
features.tabular, panel.align |
los 3 | PROPOSED | — |
forecast.catboost, forecast.statsforecast, forecast.nhits, forecast.chronos2, forecast.timesfm25 |
los 3 | PROPOSED | — |
source.alpha, source.risk, source.portfolio, source.target, source.prices |
los 3 | PROPOSED | — |
allocation.mean_variance, constraints.continuous, discrete.cp_sat, constraints.discrete |
los 3 | PROPOSED | — |
replay.feed, source.redpanda |
sólo STREAM | PROPOSED | — |
Así se llega al recuento de docs/v2/MOTORES.md: 22 operadores, 2 IMPLEMENTED y 20 PROPOSED.
Si source.snapshot y result.persist son PROPOSED, ¿cómo se ejecuta el Flow A?
El binding del registry es la vía genérica de un operador. En el primer vertical, la lectura del
snapshot y la persistencia las hace directamente el pipeline del runtime
(finaz_runtime/core/pipeline.py: snapshot → clock → EMA → persist → publish), con los paquetes
finaz_storage/finaz_data y finaz_control. Por eso el plan compilado muestra
"proposed.source.snapshot" como binding: el registro es honesto sobre qué operador tiene ya una
función genérica enlazada y cuál no. Otros operadores PROPOSED (por ejemplo forecast.nhits)
tienen adaptadores reales que se ejecutan como jobs (capítulo 5), pero todavía no como nodos
de un flow.
Analogía: la carta del restaurante
El registry es la carta. Si pides un plato que no está en la carta, la cocina no improvisa: rechaza la comanda. Algunos platos están en la carta con la nota «próximamente» (PROPOSED): puedes ver su descripción, pero el compilador sólo encadena lo que la carta permite.
3. El compilador: de FlowSpec a CompiledPlan¶
finaz_flow.compilar_modo(flujo, modo) sigue el orden del capítulo 06 del pack:
flowchart LR
A[parse / schema] --> B[límites de tamaño]
B --> C[registry y refs<br/>operador+versión pineados]
C --> D[DAG y puertos<br/>ciclos, tipos]
D --> E[shape / universo]
E --> F[unidades, target,<br/>horizonte, reloj]
F --> G[roles temporales<br/>y guardas de modo]
G --> H[capacidades y recursos]
H --> I[fronteras y estado]
I --> J[optimizaciones legales<br/>CSE / fusión]
J --> K[hash y validación<br/>del plan]
Cada error identifica nodo, puerto y regla. Veamos las guardas por modo, que son la parte más pedagógica:
| Modo | Qué exige el compilador |
|---|---|
| BATCH | El operador declara BATCH en supported_modes. Es el único modo que admite semántica NONCAUSAL_RESEARCH_ONLY (investigación, nunca «como se conocía»). |
| REPLAY | Soporte REPLAY; rechaza operadores no causales (NONCAUSAL_RESEARCH_ONLY, LABEL_OUTCOME) y no deterministas. El plan conserva semilla y dominio simulado. |
| STREAM | Lo de REPLAY y además: todo operador con estado debe declarar state_schema_hash (si no, CHECKPOINT_INCOMPATIBLE), porque en stream hay que poder hacer checkpoint y restaurar. |
Y las optimizaciones sólo cuando son demostrablemente legales:
- CSE (reglas CSE-1..3): dos nodos se fusionan sólo si coinciden operador, versión, parámetros canónicos, entradas, reloj, universo y políticas. El estado de una EMA nunca se comparte entre orígenes distintos.
- Fusión (reglas F-1..4): sólo entre operadores sin estado, deterministas, causales y sin barrera
de disponibilidad/commit.
result.persist,source.*,replay.feedysource.redpandanunca se fusionan. La fusión elimina transporte, no cambia orden ni semántica.
El CompiledPlan real del Flow A¶
Compilando el ejemplo del pack en modo BATCH, el resultado (recortado) es:
{
"kind": "CompiledPlan",
"plan_id": "fixture.flow.01_A_feature.plan.batch.v1",
"flow_hash": "sha256:a6146520d429207df049ea51050ac15898491d2af90f9bff9e6d35d8b60524ef",
"mode": "BATCH",
"compiler_version": "finaz_flow/1.0.0 (AG-012)",
"operator_bindings": {
"source": "proposed.source.snapshot",
"clock": "finaz_clock.project.project_batch",
"ema": "finaz_features.operators.ema.ema_sma_seed_reject_gap",
"persist": "proposed.result.persist"
},
"stages": {"uri": "artifact://sha256/5f8c2fc0…", "rows": 4, "format": "JSON", …},
"checkpoint_schema_hash": "sha256:27b3eaa9…",
"capability_certificate_hashes": ["sha256:a2524b89…", "…"],
"optimization_passes": [],
"plan_hash": "sha256:db758f39…"
}
Fíjate en tres cosas:
stagesno va en línea: es una referencia a un artefacto direccionado por contenido (artifact://sha256/…), con su hash y su número de filas (4 etapas, una por nodo).optimization_passesestá vacío: en Flow A no hay nada legal que fusionar (todos los pares tocan unsource.*, un operador con estado oresult.persist).plan_hashsella todo: si cambias un parámetro, cambia el hash y es otro plan.
El mismo FlowSpec compila también en REPLAY y STREAM (los tres modos devuelven OK).
# POST /platform/v1/plans compila un plan (capacidad write:runs)
curl -sS -X POST http://127.0.0.1:18300/platform/v1/plans \
-H "Authorization: Bearer <TOKEN>" \
-H "Idempotency-Key: <CLAVE>" \
-H "Content-Type: application/json" \
--data @flowspec_con_mode.json
El cuerpo debe ser el FlowSpec completo (más mode). Probado en el servidor el 2026-09-22
con sólo flow_id + flow_version, la API responde 404 con un mensaje claro:
{"code":"REFERENCE_NOT_FOUND",
"message":"referencia sin resolver: compilar exige el FlowSpec completo (la resolución la hace quien admite)",
"details":{"flow_id":"fixture.flow.01_A_feature","flow_version":1}, …}
Es decir: resolver flow_id:vN a partir del repositorio es trabajo de la admisión de runs
(POST /runs), no del compilador expuesto.
Lo que el compilador rechaza (ejemplos reales)¶
Modificando el Flow A en memoria y compilando en BATCH:
# Operador fuera de la allowlist:
FlowError [OPERATOR_UNKNOWN] operador fuera del allowlist en nodo ema: feature.python.arbitrario
# Arista persist → source (cierra un ciclo):
FlowError [DAG_CYCLE] el flujo contiene un ciclo entre: clock,ema,persist,source
Otros rechazos documentados en la suite de AG-012: modo no pedido o no soportado
(UNSUPPORTED_MODE), tipos de puerto incompatibles, y que una etiqueta (LABEL) alimente una
decisión.
4. El ClockEngine: ¿qué hora es para el mercado?¶
Un timestamp UTC no dice por sí solo «qué barra es». El ClockEngine (finaz_clock, AG-010)
traduce marcas de tiempo a coordenadas económicas: índice de barra dentro de la sesión, sesión a
la que pertenece, o barra de volumen/valor.
Analogía: el reloj de un partido
El reloj de pared sigue corriendo en el descanso, pero el reloj del partido se detiene. El ClockEngine es el reloj del partido: fuera de sesión no avanza, y si llega un evento en el descanso no le asigna minuto, lo marca como GAP.
4.1 Proyección trading/sesión¶
Una SessionGrid es una lista de intervalos semiabiertos [open, close) en nanosegundos UTC. A
partir de ella, build_minute_grid construye una rejilla de minutos negociables y project_batch
asigna cada marca a su barra. Ejemplo ejecutado con dos sesiones sintéticas (s1 de 0 a 3 min,
s2 de 10 a 12 min):
ts_ns clock_index flag session_id
0 0 OK s1
60000000005 1 OK s1
120000000000 2 OK s1
240000000000 -1 GAP None ← minuto 4: fuera de sesión
600000000000 3 OK s2 ← el índice continúa: tiempo NEGOCIABLE
660000000001 4 OK s2
Observa dos reglas del pack:
- Nunca interpola en silencio: sin barra contenedora no hay índice;
clock_index = -1y flagGAP(Q021). El propio contratoClockedRowrechaza un GAP con índice ≥ 0. - Tiempo negociable: entre
s1ys2pasan 7 minutos de pared, pero el índice sólo avanza de 2 a 3. Los intervalos fuera de sesión no acumulan.
El orden de salida es el de entrada y los empates se resuelven por posición: determinista.
4.2 Relojes de volumen y de valor¶
En vez de cortar una barra cada minuto, puedes cortarla cada vez que se negocia un umbral de
volumen (acciones) o de valor (nocional). accumulate_volume(sizes, threshold) hace exactamente eso,
con política WHOLE_EVENT: un evento nunca se trocea.
Ejemplo ejecutado con umbral 100 y volúmenes [40, 30, 50, 250, 0, 20]:
| Barra | Cierra en evento | Volumen de la barra | Umbrales cruzados | Exceso (overshoot) |
|---|---|---|---|---|
| 0 | 2 (40+30+50) | 120 | 1 | 20 |
| 1 | 3 (el evento de 250) | 250 | 2 | 70 |
Y el cursor queda en cumulative=390, bars_closed=2: los eventos 4 (volumen 0, que no avanza) y 5
(20) están acumulados esperando el siguiente umbral (400).
El evento gigante
El evento de 250 cruza dos umbrales de golpe. El reloj no inventa dos barras idénticas ni parte
el evento: cierra una barra con crossed_thresholds=2 y conserva el volumen total (Q023).
Cualquier otra política (large_event_policy) exige un operator ID nuevo; hoy sólo WHOLE_EVENT
está certificada.
accumulate_value es lo mismo con nocionales. Ambos aceptan un cursor para continuar donde se
quedó el lote anterior: por eso procesar en trozos da lo mismo que procesar todo de una vez
(«paridad chunks ≡ batch», Q009). VolumeClock envuelve esto con snapshot/restore para el stream.
4.3 Causalidad en relojes avanzados¶
Dos guardas devuelven FUTURE_INFORMATION si un reloj pide futuro:
exigir_anchor_countdown: una cuenta atrás necesita un ancla versionada y conocida; un ancla futura desconocida no se puede calcular retrospectivamente.exigir_varianza_causal: una varianza full-sample usa el futuro (rechazo); y durante el warmup aún no hay semillas disponibles.
5. Flow A de principio a fin¶
Ahora juntamos todo. El Flow A del primer vertical (pack §16) es:
flowchart LR
S["source.snapshot<br/>fixture.snapshot.01<br/>manifest 55ad7c17…"] -- data --> C["clock.project<br/>fixture.clock.trading_minutes"]
C -- clocked --> E["feature.ema.sma_seed.reject_gap<br/>period = 3, field = close"]
E -- features --> P["result.persist<br/>fixture.result.ema"]
- Snapshot: se leen las barras del snapshot sellado
fixture.snapshot.01(manifestsha256:55ad7c17…). Nada fuera de ese paquete. - Clock: cada barra recibe su índice de reloj de trading (o GAP).
- EMA:
feature.ema.sma_seed.reject_gapcon periodo 3 sobreclose. - Persist: el
FeatureBatchresultante se guarda (CAS) y se publica unResultRef.
La EMA con semilla SMA, a mano¶
La semántica está portada con aritmética exacta del oráculo congelado:
- Semilla: media simple (SMA) de los primeros
periodcierres (conmath.fsum). - α = 2 / (period + 1) → con periodo 3, α = 0,5.
- Primer índice válido =
period − 1= 2.
Con la entrada [1, 2, 3, 4, 5, 6]:
| i | cierre | cálculo | EMA | flag |
|---|---|---|---|---|
| 0 | 1 | calentamiento | null |
WARMUP |
| 1 | 2 | calentamiento | null |
WARMUP |
| 2 | 3 | SMA(1,2,3) = 6/3 | 2 | — |
| 3 | 4 | 0,5·4 + 0,5·2 | 3 | — |
| 4 | 5 | 0,5·5 + 0,5·3 | 4 | — |
| 5 | 6 | 0,5·6 + 0,5·4 | 5 | — |
Resultado: [null, null, 2, 3, 4, 5] con máscara [false, false, true, true, true, true]. Es
exactamente el fixture golden tests_v2/fixtures_shared/ema_period3_golden.json, cuya fuente
declarada es finazbench/features/canonical.py:108-110. Los tests lo comparan con ==, no con
tolerancia.
>>> ema_sma_seed_reject_gap([1.0, 2.0, 3.0, 4.0, 5.0, 6.0], 3)
ResultadoEma(valores=(None, None, 2.0, 3.0, 4.0, 5.0),
mascara_valida=(False, False, True, True, True, True),
flags=(('WARMUP',), ('WARMUP',), (), (), (), ()))
¿Y el «reject_gap»?¶
El sufijo del nombre es su política de huecos (§5.2 v1): ante un valor no finito o ausente, emite
nulo con flag GAP desde el primer hueco y sin reanudar. Ejecutado con un hueco en la posición 3:
>>> ema_sma_seed_reject_gap([1.0, 2.0, 3.0, None, 5.0, 6.0], 3)
ResultadoEma(valores=(None, None, 2.0, None, None, None),
mascara_valida=(False, False, True, False, False, False),
flags=(('WARMUP',), ('WARMUP',), (), ('GAP',), ('GAP',), ('GAP',)))
Reanudar sería otro operador
Podrías pensar «que la EMA se reinicie tras el hueco». Sería otra semántica, y en el runtime una semántica
distinta es otro operator_id (Q008). Así un resultado nunca cambia de significado por
debajo de un mismo nombre.
Una sola lógica, tres modos¶
La promesa del primer vertical es que el mismo FlowSpec da el mismo resultado lógico en
BATCH, REPLAY y STREAM. El capítulo 3 muestra cómo se comprueba: el emit del stream se compara por
hash (values_hash) con un golden de referencia.
Resumen¶
- Un FlowSpec es un DAG declarativo: operadores pineados, parámetros con hashes, aristas tipadas, modos pedidos, dominio temporal y políticas. Sin SQL ni Python arbitrario.
- El registry es una allowlist de 22 operadores: 2
IMPLEMENTED(clock.project,feature.ema.sma_seed.reject_gap) y 20PROPOSED. - El compilador genera un CompiledPlan por modo, con guardas de causalidad y determinismo
(REPLAY/STREAM) y de estado checkpointable (STREAM); optimiza sólo si es legal; sella con
plan_hash. - El ClockEngine proyecta a tiempo negociable con GAP explícito y ofrece relojes de volumen/valor
WHOLE_EVENT, continuables por cursor. - El Flow A produce el golden
[null, null, 2, 3, 4, 5], bit a bit igual al oráculo del banco de backtesting.
Para practicar¶
- Calcula a mano la EMA con semilla SMA de periodo 3 para
[2, 4, 6, 8]. ¿Cuál es el primer índice válido? - En
registry/operators/operator_registry.json, localiza el descriptor defeature.ema.sma_seed.reject_gapy busca los camposstatefulystate_schema_hash. ¿Por qué el compilador los mira sólo en STREAM? - Con la tabla de volúmenes del apartado 4.2, ¿qué barra se cerraría si el siguiente evento fuera de 15? ¿Y de 10?
- Modifica una copia del Flow A cambiando
perioda 5 y compílala en BATCH. ¿Cambiaplan_hash? ¿Yflow_hash? - Explica con tus palabras por qué
result.persistnunca se fusiona con el operador anterior.