feat(team-mode): add team runtime shutdown helpers and tests
This commit is contained in:
@@ -0,0 +1,78 @@
|
|||||||
|
import { randomUUID } from "node:crypto"
|
||||||
|
import { rm } from "node:fs/promises"
|
||||||
|
|
||||||
|
import type { Message, RuntimeState } from "../types"
|
||||||
|
|
||||||
|
export const DELETABLE_MEMBER_STATUSES = new Set<RuntimeState["members"][number]["status"]>([
|
||||||
|
"completed",
|
||||||
|
"shutdown_approved",
|
||||||
|
"errored",
|
||||||
|
])
|
||||||
|
|
||||||
|
export function createShutdownMessage(from: string, to: string, kind: Message["kind"], body: string): Message {
|
||||||
|
return {
|
||||||
|
version: 1,
|
||||||
|
messageId: randomUUID(),
|
||||||
|
from,
|
||||||
|
to,
|
||||||
|
kind,
|
||||||
|
body,
|
||||||
|
timestamp: Date.now(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
export function getRuntimeMember(runtimeState: RuntimeState, memberName: string): RuntimeState["members"][number] {
|
||||||
|
const member = runtimeState.members.find((candidate) => candidate.name === memberName)
|
||||||
|
if (!member) {
|
||||||
|
throw new Error(`unknown member '${memberName}'`)
|
||||||
|
}
|
||||||
|
|
||||||
|
return member
|
||||||
|
}
|
||||||
|
|
||||||
|
export function getLeadMemberName(runtimeState: RuntimeState): string {
|
||||||
|
const leadMember = runtimeState.members.find((member) => member.agentType === "leader")
|
||||||
|
if (!leadMember) {
|
||||||
|
throw new Error(`team '${runtimeState.teamRunId}' is missing a lead member`)
|
||||||
|
}
|
||||||
|
|
||||||
|
return leadMember.name
|
||||||
|
}
|
||||||
|
|
||||||
|
export function createSendContext(
|
||||||
|
runtimeState: RuntimeState,
|
||||||
|
senderName: string,
|
||||||
|
): { isLead: boolean; activeMembers: string[] } {
|
||||||
|
const sender = getRuntimeMember(runtimeState, senderName)
|
||||||
|
return {
|
||||||
|
isLead: sender.agentType === "leader",
|
||||||
|
activeMembers: runtimeState.members.map((member) => member.name),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
export function findLatestShutdownRequestIndex(
|
||||||
|
runtimeState: RuntimeState,
|
||||||
|
memberName: string,
|
||||||
|
requesterName?: string,
|
||||||
|
): number {
|
||||||
|
for (let index = runtimeState.shutdownRequests.length - 1; index >= 0; index -= 1) {
|
||||||
|
const shutdownRequest = runtimeState.shutdownRequests[index]
|
||||||
|
if (shutdownRequest.memberId !== memberName) continue
|
||||||
|
if (requesterName !== undefined && shutdownRequest.requesterName !== requesterName) continue
|
||||||
|
return index
|
||||||
|
}
|
||||||
|
|
||||||
|
return -1
|
||||||
|
}
|
||||||
|
|
||||||
|
export async function removeWorktrees(memberPaths: Array<string | undefined>): Promise<string[]> {
|
||||||
|
const removedWorktrees: string[] = []
|
||||||
|
|
||||||
|
for (const memberPath of new Set(memberPaths)) {
|
||||||
|
if (!memberPath) continue
|
||||||
|
await rm(memberPath, { recursive: true, force: true })
|
||||||
|
removedWorktrees.push(memberPath)
|
||||||
|
}
|
||||||
|
|
||||||
|
return removedWorktrees
|
||||||
|
}
|
||||||
@@ -0,0 +1,146 @@
|
|||||||
|
import { mkdir, mkdtemp, readdir, readFile } from "node:fs/promises"
|
||||||
|
import { tmpdir } from "node:os"
|
||||||
|
import path from "node:path"
|
||||||
|
|
||||||
|
import { TeamModeConfigSchema } from "../../../config/schema/team-mode"
|
||||||
|
import type { TeamModeConfig } from "../../../config/schema/team-mode"
|
||||||
|
import { sendMessage } from "../team-mailbox/send"
|
||||||
|
import { getInboxDir, getRuntimeStateDir, resolveBaseDir } from "../team-registry/paths"
|
||||||
|
import { saveRuntimeState, transitionRuntimeState } from "../team-state-store/store"
|
||||||
|
import { MessageSchema, type RuntimeState, type TeamSpec } from "../types"
|
||||||
|
|
||||||
|
let fixtureCounter = 0
|
||||||
|
|
||||||
|
function createUuid(sequence: number): string {
|
||||||
|
return `123e4567-e89b-42d3-a456-${sequence.toString(16).padStart(12, "0")}`
|
||||||
|
}
|
||||||
|
|
||||||
|
export function createConfig(baseDir: string): TeamModeConfig {
|
||||||
|
return TeamModeConfigSchema.parse({ base_dir: baseDir })
|
||||||
|
}
|
||||||
|
|
||||||
|
export function createSpec(worktreeRoot: string): TeamSpec {
|
||||||
|
fixtureCounter += 1
|
||||||
|
|
||||||
|
return {
|
||||||
|
version: 1,
|
||||||
|
name: `team-${fixtureCounter.toString(16).padStart(8, "0")}`,
|
||||||
|
createdAt: Date.now(),
|
||||||
|
leadAgentId: "lead",
|
||||||
|
members: [
|
||||||
|
{ kind: "subagent_type", name: "lead", subagent_type: "sisyphus", backendType: "in-process", isActive: true },
|
||||||
|
{
|
||||||
|
kind: "category",
|
||||||
|
name: "member-a",
|
||||||
|
category: "deep",
|
||||||
|
prompt: "work on task a",
|
||||||
|
backendType: "in-process",
|
||||||
|
isActive: true,
|
||||||
|
worktreePath: path.join(worktreeRoot, "member-a"),
|
||||||
|
},
|
||||||
|
{
|
||||||
|
kind: "category",
|
||||||
|
name: "member-b",
|
||||||
|
category: "deep",
|
||||||
|
prompt: "work on task b",
|
||||||
|
backendType: "in-process",
|
||||||
|
isActive: true,
|
||||||
|
worktreePath: path.join(worktreeRoot, "member-b"),
|
||||||
|
},
|
||||||
|
],
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
export async function createFixture(options?: { status?: RuntimeState["status"] }): Promise<{
|
||||||
|
baseDir: string
|
||||||
|
config: TeamModeConfig
|
||||||
|
teamRunId: string
|
||||||
|
worktreePaths: string[]
|
||||||
|
}> {
|
||||||
|
fixtureCounter += 1
|
||||||
|
const baseDir = await mkdtemp(path.join(tmpdir(), `team-runtime-shutdown-${fixtureCounter}-`))
|
||||||
|
const config = createConfig(baseDir)
|
||||||
|
const worktreeRoot = path.join(baseDir, "fixture-worktrees")
|
||||||
|
const teamRunId = createUuid(fixtureCounter)
|
||||||
|
const runtimeState: RuntimeState = {
|
||||||
|
version: 1,
|
||||||
|
teamRunId,
|
||||||
|
teamName: createSpec(worktreeRoot).name,
|
||||||
|
specSource: "project",
|
||||||
|
createdAt: Date.now(),
|
||||||
|
status: options?.status ?? "active",
|
||||||
|
leadSessionId: "lead-session",
|
||||||
|
members: [
|
||||||
|
{ name: "lead", agentType: "leader", status: "pending", pendingInjectedMessageIds: [] },
|
||||||
|
{
|
||||||
|
name: "member-a",
|
||||||
|
agentType: "general-purpose",
|
||||||
|
status: "pending",
|
||||||
|
pendingInjectedMessageIds: [],
|
||||||
|
worktreePath: path.join(worktreeRoot, "member-a"),
|
||||||
|
},
|
||||||
|
{
|
||||||
|
name: "member-b",
|
||||||
|
agentType: "general-purpose",
|
||||||
|
status: "pending",
|
||||||
|
pendingInjectedMessageIds: [],
|
||||||
|
worktreePath: path.join(worktreeRoot, "member-b"),
|
||||||
|
},
|
||||||
|
],
|
||||||
|
shutdownRequests: [],
|
||||||
|
bounds: {
|
||||||
|
maxMembers: config.max_members,
|
||||||
|
maxParallelMembers: config.max_parallel_members,
|
||||||
|
maxMessagesPerRun: config.max_messages_per_run,
|
||||||
|
maxWallClockMinutes: config.max_wall_clock_minutes,
|
||||||
|
maxMemberTurns: config.max_member_turns,
|
||||||
|
},
|
||||||
|
}
|
||||||
|
await mkdir(getRuntimeStateDir(resolveBaseDir(config), teamRunId), { recursive: true })
|
||||||
|
await saveRuntimeState(runtimeState, config)
|
||||||
|
|
||||||
|
return {
|
||||||
|
baseDir,
|
||||||
|
config,
|
||||||
|
teamRunId: runtimeState.teamRunId,
|
||||||
|
worktreePaths: [path.join(worktreeRoot, "member-a"), path.join(worktreeRoot, "member-b")],
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
export async function updateMemberStatuses(
|
||||||
|
teamRunId: string,
|
||||||
|
config: TeamModeConfig,
|
||||||
|
statuses: Record<string, RuntimeState["members"][number]["status"]>,
|
||||||
|
): Promise<void> {
|
||||||
|
await transitionRuntimeState(teamRunId, (runtimeState) => ({
|
||||||
|
...runtimeState,
|
||||||
|
members: runtimeState.members.map((member) => ({
|
||||||
|
...member,
|
||||||
|
status: statuses[member.name] ?? member.status,
|
||||||
|
})),
|
||||||
|
}), config)
|
||||||
|
}
|
||||||
|
|
||||||
|
export async function readInboxMessages(teamRunId: string, memberName: string, config: TeamModeConfig) {
|
||||||
|
const inboxDir = getInboxDir(resolveBaseDir(config), teamRunId, memberName)
|
||||||
|
const fileNames = (await readdir(inboxDir)).filter((entry) => entry.endsWith(".json")).sort()
|
||||||
|
return Promise.all(fileNames.map(async (fileName) => {
|
||||||
|
const content = await readFile(path.join(inboxDir, fileName), "utf8")
|
||||||
|
return MessageSchema.parse(JSON.parse(content))
|
||||||
|
}))
|
||||||
|
}
|
||||||
|
|
||||||
|
export function createTestMessage(overrides?: Partial<Parameters<typeof sendMessage>[0]>) {
|
||||||
|
fixtureCounter += 1
|
||||||
|
|
||||||
|
return MessageSchema.parse({
|
||||||
|
version: 1,
|
||||||
|
messageId: createUuid(fixtureCounter),
|
||||||
|
from: "lead",
|
||||||
|
to: "member-a",
|
||||||
|
kind: "message",
|
||||||
|
body: "hello",
|
||||||
|
timestamp: Date.now(),
|
||||||
|
...overrides,
|
||||||
|
})
|
||||||
|
}
|
||||||
@@ -0,0 +1,396 @@
|
|||||||
|
/// <reference types="bun-types" />
|
||||||
|
|
||||||
|
import { afterEach, describe, expect, mock, spyOn, test } from "bun:test"
|
||||||
|
import { access, mkdir, rm } from "node:fs/promises"
|
||||||
|
import path from "node:path"
|
||||||
|
|
||||||
|
import { sendMessage } from "../team-mailbox/send"
|
||||||
|
import { getRuntimeStateDir, resolveBaseDir } from "../team-registry/paths"
|
||||||
|
import * as logger from "../../../shared/logger"
|
||||||
|
import * as layoutModule from "../team-layout-tmux/layout"
|
||||||
|
import * as runtimeStateStore from "../team-state-store/store"
|
||||||
|
import { loadRuntimeState, transitionRuntimeState } from "../team-state-store/store"
|
||||||
|
import {
|
||||||
|
createFixture,
|
||||||
|
createTestMessage,
|
||||||
|
readInboxMessages,
|
||||||
|
updateMemberStatuses,
|
||||||
|
} from "./shutdown-test-fixtures"
|
||||||
|
|
||||||
|
const { approveShutdown, deleteTeam, rejectShutdown, requestShutdownOfMember } = await import("./shutdown")
|
||||||
|
|
||||||
|
describe("team-runtime shutdown", () => {
|
||||||
|
const temporaryDirectories: string[] = []
|
||||||
|
|
||||||
|
afterEach(async () => {
|
||||||
|
await Promise.all(temporaryDirectories.splice(0).map(async (directoryPath) => {
|
||||||
|
await rm(directoryPath, { recursive: true, force: true })
|
||||||
|
}))
|
||||||
|
mock.restore()
|
||||||
|
})
|
||||||
|
|
||||||
|
test("refuses team deletion while non-lead members are still active", async () => {
|
||||||
|
// given
|
||||||
|
const fixture = await createFixture()
|
||||||
|
temporaryDirectories.push(fixture.baseDir)
|
||||||
|
await updateMemberStatuses(fixture.teamRunId, fixture.config, {
|
||||||
|
"member-a": "running",
|
||||||
|
"member-b": "running",
|
||||||
|
})
|
||||||
|
|
||||||
|
// when
|
||||||
|
const result = deleteTeam(fixture.teamRunId, fixture.config)
|
||||||
|
|
||||||
|
// then
|
||||||
|
await result.then(
|
||||||
|
() => { throw new Error("expected deleteTeam to reject") },
|
||||||
|
(error: unknown) => {
|
||||||
|
if (!(error instanceof Error)) throw error
|
||||||
|
expect(error.message).toBe("members still active")
|
||||||
|
},
|
||||||
|
)
|
||||||
|
const runtimeState = await loadRuntimeState(fixture.teamRunId, fixture.config)
|
||||||
|
expect(runtimeState.status).toBe("active")
|
||||||
|
expect(runtimeState.members.filter((member) => member.agentType !== "leader").map((member) => member.status)).toEqual([
|
||||||
|
"running",
|
||||||
|
"running",
|
||||||
|
])
|
||||||
|
})
|
||||||
|
|
||||||
|
test("writes shutdown requests to the target inbox and records runtime metadata", async () => {
|
||||||
|
// given
|
||||||
|
const fixture = await createFixture()
|
||||||
|
temporaryDirectories.push(fixture.baseDir)
|
||||||
|
|
||||||
|
// when
|
||||||
|
await requestShutdownOfMember(fixture.teamRunId, "member-a", "lead", fixture.config)
|
||||||
|
|
||||||
|
// then
|
||||||
|
const inboxMessages = await readInboxMessages(fixture.teamRunId, "member-a", fixture.config)
|
||||||
|
const runtimeState = await loadRuntimeState(fixture.teamRunId, fixture.config)
|
||||||
|
expect(inboxMessages).toHaveLength(1)
|
||||||
|
expect(inboxMessages[0]).toEqual(expect.objectContaining({
|
||||||
|
from: "lead",
|
||||||
|
to: "member-a",
|
||||||
|
kind: "shutdown_request",
|
||||||
|
body: "",
|
||||||
|
}))
|
||||||
|
expect(runtimeState.shutdownRequests).toEqual([
|
||||||
|
expect.objectContaining({
|
||||||
|
memberId: "member-a",
|
||||||
|
requesterName: "lead",
|
||||||
|
requestedAt: expect.any(Number),
|
||||||
|
}),
|
||||||
|
])
|
||||||
|
})
|
||||||
|
|
||||||
|
test("approves shutdown requests and notifies the lead", async () => {
|
||||||
|
// given
|
||||||
|
const fixture = await createFixture()
|
||||||
|
temporaryDirectories.push(fixture.baseDir)
|
||||||
|
await requestShutdownOfMember(fixture.teamRunId, "member-a", "lead", fixture.config)
|
||||||
|
|
||||||
|
// when
|
||||||
|
await approveShutdown(fixture.teamRunId, "member-a", "member-a", fixture.config)
|
||||||
|
|
||||||
|
// then
|
||||||
|
const runtimeState = await loadRuntimeState(fixture.teamRunId, fixture.config)
|
||||||
|
const leadInboxMessages = await readInboxMessages(fixture.teamRunId, "lead", fixture.config)
|
||||||
|
const approvedRequest = runtimeState.shutdownRequests.find((shutdownRequest) => shutdownRequest.memberId === "member-a")
|
||||||
|
expect(approvedRequest?.approvedAt).toEqual(expect.any(Number))
|
||||||
|
expect(runtimeState.members.find((member) => member.name === "member-a")?.status).toBe("shutdown_approved")
|
||||||
|
expect(leadInboxMessages.some((message) => (
|
||||||
|
message.kind === "shutdown_approved"
|
||||||
|
&& message.from === "member-a"
|
||||||
|
&& message.to === "lead"
|
||||||
|
&& message.body === "member-a"
|
||||||
|
))).toBe(true)
|
||||||
|
})
|
||||||
|
|
||||||
|
test("rejects shutdown requests and replies to the original requester", async () => {
|
||||||
|
// given
|
||||||
|
const fixture = await createFixture()
|
||||||
|
temporaryDirectories.push(fixture.baseDir)
|
||||||
|
await requestShutdownOfMember(fixture.teamRunId, "member-a", "lead", fixture.config)
|
||||||
|
|
||||||
|
// when
|
||||||
|
await rejectShutdown(fixture.teamRunId, "member-a", "not done yet", fixture.config)
|
||||||
|
|
||||||
|
// then
|
||||||
|
const runtimeState = await loadRuntimeState(fixture.teamRunId, fixture.config)
|
||||||
|
const leadInboxMessages = await readInboxMessages(fixture.teamRunId, "lead", fixture.config)
|
||||||
|
const rejectedRequest = runtimeState.shutdownRequests.find((shutdownRequest) => shutdownRequest.memberId === "member-a")
|
||||||
|
expect(rejectedRequest).toEqual(expect.objectContaining({
|
||||||
|
rejectedAt: expect.any(Number),
|
||||||
|
rejectedReason: "not done yet",
|
||||||
|
}))
|
||||||
|
expect(leadInboxMessages.some((message) => (
|
||||||
|
message.kind === "shutdown_rejected"
|
||||||
|
&& message.from === "member-a"
|
||||||
|
&& message.to === "lead"
|
||||||
|
&& message.body === "not done yet"
|
||||||
|
))).toBe(true)
|
||||||
|
})
|
||||||
|
|
||||||
|
test("deletes team runtime resources after all non-lead members are approved", async () => {
|
||||||
|
// given
|
||||||
|
const fixture = await createFixture()
|
||||||
|
temporaryDirectories.push(fixture.baseDir)
|
||||||
|
await updateMemberStatuses(fixture.teamRunId, fixture.config, {
|
||||||
|
"member-a": "shutdown_approved",
|
||||||
|
"member-b": "shutdown_approved",
|
||||||
|
})
|
||||||
|
await Promise.all(fixture.worktreePaths.map(async (worktreePath) => {
|
||||||
|
await mkdir(worktreePath, { recursive: true })
|
||||||
|
}))
|
||||||
|
// when
|
||||||
|
const result = await deleteTeam(fixture.teamRunId, fixture.config)
|
||||||
|
|
||||||
|
// then
|
||||||
|
expect(result.removedLayout).toBe(false)
|
||||||
|
expect(result.removedWorktrees.sort()).toEqual([...fixture.worktreePaths].sort())
|
||||||
|
await Promise.all(fixture.worktreePaths.map(async (worktreePath) => {
|
||||||
|
await access(worktreePath).then(
|
||||||
|
() => { throw new Error(`expected ${worktreePath} to be removed`) },
|
||||||
|
() => undefined,
|
||||||
|
)
|
||||||
|
}))
|
||||||
|
const runtimeStateDirectory = getRuntimeStateDir(resolveBaseDir(fixture.config), fixture.teamRunId)
|
||||||
|
await access(runtimeStateDirectory).then(
|
||||||
|
() => { throw new Error(`expected ${runtimeStateDirectory} to be removed`) },
|
||||||
|
() => undefined,
|
||||||
|
)
|
||||||
|
})
|
||||||
|
|
||||||
|
test("deletes team even with active members when force=true", async () => {
|
||||||
|
// given
|
||||||
|
const fixture = await createFixture()
|
||||||
|
temporaryDirectories.push(fixture.baseDir)
|
||||||
|
await updateMemberStatuses(fixture.teamRunId, fixture.config, {
|
||||||
|
"member-a": "running",
|
||||||
|
"member-b": "running",
|
||||||
|
})
|
||||||
|
await Promise.all(fixture.worktreePaths.map(async (worktreePath) => {
|
||||||
|
await mkdir(worktreePath, { recursive: true })
|
||||||
|
}))
|
||||||
|
|
||||||
|
// when
|
||||||
|
const result = await deleteTeam(fixture.teamRunId, fixture.config, undefined, undefined, { force: true })
|
||||||
|
|
||||||
|
// then
|
||||||
|
expect(result.removedLayout).toBe(false)
|
||||||
|
expect(result.removedWorktrees.sort()).toEqual([...fixture.worktreePaths].sort())
|
||||||
|
await Promise.all(fixture.worktreePaths.map(async (worktreePath) => {
|
||||||
|
await access(worktreePath).then(
|
||||||
|
() => { throw new Error(`expected ${worktreePath} to be removed`) },
|
||||||
|
() => undefined,
|
||||||
|
)
|
||||||
|
}))
|
||||||
|
const runtimeStateDirectory = getRuntimeStateDir(resolveBaseDir(fixture.config), fixture.teamRunId)
|
||||||
|
await access(runtimeStateDirectory).then(
|
||||||
|
() => { throw new Error(`expected ${runtimeStateDirectory} to be removed`) },
|
||||||
|
() => undefined,
|
||||||
|
)
|
||||||
|
})
|
||||||
|
|
||||||
|
test("force deletes a team stuck in 'creating' status", async () => {
|
||||||
|
// given
|
||||||
|
const fixture = await createFixture({ status: "creating" })
|
||||||
|
temporaryDirectories.push(fixture.baseDir)
|
||||||
|
const transitionedStatuses: string[] = []
|
||||||
|
const originalTransitionRuntimeState = runtimeStateStore.transitionRuntimeState
|
||||||
|
spyOn(runtimeStateStore, "transitionRuntimeState").mockImplementation(async (teamRunId, transition, config) => {
|
||||||
|
const currentRuntimeState = await runtimeStateStore.loadRuntimeState(teamRunId, config)
|
||||||
|
transitionedStatuses.push(transition(currentRuntimeState).status)
|
||||||
|
return await originalTransitionRuntimeState(teamRunId, transition, config)
|
||||||
|
})
|
||||||
|
await updateMemberStatuses(fixture.teamRunId, fixture.config, {
|
||||||
|
"member-a": "pending",
|
||||||
|
"member-b": "pending",
|
||||||
|
})
|
||||||
|
|
||||||
|
// when
|
||||||
|
await deleteTeam(fixture.teamRunId, fixture.config, undefined, undefined, { force: true })
|
||||||
|
|
||||||
|
// then
|
||||||
|
expect(transitionedStatuses).toContain("deleted")
|
||||||
|
const runtimeStateDirectory = getRuntimeStateDir(resolveBaseDir(fixture.config), fixture.teamRunId)
|
||||||
|
await access(runtimeStateDirectory).then(
|
||||||
|
() => { throw new Error(`expected ${runtimeStateDirectory} to be removed`) },
|
||||||
|
() => undefined,
|
||||||
|
)
|
||||||
|
})
|
||||||
|
|
||||||
|
test("force deletes a team in 'orphaned' status", async () => {
|
||||||
|
// given
|
||||||
|
const fixture = await createFixture({ status: "orphaned" })
|
||||||
|
temporaryDirectories.push(fixture.baseDir)
|
||||||
|
const transitionedStatuses: string[] = []
|
||||||
|
const originalTransitionRuntimeState = runtimeStateStore.transitionRuntimeState
|
||||||
|
spyOn(runtimeStateStore, "transitionRuntimeState").mockImplementation(async (teamRunId, transition, config) => {
|
||||||
|
const currentRuntimeState = await runtimeStateStore.loadRuntimeState(teamRunId, config)
|
||||||
|
transitionedStatuses.push(transition(currentRuntimeState).status)
|
||||||
|
return await originalTransitionRuntimeState(teamRunId, transition, config)
|
||||||
|
})
|
||||||
|
await updateMemberStatuses(fixture.teamRunId, fixture.config, {
|
||||||
|
"member-a": "running",
|
||||||
|
"member-b": "running",
|
||||||
|
})
|
||||||
|
|
||||||
|
// when
|
||||||
|
await deleteTeam(fixture.teamRunId, fixture.config, undefined, undefined, { force: true })
|
||||||
|
|
||||||
|
// then
|
||||||
|
expect(transitionedStatuses).toContain("deleted")
|
||||||
|
const runtimeStateDirectory = getRuntimeStateDir(resolveBaseDir(fixture.config), fixture.teamRunId)
|
||||||
|
await access(runtimeStateDirectory).then(
|
||||||
|
() => { throw new Error(`expected ${runtimeStateDirectory} to be removed`) },
|
||||||
|
() => undefined,
|
||||||
|
)
|
||||||
|
})
|
||||||
|
|
||||||
|
test("force removes lead member worktree if present", async () => {
|
||||||
|
// given
|
||||||
|
const fixture = await createFixture()
|
||||||
|
temporaryDirectories.push(fixture.baseDir)
|
||||||
|
const leadWorktreePath = path.join(fixture.baseDir, "fixture-worktrees", "lead")
|
||||||
|
await transitionRuntimeState(fixture.teamRunId, (runtimeState) => ({
|
||||||
|
...runtimeState,
|
||||||
|
members: runtimeState.members.map((member) => member.name === "lead"
|
||||||
|
? { ...member, worktreePath: leadWorktreePath }
|
||||||
|
: member),
|
||||||
|
}), fixture.config)
|
||||||
|
await mkdir(leadWorktreePath, { recursive: true })
|
||||||
|
await Promise.all(fixture.worktreePaths.map(async (worktreePath) => {
|
||||||
|
await mkdir(worktreePath, { recursive: true })
|
||||||
|
}))
|
||||||
|
await updateMemberStatuses(fixture.teamRunId, fixture.config, {
|
||||||
|
"member-a": "running",
|
||||||
|
"member-b": "running",
|
||||||
|
})
|
||||||
|
|
||||||
|
// when
|
||||||
|
const result = await deleteTeam(fixture.teamRunId, fixture.config, undefined, undefined, { force: true })
|
||||||
|
|
||||||
|
// then
|
||||||
|
expect(result.removedWorktrees.sort()).toEqual([leadWorktreePath, ...fixture.worktreePaths].sort())
|
||||||
|
await access(leadWorktreePath).then(
|
||||||
|
() => { throw new Error(`expected ${leadWorktreePath} to be removed`) },
|
||||||
|
() => undefined,
|
||||||
|
)
|
||||||
|
})
|
||||||
|
|
||||||
|
test("force continues cleanup when removeTeamLayout throws", async () => {
|
||||||
|
// given
|
||||||
|
const fixture = await createFixture()
|
||||||
|
temporaryDirectories.push(fixture.baseDir)
|
||||||
|
const transitionedStatuses: string[] = []
|
||||||
|
const originalTransitionRuntimeState = runtimeStateStore.transitionRuntimeState
|
||||||
|
spyOn(runtimeStateStore, "transitionRuntimeState").mockImplementation(async (teamRunId, transition, config) => {
|
||||||
|
const currentRuntimeState = await runtimeStateStore.loadRuntimeState(teamRunId, config)
|
||||||
|
transitionedStatuses.push(transition(currentRuntimeState).status)
|
||||||
|
return await originalTransitionRuntimeState(teamRunId, transition, config)
|
||||||
|
})
|
||||||
|
spyOn(layoutModule, "canVisualize").mockReturnValue(true)
|
||||||
|
spyOn(layoutModule, "removeTeamLayout").mockRejectedValue(new Error("layout failed"))
|
||||||
|
const logSpy = spyOn(logger, "log").mockImplementation(() => {})
|
||||||
|
await updateMemberStatuses(fixture.teamRunId, fixture.config, {
|
||||||
|
"member-a": "running",
|
||||||
|
"member-b": "idle",
|
||||||
|
})
|
||||||
|
await Promise.all(fixture.worktreePaths.map(async (worktreePath) => {
|
||||||
|
await mkdir(worktreePath, { recursive: true })
|
||||||
|
}))
|
||||||
|
|
||||||
|
// when
|
||||||
|
const result = await deleteTeam(
|
||||||
|
fixture.teamRunId,
|
||||||
|
fixture.config,
|
||||||
|
{ getServerUrl: () => "http://localhost" } as never,
|
||||||
|
undefined,
|
||||||
|
{ force: true },
|
||||||
|
)
|
||||||
|
|
||||||
|
// then
|
||||||
|
expect(result.removedLayout).toBe(true)
|
||||||
|
expect(transitionedStatuses).toContain("deleted")
|
||||||
|
expect(logSpy).toHaveBeenCalledWith("team delete layout cleanup failed", {
|
||||||
|
teamRunId: fixture.teamRunId,
|
||||||
|
error: "layout failed",
|
||||||
|
})
|
||||||
|
const runtimeStateDirectory = getRuntimeStateDir(resolveBaseDir(fixture.config), fixture.teamRunId)
|
||||||
|
await access(runtimeStateDirectory).then(
|
||||||
|
() => { throw new Error(`expected ${runtimeStateDirectory} to be removed`) },
|
||||||
|
() => undefined,
|
||||||
|
)
|
||||||
|
})
|
||||||
|
|
||||||
|
test("cancels team background tasks before deleting when force=true", async () => {
|
||||||
|
// given
|
||||||
|
const fixture = await createFixture()
|
||||||
|
temporaryDirectories.push(fixture.baseDir)
|
||||||
|
await updateMemberStatuses(fixture.teamRunId, fixture.config, {
|
||||||
|
"member-a": "running",
|
||||||
|
"member-b": "idle",
|
||||||
|
})
|
||||||
|
const runtimeStatusesDuringCancellation: Array<{ teamStatus: string; memberStatuses: string[] }> = []
|
||||||
|
const cancelTaskMock = mock(async () => {
|
||||||
|
const runtimeState = await loadRuntimeState(fixture.teamRunId, fixture.config)
|
||||||
|
runtimeStatusesDuringCancellation.push({
|
||||||
|
teamStatus: runtimeState.status,
|
||||||
|
memberStatuses: runtimeState.members
|
||||||
|
.filter((member) => member.agentType !== "leader")
|
||||||
|
.map((member) => member.status),
|
||||||
|
})
|
||||||
|
return true
|
||||||
|
})
|
||||||
|
const bgMgr = {
|
||||||
|
getTasksByParentSession: () => [
|
||||||
|
{ id: "team-task-a", sessionID: "session-a", parentMessageID: `team-create:${fixture.teamRunId}:member-a` },
|
||||||
|
{ id: "team-task-b", sessionID: "session-b", parentMessageID: `team-create:${fixture.teamRunId}:member-b` },
|
||||||
|
],
|
||||||
|
cancelTask: cancelTaskMock,
|
||||||
|
}
|
||||||
|
|
||||||
|
// when
|
||||||
|
await deleteTeam(fixture.teamRunId, fixture.config, undefined, bgMgr as never, { force: true })
|
||||||
|
|
||||||
|
// then
|
||||||
|
expect(cancelTaskMock).toHaveBeenCalledTimes(2)
|
||||||
|
expect(runtimeStatusesDuringCancellation).toEqual([
|
||||||
|
{ teamStatus: "active", memberStatuses: ["running", "idle"] },
|
||||||
|
{ teamStatus: "active", memberStatuses: ["running", "idle"] },
|
||||||
|
])
|
||||||
|
})
|
||||||
|
|
||||||
|
test("blocks mailbox writes while the team is deleting", async () => {
|
||||||
|
// given
|
||||||
|
const fixture = await createFixture()
|
||||||
|
temporaryDirectories.push(fixture.baseDir)
|
||||||
|
await updateMemberStatuses(fixture.teamRunId, fixture.config, {
|
||||||
|
"member-a": "shutdown_approved",
|
||||||
|
"member-b": "shutdown_approved",
|
||||||
|
})
|
||||||
|
await transitionRuntimeState(fixture.teamRunId, (runtimeState) => ({
|
||||||
|
...runtimeState,
|
||||||
|
status: "deleting",
|
||||||
|
}), fixture.config)
|
||||||
|
|
||||||
|
// when
|
||||||
|
const result = sendMessage(
|
||||||
|
createTestMessage(),
|
||||||
|
fixture.teamRunId,
|
||||||
|
fixture.config,
|
||||||
|
{ isLead: true, activeMembers: ["lead", "member-a", "member-b"] },
|
||||||
|
)
|
||||||
|
|
||||||
|
// then
|
||||||
|
await result.then(
|
||||||
|
() => { throw new Error("expected sendMessage to reject") },
|
||||||
|
(error: unknown) => {
|
||||||
|
if (!(error instanceof Error)) throw error
|
||||||
|
expect(error.message).toBe("team is deleting")
|
||||||
|
},
|
||||||
|
)
|
||||||
|
})
|
||||||
|
})
|
||||||
@@ -0,0 +1,146 @@
|
|||||||
|
import type { TeamModeConfig } from "../../../config/schema/team-mode"
|
||||||
|
import { sendMessage } from "../team-mailbox/send"
|
||||||
|
import { loadRuntimeState, transitionRuntimeState } from "../team-state-store/store"
|
||||||
|
import {
|
||||||
|
createSendContext,
|
||||||
|
createShutdownMessage,
|
||||||
|
findLatestShutdownRequestIndex,
|
||||||
|
getLeadMemberName,
|
||||||
|
getRuntimeMember,
|
||||||
|
} from "./shutdown-helpers"
|
||||||
|
export { deleteTeam } from "./delete-team"
|
||||||
|
|
||||||
|
export async function requestShutdownOfMember(
|
||||||
|
teamRunId: string,
|
||||||
|
targetMemberName: string,
|
||||||
|
requesterName: string,
|
||||||
|
config: TeamModeConfig,
|
||||||
|
): Promise<void> {
|
||||||
|
const runtimeState = await loadRuntimeState(teamRunId, config)
|
||||||
|
getRuntimeMember(runtimeState, targetMemberName)
|
||||||
|
getRuntimeMember(runtimeState, requesterName)
|
||||||
|
|
||||||
|
const existingRequestIndex = findLatestShutdownRequestIndex(runtimeState, targetMemberName, requesterName)
|
||||||
|
const existingRequest = existingRequestIndex >= 0
|
||||||
|
? runtimeState.shutdownRequests[existingRequestIndex]
|
||||||
|
: undefined
|
||||||
|
if (existingRequest && existingRequest.approvedAt === undefined && existingRequest.rejectedAt === undefined) {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
await sendMessage(
|
||||||
|
createShutdownMessage(requesterName, targetMemberName, "shutdown_request", ""),
|
||||||
|
teamRunId,
|
||||||
|
config,
|
||||||
|
createSendContext(runtimeState, requesterName),
|
||||||
|
)
|
||||||
|
|
||||||
|
await transitionRuntimeState(teamRunId, (currentRuntimeState) => {
|
||||||
|
const duplicateRequestIndex = findLatestShutdownRequestIndex(currentRuntimeState, targetMemberName, requesterName)
|
||||||
|
const duplicateRequest = duplicateRequestIndex >= 0
|
||||||
|
? currentRuntimeState.shutdownRequests[duplicateRequestIndex]
|
||||||
|
: undefined
|
||||||
|
if (duplicateRequest && duplicateRequest.approvedAt === undefined && duplicateRequest.rejectedAt === undefined) {
|
||||||
|
return currentRuntimeState
|
||||||
|
}
|
||||||
|
|
||||||
|
return {
|
||||||
|
...currentRuntimeState,
|
||||||
|
shutdownRequests: [
|
||||||
|
...currentRuntimeState.shutdownRequests,
|
||||||
|
{ memberId: targetMemberName, requesterName, requestedAt: Date.now() },
|
||||||
|
],
|
||||||
|
}
|
||||||
|
}, config)
|
||||||
|
}
|
||||||
|
|
||||||
|
export async function approveShutdown(
|
||||||
|
teamRunId: string,
|
||||||
|
memberName: string,
|
||||||
|
approverName: string,
|
||||||
|
config: TeamModeConfig,
|
||||||
|
): Promise<void> {
|
||||||
|
const runtimeState = await loadRuntimeState(teamRunId, config)
|
||||||
|
getRuntimeMember(runtimeState, approverName)
|
||||||
|
const shutdownRequestIndex = findLatestShutdownRequestIndex(runtimeState, memberName)
|
||||||
|
if (shutdownRequestIndex < 0) {
|
||||||
|
throw new Error(`shutdown request missing for '${memberName}'`)
|
||||||
|
}
|
||||||
|
|
||||||
|
const existingRequest = runtimeState.shutdownRequests[shutdownRequestIndex]
|
||||||
|
if (existingRequest?.approvedAt !== undefined) {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
const updatedRuntimeState = await transitionRuntimeState(teamRunId, (currentRuntimeState) => {
|
||||||
|
const currentRequestIndex = findLatestShutdownRequestIndex(currentRuntimeState, memberName)
|
||||||
|
if (currentRequestIndex < 0) {
|
||||||
|
throw new Error(`shutdown request missing for '${memberName}'`)
|
||||||
|
}
|
||||||
|
|
||||||
|
const currentRequest = currentRuntimeState.shutdownRequests[currentRequestIndex]
|
||||||
|
if (!currentRequest || currentRequest.approvedAt !== undefined) {
|
||||||
|
return currentRuntimeState
|
||||||
|
}
|
||||||
|
|
||||||
|
return {
|
||||||
|
...currentRuntimeState,
|
||||||
|
members: currentRuntimeState.members.map((member) => {
|
||||||
|
if (member.name !== memberName || member.status === "completed" || member.status === "errored") {
|
||||||
|
return member
|
||||||
|
}
|
||||||
|
|
||||||
|
return { ...member, status: "shutdown_approved" }
|
||||||
|
}),
|
||||||
|
shutdownRequests: currentRuntimeState.shutdownRequests.map((shutdownRequest, index) => index === currentRequestIndex
|
||||||
|
? { ...shutdownRequest, approvedAt: Date.now() }
|
||||||
|
: shutdownRequest),
|
||||||
|
}
|
||||||
|
}, config)
|
||||||
|
|
||||||
|
await sendMessage(
|
||||||
|
createShutdownMessage(approverName, getLeadMemberName(updatedRuntimeState), "shutdown_approved", memberName),
|
||||||
|
teamRunId,
|
||||||
|
config,
|
||||||
|
createSendContext(updatedRuntimeState, approverName),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
export async function rejectShutdown(
|
||||||
|
teamRunId: string,
|
||||||
|
memberName: string,
|
||||||
|
reason: string,
|
||||||
|
config: TeamModeConfig,
|
||||||
|
): Promise<void> {
|
||||||
|
const runtimeState = await loadRuntimeState(teamRunId, config)
|
||||||
|
const shutdownRequestIndex = findLatestShutdownRequestIndex(runtimeState, memberName)
|
||||||
|
if (shutdownRequestIndex < 0) {
|
||||||
|
throw new Error(`shutdown request missing for '${memberName}'`)
|
||||||
|
}
|
||||||
|
|
||||||
|
const shutdownRequest = runtimeState.shutdownRequests[shutdownRequestIndex]
|
||||||
|
if (shutdownRequest.rejectedAt !== undefined && shutdownRequest.rejectedReason === reason) {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
|
||||||
|
await sendMessage(
|
||||||
|
createShutdownMessage(memberName, shutdownRequest.requesterName, "shutdown_rejected", reason),
|
||||||
|
teamRunId,
|
||||||
|
config,
|
||||||
|
createSendContext(runtimeState, memberName),
|
||||||
|
)
|
||||||
|
|
||||||
|
await transitionRuntimeState(teamRunId, (currentRuntimeState) => {
|
||||||
|
const currentRequestIndex = findLatestShutdownRequestIndex(currentRuntimeState, memberName)
|
||||||
|
if (currentRequestIndex < 0) {
|
||||||
|
throw new Error(`shutdown request missing for '${memberName}'`)
|
||||||
|
}
|
||||||
|
|
||||||
|
return {
|
||||||
|
...currentRuntimeState,
|
||||||
|
shutdownRequests: currentRuntimeState.shutdownRequests.map((currentRequest, index) => index === currentRequestIndex
|
||||||
|
? { ...currentRequest, rejectedAt: Date.now(), rejectedReason: reason }
|
||||||
|
: currentRequest),
|
||||||
|
}
|
||||||
|
}, config)
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user