docs/plans/gdocs-liquidation-engine.mdAdminManual trigger so ops can liquidate an
account on demand from the admin GUI (replacing the manual EP3 console
workflow).auto_liquidation_enabled gate, so existing
institutional customers are never auto-liquidated even if their accounts
breach.is_close_only freeze;
auto-liquidation on MLL is a deliberate follow-up that reuses this
engine.The relevant existing surface:
rs/risk-engine is the running risk
service. It computes IM, MMR, IMA, MMA per user (via its
calculator/buying_power_updater.rs), publishes snapshots to
Redis on RISK_SNAPSHOT_UPDATED_CHANNEL, and persists them
to ClickHouse (risk_snapshots). (Not to be confused with
the separate rs/risk-engine2 library, which is embedded
into api-gateway and order-gateway for
open-orders margin checks and does no publishing.)rs/risk-monitor subscribes to that
channel for monthly-loss-limit breach detection only.
On MLL breach today it flips users.is_close_only=true and
posts a Slack alert (PR #1248 and follow-ups). A current-status
margin_status row in Postgres for liquidation-engine use
does not exist today — risk-monitor begins writing it
in PR #3, where it also doubles as MLL margin-status persistence.rs/order-gateway exposes two
admin-only primitives in use today for manual liquidations:
cancel_all_orders(user) and
admin_place_liquidation_order(user, instrument, qty, side, …).
Liquidation orders are forced IOC and stamped with an OG-generated
L-<ulid> order ID.instruments table carries
maintenance_margin_pct, used by risk-engine to
compute MMR; no new margin-model work is needed for v1.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.
A new crate rs/liquidation-engine with its own binary.
The crate internally splits functional core / imperative
shell, mirroring rs/rome (PR #1843):
state_machine.rs — pure
handle_event(state, event) → (state, Vec<Effect>). No
I/O, no clocks, no IDs, no RNG. The driver passes wall-clock
now_ns and pre-generates all IDs (incident IDs, audit
event_ids, client order IDs) and threads them in via the
event.tasks/state_machine_driver.rs — single tokio task; the
only mutator of live EngineState. Reads events from one
bounded mpsc, calls handle_event, dispatches returned
effects.tasks/trigger_listener.rs — subscribes to Redis
LIQUIDATION_TRIGGER_CHANNEL and exposes an admin REST
endpoint. Both paths normalize to
Event::TriggerReceived.tasks/risk_snapshot_listener.rs — subscribes to
RISK_SNAPSHOT_UPDATED_CHANNEL. Emits
Event::RiskSnapshotUpdated so the core can detect
resolution and stale-mark conditions.tasks/og_client.rs — wraps OG REST calls. Outcomes feed
back as Event::OrderAcked /
Event::OrderRejected. Fills are observed via the dropcopy
stream and translated to Event::OrderFilled.tasks/persistence.rs — writes Postgres
liquidation_account_state (live state) and ClickHouse
liquidation_log (audit).tasks/notifier.rs — Slack threads keyed by
(user_id, incident_id).tasks/tick.rs — periodic Event::Tick
(default 1s) so stuck in-flight orders re-evaluate.bootstrap.rs — on startup, rebuilds
EngineState and runs the bidirectional reconcile described
in §5.5.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.
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.
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
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.
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 markThe 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.
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 widthIdle 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.
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.
| 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.
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.
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 | — |
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 --> [*]
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.
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.
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:
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.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.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.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.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.
(user, instrument). Enforced by the pure core: it
will not emit PlaceLiquidationOrder if
AccountState::Liquidating.inflight already contains that
instrument.JoinSet (one task per
PlaceLiquidationOrder). Results return as events. The
driver task itself stays serial; only it touches state.select_next_action(state) → Option<Effect> resolves
one action per tick in strict priority order: cancel any prerequisite →
place the best protective order (largest current MMR across all
Liquidating accounts) → resume the best recoverable park →
park an account. Risk-reducing placement always outranks parking, so an
ineligible account waits behind any placeable sibling. Invoked on every
Tick and known-zero-fill reject; terminal fills remain
blocked pending reconciliation.Tick events
first (idempotent and recoverable on the next tick); fills, rejects, and
triggers are never dropped. The producer-side drop is recorded in
liquidation_event_channel_drops_total{kind}.og_client retries with exponential backoff (3 attempts:
100ms / 500ms / 2s). After exhaustion it emits
Event::OrderRejected { reason: TransientExhausted }.InsufficientLiquidity,
InstrumentHalted, LimitNotMarketable): a
single Event::OrderRejected flows in. The core increments
attempts[instrument].attempts[instrument] ≥ max_attempts_before_park, the core
records pre_park_probe_until = now + safety_tick and emits
PreParkProbeStarted instead of parking immediately. During
that guard, a fresh RiskSnapshotUpdated { mma ≥ 0 }
resolves the account. If the latest fresh snapshot still has
MMA < 0, or the guard expires without a fresh healthy
snapshot, the next Tick transitions to
Parked { reason: AttemptsExhausted }.Parked { reason: AttemptsExhausted } and
Parked { reason: IneligibleInstrument } resolve on a later
fresh RiskSnapshotUpdated { mma ≥ 0 }, transitioning to
Resolved and incrementing
liquidation_parked_self_resolved_total. An ineligible park
also resumes when a fresh breaching snapshot again yields a
constructible protective order for an in-scope instrument: the core
rebuilds Liquidating (preserving the incident core and
attempt history), emits AccountResumed, and re-runs
scheduling so the cancel (Live) or shadow placement (DryRun) fires in
the same transition. A differing-identity trigger reopens either
recoverable park. AutoLiquidationDisabled and
StaleMark parks are ops-only and require
POST /admin/liquidation/unpark.OrderFilled
carries one terminal cumulative quantity. A valid fill is retained as
TerminalObserved without mutating cached risk exposure or
resetting attempts. Duplicate identical observations are inert; overfill
or contradictory terminal evidence transitions to one-shot
OutcomeUnknown.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.
liquidation_account_state rows from Postgres.L- prefix.
Liquidating PG row, reconstruct
inflight from OG. Re-emit Event::Tick so the
engine resumes decision-making.L- order in OG whose user has no
Liquidating row in PG (crash between
cancel-all/place and persist), emit
EmitSlack { kind: BootstrapOrphan }, increment
liquidation_bootstrap_orphans_total{side: og_only}, and
surface in the reconcile CLI's diff. Engine does not unilaterally cancel
— ops decides.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.
Per CLAUDE.md ("any feature tied to order lifecycle must be tested across disconnect / restart / recovery"):
margin_status. Worst-case
window: time between PG write and CH write — bounded and observable via
liquidation_audit_gap_seconds.og_client retries with
backoff. Inflight orders held by OG come back as fills/rejects on the
dropcopy stream. Pending engine-to-OG POSTs that timed out are surfaced
as
Event::OrderRejected { reason: TransientExhausted }.risk-monitor does
not automatically re-publish (pubsub is
fire-and-forget); the authoritative recovery path is the
margin_status reconcile, which fires on every LE reconnect
via the same code path as §5.5 step 3. While Redis is down, existing
inflight orders continue normally (fills come from OG, not Redis).RiskSnapshotUpdated event for any user is seen for > 30s
the engine emits EmitSlack { kind: EngineHalted } and
increments liquidation_upstream_silence_total.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:
inflight_timeout for outstanding OG calls
in its per-order JoinSet to settle (their results become
events the driver still processes).PersistAccount per modified account so
PG is consistent before exit.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.
This section is the single chokepoint inventory. Every field declared in §4 has exactly one enforcement site listed here.
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.
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.
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 }
SkewViolation → core emits
AppendAuditLog { ClockSkewViolation }, emits
EmitSlack { StaleMark } at the start of a stale/skewed
episode, and refuses to act for that user this tick. Shell increments
liquidation_clock_skew_violations_total at the same
emission site.Stale → core emits
AppendAuditLog { StaleMark }, emits
EmitSlack { StaleMark } at the start of a stale/skewed
episode, and refuses to act for that user this tick.Fresh → normal MMA-based action selection.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.
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.
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.
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).
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.
| 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 |
Three sinks, one correlation key.
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.
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).
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.
Health (mirrors the order-gateway PR #1733 shape):
liquidation_triggers_total{reason, source},
liquidation_orders_placed_total{instrument},
liquidation_orders_filled_total,
liquidation_orders_rejected_total{reason},
liquidation_alerts_total{kind},
liquidation_duplicate_triggers_total,
liquidation_stale_trigger_dropped_total,
liquidation_pre_park_probe_total{outcome=resolved|parked|timeout},
liquidation_parked_self_resolved_total.liquidation_accounts_in_breach,
liquidation_accounts_parked,
liquidation_inflight_orders,
liquidation_event_channel_depth,
liquidation_killswitch_stale_seconds,
liquidation_audit_gap_seconds.liquidation_breach_to_first_action_seconds,
liquidation_breach_to_resolved_seconds,
liquidation_clock_skew_ms.Failure-mode (the counters Rome learned to need):
liquidation_log_rows_dropped_total (CH retry exhausted;
PagerDuty)liquidation_pg_persist_failures_totalliquidation_slack_failures_totalliquidation_event_channel_drops_total{kind} (only
Tick is droppable; non-zero on fills/rejects is a
bug-bug)liquidation_bootstrap_orphans_total{side: pg_only|og_only}liquidation_clock_skew_violations_totalliquidation_upstream_silence_total (>30s without
RiskSnapshotUpdated)Three message classes, all threaded per
(user_id, incident_id):
BreachDetected — incident opened.AccountParked — engine gave up; ops attention required
(carries ParkReason).EngineHalted — kill-switch flipped, upstream silence,
or stale killswitch read.cargo insta inline snapshots for canonical transitions:
single breach → cancel + place; terminal cumulative fill → retain
TerminalObserved until authoritative reconciliation;
consecutive rejects → pre-park probe → park; pre-park probe → resolved
on healthy MMA; Parked(AttemptsExhausted) → resolved on
healthy MMA; dry-run shadows actions; per-user gate blocks emission;
duplicate trigger is idempotent; instrument ineligibility parks only
that instrument.PlaceLiquidationOrder;
Parked(AttemptsExhausted) only exits on fresh healthy MMA
or admin unpark; resolved accounts emit no actions; in-flight cap per
(user, instrument) never exceeded; duplicate events are
idempotent; qty/side always matches closing direction.Tick only, shutdown drains in-flight
JoinSet within inflight_timeout, Redis wakeup
held while margin_status flips
Breached → Cleared then consumed and dropped, OG rejects
reaching threshold while healthy MMA is delayed then resolved by the
pre-park probe or parked-self-resolve path.Reuses existing test-utils seed data and the mock EP3
environment.
margin_status.margin_status clears →
StaleTriggerDropped, no incident.L- order with
no PG row; engine alerts and surfaces in CLI diff.One canonical end-to-end happy-path test in the CI gate; slower
scenarios behind --ignored.
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).
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:
EngineState matches PG snapshot at
<to>.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.
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 |
margin_status row it writes is the only authority LE reads
when the subscriber lands.auto_liquidation_enabled defaults FALSE;
institutional accounts remain protected after the engine goes Live.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).
instruments.maintenance_margin_pct is already
first-class and risk snapshots publish MMR + MMA.liquidation_limit_bps (chosen). Catastrophic-fill
protection outweighs completeness. Unfilled limits retry next tick.rs/rome.margin_status row (chosen). PR #3 introduces a
single Postgres row that doubles as MLL margin-status persistence and
the LE trigger-of-record; no separate outbox table. Pubsub gives
sub-second wakeup latency on the happy path, but LE re-reads
margin_status before opening any incident so stale Redis
messages cannot liquidate a recovered user. Bootstrap and reconnect run
the same validation path.(user, instrument) (chosen). Stuck
market parks only itself; the rest of the account continues to wind
down.Effect::EmitMetric in the core
(rejected). Mixes "decisions" with "telemetry" and forces the
core to know metric names.TriggerReason.liquidation_limit_bps
automatically on repeated LimitNotMarketable rejects before
parking.Operational artefacts to land alongside the rollout:
docs/runbooks/liquidation-parked.md — how ops
investigates a Parked account, force-closes via EP3 console
if needed, and unparks.docs/runbooks/liquidation-bootstrap-orphan.md — what to
do when liquidation_bootstrap_orphans_total{side: og_only}
is non-zero on startup.