Files
PiCloud/crates/manager-core/src/outbox_event_emitter.rs
MechaCat02 a2360a9464 refactor(outbox): make the trigger fan-out connection-scoped
`OutboxEventEmitter` resolved matching triggers and inserted outbox rows
through `Arc<dyn TriggerRepo>` / `Arc<dyn OutboxRepo>`, i.e. against the
pool — each query on whatever connection it happened to get. That makes a
transactional outbox impossible: the data write and the outbox rows can
never share a transaction, so a crash (or an outbox error) between them
silently loses the trigger event.

Take the emitter down to a connection instead of a pool:

  * `outbox_repo::insert_on(exec, row)` and `trigger_repo::list_matching_on`
    / `list_matching_shared_on(exec, ...)` are generic over `PgExecutor`, so
    the same SQL serves a pooled connection or a `&mut *tx`. The repo trait
    methods delegate to them — the SQL keeps exactly one home.
  * `emit_on` / `emit_shared_on` take a `&mut PgConnection` and run the whole
    fan-out on it. The `ServiceEventEmitter` impl acquires one pooled
    connection and calls them, so behaviour is unchanged today; a caller
    holding a transaction can now pass `&mut *tx` and have the outbox rows
    commit with the write.
  * `OutboxEventEmitter::new` takes the `PgPool` directly (it was only ever
    constructed once, in the host wiring).

Also collapses the three copy-pasted per-app match queries (kv/docs/files
differ only in the `kind` discriminator and detail table) and the three
`emit_*` bodies into one `plan()` + one match fn, so the suppression
anti-join, the chain walk, and the empty-ops-means-any-op semantic each
exist once rather than three times.

No behaviour change — pure refactor. It is the seam the transactional
write lands on next.

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-14 19:20:41 +02:00

225 lines
8.0 KiB
Rust

//! `OutboxEventEmitter` — the real `ServiceEventEmitter` behind every stateful
//! service. On each collection mutation (KV / docs / files) it looks up the
//! triggers the event matches and writes one outbox row per match; the
//! dispatcher reads them back out-of-band.
//!
//! **The fan-out is connection-scoped, not pool-scoped.** [`emit_on`] and
//! [`emit_shared_on`] take a `&mut PgConnection`, so a caller that already holds
//! a transaction can pass `&mut *tx` and have the outbox rows commit *with* the
//! data write that produced them — the transactional-outbox pattern. That closes
//! the window where a row committed but its trigger never fired because the
//! outbox insert failed (or the process died) immediately afterwards. Services
//! that write through `crate::atomic_write` take that path; the
//! `ServiceEventEmitter` trait impl below is the non-transactional entry point,
//! kept for the callers (and the tests) that only need best-effort emission.
//!
//! Non-collection `ServiceEvent` sources are silently dropped: the SDK calls
//! `events.emit(...)` unconditionally for forward compat, and every source the
//! dispatcher has an arm for is handled here.
use async_trait::async_trait;
use picloud_shared::{
DocsEventOp, EmitError, FileMeta, FilesEventOp, GroupId, KvEventOp, SdkCallCx, ServiceEvent,
ServiceEventEmitter, TriggerEvent,
};
use sqlx::{PgConnection, PgPool};
use crate::outbox_repo::{insert_on, NewOutboxRow, OutboxSourceKind};
use crate::trigger_repo::{list_matching_on, list_matching_shared_on, KvMatchRow};
pub struct OutboxEventEmitter {
pool: PgPool,
}
impl OutboxEventEmitter {
#[must_use]
pub fn new(pool: PgPool) -> Self {
Self { pool }
}
}
#[async_trait]
impl ServiceEventEmitter for OutboxEventEmitter {
async fn emit(&self, cx: &SdkCallCx, event: ServiceEvent) -> Result<(), EmitError> {
let mut conn = self
.pool
.acquire()
.await
.map_err(|e| EmitError::Unavailable(format!("acquire connection: {e}")))?;
emit_on(&mut conn, cx, &event).await
}
async fn emit_shared(
&self,
cx: &SdkCallCx,
owning_group: GroupId,
event: ServiceEvent,
) -> Result<(), EmitError> {
let mut conn = self
.pool
.acquire()
.await
.map_err(|e| EmitError::Unavailable(format!("acquire connection: {e}")))?;
emit_shared_on(&mut conn, cx, owning_group, &event).await
}
}
/// Everything the fan-out needs, derived from a `ServiceEvent` with no I/O: the
/// trigger-table coordinates to match on, and the `TriggerEvent` the handler
/// will see as `ctx.event`. `None` for a source or op the dispatcher has no arm
/// for — the caller drops the event quietly.
struct EventPlan {
source_kind: OutboxSourceKind,
/// `triggers.kind` discriminator. A literal — never user data, since it is
/// interpolated into the match SQL.
kind: &'static str,
/// Per-kind detail table. A literal, for the same reason.
detail_table: &'static str,
collection: String,
op_str: &'static str,
trigger_event: TriggerEvent,
}
fn plan(event: &ServiceEvent) -> Option<EventPlan> {
// Collection events always carry a collection; defensively skip if not.
let collection = event.collection.clone()?;
let key = event.key.clone().unwrap_or_default();
match event.source {
"kv" => {
let op = KvEventOp::from_wire(event.op)?;
Some(EventPlan {
source_kind: OutboxSourceKind::Kv,
kind: "kv",
detail_table: "kv_trigger_details",
collection: collection.clone(),
op_str: op.as_str(),
trigger_event: TriggerEvent::Kv {
op,
collection,
key,
value: event.payload.clone(),
},
})
}
"docs" => {
let op = DocsEventOp::from_wire(event.op)?;
Some(EventPlan {
source_kind: OutboxSourceKind::Docs,
kind: "docs",
detail_table: "docs_trigger_details",
collection: collection.clone(),
op_str: op.as_str(),
trigger_event: TriggerEvent::Docs {
op,
collection,
id: key,
data: event.payload.clone(),
prev_data: event.old_payload.clone(),
},
})
}
"files" => {
let op = FilesEventOp::from_wire(event.op)?;
// The payload is the `FileMeta` JSON the files service emitted —
// never the blob bytes.
let meta: FileMeta = serde_json::from_value(event.payload.clone()?).ok()?;
Some(EventPlan {
source_kind: OutboxSourceKind::Files,
kind: "files",
detail_table: "files_trigger_details",
collection: collection.clone(),
op_str: op.as_str(),
trigger_event: TriggerEvent::Files {
op,
collection,
id: meta.id.to_string(),
name: meta.name,
content_type: meta.content_type,
size: meta.size,
checksum: meta.checksum,
prev: event.old_payload.clone(),
},
})
}
_ => None,
}
}
/// Fan a collection mutation out over the writing app's own + inherited
/// triggers. Pass `&mut *tx` to have the outbox rows commit with the write.
pub(crate) async fn emit_on(
conn: &mut PgConnection,
cx: &SdkCallCx,
event: &ServiceEvent,
) -> Result<(), EmitError> {
let Some(p) = plan(event) else { return Ok(()) };
let matches = list_matching_on(
&mut *conn,
p.kind,
p.detail_table,
cx.app_id,
&p.collection,
p.op_str,
)
.await
.map_err(|e| EmitError::Unavailable(format!("trigger lookup: {e}")))?;
write_outbox_rows(conn, cx, &p, matches).await
}
/// §11.6: fan a SHARED-collection write out over the `shared = true` triggers on
/// the OWNING group. Each fires under the WRITER app (`cx.app_id`) — the same
/// "a group template runs under the firing app" model as an inherited trigger.
pub(crate) async fn emit_shared_on(
conn: &mut PgConnection,
cx: &SdkCallCx,
owning_group: GroupId,
event: &ServiceEvent,
) -> Result<(), EmitError> {
let Some(p) = plan(event) else { return Ok(()) };
let matches = list_matching_shared_on(
&mut *conn,
p.kind,
p.detail_table,
owning_group,
&p.collection,
p.op_str,
)
.await
.map_err(|e| EmitError::Unavailable(format!("trigger lookup: {e}")))?;
write_outbox_rows(conn, cx, &p, matches).await
}
async fn write_outbox_rows(
conn: &mut PgConnection,
cx: &SdkCallCx,
p: &EventPlan,
matches: Vec<KvMatchRow>,
) -> Result<(), EmitError> {
if matches.is_empty() {
return Ok(());
}
// Serialize the originating event once so the dispatcher can hand it to the
// handler as `ctx.event` without joining back to the trigger row.
let payload = serde_json::to_value(&p.trigger_event)
.map_err(|e| EmitError::Rejected(format!("event serialize: {e}")))?;
for m in matches {
insert_on(
&mut *conn,
NewOutboxRow {
app_id: cx.app_id,
source_kind: p.source_kind,
trigger_id: Some(m.id.into()),
script_id: Some(m.script_id.into()),
reply_to: None,
payload: payload.clone(),
origin_principal: cx.principal.as_ref().map(|pr| pr.user_id),
trigger_depth: cx.trigger_depth.saturating_add(1),
root_execution_id: Some(cx.root_execution_id),
},
)
.await
.map_err(|e| EmitError::Unavailable(format!("outbox insert: {e}")))?;
}
Ok(())
}