feat(team-mode): add team runtime status query with tests
This commit is contained in:
@@ -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<string> {
|
||||||
|
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<string> {
|
||||||
|
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<BackgroundManager, "getTasksByParentSession"> & {
|
||||||
|
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)
|
||||||
|
})
|
||||||
|
})
|
||||||
@@ -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<unknown>
|
||||||
|
}
|
||||||
|
|
||||||
|
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<TeamStatus> {
|
||||||
|
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),
|
||||||
|
}
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user