`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>
63 lines
2.1 KiB
Rust
63 lines
2.1 KiB
Rust
//! `PrincipalResolver` — turns a `registered_by_principal` user id from
|
|
//! a trigger row into the `Principal` the dispatcher passes through to
|
|
//! the executor. Per design notes §4, a trigger execution runs as the
|
|
//! user that registered the trigger; the original event's caller is
|
|
//! recorded elsewhere (on the outbox row, for forensics) and does not
|
|
//! become the execution principal.
|
|
|
|
use async_trait::async_trait;
|
|
use picloud_shared::{AdminUserId, Principal};
|
|
|
|
use crate::admin_user_repo::{AdminUserRepository, AdminUserRepositoryError};
|
|
|
|
#[derive(Debug, thiserror::Error)]
|
|
pub enum PrincipalResolverError {
|
|
#[error("user not found: {0}")]
|
|
NotFound(AdminUserId),
|
|
#[error("user is inactive: {0}")]
|
|
Inactive(AdminUserId),
|
|
#[error("admin user repo error: {0}")]
|
|
Backend(String),
|
|
}
|
|
|
|
#[async_trait]
|
|
pub trait PrincipalResolver: Send + Sync {
|
|
async fn resolve(&self, user_id: AdminUserId) -> Result<Principal, PrincipalResolverError>;
|
|
}
|
|
|
|
pub struct AdminPrincipalResolver {
|
|
users: std::sync::Arc<dyn AdminUserRepository>,
|
|
}
|
|
|
|
impl AdminPrincipalResolver {
|
|
#[must_use]
|
|
pub fn new(users: std::sync::Arc<dyn AdminUserRepository>) -> Self {
|
|
Self { users }
|
|
}
|
|
}
|
|
|
|
#[async_trait]
|
|
impl PrincipalResolver for AdminPrincipalResolver {
|
|
async fn resolve(&self, user_id: AdminUserId) -> Result<Principal, PrincipalResolverError> {
|
|
let row = self
|
|
.users
|
|
.get(user_id)
|
|
.await
|
|
.map_err(|e: AdminUserRepositoryError| PrincipalResolverError::Backend(e.to_string()))?
|
|
.ok_or(PrincipalResolverError::NotFound(user_id))?;
|
|
if !row.is_active {
|
|
return Err(PrincipalResolverError::Inactive(user_id));
|
|
}
|
|
Ok(Principal {
|
|
user_id,
|
|
instance_role: row.instance_role,
|
|
// Trigger executions are cookie-session-style (no API key
|
|
// scope restriction). Per-app permissions are evaluated
|
|
// via `authz::can` against the `app_id` of the resource
|
|
// the script touches, exactly like an admin invocation.
|
|
scopes: None,
|
|
app_binding: None,
|
|
})
|
|
}
|
|
}
|