Skip to main content

backbone_integrations/infrastructure/persistence/
integration_event_repository.rs

1//! Repository for IntegrationEvent entities
2//!
3//! Originally generated by metaphor-schema; now **user-owned** — this exact path is declared under
4//! `user_owned` in `metaphor.codegen.yaml`, so the generator skips it wholesale. The custom methods
5//! below hold the hand-written IntegrationEvent SQL — the (connector, business_key) dedup claim that
6//! stops a webhook retry re-applying a payment, its re-read, and the mapped/ignored/failed transitions
7//! on both the receive and the retry path (4-layer rule: services orchestrate, repos hold SQL).
8//!
9//! Tenancy (ADR-0029): the table carries no company column and the module invents no scope —
10//! every statement runs on the request-dedicated connection when the composing service bound
11//! one, else plainly on the pool; the COMPOSING service's tenancy decorator owns org scoping
12//! (org_unit_id + RLS), so an unscoped write fails closed on a decorated deployment.
13//!
14//! Thin newtype over `backbone_orm::GenericCrudRepository<IntegrationEvent, backbone_orm::SoftDelete>`.
15//! All standard CRUD methods are available via `Deref`.
16
17use anyhow::Result;
18use sqlx::{PgPool, Row};
19use uuid::Uuid;
20
21use backbone_orm::{company_scope, org_scope};
22
23use crate::domain::entity::IntegrationEvent;
24
25/// Table name for IntegrationEvent entities
26pub const TABLE_NAME: &str = "integrations.integration_events";
27
28/// Repository for IntegrationEvent entities.
29///
30/// All standard CRUD, soft-delete, pagination, and bulk methods are
31/// provided automatically via `Deref` to `backbone_orm::GenericCrudRepository`.
32pub struct IntegrationEventRepository(
33    backbone_orm::GenericCrudRepository<IntegrationEvent, backbone_orm::SoftDelete>,
34);
35
36impl std::ops::Deref for IntegrationEventRepository {
37    type Target = backbone_orm::GenericCrudRepository<IntegrationEvent, backbone_orm::SoftDelete>;
38    fn deref(&self) -> &Self::Target { &self.0 }
39}
40
41impl IntegrationEventRepository {
42    /// Create a new repository instance.
43    pub fn new(pool: PgPool) -> Self {
44        Self(backbone_orm::GenericCrudRepository::new(pool, TABLE_NAME))
45    }
46}
47
48/// The exact row an inbound-event dedup claim writes.
49///
50/// Mirrors the raw column shape rather than the `IntegrationEvent` entity: `status` is the literal
51/// `'received'`, and `raw` binds to the `payload` column — the provider's raw notification TEXT, as the
52/// original write did. Note the dedup grain is (connector_id, business_key), NOT `external_id`:
53/// `external_id` varies per notification, while `business_key` is the stable business-action identity.
54pub struct NewEventRow<'a> {
55    pub id: Uuid,
56    pub connector_id: Uuid,
57    pub event_type: &'a str,
58    pub external_id: &'a str,
59    pub business_key: &'a str,
60    pub raw: &'a str,
61}
62
63/// A redelivered event's current state, as seen by the dedup re-read.
64pub struct EventOutcomeRow {
65    pub id: Uuid,
66    pub status: String,
67    pub mapped_ref_id: Option<Uuid>,
68}
69
70/// One failed event, as the failure report and the retry loop read it. `payload` is the raw JSON TEXT
71/// off the column — the caller parses it, so a malformed payload is the service's call, not a repo panic.
72pub struct FailedEventRow {
73    pub id: Uuid,
74    pub event_type: String,
75    pub external_id: String,
76    pub business_key: String,
77    pub error_detail: Option<String>,
78    pub payload: String,
79}
80
81/// Hand-written IntegrationEvent SQL. Lives here (not in the write service) per the module's 4-layer
82/// rule: services orchestrate and own the unit of work, repositories hold the SQL.
83impl IntegrationEventRepository {
84    /// Claim the (connector, business_key) dedup slot. `Ok(None)` = this business action was already
85    /// received (a webhook retry, or a second notification for the same action) and the caller must
86    /// re-read the original rather than re-map it. A NEW business action does not conflict.
87    ///
88    /// Runs outside any transaction; it rides the request-dedicated connection when the composing
89    /// service bound one — the decorator's org fence then keeps the dedup idempotent within the
90    /// tenant even off the request path — else runs plainly on the pool (a decorated deployment
91    /// refuses it, fail closed; ADR-0029).
92    pub async fn claim_event(
93        &self,
94        pool: &PgPool,
95        e: &NewEventRow<'_>,
96    ) -> Result<Option<Uuid>, sqlx::Error> {
97        company_scope::fetch_optional_scalar_scoped(
98            pool,
99            sqlx::query_scalar(
100                r#"INSERT INTO integrations.integration_events
101                     (id, connector_id, event_type, external_id, business_key, status, payload)
102                   VALUES ($1,$2,$3,$4,$5,'received'::integration_status,$6)
103                   ON CONFLICT (connector_id, business_key) DO NOTHING
104                   RETURNING id"#,
105            )
106            .bind(e.id).bind(e.connector_id).bind(e.event_type)
107            .bind(e.external_id).bind(e.business_key).bind(e.raw),
108        )
109        .await
110    }
111
112    /// Re-read the original after a losing dedup claim — the row is known to exist. Same connection
113    /// discipline as [`Self::claim_event`].
114    pub async fn fetch_by_business_key(
115        &self,
116        pool: &PgPool,
117        connector_id: Uuid,
118        business_key: &str,
119    ) -> Result<EventOutcomeRow, sqlx::Error> {
120        let r = org_scope::fetch_one_row_scoped(
121            pool,
122            sqlx::query(
123                r#"SELECT id, status::text AS status, mapped_ref_id FROM integrations.integration_events
124                   WHERE connector_id=$1 AND business_key=$2"#,
125            )
126            .bind(connector_id).bind(business_key),
127        )
128        .await?;
129        Ok(EventOutcomeRow {
130            id: r.get("id"), status: r.get("status"), mapped_ref_id: r.get("mapped_ref_id"),
131        })
132    }
133
134    /// Record a successful mapping to an internal action. State-guarded on `received`.
135    ///
136    /// Takes the CALLER'S connection so this and the outbox stage commit as one unit. The caller has
137    /// already bound the ambient org scope on it when one is present — don't re-bind here.
138    pub async fn mark_mapped(
139        &self,
140        conn: &mut sqlx::PgConnection,
141        event_id: Uuid,
142        mapped_ref_type: &str,
143        mapped_ref_id: Uuid,
144    ) -> Result<(), sqlx::Error> {
145        sqlx::query(
146            r#"UPDATE integrations.integration_events
147               SET status='mapped'::integration_status, mapped_ref_type=$2, mapped_ref_id=$3
148               WHERE id=$1 AND status='received'::integration_status"#,
149        )
150        .bind(event_id).bind(mapped_ref_type).bind(mapped_ref_id)
151        .execute(conn)
152        .await?;
153        Ok(())
154    }
155
156    /// Record an event the target intentionally ignored. State-guarded on `received`; same
157    /// caller-owned-tx contract as [`Self::mark_mapped`].
158    pub async fn mark_ignored(
159        &self,
160        conn: &mut sqlx::PgConnection,
161        event_id: Uuid,
162        reason: &str,
163    ) -> Result<(), sqlx::Error> {
164        sqlx::query(
165            r#"UPDATE integrations.integration_events SET status='ignored'::integration_status, error_detail=$2
166               WHERE id=$1 AND status='received'::integration_status"#,
167        )
168        .bind(event_id).bind(reason)
169        .execute(conn)
170        .await?;
171        Ok(())
172    }
173
174    /// Record a mapping rejection. State-guarded on `received`; same caller-owned-tx contract as
175    /// [`Self::mark_mapped`].
176    pub async fn mark_failed(
177        &self,
178        conn: &mut sqlx::PgConnection,
179        event_id: Uuid,
180        error_detail: &str,
181    ) -> Result<(), sqlx::Error> {
182        sqlx::query(
183            r#"UPDATE integrations.integration_events SET status='failed'::integration_status, error_detail=$2
184               WHERE id=$1 AND status='received'::integration_status"#,
185        )
186        .bind(event_id).bind(error_detail)
187        .execute(conn)
188        .await?;
189        Ok(())
190    }
191
192    /// A connector's FAILED events — the operator's failure report and the retry loop's work list.
193    ///
194    /// ID-only: the connector id alone identifies the set, so the read rides the request-dedicated
195    /// connection when the composing service bound one — the decorator's org fence then hides
196    /// another tenant's connector — else runs plainly on the pool (ADR-0029).
197    pub async fn fetch_failed(
198        &self,
199        pool: &PgPool,
200        connector_id: Uuid,
201    ) -> Result<Vec<FailedEventRow>, sqlx::Error> {
202        let rows = org_scope::fetch_all_rows_scoped(
203            pool,
204            sqlx::query(
205                r#"SELECT id, event_type, external_id, business_key, error_detail, payload
206                   FROM integrations.integration_events
207                   WHERE connector_id=$1 AND status='failed'::integration_status
208                   ORDER BY (metadata->>'created_at') NULLS FIRST"#,
209            )
210            .bind(connector_id),
211        )
212        .await?;
213        Ok(rows
214            .iter()
215            .map(|r| FailedEventRow {
216                id: r.get("id"), event_type: r.get("event_type"), external_id: r.get("external_id"),
217                business_key: r.get("business_key"), error_detail: r.get("error_detail"),
218                payload: r.get("payload"),
219            })
220            .collect())
221    }
222
223    /// The RETRY path's `failed → mapped` transition, clearing the stale error. Distinct from
224    /// [`Self::mark_mapped`]: state-guarded on `failed`, not `received`. Returns the rows affected so the
225    /// caller only stages the event (and counts it) when the transition really happened.
226    ///
227    /// Takes the CALLER'S connection so the transition and the outbox stage commit as one unit — and so
228    /// the caller can roll back when it loses. The caller has already bound the ambient org scope on
229    /// it when one is present — don't re-bind here.
230    pub async fn retry_mark_mapped(
231        &self,
232        conn: &mut sqlx::PgConnection,
233        event_id: Uuid,
234        mapped_ref_type: &str,
235        mapped_ref_id: Uuid,
236    ) -> Result<u64, sqlx::Error> {
237        let done = sqlx::query(
238            r#"UPDATE integrations.integration_events
239               SET status='mapped'::integration_status, mapped_ref_type=$2, mapped_ref_id=$3, error_detail=NULL
240               WHERE id=$1 AND status='failed'::integration_status"#,
241        )
242        .bind(event_id).bind(mapped_ref_type).bind(mapped_ref_id)
243        .execute(conn)
244        .await?;
245        Ok(done.rows_affected())
246    }
247
248    /// The RETRY path's `failed → ignored` transition (the target now says this event is a no-op).
249    /// State-guarded on `failed`.
250    ///
251    /// Runs outside any transaction; it rides the request-dedicated connection when the composing
252    /// service bound one, else runs plainly on the pool (ADR-0029).
253    pub async fn retry_mark_ignored(
254        &self,
255        pool: &PgPool,
256        event_id: Uuid,
257        reason: &str,
258    ) -> Result<(), sqlx::Error> {
259        company_scope::execute_scoped(
260            pool,
261            sqlx::query(
262                r#"UPDATE integrations.integration_events SET status='ignored'::integration_status, error_detail=$2
263                   WHERE id=$1 AND status='failed'::integration_status"#,
264            )
265            .bind(event_id).bind(reason),
266        )
267        .await?;
268        Ok(())
269    }
270
271    /// Refresh a still-failing event's error after a retry attempt — the status deliberately stays
272    /// `failed` so the next retry picks it up again. Same connection discipline as
273    /// [`Self::retry_mark_ignored`].
274    pub async fn set_error_detail(
275        &self,
276        pool: &PgPool,
277        event_id: Uuid,
278        error_detail: &str,
279    ) -> Result<(), sqlx::Error> {
280        company_scope::execute_scoped(
281            pool,
282            sqlx::query(
283                "UPDATE integrations.integration_events SET error_detail=$2 WHERE id=$1 AND status='failed'::integration_status")
284                .bind(event_id).bind(error_detail),
285        )
286        .await?;
287        Ok(())
288    }
289}
290
291backbone_core::impl_crud_repository!(IntegrationEventRepository, IntegrationEvent, soft_delete);