Skip to content

Design Projection Read Model Engine

Kadyapam edited this page Jul 6, 2026 · 3 revisions

Design: Projection / Read-Model Engine (Phase 7)

Status: Phase 7 shadow merged (ehdb#243, e0f1c0f; worker#157, eadc3a5); Phase 9 tier-2 PRIMARY-serve merged (ehdb#248, d08013c; worker#162, a56583c, released v5.62.0 36875e3) (2026-07-05). The worker dual-materializes the engine as a disabled-by-default shadow and now, under NOETL_EHDB_PROJECTION=primary, serves the read-models authoritatively from EHDB in place of the PostgreSQL materializer — reversible via two levers (see Roadmap Phase 9 tier 2). The prod/GKE cutover is a separate later step gated on the user; nothing in prod changed (NOETL_EHDB_PROJECTION defaults off, prod stays on the incumbent).

Tracks: noetl/ehdb#241 (completion program) · RFC (decided) · Roadmap Phase 7 · builds on Design: Event-Log Core Engine (Phase 6).

Phase 6 put EHDB's event-log engine underneath the append-only noetl.event producer path. Phase 7 puts the projection / read-model engine on top of that log's tail — the engine that builds and serves the materialized read-models the control plane queries, retiring the second half of the event-sourcing bottleneck: the PostgreSQL materializer and its projected state tables (materializer soak noetl/ai-meta#104, sole-writer noetl/ai-meta#103, multi-replica coherence noetl/ai-meta#115).

The boundary this preserves

A projection is a derived read-model, never an event author. The engine is built by consuming the append-only event log; it never writes an event back. The #103 sole-writer invariant is preserved — the Postgres materializer stays the only noetl.event writer during the shadow; the EHDB projection engine only reads the log and materializes a separate read-model store. Event authorship stays with the gateway/server producer path.

Platform read-models only — never business data. The engine materializes NoETL's own platform read-models (execution / event / runtime state). Tenant/domain data stays in external systems reached via playbook connectors and never flows through this engine.

What the Postgres materializer produces today (the contract to mirror)

The path Phase 7 replaces folds the noetl.event log into three read-model surfaces the control plane queries:

Read-model Written by Idempotency Query contract
noetl.event (event read-model) POST /api/internal/events/project → project_events (single-statement batch INSERT … SELECT jsonb_array_elements) ON CONFLICT (event_id) DO NOTHING — re-projecting an event_id is a no-op get-by-event_id; list-by-execution_id; ordered by event_id
projection_snapshot (execution-state read-model) POST /api/internal/projection/advance → advance_snapshot(execution_id) (bounded rebuild, no dispatch) monotonic versioned upsert — a redelivered batch is a forward no-op get folded state (status, current node) by execution_id
consumer offset event_stream.position upsert (ON CONFLICT (name) DO UPDATE) last-write-wins per consumer name resume the drain from the durable cursor

The materializer's consumer model is ack-after-materialize on a durable JetStream pull consumer (worker src/materializer.rs): drain a bounded batch with deferred ack → POST events/project → ack only on 2xx; idempotent ON CONFLICT makes a redelivered batch a no-op double-write. The Phase 7 engine reproduces the same shape over the Phase-6 EventLogDriver tail.

Semantics contract (match the Postgres materializer)

Property Postgres materializer today EHDB projection engine
Event read-model noetl.event row per event_id one materialized record per event_id, keyed + de-duped
Idempotent insert ON CONFLICT (event_id) DO NOTHING event_id dedup guard + global_sequence <= checkpoint skip
Execution state projection_snapshot (monotonic version) folded ExecutionStateView (monotonic fold)
Durable offset event_stream.position ProjectionCheckpoint.applied_through_sequence
Ordering key event-log order (event_id / stream seq) Phase-6 global sequence (gapless, monotonic)
Replay is truth rebuild snapshot from noetl.event rebuild read-models by replaying the projection store / re-applying the log

What the engine materializes

The engine composes the append-only stream primitives (the same substrate Phase 6 uses) over one noetl_projection_log store per (tenant, namespace), scoping each materialized record to its execution by subject noetl.projection.exec.<execution_id>. From that store it serves:

  • Event read-model (EventReadModelView) — the projected columns of one event (global_sequence, event_id, execution_id, event_type, node_name, status, prev_event_id), keyed on event_id. The ON CONFLICT (event_id) DO NOTHING twin.
  • Execution-state read-model (ExecutionStateView) — the per-execution fold: derived status, current_node, event_count, first/last_global_sequence, last_event_id, terminal + terminal_event_type. The projection_snapshot twin. Status is derived from the terminal event when one is folded (both dotted playbook.completed and underscore playbook_completed spellings match, as in the worker state_materializer and server event_write::is_terminal), else the latest event's status, else running.
  • Consumer checkpoint (ProjectionCheckpoint) — the highest event-log global_sequence materialized through, plus the total applied count. The event_stream.position twin.

Idempotent, exactly-once, replay-safe apply

Apply is idempotent keyed on the Phase-6 global sequence — the gapless, monotonic sequence the event-log engine assigns:

  1. The engine persists a checkpoint = the highest global sequence materialized so far (rebuilt from the store on open).
  2. An incoming event whose global_sequence <= checkpoint is skipped (skipped_below_checkpoint) — an at-least-once redelivery / replay is a no-op. This is the exactly-once materialization guarantee.
  3. An already-materialized event_id is de-duped (duplicates) as a second guard (the ON CONFLICT (event_id) twin), belt-and-suspenders even if sequences ever realign.
  4. The execution-state fold is monotonic — it only advances by the highest event seen — so re-folding a superset is a forward no-op.

The batch is applied in ascending global_sequence order (sorted defensively so exactly-once holds even on a shuffled batch, though the Phase-6 tail delivers in order). A batch is committed only when something was actually materialized — an all-skipped batch is a pure no-op, never an empty transaction.

Rebuild-from-log is deterministic and bounded: the read-models are the replay-fold of the projection store, so the same event set yields identical read-models regardless of batch boundaries (a rebuild_from_event_log_is_deterministic unit test drives four single-event batches vs one four-event batch and asserts byte-equal read-models + equal checkpoints).

Consumer / checkpoint model (the sole-writer rule, #103, preserved)

The projection engine is a read-model projector off the event log — a sibling of the Postgres materializer, not a competitor. Under the shadow it runs alongside the incumbent on its own consumer cursor: it tails the Phase-6 durable consumer, materializes, then advances its own checkpoint. It never writes noetl.event, never touches the incumbent's projection_snapshot / event_stream.position, and never blocks the authoritative materializer. The #103 sole-writer of noetl.event is untouched; the projection store is a separate, derived fabric.

Driver interface (Phase 10-ready)

The engine is exposed behind the ProjectionDriver trait so the projection tier is driver-selectable (the RFC's tunable-backend surface, Roadmap Phase 10):

pub trait ProjectionDriver {
    fn driver_name(&self) -> &'static str;
    fn apply(&self, req: &ProjectionApplyRequest) -> Result<ProjectionApplyOutcome>;
    fn read_execution_state(&self, execution_id: &str) -> Result<ProjectionReadExecutionOutcome>;
    fn read_event(&self, event_id: i64) -> Result<ProjectionReadEventOutcome>;
    fn list_executions(&self, limit: usize) -> Result<ProjectionListExecutionsOutcome>;
    fn checkpoint(&self, consumer: &str) -> Result<ProjectionCheckpoint>;
}

LocalReferenceProjectionEngine is the EHDB implementation. A Postgres-materializer driver implementing the same trait keeps the tier selectable back to the incumbent — callers program against the trait. The read_execution_state / read_event / list_executions methods are the read-serving interface Phase 9 cuts over to; nothing serves reads from them in this phase (Phase 9's cutover flips the control plane to read here). All methods are &self: durable state lives in the on-disk store, opened + dropped per op (bounded / stateless), matching the rest of the integration.

ProjectionEventInput::from_event_log_record bridges a Phase-6 EventLogRecordView (from the tail / scan_global surface) into an apply batch — the seam the Phase-6 design's "How projections consume it" section set up.

Shadow / dual-materialize validation strategy

The engine ships disabled by default. The worker wires it as a shadow (worker-rust src/ehdb/projection.rs) behind a driver-selection flag:

NOETL_EHDB_PROJECTION = off | shadow | primary        (default: off)
  • off — strict no-op. No engine opened, no metric recorded; the worker's /metrics and behaviour are byte-identical to a build without the wiring.
  • shadow — each event-log tail batch is dual-materialized into the EHDB projection engine alongside the authoritative Postgres materializer, and the EHDB read-models are compared against the incumbent's output. Reads are never served from EHDB; the authoritative materializer is untouched.
  • primary — EHDB serves the read-models authoritatively (Phase 9 tier 2): the read-model queries the control plane makes (list_executions, per-execution read_execution_state, read_event) are served by the EHDB engine in place of the PostgreSQL materializer, while the served read-models are dual-run parity-checked against the incumbent (served_primary, or primary_divergence on divergence). PRIMARY_SERVE_ACTIVATED = true so this build can serve; whether it does is the runtime NOETL_EHDB_PROJECTION choice, so the cutover is reversible without a redeploy (flip back to shadow/off restores the materializer read path with zero data loss; setting PRIMARY_SERVE_ACTIVATED = false is the structural kill switch). See Roadmap Phase 9 tier 2 for the rollback procedure.

Parity checks (compare_projection_parity)

The shadow compares the EHDB read-models against the authoritative materializer's output and reports a single divergence reason on the first failure:

  • Key parity — the EHDB read-model has exactly the authoritative execution set (no missing / extra execution rows).
  • Value parity — every shared execution's derived status, terminal flag, and event_count match the authoritative materializer.
  • Checkpoint parity / lag — the EHDB checkpoint has caught up to (>=) the authoritative offset; the lag is surfaced when it hasn't.

Comparison is against AuthoritativeExecutionState — the minimal secret-free projection of the incumbent read-model the worker can observe (execution id, status, event count, terminal). A divergence produces the parity_mismatch outcome — degraded but non-fatal: it surfaces on metrics without disturbing the authoritative path. This is the acceptance evidence: a NoETL flow drives with EHDB dual-materializing and parity holds with zero divergence, while EHDB never serves a read.

Observability

Secret-free noetl_ehdb_projection_* metrics (op-family shape, labels operation + outcome only — never a log path, payload, subject, or execution id): parity mismatches, materialize lag (checkpoint lag behind the authoritative offset), and apply latency. Disabled ⇒ no lines rendered (byte-identical /metrics).

Config surface

Env var Default Meaning
NOETL_EHDB_PROJECTION off Projection driver selection: off / shadow / primary. Unrecognised ⇒ off (fail-safe).
NOETL_EHDB_PROJECTION_MAX_BATCH 4096 Per-apply event-batch cap (MAX_APPLY_BATCH). Over-cap ⇒ rejected.

The projection wiring also respects the existing EHDB contract (NOETL_EHDB_ENABLED + NOETL_EHDB_MODE=local_reference + NOETL_EHDB_CLIENT_ROLE data-plane + NOETL_EHDB_LOCAL_REFERENCE_LOG). Control-plane roles (gateway/api/server) are refused before any engine opens.

CLI

ehdb-local-reference (engine-side) gains projection-apply / projection-read-exec / projection-read-event / projection-list / projection-checkpoint / projection-from-eventlog (the Phase-6 tail → Phase-7 apply bridge) / projection-suite, with distinct exit codes (0 ok / 3 rejected / 4 invalid / 5 unavailable). ehdb-selfcheck (worker-side, shipped in the worker image) gains the projection shadow verbs for in-image kind validation, mirroring the selfcheck exit-code contract.

Assumptions where the current path is ambiguous

  • Read-model shape. The Postgres projection_snapshot carries a richer, orchestrator-specific rebuild than the minimal derived state the reference engine folds (status, current node, counts, terminal). The reference engine materializes the queryable contract — the fields the control plane reads to answer "what is the state of execution X" — not the internal snapshot blob. The parity comparison is scoped to the observable fields (AuthoritativeExecutionState); a production driver widens the fold to the full snapshot columns when the primary cutover needs them.
  • Authoritative offset alignment. The worker cannot always observe a 1-based authoritative offset comparable to the EHDB checkpoint (the incumbent's event_stream.position is a JetStream stream sequence, not a 1-based projector counter). The parity check therefore treats the checkpoint-lag comparison as optional (skipped when the authoritative offset is not surfaced) and always enforces key + value parity.
  • One projection store per (tenant, namespace). The design assumes a single logical projection store, matching the single-materializer model today. Sharded projection (per-shard store, aligned with server shard publish noetl/ai-meta#166) is a later concern, not this slice.
  • Terminal taxonomy. Both dotted and underscore spellings of the terminal event types are matched (the WAL carries both). A missed terminal leaves an execution running (eventually re-folded), never wrong.

What remains before Phase 8

  • Production projection store format (segmented + indexed) and compaction for the primary-serve path.
  • Richer execution-state fold (full projection_snapshot column parity) for the read-cutover.
  • Sharded / multi-store projection aligned with server shard publish.
  • A Postgres-materializer ProjectionDriver implementation for the tunable surface (Phase 10) so dual-run can select either engine per deployment.
  • Primary read-serving cutover off Postgres — a separately-gated step (Phase 9) with dual-run verify + documented rollback, kind before GKE.
  • Phase 8 is the KV/state + object/blob + vector engines; Phase 7 completing is the projection tier only.

Related

Clone this wiki locally