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);