Skip to main content

backbone_integrations/infrastructure/persistence/
integration_connector_repository.rs

1//! Repository for IntegrationConnector entities
2//!
3//! Originally generated by metaphor-schema; now **user-owned** — this exact path is declared under
4//! `user_owned` in `metaphor.codegen.yaml`, so the generator skips it wholesale. The custom methods
5//! below hold the hand-written IntegrationConnector SQL — the registration and the ID-only reads a
6//! webhook's admission uses (4-layer rule: services orchestrate, repos hold SQL).
7//!
8//! Tenancy (ADR-0029): the table carries no company column and the module invents no scope —
9//! every statement here runs on the request-dedicated connection when the composing service
10//! bound one, else plainly on the pool; the COMPOSING service's tenancy decorator owns org
11//! scoping (org_unit_id + RLS + the per-unit provider unique), so an unscoped write fails
12//! closed on a decorated deployment.
13//!
14//! Thin newtype over `backbone_orm::GenericCrudRepository<IntegrationConnector, backbone_orm::SoftDelete>`.
15//! All standard CRUD methods are available via `Deref`.
16
17use anyhow::Result;
18use sqlx::{PgPool, Row};
19use uuid::Uuid;
20
21use backbone_orm::{company_scope, org_scope};
22
23use crate::domain::entity::IntegrationConnector;
24
25/// Table name for IntegrationConnector entities
26pub const TABLE_NAME: &str = "integrations.integration_connectors";
27
28/// Repository for IntegrationConnector entities.
29///
30/// All standard CRUD, soft-delete, pagination, and bulk methods are
31/// provided automatically via `Deref` to `backbone_orm::GenericCrudRepository`.
32pub struct IntegrationConnectorRepository(
33    backbone_orm::GenericCrudRepository<IntegrationConnector, backbone_orm::SoftDelete>,
34);
35
36impl std::ops::Deref for IntegrationConnectorRepository {
37    type Target = backbone_orm::GenericCrudRepository<IntegrationConnector, backbone_orm::SoftDelete>;
38    fn deref(&self) -> &Self::Target { &self.0 }
39}
40
41impl IntegrationConnectorRepository {
42    /// Create a new repository instance.
43    pub fn new(pool: PgPool) -> Self {
44        Self(backbone_orm::GenericCrudRepository::new(pool, TABLE_NAME))
45    }
46}
47
48/// The exact row a connector registration writes.
49///
50/// Mirrors the raw column shape rather than the `IntegrationConnector` entity: `status` is the
51/// literal `'active'`, and `kind`/`direction` bind as `&str` with DB-side casts (`$4::connector_kind`,
52/// `$5::connector_direction`) so an unknown value fails as a DB error rather than a deserialize panic.
53pub struct NewConnectorRow<'a> {
54    pub id: Uuid,
55    pub provider: &'a str,
56    pub kind: &'a str,
57    pub direction: &'a str,
58}
59
60/// What an inbound event needs to know about its connector before it is admitted.
61pub struct ConnectorGateRow {
62    pub kind: String,
63    pub status: String,
64}
65
66/// What a retry run needs about its connector.
67pub struct ConnectorRetryRow {
68    pub kind: String,
69}
70
71/// Hand-written IntegrationConnector SQL. Lives here (not in the write service) per the module's
72/// 4-layer rule: services orchestrate and own the unit of work, repositories hold the SQL.
73impl IntegrationConnectorRepository {
74    /// Register a connector to an external provider.
75    ///
76    /// A write outside any transaction. Tenancy (ADR-0029): the module invents no scope — the
77    /// statement rides the request-dedicated connection when the composing service bound one
78    /// (the decorator's org fence and acting-unit default then apply), else runs plainly on
79    /// the pool, where an org-scoped deployment refuses it (fail closed).
80    ///
81    /// Returns the raw `sqlx::Error` deliberately: the caller inspects it for a unique violation to turn
82    /// a duplicate provider into a domain error.
83    pub async fn insert_connector(
84        &self,
85        pool: &PgPool,
86        c: &NewConnectorRow<'_>,
87    ) -> Result<(), sqlx::Error> {
88        company_scope::execute_scoped(
89            pool,
90            sqlx::query(
91                r#"INSERT INTO integrations.integration_connectors
92                     (id, provider, kind, direction, status)
93                   VALUES ($1,$2,$3::connector_kind,$4::connector_direction,'active')"#,
94            )
95            .bind(c.id).bind(c.provider).bind(c.kind).bind(c.direction),
96        )
97        .await?;
98        Ok(())
99    }
100
101    /// Read the connector an inbound event names. `Ok(None)` = no such connector in scope.
102    ///
103    /// ID-only, outside any transaction: it rides the request-dedicated connection when the
104    /// composing service bound one — the decorator's org fence then hides other tenants' rows —
105    /// else runs plainly on the pool (a decorated deployment yields nothing unscoped, fail closed).
106    pub async fn fetch_gate(
107        &self,
108        pool: &PgPool,
109        connector_id: Uuid,
110    ) -> Result<Option<ConnectorGateRow>, sqlx::Error> {
111        let row = org_scope::fetch_optional_row_scoped(
112            pool,
113            sqlx::query(
114                r#"SELECT kind::text AS kind, status::text AS status FROM integrations.integration_connectors
115                   WHERE id=$1 AND (metadata->>'deleted_at') IS NULL"#,
116            )
117            .bind(connector_id),
118        )
119        .await?;
120        Ok(row.map(|r| ConnectorGateRow { kind: r.get("kind"), status: r.get("status") }))
121    }
122
123    /// Read the connector a retry run re-drives.
124    ///
125    /// ID-only: identified by the connector id alone, so this read rides the request-dedicated
126    /// connection when one is bound.
127    pub async fn fetch_for_retry(
128        &self,
129        pool: &PgPool,
130        connector_id: Uuid,
131    ) -> Result<Option<ConnectorRetryRow>, sqlx::Error> {
132        let row = org_scope::fetch_optional_row_scoped(
133            pool,
134            sqlx::query("SELECT kind::text AS kind FROM integrations.integration_connectors WHERE id=$1")
135                .bind(connector_id),
136        )
137        .await?;
138        Ok(row.map(|r| ConnectorRetryRow { kind: r.get("kind") }))
139    }
140}
141
142backbone_core::impl_crud_repository!(IntegrationConnectorRepository, IntegrationConnector, soft_delete);