//! Persistence for deployed teams — a baseline topology staffed with real //! workspace claws. `teams` holds the topology graph; `team_members` is the //! durable node→claw binding. use cm_domain::WorkspaceId; use serde_json::Value; use sqlx::PgPool; use time::OffsetDateTime; use uuid::Uuid; use crate::DbError; /// A team row summary (list view). pub struct TeamSummary { pub id: Uuid, pub name: String, pub kind: String, pub status: String, pub created_at: OffsetDateTime, } /// A full team (graph + metadata). pub struct Team { pub id: Uuid, pub name: String, pub kind: String, pub graph: Value, pub status: String, pub created_at: OffsetDateTime, } /// A node→claw binding within a team. pub struct TeamMember { pub node_id: String, pub claw_id: Uuid, pub role: String, } /// Insert a team (the topology graph). Members are added separately. /// Lifecycle defaults to `permanent`; use `insert_ephemeral_team` for the /// scheduled / triggered flows. pub async fn insert_team( pool: &PgPool, id: Uuid, workspace_id: WorkspaceId, name: &str, kind: &str, graph: &Value, ) -> Result<(), DbError> { insert_team_with_lifecycle(pool, id, workspace_id, name, kind, graph, "permanent").await } /// Insert a team with an explicit `lifecycle` (`permanent` | `ephemeral`). /// Ephemeral teams are torn down after the last in-flight run terminates — /// see cm-api::topology_worker::maybe_teardown_ephemeral_team. pub async fn insert_team_with_lifecycle( pool: &PgPool, id: Uuid, workspace_id: WorkspaceId, name: &str, kind: &str, graph: &Value, lifecycle: &str, ) -> Result<(), DbError> { sqlx::query!( "INSERT INTO teams (id, workspace_id, name, kind, graph, lifecycle) VALUES ($1, $2, $3, $4, $5, $6)", id, workspace_id.as_uuid(), name, kind, graph, lifecycle, ) .execute(pool) .await?; Ok(()) } /// Rebuild a team's topology (kind + graph) in place; node→claw bindings (which /// key off stable node ids `n0..`) are left untouched. Non-macro query so it /// needs no offline sqlx cache entry. pub async fn set_topology( pool: &PgPool, id: Uuid, workspace_id: WorkspaceId, kind: &str, graph: &Value, ) -> Result<(), DbError> { let res = sqlx::query("UPDATE teams SET kind = $3, graph = $4 WHERE id = $1 AND workspace_id = $2") .bind(id) .bind(workspace_id.as_uuid()) .bind(kind) .bind(graph) .execute(pool) .await?; if res.rows_affected() == 0 { return Err(DbError::NotFound); } Ok(()) } /// Delete a team and its node→claw bindings. /// Rename a team in-place. Only the display name changes; the topology /// graph + member bindings are untouched. pub async fn rename_team( pool: &PgPool, id: Uuid, workspace_id: WorkspaceId, name: &str, ) -> Result<(), DbError> { let res = sqlx::query("UPDATE teams SET name = $3 WHERE id = $1 AND workspace_id = $2") .bind(id) .bind(workspace_id.as_uuid()) .bind(name) .execute(pool) .await?; if res.rows_affected() == 0 { return Err(DbError::NotFound); } Ok(()) } pub async fn delete_team( pool: &PgPool, id: Uuid, workspace_id: WorkspaceId, ) -> Result<(), DbError> { sqlx::query("DELETE FROM team_members WHERE team_id = $1") .bind(id) .execute(pool) .await?; let res = sqlx::query("DELETE FROM teams WHERE id = $1 AND workspace_id = $2") .bind(id) .bind(workspace_id.as_uuid()) .execute(pool) .await?; if res.rows_affected() == 0 { return Err(DbError::NotFound); } Ok(()) } /// Bind a claw to a topology node within a team. pub async fn add_member( pool: &PgPool, team_id: Uuid, node_id: &str, claw_id: Uuid, role: &str, ) -> Result<(), DbError> { sqlx::query!( "INSERT INTO team_members (team_id, node_id, claw_id, role) VALUES ($1, $2, $3, $4)", team_id, node_id, claw_id, role, ) .execute(pool) .await?; Ok(()) } /// The most recent teams for a workspace, newest first. pub async fn list_for_workspace( pool: &PgPool, workspace_id: WorkspaceId, limit: i64, ) -> Result, DbError> { let rows = sqlx::query!( "SELECT id, name, kind, status, created_at FROM teams WHERE workspace_id = $1 ORDER BY created_at DESC LIMIT $2", workspace_id.as_uuid(), limit, ) .fetch_all(pool) .await?; Ok(rows .into_iter() .map(|r| TeamSummary { id: r.id, name: r.name, kind: r.kind, status: r.status, created_at: r.created_at, }) .collect()) } /// A single team, workspace-scoped. pub async fn get_team(pool: &PgPool, id: Uuid, workspace_id: WorkspaceId) -> Result { let row = sqlx::query!( "SELECT id, name, kind, graph, status, created_at FROM teams WHERE id = $1 AND workspace_id = $2", id, workspace_id.as_uuid(), ) .fetch_optional(pool) .await? .ok_or(DbError::NotFound)?; Ok(Team { id: row.id, name: row.name, kind: row.kind, graph: row.graph, status: row.status, created_at: row.created_at, }) } /// Distinct agent ids bound to a team via `team_members`. Used by the /// cascade-reap path so deleting a team also purges the agents inside it. pub async fn agents_of_team(pool: &PgPool, team_id: Uuid) -> Result, DbError> { let rows = sqlx::query!( "SELECT DISTINCT claw_id FROM team_members WHERE team_id = $1", team_id, ) .fetch_all(pool) .await?; Ok(rows.into_iter().map(|r| r.claw_id).collect()) } /// The node→claw bindings for a team. pub async fn members_for_team(pool: &PgPool, team_id: Uuid) -> Result, DbError> { let rows = sqlx::query!( "SELECT node_id, claw_id, role FROM team_members WHERE team_id = $1 ORDER BY node_id", team_id, ) .fetch_all(pool) .await?; Ok(rows .into_iter() .map(|r| TeamMember { node_id: r.node_id, claw_id: r.claw_id, role: r.role, }) .collect()) }