Skip to content

Design Event Log Core Engine

Kadyapam edited this page Jul 7, 2026 · 2 revisions

Design: Event-Log Core Engine (Phase 6)

Status: design + shadow slice landed (2026-07-05). The engine slice is merged (ehdb-reference::eventlog); the worker mirrors it as a disabled-by-default shadow. Cutover to serving the log from EHDB is NOT part of this phase — it is a later, separately-gated step.

Tracks: noetl/ehdb#241 (completion program) · RFC (decided) · Roadmap Phase 6.

This is the headline fix for NoETL's event-sourcing bottleneck: EHDB's event-log core engine becomes the durable persistence + ordering + serving layer underneath the append-only noetl.event producer path, replacing the NATS JetStream + PostgreSQL log-and-store path that is today's scaling pressure point (off-server state builder, unbounded WAL index noetl/ai-meta#166, materializer soak noetl/ai-meta#104, multi-replica affinity noetl/ai-meta#115).

The boundary this preserves

Event authorship is unchanged. The gateway/server remain the gatekeeper of what enters the log through the append-only producer path. This engine is the disk-and-index under an append the producer already authorized — it does not decide what is appended, it persists, orders, and serves it. Non-log data-plane roles never fabricate events.

Platform event log only — never business data. The engine serves NoETL's own platform event log. Tenant/domain data stays in external systems (Postgres/Snowflake/Kafka/…) reached via playbook connectors and never flows through this engine.

Semantics contract (match the JetStream + Postgres path)

The engine reproduces the guarantees the current path gives so the driver is swappable underneath without changing observable behaviour:

Property JetStream + Postgres today EHDB event-log engine
Global ordering JetStream stream sequence / Postgres ordering column StreamSequence on one canonical noetl_event_log stream — monotonic, gapless from 1
Per-execution scope events carry execution_id; read WHERE execution_id = ? subject noetl.event.exec.<execution_id>; per-execution read = subject-filtered replay
Tail / subscribe JetStream durable consumer durable consumer with ack cursor persisted in the transaction log
Offset / ack consumer ack, at-least-once explicit ack-after-materialize; cursor never moves backward
Append-only + immutable append-only stream / insert-only table records never mutated; KeepAll retention
Replay is truth rebuild state by replaying the log state is a replay of the on-disk transaction log

Ordering + sequence assignment

All events land on one canonical stream (noetl_event_log) per (tenant, namespace). Using a single stream makes its StreamSequence the global event-log sequence: next = record_count + 1, monotonic and gapless. Per-execution order is preserved because appends within an execution take increasing global sequences and the subject filter keeps them in sequence order — exactly the JetStream model (one stream, global seq, per-subject filter).

On-disk / segment format

The engine reuses the existing EHDB append/read primitives rather than a new format: the durable substrate is the LocalJsonlTransactionLog (one fsynced JSONL TransactionRecord per commit), replayed on open into the in-memory InMemoryStreamLog projection. Each event is a StreamMutation::Publish { stream, subject, payload, sequence } inside a single atomic commit; the stream is created on first use with KeepAll retention. This is the same segment-and-index shape the Phase C/D data-plane and event-stream helpers already use — the event-log engine is a purpose-scoped façade over it, not a second storage stack.

Production hardening (segmented append files + a sparse offset index for O(log n) sequence seeks, compaction, and cross-node replication) is deferred to the primary-serve phase; the local-reference substrate is the reference implementation of the contract, not the production disk format.

Driver interface (Phase 10-ready)

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

pub trait EventLogDriver {
    fn driver_name(&self) -> &'static str;
    fn append(&self, req: &EventLogAppendRequest) -> Result<EventLogAppendOutcome>;
    fn scan_global(&self, req: &EventLogScanRequest) -> Result<EventLogScanOutcome>;
    fn read_execution(&self, req: &EventLogReadExecutionRequest) -> Result<EventLogReadExecutionOutcome>;
    fn tail(&self, req: &EventLogTailRequest) -> Result<EventLogTailOutcome>;
    fn ack(&self, req: &EventLogAckRequest) -> Result<EventLogAckOutcome>;
}

LocalReferenceEventLogDriver is the EHDB implementation. A JetStream+Postgres driver implementing the same trait is what keeps the tier selectable back to the incumbent — callers program against the trait, not the concrete engine. All methods are &self: the engine is stateless per call (durable state lives in the on-disk log, opened + dropped per op), matching the bounded/stateless discipline of the rest of the integration.

How projections consume it (Phase 7 hook)

Phase 7's projection engine drives off this engine's serving surface: tail a durable consumer over the canonical stream, materialize each batch into a read-model, then ack the cursor. Because the cursor is durable and ack is explicit-after-materialize, a projector restart resumes exactly where it left off (at-least-once; the materialize step is idempotent on global_sequence). read_execution gives a projection a bounded per-execution rebuild without scanning the whole log. No Phase 7 code lands here — this documents the seam it will attach to.

Shadow / dual-write validation strategy

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

NOETL_EHDB_EVENTLOG = 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 already-authored platform event is mirrored into the EHDB engine alongside the existing JetStream+Postgres path and the mirror is compared against the authoritative log. Reads are never served from EHDB; the authoritative producer path is untouched.
  • primary — recognised but NOT activated in this phase. Requesting it is refused with a distinct outcome and the worker stays on the incumbent path. A compile-time PRIMARY_SERVE_ACTIVATED = false makes serving structurally impossible for this build; flipping it is the later, separately-gated cutover.

Parity checks (compare_shadow_parity)

Per mirrored event the shadow compares three properties and reports a single divergence reason on the first failure:

  • Sequence parity — EHDB's assigned sequence equals the authoritative sequence when it is known and comparable (a controlled selfcheck drive, or a JetStream stream-seq mirrored from origin). When the authoritative sequence is not a 1-based value aligned with the EHDB stream, the raw-value check is skipped and the shadow relies on count + ordering parity (the safe default).
  • Count parity — the EHDB record count advances by exactly one per mirror (no gap, no double-write), tracked with process-local contiguity bookkeeping seeded on the first mirror.
  • Ordering — the EHDB sequence is strictly greater than the previous mirror's (monotonic).

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 for the phase: a NoETL flow drives with EHDB shadowing the log and parity holds with zero divergence, while EHDB never serves a read.

Observability

Secret-free noetl_ehdb_eventlog_* metrics (op-family shape, labels operation + outcome only — never a log path, payload, subject, or stream): noetl_ehdb_eventlog_ops_total, noetl_ehdb_eventlog_last_ok, noetl_ehdb_eventlog_last_degraded (flags a parity mismatch or engine hiccup), and noetl_ehdb_eventlog_last_duration_seconds (mirror/append latency). Disabled ⇒ no lines rendered (byte-identical /metrics).

Config surface

Env var Default Meaning
NOETL_EHDB_EVENTLOG off Event-log driver selection: off / shadow / primary. Unrecognised ⇒ off (fail-safe).
NOETL_EHDB_EVENTLOG_MAX_PAYLOAD_BYTES 262144 Per-event payload cap (clamped to 1 MiB ceiling). Over-cap ⇒ rejected.

The event-log 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 eventlog-append / eventlog-scan / eventlog-read-exec / eventlog-tail / eventlog-ack / eventlog-suite, with distinct exit codes (0 ok / 3 rejected / 4 invalid / 5 unavailable). ehdb-selfcheck (worker-side, shipped in the worker image) gains mirror-eventlog and eventlog-suite for in-image kind validation, mirroring the selfcheck exit-code contract (0 ok / 3 rejected / 4 guard-refused-or-invalid / 5 degraded).

Assumptions where the current path is ambiguous

  • Authoritative sequence alignment. The worker cannot always observe a 1-based authoritative sequence comparable to the EHDB stream sequence (Postgres ordering is not a 1-based counter surfaced to the worker). The shadow therefore treats raw sequence-value parity as optional and always enforces count + ordering parity. Where a JetStream stream sequence is available and mirrored from origin, the raw check applies.
  • One canonical stream per (tenant, namespace). The design assumes a single logical event-log stream carries global ordering, matching the single noetl_events JetStream stream today. Sharded ordering (multi-stream, per-shard sequence) is a later concern aligned with the server-side shard-publish work (noetl/ai-meta#166 Phase 5), not this slice.
  • Retention. The event log is source-of-truth, so the engine uses KeepAll. Tiered GC / truncation is a separate, later concern (Postgres noetl.event tier-GC lineage) and is deliberately not a knob here.

Runtime integration — the live event-append hook

The shadow mirror is wired into the worker's live event path as of worker v5.67.0 (noetl/worker#167), closing the gap the 2026-07-06 live-in-kind session surfaced: before it, the only runtime EHDB hooks were the readiness preflight and the /metrics renderer, so the per-tier mirror engines ran only via ehdb-selfcheck + tests — a real drive did NOT dual-write into EHDB.

The seam

The hook lives at the worker's single authoritative event-append chokepoint: ControlPlaneClient::emit_event (repos/worker/src/client/control_plane.rs). Every worker event path funnels through it — EventEmitter::emit, emit_event_with_retry, spool_runtime, subscription, and plugin. After the control plane durably accepts an event (POST /api/events returns success), the hook mirrors that just-authored event into the derived EHDB fabric via the same eventlog::mirror_event shadow dual-write + parity path the selfcheck exercises. Hooking the chokepoint (rather than one high-level emitter) means all live events are mirrored regardless of which path emitted them.

Two functions in src/ehdb/eventlog.rs carry it:

  • runtime_hook_env(env) -> Option<EnvMap> — resolves the arming state once at client construction (the process env is immutable for the worker's lifetime, so the per-event path re-collects nothing). Returns Some(env) — "mirror every live event" — only when NOETL_EHDB_ENABLED is truthy and NOETL_EHDB_EVENTLOG=shadow and the resolved contract is a data-plane role (worker/playbook/system) with a local-reference log. Every other case (disabled, tier off/primary, control-plane role, malformed contract) returns None — a strict per-event no-op.
  • mirror_live_event(env, execution_id, payload) — the per-event call. Panic-isolated (catch_unwind) and best-effort: engine errors surface as metered non-ok outcomes and a panic is caught and returned as Unavailable, so a mirror failure never propagates into the authoritative event path. Passes authoritative_sequence = None (the worker doesn't know the server-assigned global sequence at emit time), so parity relies on the engine's count + monotonic-order invariant.

Invariants preserved

  • Disabled-by-default / byte-identical. No EHDB, tier off/ primary, or a control-plane role ⇒ the hook is a strict no-op and /metrics is unchanged.
  • Shadow never serves. The live hook only mirrors; the incumbent JetStream+Postgres path stays authoritative. Event authorship and event-log semantics are untouched — the gateway/server remain the gatekeeper of what is appended.
  • Live metrics. A real drive now advances noetl_ehdb_eventlog_ops_total{operation="mirror",…} and the noetl_ehdb_eventlog_last_* gauges — not just the selfcheck.

Remaining per-tier runtime hooks (follow-up slices)

The event-log tier is wired first. The other four tier mirrors plug into the same emit_event seam (or the analogous state/projection write sites) in follow-up slices:

Tier Engine module Live write site to hook
projection ehdb/projection.rs read-model materialization on event apply
kv / state ehdb/kv.rs state-cache / KV write path
object / blob ehdb/object.rs result/object store put
vector ehdb/vector.rs embedding/retrieval ingest path

Each is additive, shadow-only, disabled-by-default, and error-isolated, mirroring the event-log hook's shape.

What remains before Phase 7

  • Production disk format (segmented append files + sparse offset index) and compaction for the primary-serve path.
  • Sharded / multi-stream global ordering aligned with server shard publish.
  • The projection engine (Phase 7) attaching to the tail/ack/ read_execution serving surface.
  • A JetStream+Postgres EventLogDriver implementation for the tunable surface (Phase 10) so dual-run can select either engine per deployment.
  • Primary-serve cutover — a separately-gated step with dual-run verify + documented rollback, kind before GKE.

Related

Clone this wiki locally