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 IntegrationConnector 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 IntegrationConnector SQL — the registration and the ID-only reads a
//! webhook's admission uses (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 here 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 + the per-unit provider unique), so an unscoped write fails
//! closed on a decorated deployment.
//!
//! Thin newtype over `backbone_orm::GenericCrudRepository<IntegrationConnector, 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::IntegrationConnector;

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

/// Repository for IntegrationConnector entities.
///
/// All standard CRUD, soft-delete, pagination, and bulk methods are
/// provided automatically via `Deref` to `backbone_orm::GenericCrudRepository`.
pub struct IntegrationConnectorRepository(
    backbone_orm::GenericCrudRepository<IntegrationConnector, backbone_orm::SoftDelete>,
);

impl std::ops::Deref for IntegrationConnectorRepository {
    type Target = backbone_orm::GenericCrudRepository<IntegrationConnector, backbone_orm::SoftDelete>;
    fn deref(&self) -> &Self::Target { &self.0 }
}

impl IntegrationConnectorRepository {
    /// Create a new repository instance.
    pub fn new(pool: PgPool) -> Self {
        Self(backbone_orm::GenericCrudRepository::new(pool, TABLE_NAME))
    }
}

/// The exact row a connector registration writes.
///
/// Mirrors the raw column shape rather than the `IntegrationConnector` entity: `status` is the
/// literal `'active'`, and `kind`/`direction` bind as `&str` with DB-side casts (`$4::connector_kind`,
/// `$5::connector_direction`) so an unknown value fails as a DB error rather than a deserialize panic.
pub struct NewConnectorRow<'a> {
    pub id: Uuid,
    pub provider: &'a str,
    pub kind: &'a str,
    pub direction: &'a str,
}

/// What an inbound event needs to know about its connector before it is admitted.
pub struct ConnectorGateRow {
    pub kind: String,
    pub status: String,
}

/// What a retry run needs about its connector.
pub struct ConnectorRetryRow {
    pub kind: String,
}

/// Hand-written IntegrationConnector 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 IntegrationConnectorRepository {
    /// Register a connector to an external provider.
    ///
    /// A write outside any transaction. Tenancy (ADR-0029): the module invents no scope — the
    /// statement rides the request-dedicated connection when the composing service bound one
    /// (the decorator's org fence and acting-unit default then apply), else runs plainly on
    /// the pool, where an org-scoped deployment refuses it (fail closed).
    ///
    /// Returns the raw `sqlx::Error` deliberately: the caller inspects it for a unique violation to turn
    /// a duplicate provider into a domain error.
    pub async fn insert_connector(
        &self,
        pool: &PgPool,
        c: &NewConnectorRow<'_>,
    ) -> Result<(), sqlx::Error> {
        company_scope::execute_scoped(
            pool,
            sqlx::query(
                r#"INSERT INTO integrations.integration_connectors
                     (id, provider, kind, direction, status)
                   VALUES ($1,$2,$3::connector_kind,$4::connector_direction,'active')"#,
            )
            .bind(c.id).bind(c.provider).bind(c.kind).bind(c.direction),
        )
        .await?;
        Ok(())
    }

    /// Read the connector an inbound event names. `Ok(None)` = no such connector in scope.
    ///
    /// ID-only, outside any transaction: it rides the request-dedicated connection when the
    /// composing service bound one — the decorator's org fence then hides other tenants' rows —
    /// else runs plainly on the pool (a decorated deployment yields nothing unscoped, fail closed).
    pub async fn fetch_gate(
        &self,
        pool: &PgPool,
        connector_id: Uuid,
    ) -> Result<Option<ConnectorGateRow>, sqlx::Error> {
        let row = org_scope::fetch_optional_row_scoped(
            pool,
            sqlx::query(
                r#"SELECT kind::text AS kind, status::text AS status FROM integrations.integration_connectors
                   WHERE id=$1 AND (metadata->>'deleted_at') IS NULL"#,
            )
            .bind(connector_id),
        )
        .await?;
        Ok(row.map(|r| ConnectorGateRow { kind: r.get("kind"), status: r.get("status") }))
    }

    /// Read the connector a retry run re-drives.
    ///
    /// ID-only: identified by the connector id alone, so this read rides the request-dedicated
    /// connection when one is bound.
    pub async fn fetch_for_retry(
        &self,
        pool: &PgPool,
        connector_id: Uuid,
    ) -> Result<Option<ConnectorRetryRow>, sqlx::Error> {
        let row = org_scope::fetch_optional_row_scoped(
            pool,
            sqlx::query("SELECT kind::text AS kind FROM integrations.integration_connectors WHERE id=$1")
                .bind(connector_id),
        )
        .await?;
        Ok(row.map(|r| ConnectorRetryRow { kind: r.get("kind") }))
    }
}

backbone_core::impl_crud_repository!(IntegrationConnectorRepository, IntegrationConnector, soft_delete);