Skip to main content

backbone_integrations/application/service/
integrations_write_service.rs

1//! The hand-authored integrations write path (user-owned; survives regen).
2//!
3//! The connector hub: receive an inbound provider event **idempotently** on (connector, external_id) —
4//! providers deliver webhooks at-least-once, so a retry must not re-map (double-apply a payment / create a
5//! duplicate order) — map it to an internal action via a `TargetPort`, or record it as intentionally
6//! ignored. Posts NO GL. Integrations reaches a module only through its public contract (the port).
7//!
8//! Tenancy (ADR-0029): the module is tenant-agnostic — its tables carry no company column and this
9//! service invents no scope. Transactions bind the ambient org scope when the composing service
10//! resolved one; standalone statements ride the request-dedicated connection when one is bound and
11//! run plainly on the pool otherwise, so a decorated deployment's org fence (org_unit_id + RLS +
12//! the decorator's per-unit uniques) enforces isolation and an unscoped write fails closed. The
13//! `company_id` that remains on the boundary shapes here is the documented legacy twin: the
14//! webhook names it, and it routes to the two still-company-shaped edges — the `TargetPort`
15//! (target modules may still be company-fenced) and the outbox/event payloads (the relay keeps a
16//! tenant column).
17
18use backbone_orm::org_scope;
19use chrono::Utc;
20use sqlx::PgPool;
21use uuid::Uuid;
22
23use crate::infrastructure::persistence::{
24    IntegrationConnectorRepository, IntegrationEventRepository, NewConnectorRow, NewEventRow,
25};
26
27use super::integrations_events::*;
28use super::integrations_ports::*;
29
30#[derive(Debug, thiserror::Error)]
31pub enum IntegrationError {
32    #[error("db: {0}")]
33    Db(#[from] sqlx::Error),
34    #[error("not found: {0}")]
35    NotFound(&'static str),
36    #[error("invalid state: {0}")]
37    InvalidState(&'static str),
38    #[error("invalid input: {0}")]
39    Invalid(String),
40    #[error("mapping rejected: {0}")]
41    MappingRejected(String),
42}
43
44pub struct NewConnector {
45    pub provider: String,
46    pub kind: String,      // payment_gateway | marketplace | bank_feed | courier
47    pub direction: String, // inbound | outbound | both
48}
49
50/// An inbound provider event (already parsed into `payload`).
51pub struct InboundEvent {
52    /// Legacy tenancy twin (ADR-0029): the module's own tables carry no company column, but the
53    /// event routes to two still-company-shaped edges — the `TargetPort` request to the target
54    /// module and the outbox/event payloads (the cross-tenant relay keys on it). An unknown or
55    /// stale value fails closed at those edges, never inside this module.
56    pub company_id: Uuid,
57    pub connector_id: Uuid,
58    pub event_type: String,
59    /// The provider's raw notification id — audit only (NOT the dedup key; it varies per notification).
60    pub external_id: String,
61    /// The BUSINESS action identity (order/transaction ref + terminal state, e.g. "SO-9:settled") — the
62    /// dedup key, stable across the multiple notifications a provider sends per order.
63    pub business_key: String,
64    pub raw: String,
65    pub payload: serde_json::Value,
66}
67
68#[derive(Debug, Clone, PartialEq)]
69pub struct ReceiveOutcome {
70    pub event_id: Uuid,
71    pub status: String, // mapped | ignored | failed | duplicate
72    pub mapped_ref_id: Option<Uuid>,
73    pub duplicate: bool,
74}
75
76/// A failed event, surfaced so an operator can see + retry an unbooked action without touching the ledger.
77#[derive(Debug, Clone, PartialEq)]
78pub struct FailedEvent {
79    pub event_id: Uuid,
80    pub event_type: String,
81    pub external_id: String,
82    pub business_key: String,
83    pub error_detail: Option<String>,
84}
85
86pub struct IntegrationsWriteService {
87    pool: PgPool,
88    connectors: IntegrationConnectorRepository,
89    events: IntegrationEventRepository,
90}
91
92impl IntegrationsWriteService {
93    pub fn new(pool: PgPool) -> Self {
94        let connectors = IntegrationConnectorRepository::new(pool.clone());
95        let events = IntegrationEventRepository::new(pool.clone());
96        Self { pool, connectors, events }
97    }
98
99    /// Register a connector to an external provider.
100    pub async fn register_connector(&self, c: NewConnector) -> Result<Uuid, IntegrationError> {
101        if c.provider.trim().is_empty() {
102            return Err(IntegrationError::Invalid("connector needs a provider".into()));
103        }
104        let id = Uuid::new_v4();
105        // Tenancy (ADR-0029): no scope is invented here — the insert rides the composing service's
106        // request scope when one is bound (the decorator's per-unit provider unique then arbitrates
107        // duplicates), and fails closed on an org-scoped deployment otherwise.
108        let r = self.connectors.insert_connector(&self.pool, &NewConnectorRow {
109            id,
110            provider: &c.provider,
111            kind: &c.kind,
112            direction: &c.direction,
113        })
114        .await;
115        match r {
116            Ok(_) => Ok(id),
117            Err(e) if e.as_database_error().map(|d| d.is_unique_violation()).unwrap_or(false) =>
118                Err(IntegrationError::Invalid("a connector for this provider already exists".into())),
119            Err(e) => Err(e.into()),
120        }
121    }
122
123    /// Receive an inbound event: dedup on (connector, external_id), map it to an internal action via the
124    /// `TargetPort`, and record the outcome. A retried webhook returns the original with `duplicate=true` —
125    /// it never re-maps. Emits `IntegrationEventMapped` / `IntegrationEventFailed` / `IntegrationEventIgnored`.
126    pub async fn receive_event(
127        &self,
128        e: InboundEvent,
129        mapper: &dyn TargetPort,
130        events: &dyn IntegrationEventSink,
131    ) -> Result<ReceiveOutcome, IntegrationError> {
132        if e.external_id.trim().is_empty() {
133            return Err(IntegrationError::Invalid("an inbound event needs an external_id".into()));
134        }
135        if e.business_key.trim().is_empty() {
136            return Err(IntegrationError::Invalid("an inbound event needs a business_key (the order/transaction ref + state)".into()));
137        }
138        // Tenancy (ADR-0029): the reads below ride the composing service's request scope when one is
139        // bound — the decorator's org fence then hides another tenant's connector — and run plainly
140        // on the pool otherwise (a decorated deployment yields nothing unscoped, fail closed).
141        // The connector must exist and be active.
142        let conn = self
143            .connectors
144            .fetch_gate(&self.pool, e.connector_id)
145            .await?
146            .ok_or(IntegrationError::NotFound("connector"))?;
147        if conn.status != "active" {
148            return Err(IntegrationError::InvalidState("connector is not active"));
149        }
150        let connector_kind = conn.kind;
151
152        // Claim the (connector, business_key) dedup slot — a webhook retry OR a second notification for the
153        // same business action conflicts here; a new business action does not.
154        let inserted = self
155            .events
156            .claim_event(&self.pool, &NewEventRow {
157                id: Uuid::new_v4(),
158                connector_id: e.connector_id,
159                event_type: &e.event_type,
160                external_id: &e.external_id,
161                business_key: &e.business_key,
162                raw: &e.raw,
163            })
164            .await?;
165
166        let Some(event_id) = inserted else {
167            let row = self
168                .events
169                .fetch_by_business_key(&self.pool, e.connector_id, &e.business_key)
170                .await?;
171            return Ok(ReceiveOutcome {
172                event_id: row.id, status: row.status,
173                mapped_ref_id: row.mapped_ref_id, duplicate: true,
174            });
175        };
176
177        let req = MapRequest {
178            company_id: e.company_id, connector_kind, event_type: e.event_type.clone(),
179            external_id: e.external_id.clone(),
180            // Forward the BUSINESS key (not the per-notification event id) so the target dedups the effect
181            // on an order-scoped key — the second layer against a double-applied payment.
182            idempotency_key: e.business_key.clone(), payload: e.payload.clone(),
183        };
184        match mapper.map(&req).await {
185            Ok(MapOutcome::Mapped(mref)) => {
186                let ev = IntegrationEvent::IntegrationEventMapped(IntegrationEventMapped {
187                    event_id, company_id: e.company_id, connector_id: e.connector_id, event_type: e.event_type.clone(),
188                    external_id: e.external_id.clone(),
189                    internal_ref_type: mref.internal_ref_type.clone(), internal_ref_id: mref.internal_ref_id,
190                });
191                let mut tx = self.pool.begin().await?;
192                bind_ambient_org_scope(&mut tx).await?;
193                self.events
194                    .mark_mapped(&mut tx, event_id, &mref.internal_ref_type, mref.internal_ref_id)
195                    .await?;
196                stage(&mut tx, &ev).await?;
197                tx.commit().await?;
198                events.publish(&ev);
199                Ok(ReceiveOutcome { event_id, status: "mapped".into(), mapped_ref_id: Some(mref.internal_ref_id), duplicate: false })
200            }
201            Ok(MapOutcome::Ignored(reason)) => {
202                let ev = IntegrationEvent::IntegrationEventIgnored {
203                    event_id, company_id: e.company_id, connector_id: e.connector_id, external_id: e.external_id.clone(), reason: reason.clone(),
204                };
205                let mut tx = self.pool.begin().await?;
206                bind_ambient_org_scope(&mut tx).await?;
207                self.events.mark_ignored(&mut tx, event_id, &reason).await?;
208                stage(&mut tx, &ev).await?;
209                tx.commit().await?;
210                events.publish(&ev);
211                Ok(ReceiveOutcome { event_id, status: "ignored".into(), mapped_ref_id: None, duplicate: false })
212            }
213            Err(rej) => {
214                let ev = IntegrationEvent::IntegrationEventFailed {
215                    event_id, company_id: e.company_id, connector_id: e.connector_id,
216                    external_id: e.external_id.clone(), reason: rej.code.clone(),
217                };
218                let mut tx = self.pool.begin().await?;
219                bind_ambient_org_scope(&mut tx).await?;
220                self.events.mark_failed(&mut tx, event_id, &rej.message).await?;
221                stage(&mut tx, &ev).await?;
222                tx.commit().await?;
223                events.publish(&ev);
224                Ok(ReceiveOutcome { event_id, status: "failed".into(), mapped_ref_id: None, duplicate: false })
225            }
226        }
227    }
228
229    /// The failure report for a connector — the failed events with their keys + error — so an operator can
230    /// see WHICH provider events (e.g. settled payments) failed to book, WITHOUT querying the private ledger
231    /// (completeness council 2026-07-11).
232    pub async fn failures(&self, connector_id: Uuid) -> Result<Vec<FailedEvent>, IntegrationError> {
233        // Tenancy (ADR-0029), ID-only pattern: the connector id alone identifies the report, so the read
234        // rides the composing service's request-dedicated connection when one is bound — the decorator's
235        // org fence then hides another tenant's connector — and runs plainly on the pool otherwise.
236        let rows = self.events.fetch_failed(&self.pool, connector_id).await?;
237        Ok(rows.into_iter().map(|r| FailedEvent {
238            event_id: r.id, event_type: r.event_type, external_id: r.external_id,
239            business_key: r.business_key, error_detail: r.error_detail,
240        }).collect())
241    }
242
243    /// Re-drive a connector's FAILED events through the target after the cause is fixed — the recovery path
244    /// the dedup would otherwise weld shut (a re-delivered webhook dedups and never re-maps). Re-invokes the
245    /// `TargetPort` under the SAME business-key idempotency the target enforces, so it can't double-apply; a
246    /// still-`failed` event that now maps transitions `failed → mapped`. Returns the number newly mapped
247    /// (completeness council 2026-07-11).
248    ///
249    /// `company_id` is the legacy tenancy twin (ADR-0029) — the module's tables carry no company column,
250    /// but the re-mapped events still route to the `TargetPort` and the outbox/event payloads, which are
251    /// company-shaped; the CALLER (the composing host, which knows the tenant) names it, and an unknown
252    /// value fails closed at those edges, never inside this module.
253    pub async fn retry_failed(
254        &self,
255        company_id: Uuid,
256        connector_id: Uuid,
257        mapper: &dyn TargetPort,
258        events: &dyn IntegrationEventSink,
259    ) -> Result<usize, IntegrationError> {
260        // Tenancy (ADR-0029), ID-only pattern: identified by the connector id alone, so this first read
261        // rides the request-dedicated connection when one is bound.
262        let conn = self
263            .connectors
264            .fetch_for_retry(&self.pool, connector_id)
265            .await?
266            .ok_or(IntegrationError::NotFound("connector"))?;
267        let connector_kind = conn.kind;
268
269        let rows = self.events.fetch_failed(&self.pool, connector_id).await?;
270
271        let mut mapped = 0usize;
272        for row in &rows {
273            let event_id = row.id;
274            let business_key = row.business_key.clone();
275            let req = MapRequest {
276                company_id, connector_kind: connector_kind.clone(), event_type: row.event_type.clone(),
277                external_id: row.external_id.clone(), idempotency_key: business_key.clone(),
278                payload: serde_json::from_str(&row.payload).unwrap_or(serde_json::Value::Null),
279            };
280            match mapper.map(&req).await {
281                Ok(MapOutcome::Mapped(mref)) => {
282                    let ev = IntegrationEvent::IntegrationEventMapped(IntegrationEventMapped {
283                        event_id, company_id, connector_id, event_type: row.event_type.clone(),
284                        external_id: row.external_id.clone(),
285                        internal_ref_type: mref.internal_ref_type.clone(), internal_ref_id: mref.internal_ref_id,
286                    });
287                    let mut tx = self.pool.begin().await?;
288                    bind_ambient_org_scope(&mut tx).await?;
289                    let m = self
290                        .events
291                        .retry_mark_mapped(&mut tx, event_id, &mref.internal_ref_type, mref.internal_ref_id)
292                        .await?;
293                    if m == 1 {
294                        stage(&mut tx, &ev).await?;
295                        tx.commit().await?;
296                        events.publish(&ev);
297                        mapped += 1;
298                    } else {
299                        tx.rollback().await?;
300                    }
301                }
302                Ok(MapOutcome::Ignored(reason)) => {
303                    self.events.retry_mark_ignored(&self.pool, event_id, &reason).await?;
304                }
305                Err(rej) => {
306                    self.events.set_error_detail(&self.pool, event_id, &rej.message).await?;
307                }
308            }
309        }
310        Ok(mapped)
311    }
312}
313
314/// Bind the ambient org scope of the current request onto a fresh transaction, when the composing
315/// service resolved one (ADR-0029). The decorator's org fence (org_unit_id + RLS + the acting-unit
316/// default) then applies inside the transaction; with no scope bound the transaction runs unscoped
317/// and a decorated deployment refuses its writes (fail closed).
318async fn bind_ambient_org_scope(tx: &mut sqlx::Transaction<'_, sqlx::Postgres>) -> Result<(), IntegrationError> {
319    if let Some(scope) = org_scope::current_org_scope() {
320        org_scope::bind_org_scope_on(&mut **tx, &scope).await?;
321    }
322    Ok(())
323}
324
325async fn stage(tx: &mut sqlx::Transaction<'_, sqlx::Postgres>, event: &IntegrationEvent) -> Result<(), IntegrationError> {
326    let (etype, agg_id) = match event {
327        IntegrationEvent::IntegrationEventMapped(m) => ("IntegrationEventMapped", m.event_id),
328        IntegrationEvent::IntegrationEventFailed { event_id, .. } => ("IntegrationEventFailed", *event_id),
329        IntegrationEvent::IntegrationEventIgnored { event_id, .. } => ("IntegrationEventIgnored", *event_id),
330    };
331    let payload = serde_json::to_value(event).map_err(|e| IntegrationError::Invalid(e.to_string()))?;
332    let company_id: Uuid = payload
333        .get("company_id")
334        .and_then(|v| v.as_str())
335        .ok_or_else(|| IntegrationError::Invalid("integration event missing company_id".into()))?
336        .parse()
337        .map_err(|e| IntegrationError::Invalid(format!("company_id parse: {e}")))?;
338    let record = backbone_outbox::OutboxRecord::new(
339        etype, "IntegrationEvent", agg_id.to_string(), company_id, payload, Utc::now(),
340    );
341    backbone_outbox::outbox::stage(&mut **tx, "integrations", &record)
342        .await.map_err(|e| IntegrationError::Invalid(format!("outbox stage: {e}")))?;
343    Ok(())
344}