backbone-integrations 0.6.0

Integration registry: connectors, integration accounts and an idempotent inbound event lane, with one OAuth flow (HMAC-bound state, PKCE)
Documentation
//! Repository for IntegrationEvent entities
//!
//! Originally generated by metaphor-schema; now **user-owned** — this exact path is declared under
//! `user_owned` in `metaphor.codegen.yaml`, so the generator skips it wholesale. The custom methods
//! below hold the hand-written IntegrationEvent SQL — the (connector, business_key) dedup claim that
//! stops a webhook retry re-applying a payment, its re-read, and the mapped/ignored/failed transitions
//! on both the receive and the retry path (4-layer rule: services orchestrate, repos hold SQL).
//!
//! Tenancy (ADR-0029): the table carries no company column and the module invents no scope —
//! every statement runs on the request-dedicated connection when the composing service bound
//! one, else plainly on the pool; the COMPOSING service's tenancy decorator owns org scoping
//! (org_unit_id + RLS), so an unscoped write fails closed on a decorated deployment.
//!
//! Thin newtype over `backbone_orm::GenericCrudRepository<IntegrationEvent, backbone_orm::SoftDelete>`.
//! All standard CRUD methods are available via `Deref`.

use anyhow::Result;
use sqlx::{PgPool, Row};
use uuid::Uuid;

use backbone_orm::{company_scope, org_scope};

use crate::domain::entity::IntegrationEvent;

/// Table name for IntegrationEvent entities
pub const TABLE_NAME: &str = "integrations.integration_events";

/// Repository for IntegrationEvent entities.
///
/// All standard CRUD, soft-delete, pagination, and bulk methods are
/// provided automatically via `Deref` to `backbone_orm::GenericCrudRepository`.
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 {
    /// Create a new repository instance.
    pub fn new(pool: PgPool) -> Self {
        Self(backbone_orm::GenericCrudRepository::new(pool, TABLE_NAME))
    }
}

/// The exact row an inbound-event dedup claim writes.
///
/// Mirrors the raw column shape rather than the `IntegrationEvent` entity: `status` is the literal
/// `'received'`, and `raw` binds to the `payload` column — the provider's raw notification TEXT, as the
/// original write did. Note the dedup grain is (connector_id, business_key), NOT `external_id`:
/// `external_id` varies per notification, while `business_key` is the stable business-action identity.
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,
}

/// A redelivered event's current state, as seen by the dedup re-read.
pub struct EventOutcomeRow {
    pub id: Uuid,
    pub status: String,
    pub mapped_ref_id: Option<Uuid>,
}

/// One failed event, as the failure report and the retry loop read it. `payload` is the raw JSON TEXT
/// off the column — the caller parses it, so a malformed payload is the service's call, not a repo panic.
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,
}

/// Hand-written IntegrationEvent SQL. Lives here (not in the write service) per the module's 4-layer
/// rule: services orchestrate and own the unit of work, repositories hold the SQL.
impl IntegrationEventRepository {
    /// Claim the (connector, business_key) dedup slot. `Ok(None)` = this business action was already
    /// received (a webhook retry, or a second notification for the same action) and the caller must
    /// re-read the original rather than re-map it. A NEW business action does not conflict.
    ///
    /// Runs outside any transaction; it rides the request-dedicated connection when the composing
    /// service bound one — the decorator's org fence then keeps the dedup idempotent within the
    /// tenant even off the request path — else runs plainly on the pool (a decorated deployment
    /// refuses it, fail closed; ADR-0029).
    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
    }

    /// Re-read the original after a losing dedup claim — the row is known to exist. Same connection
    /// discipline as [`Self::claim_event`].
    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"),
        })
    }

    /// Record a successful mapping to an internal action. State-guarded on `received`.
    ///
    /// Takes the CALLER'S connection so this and the outbox stage commit as one unit. The caller has
    /// already bound the ambient org scope on it when one is present — don't re-bind here.
    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(())
    }

    /// Record an event the target intentionally ignored. State-guarded on `received`; same
    /// caller-owned-tx contract as [`Self::mark_mapped`].
    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(())
    }

    /// Record a mapping rejection. State-guarded on `received`; same caller-owned-tx contract as
    /// [`Self::mark_mapped`].
    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(())
    }

    /// A connector's FAILED events — the operator's failure report and the retry loop's work list.
    ///
    /// ID-only: the connector id alone identifies the set, so the read rides the request-dedicated
    /// connection when the composing service bound one — the decorator's org fence then hides
    /// another tenant's connector — else runs plainly on the pool (ADR-0029).
    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())
    }

    /// The RETRY path's `failed → mapped` transition, clearing the stale error. Distinct from
    /// [`Self::mark_mapped`]: state-guarded on `failed`, not `received`. Returns the rows affected so the
    /// caller only stages the event (and counts it) when the transition really happened.
    ///
    /// Takes the CALLER'S connection so the transition and the outbox stage commit as one unit — and so
    /// the caller can roll back when it loses. The caller has already bound the ambient org scope on
    /// it when one is present — don't re-bind here.
    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())
    }

    /// The RETRY path's `failed → ignored` transition (the target now says this event is a no-op).
    /// State-guarded on `failed`.
    ///
    /// Runs outside any transaction; it rides the request-dedicated connection when the composing
    /// service bound one, else runs plainly on the pool (ADR-0029).
    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(())
    }

    /// Refresh a still-failing event's error after a retry attempt — the status deliberately stays
    /// `failed` so the next retry picks it up again. Same connection discipline as
    /// [`Self::retry_mark_ignored`].
    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);