use cm_domain::{AgentId, WorkspaceId}; use sqlx::PgPool; use uuid::Uuid; use crate::DbError; /// A connected app (spec §14 AppConnection). The credential lives in the /// broker's encrypted store; this row only carries the reference. #[derive(Debug, Clone, serde::Serialize, serde::Deserialize)] pub struct AppConnection { pub id: Uuid, pub workspace_id: Uuid, pub agent_id: Option, pub provider: String, pub auth_type: String, pub status: String, pub secret_ref: Option, } pub async fn insert( pool: &PgPool, workspace_id: WorkspaceId, agent_id: Option, provider: &str, auth_type: &str, secret_ref: Uuid, ) -> Result { let id = Uuid::now_v7(); sqlx::query!( "INSERT INTO app_connections (id, workspace_id, agent_id, provider, auth_type, status, secret_ref) VALUES ($1, $2, $3, $4, $5, 'connected', $6)", id, workspace_id.as_uuid(), agent_id.map(|a| a.as_uuid()), provider, auth_type, secret_ref, ) .execute(pool) .await?; Ok(AppConnection { id, workspace_id: workspace_id.as_uuid(), agent_id: agent_id.map(|a| a.as_uuid()), provider: provider.to_owned(), auth_type: auth_type.to_owned(), status: "connected".into(), secret_ref: Some(secret_ref), }) } /// Connections visible to an agent: its own plus workspace-wide ones. pub async fn list_for_agent( pool: &PgPool, workspace_id: WorkspaceId, agent_id: AgentId, ) -> Result, DbError> { let rows = sqlx::query_as!( AppConnection, r#"SELECT id, workspace_id, agent_id, provider, auth_type, status, secret_ref FROM app_connections WHERE workspace_id = $1 AND (agent_id IS NULL OR agent_id = $2) AND status = 'connected' ORDER BY created_at"#, workspace_id.as_uuid(), agent_id.as_uuid(), ) .fetch_all(pool) .await?; Ok(rows) } /// Workspace-wide connections (agent_id IS NULL) — the global /apps page. pub async fn list_for_workspace( pool: &PgPool, workspace_id: WorkspaceId, ) -> Result, DbError> { let rows = sqlx::query_as!( AppConnection, r#"SELECT id, workspace_id, agent_id, provider, auth_type, status, secret_ref FROM app_connections WHERE workspace_id = $1 AND agent_id IS NULL AND status = 'connected' ORDER BY created_at"#, workspace_id.as_uuid(), ) .fetch_all(pool) .await?; Ok(rows) } /// The agent's live connection for one provider, if any. pub async fn find_provider( pool: &PgPool, workspace_id: WorkspaceId, agent_id: AgentId, provider: &str, ) -> Result, DbError> { let row = sqlx::query_as!( AppConnection, r#"SELECT id, workspace_id, agent_id, provider, auth_type, status, secret_ref FROM app_connections WHERE workspace_id = $1 AND (agent_id IS NULL OR agent_id = $2) AND provider = $3 AND status = 'connected' ORDER BY created_at DESC LIMIT 1"#, workspace_id.as_uuid(), agent_id.as_uuid(), provider, ) .fetch_optional(pool) .await?; Ok(row) } /// Fetch a single connection by id, workspace-scoped so cross-tenant reads /// return NotFound. Used when another table (e.g. repo_connections) needs /// to resolve the secret_ref of the credential it points at. pub async fn get( pool: &PgPool, id: Uuid, workspace_id: WorkspaceId, ) -> Result { let row = sqlx::query_as!( AppConnection, r#"SELECT id, workspace_id, agent_id, provider, auth_type, status, secret_ref FROM app_connections WHERE id = $1 AND workspace_id = $2"#, id, workspace_id.as_uuid(), ) .fetch_optional(pool) .await? .ok_or(DbError::NotFound)?; Ok(row) } pub async fn disconnect(pool: &PgPool, id: Uuid) -> Result<(), DbError> { let result = sqlx::query!( "UPDATE app_connections SET status = 'disconnected' WHERE id = $1", id, ) .execute(pool) .await?; if result.rows_affected() == 0 { return Err(DbError::NotFound); } Ok(()) }