Liquidation Engine — Design Spec

1. Goals & Non-Goals

Goals (v1)

  1. Continuously detect when a user's maintenance-margin-available (MMA) becomes negative (a "margin breach").
  2. On breach, automatically cancel the user's open orders and liquidate positions until the account is no longer in breach or no positions remain.
  3. Support an AdminManual trigger so ops can liquidate an account on demand from the admin GUI (replacing the manual EP3 console workflow).
  4. Be safe to ship: dark by default behind a global kill-switch and a per-user auto_liquidation_enabled gate, so existing institutional customers are never auto-liquidated even if their accounts breach.
  5. Provide a full audit log of every decision and side effect, deterministically replayable from data alone — no real-world side effects on replay.
  6. Provide live operational visibility (Prometheus metrics, Slack threads per incident, admin GUI state page, reconciliation CLI).

Non-goals (v1)

2. Background

The relevant existing surface:

Missing piece: the orchestration. Detecting MMA-negative edges, deciding what to cancel and liquidate, executing it, retrying on failure, parking accounts that cannot be unwound, and recording everything for audit.

3. Architecture

A new crate rs/liquidation-engine with its own binary. The crate internally splits functional core / imperative shell, mirroring rs/rome (PR #1843):

risk-monitor is the canonical breach detector. It is extended (small change) to publish a trigger when it observes a user's MMA transition from ≥ 0 to < 0. The liquidation engine never recomputes MM itself.

3.1 Transports and trust model

Two ingress paths feed trigger_listener:

Ingress Protocol Trust enforcement
Risk-monitor → engine Redis pubsub LIQUIDATION_TRIGGER_CHANNEL Redis is treated as a trusted internal wakeup bus, not as breach authority. The payload carries published_by ∈ {risk-monitor}; engine logs/metrics it, then validates margin_status in Postgres before emitting TriggerReceived. Operators with Redis credentials can forge wakeups, but cannot forge a validated margin breach without Postgres access — that is the accepted threat model.
Admin GUI / ops → engine HTTPS POST /admin/liquidate on the engine's admin port Behind the same admin-auth middleware as the rest of the admin REST surface (admin-token + allowlist). The engine binds the admin port on the internal interface only.

The engine never accepts triggers from any other source. Any inbound event from an unexpected source is dropped and counted.

3.2 Sequence diagram

sequenceDiagram
    autonumber
    participant RM as risk-monitor
    participant Redis
    participant LE as liquidation-engine<br/>(driver + core)
    participant OG as order-gateway
    participant PG as postgres
    participant CH as clickhouse
    participant SL as slack
    participant Admin as admin GUI

    Note over RM: detects MMA crossing < 0
    RM->>PG: upsert margin_status (breached)
    RM->>Redis: PUBLISH LIQUIDATION_TRIGGER_CHANNEL<br/>{user, reason: MarginBreach, scope: WholeUser, published_by: risk-monitor}
    Redis-->>LE: wakeup event
    LE->>PG: SELECT margin_status WHERE user=… AND status='breached'
    alt margin_status still breached
        LE->>LE: handle_event(state, TriggerReceived{breach_epoch}, now)<br/>→ effects: [PersistAccount, AppendAuditLog, CancelAllOrders, …]
    else margin_status already cleared
        LE->>CH: log StaleTriggerDropped
    end
    LE->>PG: upsert liquidation_account_state (Liquidating)
    LE->>CH: insert liquidation_log (TriggerReceived)
    LE->>SL: thread "BreachDetected user=…"
    LE->>OG: POST cancel_all_orders(user)
    OG-->>LE: ack
    LE->>LE: handle_event(state, OrdersCancelled, now)
    LE->>OG: POST admin_place_liquidation_order<br/>(user, instrument, qty=min(pos, increment), IOC, limit=mark±bps)
    OG-->>LE: order_id ack
    LE->>CH: insert liquidation_log (LiquidationOrderPlaced)

    alt Fill
        OG-->>LE: order_filled (dropcopy stream)
        LE->>LE: handle_event(state, OrderFilled, now) → retain TerminalObserved
        LE->>CH: log OrderFilled
    else Reject / no liquidity
        OG-->>LE: order_rejected
        LE->>LE: handle_event(state, OrderRejected, now)<br/>→ attempts++, if ≥ N → PreParkProbe
        alt fresh MMA ≥ 0 before guard expires
            LE->>LE: handle_event(state, RiskSnapshotUpdated, now) → Resolved
        else still breached or guard timeout
            LE->>SL: alert AccountParked
        end
    end

    Redis-->>LE: RISK_SNAPSHOT_UPDATED (user, MMA ≥ 0)
    alt cancel not pending and no blocking attempt
        LE->>LE: handle_event(state, RiskSnapshotUpdated, now) → Resolved
        LE->>PG: liquidation_account_state = Resolved
        LE->>CH: log AccountResolved
        LE->>SL: thread "Resolved user=…"
    else TerminalObserved or OutcomeUnknown remains
        LE->>LE: retain Liquidating until authoritative reconciliation
    end

    Note over Admin,LE: Manual trigger path
    Admin->>LE: POST /admin/liquidate {user, scope} (admin-auth)
    LE->>LE: same Event::TriggerReceived path

4. Wire and storage shapes

This section declares shapes only. Where each field's invariants are enforced is described in §5 (runtime behavior) and §6 (safety controls). Every field declared here has exactly one declared enforcement site; no field is accepted that is not enforced.

4.1 Core types (state_machine.rs, pure)

Domain newtypes — single declaration, used everywhere:

pub struct UserId(pub String);                       // existing repo type, re-used
pub struct InstrumentId(pub String);                 // existing repo type, re-used
pub struct IncidentId(pub Ulid);
pub struct AuditEventId(pub Ulid);
pub struct ClientOrderId(pub Ulid);                  // generated by driver, threaded into core
pub struct ExchangeOrderId(pub String);              // OG-stamped, "L-<ulid>"
pub struct Ns(pub i64);                              // wall-clock nanoseconds (driver-supplied)
pub struct Qty(NonZeroU64);                          // magnitude only; closing direction lives in Side
pub struct Mma(pub Decimal);                         // signed; < 0 means breach
pub struct LimitBps(NonZeroU16);                     // protective limit width around mark

The pure core uses these everywhere. String / i64 / Decimal never appear directly on event or effect signatures except where wrapped above.

Account status is a sum type so each variant carries exactly the fields it needs. Statuses like "Liquidating with no incident" are unrepresentable:

pub enum AccountState {
    Idle {
        user_id: UserId,
    },
    Liquidating {
        user_id: UserId,
        incident_id: IncidentId,
        trigger_reason: TriggerReason,
        breach_epoch: Option<Ns>,                  // set for validated MarginBreach triggers
        scope: Scope,
        opened_at: Ns,
        last_action_at: Ns,
        attempts: BTreeMap<InstrumentId, u32>,        // per-instrument rejects; retained with terminal attempts
        inflight: BTreeMap<InstrumentId, InflightOrder>,
        pre_park_probe_until: Option<Ns>,             // set after attempts[instr] reaches N
        last_mma: Option<Mma>,
        last_snapshot_at: Option<Ns>,
        auto_liquidation_enabled: bool,               // snapshotted at trigger time (see §6)
    },
    Parked {
        user_id: UserId,
        incident_id: IncidentId,
        parked_at: Ns,
        reason: ParkReason,
    },
    Resolved {
        user_id: UserId,
        incident_id: IncidentId,
        resolved_at: Ns,
    },
}

pub struct InflightOrder {
    pub client_order_id: ClientOrderId,
    pub observation: InflightOrderObservation,
    pub instrument: InstrumentId,
    pub qty: Qty,
    pub side: Side,
    pub limit_price: Decimal,
    pub placed_at: Ns,
}

pub enum InflightOrderObservation {
    Submitted,
    Acknowledged { exchange_order_id: ExchangeOrderId },
    TerminalObserved {
        exchange_order_id: Option<ExchangeOrderId>,
        cumulative_fill: Qty,
    },
    OutcomeUnknown {
        exchange_order_id: Option<ExchangeOrderId>,
        terminal_observation: Option<TerminalOrderObservation>,
        alerted_at: Ns,
        reason: OrderOutcomeUnknownReason,
    },
}

pub enum TriggerReason { MarginBreach, AdminManual /* future: MonthlyLossBreach, ComplianceFreeze, InstrumentDelisting */ }
pub enum Scope { WholeUser, Instrument(InstrumentId) }
pub enum ParkReason { AttemptsExhausted, IneligibleInstrument, StaleMark, AutoLiquidationDisabled }
impl ParkReason {
    pub fn is_recoverable(&self) -> bool {
        matches!(self, Self::AttemptsExhausted | Self::IneligibleInstrument)
    }
}
pub enum RejectReason { InsufficientLiquidity, InstrumentHalted, TransientExhausted, LimitNotMarketable, Other(String) }
pub enum TriggerSource { RiskMonitor, AdminApi, ReconcileBootstrap }
pub enum TriggerValidationSource { MarginStatus, AdminAuth }
pub enum Freshness { Fresh, Stale, SkewViolation }
pub enum ValidationStatus { Validated { breach_epoch: Ns }, StaleMarginStatus }

Events and effects:

pub enum Event {
    TriggerReceived {
        user: UserId,
        reason: TriggerReason,
        scope: Scope,
        incident_id: IncidentId,
        auto_liquidation_enabled: bool,           // driver snapshots from PG users row at ingress (§6.2)
        published_by: TriggerSource,              // RiskMonitor | AdminApi
        breach_epoch: Option<Ns>,                 // margin_status.updated_at for MarginBreach, None for AdminManual
        validated_by: TriggerValidationSource,    // MarginStatus for risk-monitor/reconcile, AdminAuth for manual
        now: Ns,
    },
    RiskSnapshotUpdated {
        user: UserId,
        mma: Mma,
        snapshot_at: Ns,                          // publisher's clock
        freshness: Freshness,                     // shell labels Fresh | Stale | SkewViolation (§6.3); core decides
        now: Ns,
    },
    TriggerCandidateValidated {
        user: UserId,
        status: ValidationStatus,                 // Validated{breach_epoch} | StaleMarginStatus
        reason: TriggerReason,
        scope: Scope,
        now: Ns,
    },
    OrdersCancelled { user: UserId, now: Ns },
    OrderAcked { user: UserId, client_order_id: ClientOrderId, exchange_order_id: ExchangeOrderId, now: Ns },
    OrderFilled { user: UserId, client_order_id: ClientOrderId, qty_filled: Qty, now: Ns },
    OrderRejected { user: UserId, client_order_id: ClientOrderId, reason: RejectReason, now: Ns },
    Tick { now: Ns },
    ConfigUpdated { config: EngineConfig, now: Ns },
    KillSwitchChanged { mode: KillSwitch, now: Ns },
}

pub enum Effect {
    CancelAllOrders { user: UserId },
    PlaceLiquidationOrder {
        user: UserId,
        instrument: InstrumentId,
        qty: Qty,
        side: Side,
        limit_price: Decimal,
        client_order_id: ClientOrderId,           // pre-generated by driver
    },
    PersistAccount { snapshot: AccountState },
    AppendAuditLog { entry: LiquidationLogEntry },// includes pre-generated event_id (§5.2)
    EmitSlack { kind: SlackKind, user: UserId, incident_id: IncidentId, payload: SlackPayload },
}

pub fn handle_event(state: EngineState, event: Event) -> (EngineState, Vec<Effect>);

Metrics are not an effect. They are incremented synchronously by the task best placed to observe the underlying signal (driver counts events handled, og_client counts OG outcomes, persistence counts insert failures, etc.). Keeping metrics out of the effect stream avoids mixing "decisions" with "telemetry" inside the core.

4.2 Postgres schema (additive, Atlas-managed db/postgres/1.sql)

CREATE TABLE liquidation_account_state (
    user_id        TEXT PRIMARY KEY,
    status         TEXT NOT NULL CHECK (status IN ('Liquidating','Parked','Resolved')),
    incident_id    TEXT NOT NULL,
    trigger_reason TEXT NOT NULL,
    scope          JSONB NOT NULL,
    opened_at      TIMESTAMPTZ NOT NULL,
    last_action_at TIMESTAMPTZ NOT NULL,
    attempts       JSONB NOT NULL DEFAULT '{}'::jsonb,   -- {instrument_id: count}
    pre_park_probe_until TIMESTAMPTZ,
    last_mma       NUMERIC,
    park_reason    TEXT,
    updated_at     TIMESTAMPTZ NOT NULL
);

ALTER TABLE users
    ADD COLUMN auto_liquidation_enabled BOOLEAN NOT NULL DEFAULT FALSE;

ALTER TABLE instruments
    ADD COLUMN liquidation_increment_qty BIGINT,                       -- NULL = instrument not eligible
    ADD COLUMN liquidation_limit_bps     INT NOT NULL DEFAULT 50;      -- protective limit width

Idle is the absence of a row, not a stored value. The engine never writes an Idle row. On each Tick, the core removes the lowest-ID resolved account and its current exposure cache entry, then emits ClearAccount; normal liquidation scheduling still runs on that tick. All schema additions are nullable or defaulted; no existing service breaks on migration.

4.3 ClickHouse schema (db/clickhouse/init.sql)

CREATE TABLE liquidation_log (
    event_id        String,                              -- ULID, generated by driver; CH dedup key
    ts              DateTime64(9),
    user_id         String,
    incident_id     String,
    event_type      LowCardinality(String),
    trigger_reason  LowCardinality(String),
    instrument_id   Nullable(String),
    client_order_id Nullable(String),
    exchange_order_id Nullable(String),
    payload         String,                              -- JSON
    engine_version  LowCardinality(String)
) ENGINE = ReplacingMergeTree
ORDER BY (user_id, ts, event_id);

ReplacingMergeTree dedups on the full ORDER BY tuple (user_id, ts, event_id), not on event_id alone. The driver freezes both event_id (ULID) and ts (= event.now — the same Ns value threaded into the originating Event) before the first insert attempt, and reuses both on every retry; the writer never stamps ts itself. Without freezing ts a retry using now() would land on a different ORDER BY tuple and silently double-write. Combined with bounded-retry retention in the writer (§5.2), this avoids both silent loss and double-writes.

4.4 Redis keys and channels

Key / channel Purpose Owner
LIQUIDATION_TRIGGER_CHANNEL (pubsub) Risk-monitor → engine wakeup for margin-status validation risk-monitor (writer), engine (reader)
RISK_SNAPSHOT_UPDATED_CHANNEL (pubsub, existing) Snapshot updates; engine uses for resolution + freshness risk-engine (writer), many readers
liquidation_engine:kill_switch (string) Source of truth for engine mode (Live / DryRun / Off) ops (writer), engine (reader)

Redis pubsub is best-effort and is not the trigger of record. margin_status is the breach authority: every Redis wakeup is re-validated against Postgres in the shell, then turned into Event::TriggerCandidateValidated { user, status: Validated | StaleMarginStatus, breach_epoch }. The validated branch is what becomes Event::TriggerReceived inside the core; the stale branch causes the core to emit AppendAuditLog { StaleTriggerDropped }. Pre-event audit emission does not exist — every audit row originates from handle_event, which keeps the §8.5 effect-sequence-equality contract intact. liquidation_stale_trigger_dropped_total is incremented at the same site that emits the audit.

4.5 Engine config

pub struct EngineConfig {
    pub max_attempts_before_park: u32,            // default 5
    pub inflight_timeout: Duration,               // default 5s; callback outcome becomes unknown
    pub mark_freshness_max: Duration,             // default 3s; stale-mark refusal
    pub safety_tick: Duration,                    // default 1s
    pub clock_skew_tolerance: Duration,           // default 500ms; bounds engine-vs-publisher skew (§6.3)
}

v1 ships exactly the fields above. Rate limits and notional caps are v1.1 follow-ups (§11); they are deliberately absent from this struct to avoid wire fields without enforcement.

5. Runtime behavior (imperative shell)

5.0 Lifecycle at a glance

The table below is the contract the pure core honors; §5.1–§5.7 are the runtime mechanics that deliver it. Every row is an Event the core sees; every column is an effect channel the driver dispatches to in the order §5.2 defines. All audit / slack emission is produced inside handle_event in response to an event; the shell never writes an audit row directly. Pre-event filtering would break §8.5's effect-sequence equality, so freshness, clock-skew, and stale-margin-status are expressed as event payload fields (Freshness, ValidationStatus) and let the core decide what to emit. next denotes select_next_action(state) → Option<Effect> (§5.3): emitted only if an eligible instrument exists, the kill-switch is Live, the mark is fresh (§6.3), and auto_liquidation_enabled was true at trigger ingress (§6.2).

Event PersistAccount (PG) RemovePersist CancelAllOrders (OG) PlaceLiquidationOrder (OG) AppendAuditLog (CH) EmitSlack
TriggerCandidateValidated (StaleMarginStatus) — — — — StaleTriggerDropped —
TriggerReceived (validated MarginBreach) Liquidating — ✓ next TriggerReceived BreachDetected
TriggerReceived (AdminManual) Liquidating — ✓ next TriggerReceived BreachDetected
TriggerReceived (duplicate, §6.5) — — — — DuplicateTriggerDropped —
TriggerReceived (auto_liquidation_enabled=false) Parked (AutoLiquidationDisabled) — — — AccountParked AccountParked
OrdersCancelled Liquidating — — next OrdersCancelled —
OrderAcked Liquidating (inflight[instr]=Acknowledged) — — — OrderAcked —
OrderFilled (valid terminal cumulative quantity) Liquidating (inflight[instr]=TerminalObserved) — — — OrderFilled —
Tick (Submitted/Acknowledged callback timeout) Liquidating (inflight[instr]=OutcomeUnknown) — — — OrderOutcomeUnknown OrderOutcomeUnknown
conflicting Ack, overfill, or contradictory terminal callback Liquidating (inflight[instr]=OutcomeUnknown) — — — OrderOutcomeUnknown OrderOutcomeUnknown
OrderRejected (attempts < N) Liquidating (attempts[instr]++) — — next OrderRejected —
OrderRejected (attempts ≥ N) Liquidating (pre_park_probe_until=now+safety_tick) — — — PreParkProbeStarted —
Tick (pre-park probe, MMA still < 0 or timeout) Parked (AttemptsExhausted) — — — AccountParked AccountParked
RiskSnapshotUpdated (MMA ≥ 0, Liquidating with cancel not pending and no attempt, or Parked(AttemptsExhausted/IneligibleInstrument)) Resolved ✓ (next tick) — — AccountResolved Resolved
RiskSnapshotUpdated (freshness=Stale, §6.3) — — — — StaleMark StaleMark
RiskSnapshotUpdated (freshness=SkewViolation, §6.3) — — — — ClockSkewViolation StaleMark
Tick — — — next — —
KillSwitchChanged (Off) — — — — KillSwitchChanged EngineHalted
KillSwitchChanged (DryRun | Live) — — — — KillSwitchChanged —
ConfigUpdated — — — — ConfigUpdated —

Account lifecycle

stateDiagram-v2
    [*] --> Idle
    Idle --> Liquidating : Validated TriggerReceived / CancelAllOrders, EmitSlack BreachDetected
    Liquidating --> Liquidating : OrderFilled / retain TerminalObserved
    Liquidating --> Liquidating : callback timeout or invariant violation / retain OutcomeUnknown, alert once
    Liquidating --> Liquidating : OrderRejected attempts lt N / attempts++
    Liquidating --> Liquidating : OrderRejected attempts ge N / start PreParkProbe
    Liquidating --> Parked : PreParkProbe timeout while MMA lt 0
    Liquidating --> Parked : IneligibleInstrument or StaleMark or AutoLiquidationDisabled
    Liquidating --> Resolved : RiskSnapshotUpdated MMA ge 0 and cancel not pending and no attempt / EmitSlack Resolved
    Parked --> Resolved : RiskSnapshotUpdated MMA ge 0 and reason=AttemptsExhausted or IneligibleInstrument
    Parked --> Liquidating : reason=IneligibleInstrument and fresh breach becomes constructible
    Parked --> Liquidating : admin unpark
    Resolved --> Idle : next Tick / ClearAccount
    Resolved --> [*]

Per-instrument inflight order lifecycle

stateDiagram-v2
    [*] --> Submitted : PlaceLiquidationOrder
    Submitted --> Acknowledged : OrderAcked
    Submitted --> TerminalObserved : terminal cumulative OrderFilled
    Acknowledged --> TerminalObserved : terminal cumulative OrderFilled
    Submitted --> OutcomeUnknown : callback timeout, overfill, or invariant violation
    Acknowledged --> OutcomeUnknown : callback timeout, conflicting Ack, or invariant violation
    TerminalObserved --> OutcomeUnknown : contradictory terminal callback or conflicting Ack
    Submitted --> [*] : known-zero-fill OrderRejected
    Acknowledged --> [*] : known-zero-fill OrderRejected
    OutcomeUnknown --> TerminalObserved : matching late terminal fill after callback timeout
    OutcomeUnknown --> [*] : matching late known-zero-fill reject after callback timeout

Account-level transitions describe the row in liquidation_account_state (PG) and the audit stream in liquidation_log (CH). The per-instrument lifecycle is durable bookkeeping inside AccountState::Liquidating.inflight. TerminalObserved and OutcomeUnknown remain inflight and block placement, parking, and resolution until authoritative reconciliation supplies causal risk evidence; elapsed time alone never drops an attempt. The two diagrams compose: an AttemptsExhausted account parks only after a one-safety-tick pre-park probe fails to observe fresh healthy MMA. Parking is account-level and always the lowest-priority action: an account parks only after every eligible sibling instrument has drained and no attempt blocks progress. IneligibleInstrument resolves on a fresh healthy snapshot, or self-heals back to Liquidating when a fresh breaching snapshot again yields a constructible protective order. StaleMark and AutoLiquidationDisabled are ops-only until admin unpark.

5.1 Driver loop

loop {
    tokio::select! {
        biased;
        _ = shutdown.cancelled() => break,
        Some(event) = events_rx.recv() => {
            let (next, effects) = handle_event(state.clone(), event);
            state = next;
            for effect in effects {
                dispatch(effect).await; // see §5.2
            }
        }
    }
}

Only the driver task mutates state. All other tasks own no engine state — they only push events.

5.2 Effect dispatch

Effects emitted by the core within a single tick are executed in the order they appear in the returned Vec<Effect>. The core emits them in a deliberate order so that durable persistence always precedes any externally observable action:

  1. PersistAccount — Postgres upsert. On failure the driver backs off (100ms / 500ms / 2s) and re-enqueues the originating event; no later effect runs. The liquidation_pg_persist_failures_total counter increments on every failure.
  2. AppendAuditLog — ClickHouse insert. The driver pre-generates event_id (ULID) and threads event.now (the Ns driving the originating Event) into the entry before the first insert attempt; both feed the row's event_id and ts columns, so retries collide on the full ORDER BY tuple and ReplacingMergeTree (§4.3) drops the duplicate on merge. The writer never stamps ts itself. The writer batches inserts and retains the batch on insert failure with bounded retry; after retry exhaustion the rows go to a local on-disk spool file, liquidation_log_rows_dropped_total increments, and PagerDuty fires. The writer does not set wait_for_async_insert=0.
  3. CancelAllOrders — POSTed to OG. Result is fed back as Event::OrdersCancelled (or OrderRejected for the cancel-all call). Driver does not synchronously block subsequent effects; later effects in this tick may run concurrently against different users.
  4. PlaceLiquidationOrder — kill-switch and dry-run gates checked in the shell, not in the core (§6.1). Off / DryRun produces a slack notice and an audit row but no OG call; Live POSTs to OG. Either way, the eventual outcome arrives as an event.
  5. EmitSlack — best-effort; failure is logged and counted (liquidation_slack_failures_total), never blocks state progress.

Between writing to PG and writing to CH the engine can crash. PG is authoritative for live state; CH is the append-only audit history. A momentary CH gap is acceptable but is monitored via liquidation_audit_gap_seconds (max age of PG updated_at with no corresponding CH row for the same (user, ts)); a sustained gap pages.

5.3 Concurrency and ordering

5.4 Retry and park policy

5.5 Restart recovery and reconciliation

bootstrap.rs runs three steps on startup, in order, and emits structured Events so the rest of the engine sees recovery exactly as it would steady-state operation.

  1. Load live state. Read all liquidation_account_state rows from Postgres.
  2. Reconcile bidirectionally with order-gateway. Query OG for every open order with the L- prefix.
  3. Recover missed triggers via margin_status. Run the same validation path used by the Redis wakeup listener: read the margin_status table (introduced in PR #3 as the durable trigger-of-record) for users currently in a breached state with no Liquidating row in liquidation_account_state, then emit synthetic Event::TriggerReceived { reason: MarginBreach, published_by: ReconcileBootstrap, breach_epoch: margin_status.updated_at, validated_by: MarginStatus } for each. This is the durable replacement for Redis pubsub reliability: pubsub wakes the engine, while margin_status is the trigger authority.

A standalone liquidation-reconcile subcommand (same crate) runs the same diff against PG + OG + CH on demand without starting the engine, for use during incidents or manual ops audits.

ClickHouse liquidation_log is never consulted as live-state truth. It is the audit/replay source only.

5.6 Disconnect matrix

Per CLAUDE.md ("any feature tied to order lifecycle must be tested across disconnect / restart / recovery"):

5.7 Shutdown

The binary owns a top-level tokio_util::CancellationToken and a JoinSet of every task spawned in §3. SIGTERM cancels the token. Tasks observe cancellation in their tokio::select!, drain pending work, and exit. The driver in particular:

  1. Stops accepting new events from the mpsc.
  2. Waits up to inflight_timeout for outstanding OG calls in its per-order JoinSet to settle (their results become events the driver still processes).
  3. Issues one final PersistAccount per modified account so PG is consistent before exit.
  4. Returns; the binary's main joins the JoinSet and exits.

Any task panic propagates through the JoinSet and aborts the process — there are no detached tasks. This is the same supervisor pattern referenced from trade-engine2's state monitor.

6. Safety controls and invariant enforcement

This section is the single chokepoint inventory. Every field declared in §4 has exactly one enforcement site listed here.

6.1 Kill-switch (global)

Redis key liquidation_engine:kill_switch is the source of truth; values are Live | DryRun | Off. The driver subscribes to keyspace notifications and emits Event::KillSwitchChanged so the core mirrors the current mode for replay determinism. Just before dispatching each PlaceLiquidationOrder, the shell re-reads the Redis key (belt-and-braces); if Redis is down it uses the last known value and flags liquidation_killswitch_stale_seconds. DryRun runs the full pipeline but stops at OG. Off suppresses the effect entirely and surfaces a slack notice. Default at first deploy is Off.

6.2 Per-user gate

users.auto_liquidation_enabled defaults FALSE. The driver snapshots this flag into the event at trigger time after any required margin_status validation (trigger_listener does the PG reads before pushing Event::TriggerReceived). The core then refuses to emit PlaceLiquidationOrder for an account whose snapshotted flag is false, regardless of TriggerReason. Snapshotting at ingress keeps the core pure (no PG reads inside handle_event) and replay deterministic.

6.3 Mark freshness and clock skew

Labelled in the shell, decided in the core. The shell computes Freshness and stamps it onto the event; the core owns every audit / slack row that follows, so §8.5 effect-sequence equality holds on replay. Per snapshot:

age  = max(0, now_ns(engine) - snapshot.ts_ns)
skew = snapshot.ts_ns - now_ns(engine)
freshness = if skew > clock_skew_tolerance { SkewViolation }
            else if age > mark_freshness_max { Stale }
            else                              { Fresh }

Every stale or skewed snapshot remains audited. EmitSlack { StaleMark } is edge-triggered once per account across a contiguous stale/skewed episode; an accepted newer fresh snapshot resets the alert latch. Replayed or out-of-order fresh snapshots do not reset it. The latch is in-memory, so restart may emit one new alert for an episode already in progress.

liquidation_clock_skew_ms histogram records the signed skew on every snapshot so drifts are visible before they become refusals.

6.4 Instrument eligibility

At trigger ingress the shell joins instruments and rejects triggers whose scope = Instrument(id) references a non-existent or liquidation_increment_qty IS NULL instrument with Alert { UnknownOrIneligibleInstrument }. For WholeUser scope, eligibility is checked per-instrument at select_next_action time, but the park model is account-level (KISS): eligible in-scope instruments keep being placed, and the account parks with reason: IneligibleInstrument only once no in-scope instrument yields a constructible protective order and no order is still inflight. Per-instrument attempts/exhaustion bookkeeping is retained inside Liquidating, but there is no per-instrument park state. Because the reason is recoverable, the account self-heals back to Liquidating if a later fresh breaching snapshot makes any in-scope instrument constructible again.

6.5 Trigger deduplication

For validated MarginBreach triggers, dedup is keyed by (user, reason, breach_epoch) while account status is Liquidating or Parked. Redis wakeups whose margin_status row is no longer Breached are dropped before they become TriggerReceived, audited as StaleTriggerDropped, and counted in liquidation_stale_trigger_dropped_total. Admin manual triggers do not have a breach epoch; they keep incident-based dedup against the current Liquidating or Parked account. Duplicate validated triggers do not open a second incident and are counted in liquidation_duplicate_triggers_total.

6.6 Order shape

Qty is NonZeroU64 magnitude by type; Side carries direction. LimitBps is NonZeroU16. The core computes the protective limit as mark ± liquidation_limit_bps, then rounds buys up and sells down to the instrument tick so the persisted inflight state, audit record, and order effect carry the same executable price. The core also asserts the emitted Side matches the closing direction of the user's current position and otherwise refuses to emit (a positive-test invariant for proptest).

6.7 Admin REST authentication

POST /admin/liquidate and POST /admin/liquidation/unpark are mounted behind the existing admin-auth middleware (admin-token + allowlist). The engine binds its admin port on the internal interface only. Unauthorized requests are rejected at the middleware layer and never reach the trigger listener.

Summary table

Field / claim Enforcement site
users.auto_liquidation_enabled = false ⇒ no orders §6.2 — snapshotted at ingress, checked in core
Mark freshness ≤ mark_freshness_max §6.3 — shell, before emitting RiskSnapshotUpdated
Engine and publisher clocks within clock_skew_tolerance §6.3 — shell
Instrument has liquidation_increment_qty IS NOT NULL §6.4 — shell at trigger ingress; core at action selection
Qty != 0 (magnitude), LimitBps != 0, tick-aligned limit, Side matches closing direction §6.6 — types + core invariant
Duplicate TriggerReceived is idempotent §6.5 — core
Stale Redis wakeup cannot open an incident §6.5 — shell validates margin_status before event emission
AttemptsExhausted park is recoverable by healthy MMA §5.4 — pre-park probe + recoverable park reason
IneligibleInstrument park self-heals when eligibility returns §5.4 / §6.4 — resume on fresh breach with constructible order
Inflight orders are never orphaned by parking §5.3 — any inflight order pins the account in Liquidating
PlaceLiquidationOrder gated by kill-switch §6.1 — shell, double-checked at dispatch
Admin REST callers are authenticated §6.7 — middleware
Trigger source is identified §3.1 — published_by on event

7. Observability

Three sinks, one correlation key.

7.1 Correlation

Every log line, metric exemplar, ClickHouse row, and Slack message carries incident_id. The driver opens a tracing::Span named incident=<ulid> for every event processed inside a known incident, and shell tasks attach the same span to their child work. The crate uses tracing (not log) so spans propagate automatically. The incident_id is also a label on histogram exemplars where the metrics backend supports it.

7.2 ClickHouse audit log

Append-only event stream. Enables full incident replay through the pure core (see §8.5). Event types: TriggerReceived, OrdersCancelled, LiquidationOrderPlaced, OrderAcked, OrderFilled, OrderRejected, AccountResolved, AccountParked, Alert. Replay never produces external side effects (§8.5).

7.3 Postgres live state

Authoritative current state for Liquidating | Parked | Resolved. Surfaces in the admin GUI. Updated by the imperative shell after every driver tick. No engine state lives anywhere else.

7.4 Prometheus metrics

Health (mirrors the order-gateway PR #1733 shape):

Failure-mode (the counters Rome learned to need):

7.5 Slack

Three message classes, all threaded per (user_id, incident_id):

8. Testing strategy

8.1 Pure core — property + snapshot

8.2 Driver — actor in isolation

8.3 Integration — full crate + real OG + mock EP3

Reuses existing test-utils seed data and the mock EP3 environment.

One canonical end-to-end happy-path test in the CI gate; slower scenarios behind --ignored.

8.4 Disconnect / restart matrix

Per CLAUDE.md, covered in the driver test layer with mocked transports: client disconnect mid-fill, OG restart, engine restart, Redis flap (recovers triggers via margin_status).

8.5 Replay determinism

Separate subcommand liquidation-engine replay --from <ts> --to <ts> reads liquidation_log rows, threads them through handle_event with an effect-recording sink (not the live shell), and asserts:

Effect-sequence equality holds because every audit row originates inside handle_event in response to an event — there is no pre-event audit emission anywhere in the shell (§4.4 stale-Redis, §6.3 stale/skew snapshots all turn into events the core sees, never into shell-side audit writes). This binary path cannot be configured to dispatch effects through the live shell. The replay sink is a different type from the live shell dispatcher; mis-wiring is a compile error. This protects goal #5 from accidental re-execution.

9. Rollout plan — 9 PRs

Each PR is ≤ ~1k LOC. The app is fully functional after every merge. No PR depends on a future PR's behavior; the engine is dark until step 9.

PR #5 of the original plan (integration tests + metrics + reconcile CLI) is split into PR #5 and PR #6 here. Metrics + integration tests stay together — metrics must land with the code that emits them so the engine is never deployed without observability, and integration tests are the contract those metrics live or die by. The reconcile + replay subcommands are ops-facing diagnostics with a separate audience and a separate determinism contract (§8.5 — replay sink is a compile-time-distinct type from the live shell dispatcher); they ship as their own PR.

# PR Scope LOC est.
1 chore(db): liquidation engine schemas Postgres liquidation_account_state + nullable columns on users / instruments. ClickHouse liquidation_log with ReplacingMergeTree. New sdk-internal types: domain newtypes, Event, Effect, TriggerReason, Scope, EngineConfig, LiquidationLogEntry. No runtime use. ~450
2 feat(liquidation-engine): pure core New crate rs/liquidation-engine (lib only). state_machine.rs with handle_event. Inline insta snapshots + proptest invariants. ~850
3 feat(risk-monitor): publish MarginBreach wakeups Risk-monitor adds MMA-edge detection on top of the existing snapshot consumer. Introduces the margin_status table (PG, additive) and starts upserting it on every detected breach, then publishes LIQUIDATION_TRIGGER_CHANNEL on MMA: ≥ 0 → < 0 — idempotent per-user via the margin_status row that this PR has just written. The row is the durable trigger of record; pubsub is only a wakeup. No subscriber yet → no-op end-to-end. ~500
4 feat(liquidation-engine): imperative shell + binary tasks/*, bootstrap.rs with bidirectional reconcile and shared margin_status validation for Redis wakeups + bootstrap, main.rs with JoinSet + CancellationToken. Includes pre-park probe handling and AttemptsExhausted self-resolve. Admin REST behind admin-auth middleware. Driver tests with mocked transports. Deploy manifest with kill_switch=Off env-baked. ~950
5 feat(liquidation-engine): metrics + integration tests Prometheus counters / gauges / histograms wired into the task that observes each signal (driver counts events, og_client counts OG outcomes, persistence counts insert failures, etc.), including stale-trigger drops, pre-park probe outcomes, and parked self-resolve. /metrics endpoint. Slack message templates (BreachDetected, AccountParked, EngineHalted) threaded per (user_id, incident_id). End-to-end integration tests against mock EP3 (happy path, OG-reject → pre-park probe → parked, admin manual trigger incl. unauthenticated rejection, restart mid-incident, stale risk-snapshot, stale Redis wakeup drop, bootstrap orphan). Disconnect-matrix tests per CLAUDE.md (engine crash, OG restart, Redis flap, risk-monitor crash) plus the two TLA counterexample regressions. Engine still dark. ~500
6 feat(liquidation-engine): reconcile + replay subcommands liquidation-engine reconcile — diffs PG ↔︎ OG ↔︎ CH on demand, exits non-zero on drift, used during incidents and ops audits. liquidation-engine replay --from <ts> --to <ts> — reads liquidation_log rows, threads them through handle_event with an effect-recording sink that is a different type from the live shell dispatcher (compile error on mis-wiring; protects goal #5 from accidental re-execution). Final EngineState and effect sequence asserted against snapshot. Engine still dark. ~300
7 feat(admin-gui): liquidation status + manual trigger Wire the existing "Liquidate Account" UI placeholder (UserOverview.tsx:184) to POST /admin/liquidate. New admin page: liquidation account state list, unpark button. Per-user auto_liquidation_enabled toggle. ~600
8 chore(liquidation-engine): flip to DryRun, soak Kill-switch = DryRun in staging + prod. Runbook (docs/runbooks/liquidation-parked.md) and doc updates. Soak ≥ 1 week. ~150
9 chore(liquidation-engine): enable for retail cohort Flip kill-switch to Live. Backfill auto_liquidation_enabled=true for first retail accounts. Populate instruments.liquidation_increment_qty for in-scope products. Final runbook updates. ~100

Invariants across every merge

Parallelizable

PR #2 and PR #3 are independent of each other once PR #1's types land. PR #6 and PR #7 can land in either order after PR #5; PR #6 only requires the shell + binary from PR #4 (its CLI subcommands can land before metrics-and-tests if needed, but landing PR #5 first means PR #6's reconcile CLI gets observed via the same metric surface).

10. Alternatives considered

10.1 Engine location

10.2 Maintenance-margin source

10.3 Ordering across breached accounts

10.4 Order type

10.5 Architecture pattern

10.6 Re-evaluation cadence

10.7 Trigger durability

10.8 Attempts counter scope

10.9 Metrics as effects

11. Open questions and follow-ups

Operational artefacts to land alongside the rollout: