# Implementation Plan — Event Outbox (E4) Continuation plan for finishing the outbox on the EC2. Branch: `feat/event-outbox-e4`. Concept: `docs/event-pipeline-design-summary.md`. Backlog: `docs/enhancements.md` → E4. ## Reconciliation (2026-10-05) — the original plan was partly stale After steps 1–2 were committed, the branch moved on to E7 (subscriber/ subscription model + CIM dispatch). Two assumptions in the first draft of this plan no longer hold, so the design was re-grounded against the live code: - **Transport is HTTP, not Kafka.** The relay was to publish to Kafka `dori.alerts.dispatch` for a CIM Kafka consumer — but that consumer never existed (CIM only consumes `topic_incidents_created`). The live alert path is `writer → INSERT inference_event → trigger pg_notify('event_insert') → DIM listener → forward_event() HTTP /event_process/ → CIM process_event`. The relay therefore targets **CIM over HTTP + the WebSocket fan-out**, and lives in **DIM** (which already owns `forward_event` + the WS registry). - **Alerting is decided by subscriptions, not event-type flags.** `is_trigger`/ `is_message` were dropped (migration 023). "Does it alert" = "an active `subscription` matches" (CIM `dispatch_via_routes`). The gate is now `is_notification OR has_matching_subscription` — exactly what the current trigger (`create_notify_trigger.subscriber.sql`, migration 025) computes. ### Staging (decided with Ravi) - **Stage 1 — durable delivery (this change).** Keep the trigger's enrichment + gate, but have it write a durable `event_outbox` row (+ `NOTIFY 'outbox'`) instead of `pg_notify`-ing the payload. DIM's relay drains the table and delivers to CIM + WS. Zero writer changes, identical payload, no regression, independently verifiable on EC2. This delivers E4's core value (no lost alerts). - **Stage 2 — routing in the writer (later, after Stage 1 verifies).** Move the enrichment + gate into an app `emit_event()`, swap the writers, and drop the trigger (the folded-in E5). Plan steps 4/6/7 below. ## Status | Step | State | |---|---| | Design summary + backlog | ✅ committed (`c7058b9`) | | 1. Migration `019_event_outbox` (table + indexes) | ✅ `alembic/versions/019_event_outbox.py` | | 2. ORM `DEventOutbox` | ✅ end of `common/dori_model/dModel.py` | | **Stage 1** | | | 3. Relay (`run_outbox_relay`, in DIM, HTTP+WS) | ✅ `services/dim/src/service/outbox_relay.py` | | 3a. Trigger writes outbox row + `NOTIFY 'outbox'` | ✅ migration `026_outbox_trigger` + `create_notify_trigger.outbox.sql` | | 3b. DIM lifespan: relay replaces listener | ✅ `services/dim/src/main.py` | | 3c. `forward_event(raise_on_error=)` for retry | ✅ `services/dim/src/service/cim_client.py` | | 3d. Relay unit tests (backoff/retry/dead-letter) | ✅ `unit_tests/test_outbox_relay.py` | | 3e. EC2 verify: migrate → inject → dispatched; CIM-down → pending→retry(2/4/8s)→recover | ✅ 2026-10-05 on mapp (rows 1-3 dispatched; no loss) | > **Fixed during 3e:** `_RETRY_SQL` used a bare `$2 >= $4`, which asyncpg > rejects (`AmbiguousParameterError`) so the backoff UPDATE never ran and the > relay crash-looped on every CIM failure (loss-free only by rollback luck). > Now cast (`$2::int` etc.); regression-guarded in the unit test. > > **Retention TODO (Stage 2):** `dispatched` rows are never pruned — add a > reaper (e.g. delete dispatched older than N days) before customer-event volume > lands in the outbox (step 6). | **Stage 2** (verified on mapp 2026-10-05) | | | 4. Writer swap → shared app `emit_event()` (enrich+gate+insert+outbox) | ✅ `common/dori_utils/events.py`; monitor (`mlc_check` both modes + `system_event_service`) + analytics (`inference_ingest`) | | 5. CIM dedupe on `alert_id` (= event_id) for at-least-once | ✅ `dori.event_dispatch` (migration 027) + CIM `_already_dispatched`/`_mark_dispatched` | | 6. Route customer events (`inference_ingest.py`) through `emit_event()` | ✅ | | 7. Drop the trigger (last, after verify) | ✅ migration 028 (two-phase: 027 dedupe+double-write, then 028 drop); trigger count now 0 | | + Reaper for dispatched/dead rows | ✅ DIM relay `_maybe_reap`, `outbox_retention_days` | | + Two editable lanes on event_type: `is_notification` (UI) + **`is_alert`** (CIM) | ✅ migration 029; emit_event lights both from flags + payload `isNotification`/`isAlert`; relay gates WS vs CIM; flags API + UIs (/fe/ Events toggles, ICD Events chip) | | 8. ICD/COD UI: remove "Contact Points"→**Subscribers** and "Notification Routes"→**Subscriptions** (smartvision_va ContactPointsICD/NotificationRoutesICD still call dead `contact_point`/`notification_route`; repoint to `/ops/subscriber` + `/ops/subscriptions`, per-company) | ⬜ TODO | ### Stage 2 cutover (as executed, 2026-10-05) 1. Migration **027** adds `dori.event_dispatch`; deploy the `emit_event` writers (monitor/analytics) + CIM dedupe. Trigger still runs → each event gets TWO outbox rows (trigger + app, same `event_id`); CIM dedupes → one alert, no loss. Verified: MLC_DOWN → rows 4&5 same id → CIM "dispatched" then "duplicate". 2. Migration **028** drops the trigger → `emit_event` sole writer, one row/event. Verified: MLC_DOWN → single row 6, CIM dispatched, no duplicate. Rollback: `alembic downgrade 026_outbox_trigger` recreates the trigger (from `create_notify_trigger.outbox.sql`) + drops `event_dispatch`; redeploy prior images. ## Stage 1 EC2 verification (step 3e) — do this before any Stage 2 coding ``` # on the EC2 checkout of this branch alembic upgrade head # applies 026; trigger now writes event_outbox python -c "from dori_model.dModel import DEventOutbox; print(DEventOutbox.__table__)" # restart DIM so the relay (not the old listener) is running ``` Then exercise the pipeline and confirm the durable hop: - Trigger an MLC_DOWN (monitor active-poll, or inject via edge_simulator as in `unit_tests/test_dim.py::test_ws_fanout_end_to_end`). - `SELECT id,status,attempts,dispatched_at FROM dori.event_outbox ORDER BY id DESC LIMIT 5;` → the row appears `pending` then flips to `dispatched`. - SMS lands at DORI_ADMIN; dashboard WS still receives the event. - Kill CIM briefly, fire another event → row stays `pending`, `attempts` climbs, `next_attempt_at` advances; bring CIM back → it delivers. Nothing lost. Do NOT start Stage 2 (writer swap / trigger drop) until this passes. ## Stage 1 — Relay (DONE) — `services/dim/src/service/outbox_relay.py` A loop that drains the outbox and delivers over HTTP+WS (NOT Kafka — see Reconciliation). Shape as built: - On a `NOTIFY 'outbox'` wake OR a fallback poll (`settings.outbox_poll_seconds`): claim a batch — `SELECT id, event_id, attempts, payload FROM dori.event_outbox WHERE status='pending' AND next_attempt_at <= now() ORDER BY id FOR UPDATE SKIP LOCKED LIMIT :batch` - For each row: `registry.fan_out(deviceId, payload)` (ephemeral UI lane, best-effort) + `forward_event(payload, raise_on_error=True)` (durable alert lane), then `UPDATE ... SET status='dispatched', dispatched_at=now()`. - On forward failure: `attempts+1`, `next_attempt_at = now() + backoff` (exp, cap `outbox_backoff_cap_seconds`), `last_error=...`, `status='dead'` once `attempts+1 >= outbox_max_attempts`. - Commit per batch; SINGLE relay worker preserves per-aggregate ordering. - Runs in DIM's lifespan, replacing `run_listener_loop`. ### Everything below is Stage 2 (deferred until step 3e passes) - Ordering: single relay worker, or shard by `hash(aggregate_key)`, so MLC_UP never overtakes MLC_DOWN for the same mlc. LISTEN helper: reuse DIM's asyncpg listen pattern (`services/dim/src/service/listener.py` `run_listener_loop`) but channel `outbox`; the payload is empty (a nudge) — the relay reads rows, not the notify. ## Step 4 — Writer swap (monitor) In `services/monitor/src/utils/mlc_check.py`, replace `_insert_system_event` with `_emit_system_event` that, in ONE transaction: 1. `INSERT INTO inference.inference_event (...) RETURNING id` → event_id. 2. Look up the event-type flags for `event_type` in DORI_SYSTEM_SERVICE_BUNDLE (cache the tiny `event_type` catalog in memory; refresh on TTL). Flags: `is_notification` (UI), `is_trigger` (system alert), `is_message` (customer). 3. If `is_trigger` (needs alert): `INSERT INTO dori.event_outbox (event_id, aggregate_key='mlc:', topic=settings.topic_alerts_dispatch, payload={alert_id: event_id, channel:'system', event_type, **event_data})` then `NOTIFY outbox`. 4. If `is_notification`: `pg_notify('ui_event', payload)` (ephemeral UI path). 5. One `commit`. Also apply the same to `services/monitor/src/service/system_event_service.py` `process_system_event` (same INSERT today). And replace the Mode-1 dual write in `_watch_transitions` (the bare `producer.send` with no DB row) with the same event+outbox insert — that removes a live at-most-once hole. **Critical:** route ALL three event writers (`analytics_app/.../inference_ingest.py`, the two monitor writers) through one shared `emit_event()` helper, or a direct insert silently skips notifications. Add a test asserting nothing else inserts into `inference_event`. ## Step 5 — Consumer + dedupe The `dori.alerts.dispatch` consumer is CIM's dispatcher. Ensure it dedupes on `alert_id` (= event_id) because the outbox is at-least-once (a crash after the Kafka send but before the `dispatched` UPDATE re-sends). CIM already has `check_message_log_exists` cooloff for customer messages — extend the same idea to system alerts keyed on alert_id. ## Step 6 — Customer events Extend the shared `emit_event()` helper to `services/analytics_app/src/utils/inference_ingest.py` (the camera-event firehose). Same two-lane decision; gate the outbox on `is_message` for customer events (vs `is_trigger` for system). ## Step 7 — Drop the trigger (LAST) Only after the relay path is verified end-to-end (system + customer): - New migration `020_drop_notify_trigger`: `DROP TRIGGER event_process`, `DROP FUNCTION notify_process`. Keep it reversible (recreate from `alembic/sql/create_notify_trigger.std.sql` in downgrade). - Keep `pg_notify('ui_event')` emission in the writer for the UI lane. ## Deploy / test on EC2 - Build remotely: `deploy-scripts/build-remote-ec2.sh` (per repo memory — build on the EC2, not locally). - Verify: write a test MLC_DOWN → row lands in `dori.event_outbox` → relay marks it `dispatched` → SMS arrives at DORI_ADMIN → row count/retry sane. - The prod-DB guard in `unit_tests/test_system_event_alerts.py` already refuses to run against prod `dori_dash`; extend those tests for the outbox path. ## Open decisions (ask Ravi) - Relay as part of monitor vs its own service. - Whether to consolidate `is_trigger` + `is_message` now or wait for E7. - Whether to keep DIM in the UI path or have the relay push UI too (E7 question).