use anyhow::Result;
use sqlx::{PgPool, Row};
use uuid::Uuid;
use backbone_orm::{company_scope, org_scope};
use crate::domain::entity::IntegrationEvent;
pub const TABLE_NAME: &str = "integrations.integration_events";
pub struct IntegrationEventRepository(
backbone_orm::GenericCrudRepository<IntegrationEvent, backbone_orm::SoftDelete>,
);
impl std::ops::Deref for IntegrationEventRepository {
type Target = backbone_orm::GenericCrudRepository<IntegrationEvent, backbone_orm::SoftDelete>;
fn deref(&self) -> &Self::Target { &self.0 }
}
impl IntegrationEventRepository {
pub fn new(pool: PgPool) -> Self {
Self(backbone_orm::GenericCrudRepository::new(pool, TABLE_NAME))
}
}
pub struct NewEventRow<'a> {
pub id: Uuid,
pub connector_id: Uuid,
pub event_type: &'a str,
pub external_id: &'a str,
pub business_key: &'a str,
pub raw: &'a str,
}
pub struct EventOutcomeRow {
pub id: Uuid,
pub status: String,
pub mapped_ref_id: Option<Uuid>,
}
pub struct FailedEventRow {
pub id: Uuid,
pub event_type: String,
pub external_id: String,
pub business_key: String,
pub error_detail: Option<String>,
pub payload: String,
}
impl IntegrationEventRepository {
pub async fn claim_event(
&self,
pool: &PgPool,
e: &NewEventRow<'_>,
) -> Result<Option<Uuid>, sqlx::Error> {
company_scope::fetch_optional_scalar_scoped(
pool,
sqlx::query_scalar(
r#"INSERT INTO integrations.integration_events
(id, connector_id, event_type, external_id, business_key, status, payload)
VALUES ($1,$2,$3,$4,$5,'received'::integration_status,$6)
ON CONFLICT (connector_id, business_key) DO NOTHING
RETURNING id"#,
)
.bind(e.id).bind(e.connector_id).bind(e.event_type)
.bind(e.external_id).bind(e.business_key).bind(e.raw),
)
.await
}
pub async fn fetch_by_business_key(
&self,
pool: &PgPool,
connector_id: Uuid,
business_key: &str,
) -> Result<EventOutcomeRow, sqlx::Error> {
let r = org_scope::fetch_one_row_scoped(
pool,
sqlx::query(
r#"SELECT id, status::text AS status, mapped_ref_id FROM integrations.integration_events
WHERE connector_id=$1 AND business_key=$2"#,
)
.bind(connector_id).bind(business_key),
)
.await?;
Ok(EventOutcomeRow {
id: r.get("id"), status: r.get("status"), mapped_ref_id: r.get("mapped_ref_id"),
})
}
pub async fn mark_mapped(
&self,
conn: &mut sqlx::PgConnection,
event_id: Uuid,
mapped_ref_type: &str,
mapped_ref_id: Uuid,
) -> Result<(), sqlx::Error> {
sqlx::query(
r#"UPDATE integrations.integration_events
SET status='mapped'::integration_status, mapped_ref_type=$2, mapped_ref_id=$3
WHERE id=$1 AND status='received'::integration_status"#,
)
.bind(event_id).bind(mapped_ref_type).bind(mapped_ref_id)
.execute(conn)
.await?;
Ok(())
}
pub async fn mark_ignored(
&self,
conn: &mut sqlx::PgConnection,
event_id: Uuid,
reason: &str,
) -> Result<(), sqlx::Error> {
sqlx::query(
r#"UPDATE integrations.integration_events SET status='ignored'::integration_status, error_detail=$2
WHERE id=$1 AND status='received'::integration_status"#,
)
.bind(event_id).bind(reason)
.execute(conn)
.await?;
Ok(())
}
pub async fn mark_failed(
&self,
conn: &mut sqlx::PgConnection,
event_id: Uuid,
error_detail: &str,
) -> Result<(), sqlx::Error> {
sqlx::query(
r#"UPDATE integrations.integration_events SET status='failed'::integration_status, error_detail=$2
WHERE id=$1 AND status='received'::integration_status"#,
)
.bind(event_id).bind(error_detail)
.execute(conn)
.await?;
Ok(())
}
pub async fn fetch_failed(
&self,
pool: &PgPool,
connector_id: Uuid,
) -> Result<Vec<FailedEventRow>, sqlx::Error> {
let rows = org_scope::fetch_all_rows_scoped(
pool,
sqlx::query(
r#"SELECT id, event_type, external_id, business_key, error_detail, payload
FROM integrations.integration_events
WHERE connector_id=$1 AND status='failed'::integration_status
ORDER BY (metadata->>'created_at') NULLS FIRST"#,
)
.bind(connector_id),
)
.await?;
Ok(rows
.iter()
.map(|r| FailedEventRow {
id: r.get("id"), event_type: r.get("event_type"), external_id: r.get("external_id"),
business_key: r.get("business_key"), error_detail: r.get("error_detail"),
payload: r.get("payload"),
})
.collect())
}
pub async fn retry_mark_mapped(
&self,
conn: &mut sqlx::PgConnection,
event_id: Uuid,
mapped_ref_type: &str,
mapped_ref_id: Uuid,
) -> Result<u64, sqlx::Error> {
let done = sqlx::query(
r#"UPDATE integrations.integration_events
SET status='mapped'::integration_status, mapped_ref_type=$2, mapped_ref_id=$3, error_detail=NULL
WHERE id=$1 AND status='failed'::integration_status"#,
)
.bind(event_id).bind(mapped_ref_type).bind(mapped_ref_id)
.execute(conn)
.await?;
Ok(done.rows_affected())
}
pub async fn retry_mark_ignored(
&self,
pool: &PgPool,
event_id: Uuid,
reason: &str,
) -> Result<(), sqlx::Error> {
company_scope::execute_scoped(
pool,
sqlx::query(
r#"UPDATE integrations.integration_events SET status='ignored'::integration_status, error_detail=$2
WHERE id=$1 AND status='failed'::integration_status"#,
)
.bind(event_id).bind(reason),
)
.await?;
Ok(())
}
pub async fn set_error_detail(
&self,
pool: &PgPool,
event_id: Uuid,
error_detail: &str,
) -> Result<(), sqlx::Error> {
company_scope::execute_scoped(
pool,
sqlx::query(
"UPDATE integrations.integration_events SET error_detail=$2 WHERE id=$1 AND status='failed'::integration_status")
.bind(event_id).bind(error_detail),
)
.await?;
Ok(())
}
}
backbone_core::impl_crud_repository!(IntegrationEventRepository, IntegrationEvent, soft_delete);