`OutboxEventEmitter` replaces `NoopEventEmitter` in the picloud binary's `Services` bundle. KV mutations now fan out to the outbox via `TriggerRepo::list_matching_kv` — one row per matching trigger, carrying the serialized `TriggerEvent` payload + the matching trigger's retry policy. `Dispatcher` is the single tokio task that polls the outbox every 100ms, claims due rows via FOR UPDATE SKIP LOCKED (with a batch cap), and routes each to the executor. Shares the `ExecutionGate` with sync HTTP per design notes §2 — gate saturation reschedules the row instead of dropping it. Outcome handling matches design notes §3 and §4: - reply_to.is_some() (sync HTTP): never retry. Deliver via `InboxResolver`; if the receiver was dropped, write an `abandoned_executions` row. - is_dead_letter_handler == true: never retry, never DL. On failure, annotate the original DL row with `resolution = 'handler_failed'`. Stops the recursion that would otherwise re-fire a broken handler script. - Otherwise async: bump attempt_count, reschedule with exponential backoff + ±jitter; once max_attempts is reached, write a `dead_letters` row and drop from outbox. - Trigger-depth limit: `cx.trigger_depth > max_trigger_depth` skips execution entirely (log + future metric), NEVER dead-letters. Loops are not retried via the DL chain — they're terminated. `InboxResolver` trait lands in `picloud-shared` with a `NoopInboxResolver` bootstrap that flags every delivery as `Abandoned`. Commit 6 replaces the noop with the real in-process registry in `orchestrator-core`. `AdminPrincipalResolver` builds a `Principal` from a trigger's `registered_by_principal` user id so the dispatched script executes as the trigger registrant (design notes §4). Unit tests cover backoff math (exponential/linear/constant) + jitter range + ExecError → InboxFailureKind classification + the status-code table mapping. Integration tests for the full dispatcher loop need a real Postgres + executor; reviewer runs them via the manual smoke flow in the plan / HANDBACK. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
87 lines
3.3 KiB
Rust
87 lines
3.3 KiB
Rust
//! `InboxResolver` — abstraction the dispatcher uses to deliver sync
|
|
//! HTTP results back to the orchestrator that's awaiting them on a
|
|
//! oneshot channel. Lives in `picloud-shared` because the dispatcher
|
|
//! (manager-core) and the registry impl (orchestrator-core) live in
|
|
//! different crates and need a shared trait surface.
|
|
//!
|
|
//! v1.1.1 ships an in-process implementation in `orchestrator-core`
|
|
//! that keeps a `HashMap<inbox_id, oneshot::Sender<...>>`. Cluster
|
|
//! mode (v1.3+) swaps this for a Postgres `LISTEN/NOTIFY`-based
|
|
//! resolver without touching the dispatcher code (design notes §3
|
|
//! implementation table).
|
|
//!
|
|
//! Until commit 6 wires up the real registry, `NoopInboxResolver`
|
|
//! (`Abandoned` for every attempt) keeps the dispatcher able to run.
|
|
|
|
use async_trait::async_trait;
|
|
use uuid::Uuid;
|
|
|
|
use crate::ExecResponseSummary;
|
|
|
|
/// Result of trying to hand back a sync-HTTP outcome.
|
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
|
pub enum InboxDeliveryOutcome {
|
|
/// Receiver still attached; result was delivered. Dispatcher
|
|
/// deletes the outbox row.
|
|
Delivered,
|
|
/// Receiver was dropped (orchestrator timed out). Dispatcher
|
|
/// writes an `abandoned_executions` row.
|
|
Abandoned,
|
|
}
|
|
|
|
/// Outcome shape the dispatcher delivers to the inbox. Carries enough
|
|
/// to reconstruct an HTTP response — full body via JSON, optional
|
|
/// error string when the executor reported a failure.
|
|
#[derive(Debug, Clone)]
|
|
pub enum InboxResult {
|
|
/// Successful execution. `response` is the `ExecResponse` summary
|
|
/// (status code + body + headers + logs).
|
|
Success(ExecResponseSummary),
|
|
/// Failure modes — script threw, op-budget, timeout, etc. The
|
|
/// orchestrator maps these to the design-notes §3 status codes
|
|
/// (422/502/503/504/507/500) when responding to the HTTP caller.
|
|
Failure {
|
|
kind: InboxFailureKind,
|
|
message: String,
|
|
},
|
|
}
|
|
|
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
|
pub enum InboxFailureKind {
|
|
/// Script's Rhai code threw or hit a runtime error → 502.
|
|
Runtime,
|
|
/// Wall-clock exceeded → 504.
|
|
Timeout,
|
|
/// Operation budget exceeded → 507.
|
|
OperationBudget,
|
|
/// Gate refused admission → 503.
|
|
Overloaded,
|
|
/// Script parse failure / bad-request → 422.
|
|
Validation,
|
|
/// Platform problem (executor crashed, dispatcher crashed, etc.) → 500.
|
|
Platform,
|
|
}
|
|
|
|
#[async_trait]
|
|
pub trait InboxResolver: Send + Sync {
|
|
/// Attempt to deliver `result` to the receiver registered under
|
|
/// `inbox_id`. Returns `Delivered` if the channel was alive,
|
|
/// `Abandoned` if the receiver was already dropped (the
|
|
/// orchestrator's timeout fired before the dispatcher got here).
|
|
async fn deliver(&self, inbox_id: Uuid, result: InboxResult) -> InboxDeliveryOutcome;
|
|
}
|
|
|
|
/// Bootstrap impl used before the real registry is wired in. Every
|
|
/// delivery is treated as abandoned — the dispatcher records an
|
|
/// abandoned-execution row and moves on. Replaced in `build_app` with
|
|
/// the in-process `InboxRegistry` from orchestrator-core.
|
|
#[derive(Debug, Default, Clone, Copy)]
|
|
pub struct NoopInboxResolver;
|
|
|
|
#[async_trait]
|
|
impl InboxResolver for NoopInboxResolver {
|
|
async fn deliver(&self, _inbox_id: Uuid, _result: InboxResult) -> InboxDeliveryOutcome {
|
|
InboxDeliveryOutcome::Abandoned
|
|
}
|
|
}
|