1use 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, pub direction: String, }
49
50pub struct InboundEvent {
52 pub company_id: Uuid,
57 pub connector_id: Uuid,
58 pub event_type: String,
59 pub external_id: String,
61 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, pub mapped_ref_id: Option<Uuid>,
73 pub duplicate: bool,
74}
75
76#[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 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 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 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 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 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 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 pub async fn failures(&self, connector_id: Uuid) -> Result<Vec<FailedEvent>, IntegrationError> {
233 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 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 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
314async 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}