diff --git a/src/features/team-mode/team-runtime/status.test.ts b/src/features/team-mode/team-runtime/status.test.ts new file mode 100644 index 000000000..1107aff33 --- /dev/null +++ b/src/features/team-mode/team-runtime/status.test.ts @@ -0,0 +1,143 @@ +import { afterEach, describe, expect, test } from "bun:test" +import { randomUUID } from "node:crypto" +import { mkdir, mkdtemp, rm, writeFile } from "node:fs/promises" +import { tmpdir } from "node:os" +import path from "node:path" + +import type { BackgroundManager } from "../../background-agent/manager" +import { TeamModeConfigSchema } from "../../../config/schema/team-mode" +import type { TeamModeConfig } from "../../../config/schema/team-mode" +import { createTask } from "../team-tasklist/store" +import { createTaskInput } from "../team-tasklist/test-support" +import { getInboxDir, getTasksDir, resolveBaseDir } from "../team-registry/paths" +import { createRuntimeState, saveRuntimeState } from "../team-state-store/store" +import { aggregateStatus } from "./status" + +async function createTemporaryBaseDir(): Promise { + return await mkdtemp(path.join(tmpdir(), "team-mode-status-")) +} + +function createConfig(baseDir: string): TeamModeConfig { + return TeamModeConfigSchema.parse({ base_dir: baseDir, enabled: true }) +} + +async function seedRuntimeState(baseDir: string, teamName: string, leadSessionId: string, memberSessionIds: string[]): Promise { + const config = createConfig(baseDir) + const runtimeState = await createRuntimeState( + { + version: 1, + name: teamName, + createdAt: Date.now(), + leadAgentId: "lead", + members: [ + { kind: "subagent_type", name: "lead", subagent_type: "sisyphus", backendType: "in-process", isActive: true, color: "red" }, + ...memberSessionIds.map((sessionID, index) => ({ + kind: "category" as const, + name: `member-${index + 1}`, + category: "deep" as const, + prompt: "implement task", + backendType: "in-process" as const, + isActive: true, + color: index === 0 ? "blue" : "green", + })), + ], + }, + leadSessionId, + "project", + config, + ) + const updatedRuntimeState = { + ...runtimeState, + members: runtimeState.members.map((member, index) => index === 0 ? { ...member, sessionId: leadSessionId, status: "running" as const } : { ...member, sessionId: memberSessionIds[index - 1], status: "running" as const }), + } + await saveRuntimeState(updatedRuntimeState, config) + return updatedRuntimeState.teamRunId +} + +describe("aggregateStatus", () => { + const temporaryDirectories: string[] = [] + + afterEach(async () => { + await Promise.all(temporaryDirectories.splice(0).map(async (directoryPath) => rm(directoryPath, { recursive: true, force: true }))) + }) + + test("surfaces stale locks from claims directory", async () => { + // given + const baseDir = await createTemporaryBaseDir() + temporaryDirectories.push(baseDir) + const config = createConfig(baseDir) + const teamRunId = await seedRuntimeState(baseDir, "team-gamma", "lead-3", []) + const claimsDir = path.join(getTasksDir(resolveBaseDir(config), teamRunId), "claims") + await mkdir(claimsDir, { recursive: true }) + const claimedTask = await createTask(teamRunId, createTaskInput(), config) + await writeFile(path.join(claimsDir, `${claimedTask.id}.lock`), "owner\n999999\n1\n") + + // when + const result = await aggregateStatus(teamRunId, config) + + // then + expect(result.staleLocks).toEqual([path.join(claimsDir, `${claimedTask.id}.lock`)]) + }) + + test("aggregates members plus tasks plus unread counts", async () => { + // given + const baseDir = await createTemporaryBaseDir() + temporaryDirectories.push(baseDir) + const config = createConfig(baseDir) + const teamRunId = await seedRuntimeState(baseDir, "team-alpha", "lead-1", ["session-a", "session-b"]) + const inboxDir = getInboxDir(resolveBaseDir(config), teamRunId, "member-1") + await mkdir(inboxDir, { recursive: true }) + await writeFile(path.join(inboxDir, "1.json"), JSON.stringify({ version: 1, messageId: randomUUID(), from: "lead", to: "member-1", kind: "message", body: "a", timestamp: 1 }) + "\n") + await writeFile(path.join(inboxDir, "2.json"), JSON.stringify({ version: 1, messageId: randomUUID(), from: "lead", to: "member-1", kind: "message", body: "b", timestamp: 2 }) + "\n") + await createTask(teamRunId, createTaskInput({ subject: "a" }), config) + await createTask(teamRunId, createTaskInput({ subject: "b" }), config) + await createTask(teamRunId, createTaskInput({ subject: "c" }), config) + await createTask(teamRunId, createTaskInput({ subject: "d" }), config) + + // when + const result = await aggregateStatus(teamRunId, config) + + // then + expect(result.teamName).toBe("team-alpha") + expect(result.members).toEqual([ + expect.objectContaining({ name: "lead", unreadMessages: 0 }), + expect.objectContaining({ name: "member-1", unreadMessages: 2 }), + expect.objectContaining({ name: "member-2", unreadMessages: 0 }), + ]) + expect(Object.keys(result.members[0] ?? {})).toEqual(expect.arrayContaining(["name", "unreadMessages"])) + expect(result.tasks).toEqual({ pending: 4, claimed: 0, in_progress: 0, completed: 0, deleted: 0, total: 4 }) + }) + + test("surfaces queued and running counts on same model", async () => { + // given + const baseDir = await createTemporaryBaseDir() + temporaryDirectories.push(baseDir) + const config = createConfig(baseDir) + const teamRunId = await seedRuntimeState(baseDir, "team-beta", "lead-2", []) + const backgroundManager = { + getTasksByParentSession: () => [ + { status: "running", model: { providerID: "anthropic", modelID: "claude-opus-4-7" } }, + { status: "running", model: { providerID: "anthropic", modelID: "claude-opus-4-7" } }, + { status: "running", model: { providerID: "anthropic", modelID: "claude-opus-4-7" } }, + { status: "running", model: { providerID: "anthropic", modelID: "claude-opus-4-7" } }, + { status: "running", model: { providerID: "anthropic", modelID: "claude-opus-4-7" } }, + { status: "pending", model: { providerID: "anthropic", modelID: "claude-opus-4-7" } }, + { status: "pending", model: { providerID: "anthropic", modelID: "claude-opus-4-7" } }, + { status: "pending", model: { providerID: "anthropic", modelID: "claude-opus-4-7" } }, + ], + getConcurrencyCounts: () => ({ running: 5, queued: 3 }), + listTasksByParentSession: () => [{}, {}, {}, {}], + } satisfies Pick & { + getConcurrencyCounts?: (modelOrUndefined?: string) => { running: number; queued: number } + listTasksByParentSession?: (sessionID: string) => unknown[] + } + + // when + const result = await aggregateStatus(teamRunId, config, backgroundManager) + + // then + expect(result.concurrency.runningOnSameModel).toBe(5) + expect(result.concurrency.queuedOnSameModel).toBe(3) + expect(result.concurrency.teamRunIdSpecific).toBe(4) + }) +}) diff --git a/src/features/team-mode/team-runtime/status.ts b/src/features/team-mode/team-runtime/status.ts new file mode 100644 index 000000000..edaddc914 --- /dev/null +++ b/src/features/team-mode/team-runtime/status.ts @@ -0,0 +1,155 @@ +import type { BackgroundManager } from "../../background-agent/manager" +import type { TeamModeConfig } from "../../../config/schema/team-mode" +import type { RuntimeState, Task } from "../types" +import { detectStaleLock } from "../team-state-store/locks" +import { loadRuntimeState } from "../team-state-store/store" +import { listUnreadMessages } from "../team-mailbox/inbox" +import { listTasks } from "../team-tasklist/list" +import { getTasksDir, resolveBaseDir } from "../team-registry/paths" +import { readdir } from "node:fs/promises" +import path from "node:path" + +export interface TeamStatus { + teamName: string + teamRunId: string + status: RuntimeState["status"] + leadSessionId?: string + createdAt: number + members: Array<{ + name: string + sessionId?: string + status: RuntimeState["members"][number]["status"] + color?: string + worktreePath?: string + unreadMessages: number + paneId?: string + }> + tasks: { + pending: number + claimed: number + in_progress: number + completed: number + deleted: number + total: number + } + shutdownRequests: RuntimeState["shutdownRequests"] + concurrency: { + runningOnSameModel: number + queuedOnSameModel: number + teamRunIdSpecific?: number + } + bounds: RuntimeState["bounds"] + staleLocks: string[] +} + +type ConcurrencyCounts = { + running: number + queued: number +} + +type TeamBackgroundManager = BackgroundManager & { + getConcurrencyCounts?: (modelOrUndefined?: string) => ConcurrencyCounts + listTasksByParentSession?: (sessionID: string) => Array +} + +function getPrimaryModelKey(bgMgr: TeamBackgroundManager | undefined, leadSessionId: string | undefined): string | undefined { + if (!bgMgr || !leadSessionId) return undefined + + const tasksByParent = bgMgr.getTasksByParentSession(leadSessionId) + if (tasksByParent.length === 0) return undefined + + const firstModel = tasksByParent[0]?.model + if (!firstModel) return undefined + + return `${firstModel.providerID}/${firstModel.modelID}` +} + +function countTasks(tasks: Task[]): TeamStatus["tasks"] { + const counts = { + pending: 0, + claimed: 0, + in_progress: 0, + completed: 0, + deleted: 0, + total: 0, + } + + for (const task of tasks) { + counts[task.status] += 1 + counts.total += 1 + } + + return counts +} + +function resolveConcurrencyCounts(bgMgr: TeamBackgroundManager | undefined, leadSessionId: string | undefined): ConcurrencyCounts { + if (!bgMgr || !leadSessionId) return { running: 0, queued: 0 } + + const modelKey = getPrimaryModelKey(bgMgr, leadSessionId) + const tasksByParent = bgMgr.getTasksByParentSession(leadSessionId) + const counts = bgMgr.getConcurrencyCounts?.(modelKey) + + if (counts) { + return { running: counts.running, queued: counts.queued } + } + + const running = tasksByParent.filter((task) => task.status === "running").length + const queued = tasksByParent.filter((task) => task.status === "pending").length + + return { running, queued } +} + +export async function aggregateStatus( + teamRunId: string, + config: TeamModeConfig, + bgMgr?: BackgroundManager, +): Promise { + const runtimeState = await loadRuntimeState(teamRunId, config) + const unreadCounts = await Promise.all( + runtimeState.members.map(async (member) => ({ + member, + unreadMessages: (await listUnreadMessages(teamRunId, member.name, config)).length, + })), + ) + const tasks = await listTasks(teamRunId, config) + const teamBackgroundManager: TeamBackgroundManager | undefined = bgMgr + const concurrencyCounts = resolveConcurrencyCounts(teamBackgroundManager, runtimeState.leadSessionId) + const teamRunIdSpecific = teamBackgroundManager?.listTasksByParentSession?.(runtimeState.leadSessionId ?? teamRunId)?.length + const baseDir = resolveBaseDir(config) + const claimsDir = path.join(getTasksDir(baseDir, teamRunId), "claims") + const staleLockEntries = await readdir(claimsDir, { withFileTypes: true }).catch(() => []) + const staleLockPaths = await Promise.all( + staleLockEntries + .filter((entry) => entry.isFile() && entry.name.endsWith(".lock")) + .map(async (entry) => { + const lockPath = path.join(claimsDir, entry.name) + return (await detectStaleLock(lockPath, 300_000)) ? lockPath : undefined + }), + ) + + return { + teamName: runtimeState.teamName, + teamRunId: runtimeState.teamRunId, + status: runtimeState.status, + leadSessionId: runtimeState.leadSessionId, + createdAt: runtimeState.createdAt, + members: unreadCounts.map(({ member, unreadMessages }) => ({ + name: member.name, + sessionId: member.sessionId, + status: member.status, + color: member.color, + worktreePath: member.worktreePath, + unreadMessages, + paneId: member.tmuxPaneId, + })), + tasks: countTasks(tasks), + shutdownRequests: runtimeState.shutdownRequests, + concurrency: { + runningOnSameModel: concurrencyCounts.running, + queuedOnSameModel: concurrencyCounts.queued, + teamRunIdSpecific, + }, + bounds: runtimeState.bounds, + staleLocks: staleLockPaths.filter((lockPath): lockPath is string => lockPath !== undefined), + } +}