Skip to content
77 linesCodeBlameRaw
1/**
2 * The one-time move of agents' own teams to team memberships
3 * (migrations/0010_agent_team_memberships.sql). An agent is on teams the
4 * way a person is, through identity's team_agents; agents made before
5 * that named a team themselves, in `agents.team`.
6 *
7 * Every cron run takes the agents not moved yet (none, once it is done:
8 * a partial index makes that one cheap read) and asks identity to put
9 * each on the team its old value names, by slug in any case, in its own
10 * workspace. A value that names no team is dropped and logged. Archived
11 * and personal agents are marked without joining anything: neither is on
12 * a team. Each agent is marked `team_moved_at` once identity answered for
13 * its workspace, so a second run does nothing, and identity's insert is
14 * idempotent besides.
15 */
16import type { AdoptedAgentTeam } from "@g1t/contracts";
17
18/** Agents read per run: the move finishes over a few runs in a large installation. */
19export const MOVE_BATCH = 200;
20
21export type MoveRow = { id: string; workspace_id: string; team: string; scope: string | null; archived_at: string | null };
22
23/** What identity is asked, per workspace: the agents that may join, by the team each named. */
24export function claimsByWorkspace(rows: readonly MoveRow[]): Map<string, { agent_id: string; team: string }[]> {
25 const out = new Map<string, { agent_id: string; team: string }[]>();
26 for (const row of rows) {
27 if (row.archived_at || row.scope === "personal") continue;
28 const team = row.team.trim();
29 if (!team) continue;
30 out.set(row.workspace_id, [...(out.get(row.workspace_id) ?? []), { agent_id: row.id, team }]);
31 }
32 return out;
33}
34
35export type MoveReport = { seen: number; added: number; already: number; dropped: { agent_id: string; team: string }[] };
36
37/** `adopt` is identity's `adopt_agent_teams` (index.ts passes it). */
38export async function moveAgentTeams(
39 db: D1Database,
40 adopt: (workspaceId: string, agents: { agent_id: string; team: string }[]) => Promise<AdoptedAgentTeam[]>,
41 now: Date = new Date(),
42): Promise<MoveReport> {
43 const rows = await db.prepare(
44 `SELECT id, workspace_id, team, scope, archived_at FROM agents
45 WHERE team IS NOT NULL AND team <> '' AND team_moved_at IS NULL LIMIT ?`,
46 )
47 .bind(MOVE_BATCH)
48 .all<MoveRow>();
49 const report: MoveReport = { seen: 0, added: 0, already: 0, dropped: [] };
50 if (!rows.results.length) return report;
51 const claims = claimsByWorkspace(rows.results);
52 // Agents identity is not asked about (archived, personal) are done as they are.
53 const done = new Set(rows.results.filter((row) => !claims.get(row.workspace_id)?.some((c) => c.agent_id === row.id)).map((row) => row.id));
54 for (const [workspaceId, agents] of claims) {
55 try {
56 const outcomes = await adopt(workspaceId, agents);
57 for (const outcome of outcomes) {
58 if (outcome.outcome === "added") report.added += 1;
59 else if (outcome.outcome === "already") report.already += 1;
60 else report.dropped.push({ agent_id: outcome.agent_id, team: outcome.team });
61 }
62 for (const agent of agents) done.add(agent.agent_id);
63 } catch (error) {
64 // Left unmarked: the next run asks again.
65 console.error("agents: moving agents' teams to memberships failed for a workspace", workspaceId, String(error));
66 }
67 }
68 if (done.size) {
69 await db.prepare("UPDATE agents SET team_moved_at = ? WHERE id IN (SELECT value FROM json_each(?)) AND team_moved_at IS NULL")
70 .bind(now.toISOString(), JSON.stringify([...done]))
71 .run();
72 }
73 report.seen = done.size;
74 for (const drop of report.dropped) console.log("agents: an agent's old team names no team, so it joins none", drop.agent_id, drop.team);
75 if (report.seen) console.log("agents: moved agents' teams to memberships", JSON.stringify({ seen: report.seen, added: report.added, already: report.already, dropped: report.dropped.length }));
76 return report;
77}