//! `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 { // 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, ) -> 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(()) }