Migrations 0008-0011 lay down the triggers framework's storage: - `triggers` + `kv_trigger_details` + `dead_letter_trigger_details` (Layout E, design notes §2). Parent table carries common columns including `registered_by_principal` — the dispatcher uses this to run the trigger as the user that registered it (design notes §4). - `outbox`: universal async dispatch substrate. KV/cron/pubsub/queue/ email/dead-letter all write rows in the same shape; the dispatcher claims due rows via FOR UPDATE SKIP LOCKED. `reply_to` is the NATS-style inbox id for sync HTTP (commit 6) — its presence flags "don't retry" per the design. - `dead_letters`: exact schema from design notes §4 with the four- value `resolution` CHECK constraint (`replayed | ignored | handled_by_script | handler_failed`) and partial index on unresolved rows for the dashboard badge. - `abandoned_executions`: forensic table for the dispatcher's "tried to resolve a dropped inbox" edge case (design notes §3 #9). Repo surfaces with Postgres impls behind traits so unit tests can swap in-memory backings: - `TriggerRepo` — CRUD + the `list_matching_kv` / `list_matching_dead_letter` hot paths the dispatcher uses. Includes a `collection_matches` helper that handles `*`, `prefix:*`, and exact-name globs. - `OutboxRepo` — insert + claim-due + delete + reschedule. - `DeadLetterRepo` — insert + get + list + unresolved-count + resolve + GC. - `AbandonedRepo` — insert + GC. `TriggerConfig::from_env` (new module) follows the existing `SandboxCeiling` env-loading pattern for `PICLOUD_MAX_TRIGGER_DEPTH`, `PICLOUD_TRIGGER_RETRY_*`, `PICLOUD_DEAD_LETTER_RETENTION_DAYS`, and `PICLOUD_ABANDONED_EXECUTIONS_RETENTION_DAYS`. `Capability::AppManageTriggers(AppId)` and `AppDeadLetterManage(AppId)` join the enum. Both map onto the existing `Scope::AppAdmin` per the seven-scope commitment; `role_satisfies` grants them at the `AppAdmin` per-app role. Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
129 lines
3.7 KiB
Rust
129 lines
3.7 KiB
Rust
//! `AbandonedExecutionsRepo` — forensic table written by the
|
|
//! dispatcher when it tries to resolve a sync-HTTP inbox channel
|
|
//! that's already been dropped (orchestrator timed out and gave up).
|
|
//!
|
|
//! Schema: see `migrations/0011_abandoned_executions.sql`.
|
|
//!
|
|
//! Tiny surface: insert + GC. Reading happens via direct SQL when
|
|
//! correlating the metric counter spike.
|
|
|
|
use async_trait::async_trait;
|
|
use chrono::{DateTime, Utc};
|
|
use picloud_shared::{AppId, ScriptId};
|
|
use sqlx::PgPool;
|
|
use uuid::Uuid;
|
|
|
|
#[derive(Debug, thiserror::Error)]
|
|
pub enum AbandonedRepoError {
|
|
#[error("database error: {0}")]
|
|
Db(#[from] sqlx::Error),
|
|
}
|
|
|
|
#[derive(Debug, Clone)]
|
|
pub struct NewAbandonedExecution {
|
|
pub app_id: AppId,
|
|
pub outbox_id: Uuid,
|
|
pub script_id: Option<ScriptId>,
|
|
pub inbox_id: Uuid,
|
|
pub status_code: u16,
|
|
pub result_summary: Option<String>,
|
|
}
|
|
|
|
#[async_trait]
|
|
pub trait AbandonedRepo: Send + Sync {
|
|
async fn insert(&self, row: NewAbandonedExecution) -> Result<Uuid, AbandonedRepoError>;
|
|
|
|
/// Retention sweep — deletes rows older than `older_than` up to
|
|
/// `limit` at a time.
|
|
async fn gc(&self, older_than: DateTime<Utc>, limit: i64) -> Result<u64, AbandonedRepoError>;
|
|
}
|
|
|
|
pub struct PostgresAbandonedRepo {
|
|
pool: PgPool,
|
|
}
|
|
|
|
impl PostgresAbandonedRepo {
|
|
#[must_use]
|
|
pub fn new(pool: PgPool) -> Self {
|
|
Self { pool }
|
|
}
|
|
}
|
|
|
|
const SUMMARY_CAP_BYTES: usize = 4096;
|
|
|
|
#[async_trait]
|
|
impl AbandonedRepo for PostgresAbandonedRepo {
|
|
async fn insert(&self, row: NewAbandonedExecution) -> Result<Uuid, AbandonedRepoError> {
|
|
// Truncate the summary at write-time. The forensic table
|
|
// doesn't need megabytes; the original outbox row may have
|
|
// been arbitrary size but we lose nothing useful by clipping.
|
|
let summary = row.result_summary.map(|s| truncate(s, SUMMARY_CAP_BYTES));
|
|
let (id,): (Uuid,) = sqlx::query_as(
|
|
"INSERT INTO abandoned_executions ( \
|
|
app_id, outbox_id, script_id, inbox_id, status_code, result_summary \
|
|
) VALUES ($1, $2, $3, $4, $5, $6) \
|
|
RETURNING id",
|
|
)
|
|
.bind(row.app_id.into_inner())
|
|
.bind(row.outbox_id)
|
|
.bind(row.script_id.map(ScriptId::into_inner))
|
|
.bind(row.inbox_id)
|
|
.bind(i32::from(row.status_code))
|
|
.bind(summary)
|
|
.fetch_one(&self.pool)
|
|
.await?;
|
|
Ok(id)
|
|
}
|
|
|
|
async fn gc(&self, older_than: DateTime<Utc>, limit: i64) -> Result<u64, AbandonedRepoError> {
|
|
let res = sqlx::query(
|
|
"DELETE FROM abandoned_executions \
|
|
WHERE id IN ( \
|
|
SELECT id FROM abandoned_executions \
|
|
WHERE created_at < $1 \
|
|
FOR UPDATE SKIP LOCKED \
|
|
LIMIT $2 \
|
|
)",
|
|
)
|
|
.bind(older_than)
|
|
.bind(limit)
|
|
.execute(&self.pool)
|
|
.await?;
|
|
Ok(res.rows_affected())
|
|
}
|
|
}
|
|
|
|
fn truncate(mut s: String, max_bytes: usize) -> String {
|
|
if s.len() <= max_bytes {
|
|
return s;
|
|
}
|
|
// Walk back from `max_bytes` to a UTF-8 char boundary so we never
|
|
// panic on `truncate` mid-codepoint.
|
|
let mut cut = max_bytes;
|
|
while cut > 0 && !s.is_char_boundary(cut) {
|
|
cut -= 1;
|
|
}
|
|
s.truncate(cut);
|
|
s
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
|
|
#[test]
|
|
fn truncate_respects_char_boundaries() {
|
|
// 3-byte UTF-8 chars; cap inside the middle char should walk
|
|
// back to the start.
|
|
let s = "héllo".to_string();
|
|
let t = truncate(s, 2);
|
|
assert!(t.is_char_boundary(t.len()));
|
|
assert_eq!(t, "h");
|
|
}
|
|
|
|
#[test]
|
|
fn truncate_passthrough_for_short_strings() {
|
|
assert_eq!(truncate("ok".into(), 100), "ok");
|
|
}
|
|
}
|