use anyhow::Result;
use sqlx::{PgPool, Row};
use uuid::Uuid;
use backbone_orm::{company_scope, org_scope};
use crate::domain::entity::IntegrationConnector;
pub const TABLE_NAME: &str = "integrations.integration_connectors";
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 {
pub fn new(pool: PgPool) -> Self {
Self(backbone_orm::GenericCrudRepository::new(pool, TABLE_NAME))
}
}
pub struct NewConnectorRow<'a> {
pub id: Uuid,
pub provider: &'a str,
pub kind: &'a str,
pub direction: &'a str,
}
pub struct ConnectorGateRow {
pub kind: String,
pub status: String,
}
pub struct ConnectorRetryRow {
pub kind: String,
}
impl IntegrationConnectorRepository {
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(())
}
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") }))
}
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);