use backbone_orm::org_scope;
use chrono::Utc;
use sqlx::PgPool;
use uuid::Uuid;
use crate::infrastructure::persistence::{
IntegrationConnectorRepository, IntegrationEventRepository, NewConnectorRow, NewEventRow,
};
use super::integrations_events::*;
use super::integrations_ports::*;
#[derive(Debug, thiserror::Error)]
pub enum IntegrationError {
#[error("db: {0}")]
Db(#[from] sqlx::Error),
#[error("not found: {0}")]
NotFound(&'static str),
#[error("invalid state: {0}")]
InvalidState(&'static str),
#[error("invalid input: {0}")]
Invalid(String),
#[error("mapping rejected: {0}")]
MappingRejected(String),
}
pub struct NewConnector {
pub provider: String,
pub kind: String, pub direction: String, }
pub struct InboundEvent {
pub company_id: Uuid,
pub connector_id: Uuid,
pub event_type: String,
pub external_id: String,
pub business_key: String,
pub raw: String,
pub payload: serde_json::Value,
}
#[derive(Debug, Clone, PartialEq)]
pub struct ReceiveOutcome {
pub event_id: Uuid,
pub status: String, pub mapped_ref_id: Option<Uuid>,
pub duplicate: bool,
}
#[derive(Debug, Clone, PartialEq)]
pub struct FailedEvent {
pub event_id: Uuid,
pub event_type: String,
pub external_id: String,
pub business_key: String,
pub error_detail: Option<String>,
}
pub struct IntegrationsWriteService {
pool: PgPool,
connectors: IntegrationConnectorRepository,
events: IntegrationEventRepository,
}
impl IntegrationsWriteService {
pub fn new(pool: PgPool) -> Self {
let connectors = IntegrationConnectorRepository::new(pool.clone());
let events = IntegrationEventRepository::new(pool.clone());
Self { pool, connectors, events }
}
pub async fn register_connector(&self, c: NewConnector) -> Result<Uuid, IntegrationError> {
if c.provider.trim().is_empty() {
return Err(IntegrationError::Invalid("connector needs a provider".into()));
}
let id = Uuid::new_v4();
let r = self.connectors.insert_connector(&self.pool, &NewConnectorRow {
id,
provider: &c.provider,
kind: &c.kind,
direction: &c.direction,
})
.await;
match r {
Ok(_) => Ok(id),
Err(e) if e.as_database_error().map(|d| d.is_unique_violation()).unwrap_or(false) =>
Err(IntegrationError::Invalid("a connector for this provider already exists".into())),
Err(e) => Err(e.into()),
}
}
pub async fn receive_event(
&self,
e: InboundEvent,
mapper: &dyn TargetPort,
events: &dyn IntegrationEventSink,
) -> Result<ReceiveOutcome, IntegrationError> {
if e.external_id.trim().is_empty() {
return Err(IntegrationError::Invalid("an inbound event needs an external_id".into()));
}
if e.business_key.trim().is_empty() {
return Err(IntegrationError::Invalid("an inbound event needs a business_key (the order/transaction ref + state)".into()));
}
let conn = self
.connectors
.fetch_gate(&self.pool, e.connector_id)
.await?
.ok_or(IntegrationError::NotFound("connector"))?;
if conn.status != "active" {
return Err(IntegrationError::InvalidState("connector is not active"));
}
let connector_kind = conn.kind;
let inserted = self
.events
.claim_event(&self.pool, &NewEventRow {
id: Uuid::new_v4(),
connector_id: e.connector_id,
event_type: &e.event_type,
external_id: &e.external_id,
business_key: &e.business_key,
raw: &e.raw,
})
.await?;
let Some(event_id) = inserted else {
let row = self
.events
.fetch_by_business_key(&self.pool, e.connector_id, &e.business_key)
.await?;
return Ok(ReceiveOutcome {
event_id: row.id, status: row.status,
mapped_ref_id: row.mapped_ref_id, duplicate: true,
});
};
let req = MapRequest {
company_id: e.company_id, connector_kind, event_type: e.event_type.clone(),
external_id: e.external_id.clone(),
idempotency_key: e.business_key.clone(), payload: e.payload.clone(),
};
match mapper.map(&req).await {
Ok(MapOutcome::Mapped(mref)) => {
let ev = IntegrationEvent::IntegrationEventMapped(IntegrationEventMapped {
event_id, company_id: e.company_id, connector_id: e.connector_id, event_type: e.event_type.clone(),
external_id: e.external_id.clone(),
internal_ref_type: mref.internal_ref_type.clone(), internal_ref_id: mref.internal_ref_id,
});
let mut tx = self.pool.begin().await?;
bind_ambient_org_scope(&mut tx).await?;
self.events
.mark_mapped(&mut tx, event_id, &mref.internal_ref_type, mref.internal_ref_id)
.await?;
stage(&mut tx, &ev).await?;
tx.commit().await?;
events.publish(&ev);
Ok(ReceiveOutcome { event_id, status: "mapped".into(), mapped_ref_id: Some(mref.internal_ref_id), duplicate: false })
}
Ok(MapOutcome::Ignored(reason)) => {
let ev = IntegrationEvent::IntegrationEventIgnored {
event_id, company_id: e.company_id, connector_id: e.connector_id, external_id: e.external_id.clone(), reason: reason.clone(),
};
let mut tx = self.pool.begin().await?;
bind_ambient_org_scope(&mut tx).await?;
self.events.mark_ignored(&mut tx, event_id, &reason).await?;
stage(&mut tx, &ev).await?;
tx.commit().await?;
events.publish(&ev);
Ok(ReceiveOutcome { event_id, status: "ignored".into(), mapped_ref_id: None, duplicate: false })
}
Err(rej) => {
let ev = IntegrationEvent::IntegrationEventFailed {
event_id, company_id: e.company_id, connector_id: e.connector_id,
external_id: e.external_id.clone(), reason: rej.code.clone(),
};
let mut tx = self.pool.begin().await?;
bind_ambient_org_scope(&mut tx).await?;
self.events.mark_failed(&mut tx, event_id, &rej.message).await?;
stage(&mut tx, &ev).await?;
tx.commit().await?;
events.publish(&ev);
Ok(ReceiveOutcome { event_id, status: "failed".into(), mapped_ref_id: None, duplicate: false })
}
}
}
pub async fn failures(&self, connector_id: Uuid) -> Result<Vec<FailedEvent>, IntegrationError> {
let rows = self.events.fetch_failed(&self.pool, connector_id).await?;
Ok(rows.into_iter().map(|r| FailedEvent {
event_id: r.id, event_type: r.event_type, external_id: r.external_id,
business_key: r.business_key, error_detail: r.error_detail,
}).collect())
}
pub async fn retry_failed(
&self,
company_id: Uuid,
connector_id: Uuid,
mapper: &dyn TargetPort,
events: &dyn IntegrationEventSink,
) -> Result<usize, IntegrationError> {
let conn = self
.connectors
.fetch_for_retry(&self.pool, connector_id)
.await?
.ok_or(IntegrationError::NotFound("connector"))?;
let connector_kind = conn.kind;
let rows = self.events.fetch_failed(&self.pool, connector_id).await?;
let mut mapped = 0usize;
for row in &rows {
let event_id = row.id;
let business_key = row.business_key.clone();
let req = MapRequest {
company_id, connector_kind: connector_kind.clone(), event_type: row.event_type.clone(),
external_id: row.external_id.clone(), idempotency_key: business_key.clone(),
payload: serde_json::from_str(&row.payload).unwrap_or(serde_json::Value::Null),
};
match mapper.map(&req).await {
Ok(MapOutcome::Mapped(mref)) => {
let ev = IntegrationEvent::IntegrationEventMapped(IntegrationEventMapped {
event_id, company_id, connector_id, event_type: row.event_type.clone(),
external_id: row.external_id.clone(),
internal_ref_type: mref.internal_ref_type.clone(), internal_ref_id: mref.internal_ref_id,
});
let mut tx = self.pool.begin().await?;
bind_ambient_org_scope(&mut tx).await?;
let m = self
.events
.retry_mark_mapped(&mut tx, event_id, &mref.internal_ref_type, mref.internal_ref_id)
.await?;
if m == 1 {
stage(&mut tx, &ev).await?;
tx.commit().await?;
events.publish(&ev);
mapped += 1;
} else {
tx.rollback().await?;
}
}
Ok(MapOutcome::Ignored(reason)) => {
self.events.retry_mark_ignored(&self.pool, event_id, &reason).await?;
}
Err(rej) => {
self.events.set_error_detail(&self.pool, event_id, &rej.message).await?;
}
}
}
Ok(mapped)
}
}
async fn bind_ambient_org_scope(tx: &mut sqlx::Transaction<'_, sqlx::Postgres>) -> Result<(), IntegrationError> {
if let Some(scope) = org_scope::current_org_scope() {
org_scope::bind_org_scope_on(&mut **tx, &scope).await?;
}
Ok(())
}
async fn stage(tx: &mut sqlx::Transaction<'_, sqlx::Postgres>, event: &IntegrationEvent) -> Result<(), IntegrationError> {
let (etype, agg_id) = match event {
IntegrationEvent::IntegrationEventMapped(m) => ("IntegrationEventMapped", m.event_id),
IntegrationEvent::IntegrationEventFailed { event_id, .. } => ("IntegrationEventFailed", *event_id),
IntegrationEvent::IntegrationEventIgnored { event_id, .. } => ("IntegrationEventIgnored", *event_id),
};
let payload = serde_json::to_value(event).map_err(|e| IntegrationError::Invalid(e.to_string()))?;
let company_id: Uuid = payload
.get("company_id")
.and_then(|v| v.as_str())
.ok_or_else(|| IntegrationError::Invalid("integration event missing company_id".into()))?
.parse()
.map_err(|e| IntegrationError::Invalid(format!("company_id parse: {e}")))?;
let record = backbone_outbox::OutboxRecord::new(
etype, "IntegrationEvent", agg_id.to_string(), company_id, payload, Utc::now(),
);
backbone_outbox::outbox::stage(&mut **tx, "integrations", &record)
.await.map_err(|e| IntegrationError::Invalid(format!("outbox stage: {e}")))?;
Ok(())
}