diff --git a/src/hooks/team-session-events/team-member-error-handler.test.ts b/src/hooks/team-session-events/team-member-error-handler.test.ts new file mode 100644 index 000000000..8d1418929 --- /dev/null +++ b/src/hooks/team-session-events/team-member-error-handler.test.ts @@ -0,0 +1,172 @@ +/// + +import { afterEach, describe, expect, test } from "bun:test" +import { randomUUID } from "node:crypto" +import { mkdtemp, mkdir, rm } 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 { + clearTeamSessionRegistry, + registerTeamSession, +} from "../../features/team-mode/team-session-registry" +import type { RuntimeState } from "../../features/team-mode/types" +import { loadRuntimeState, saveRuntimeState } from "../../features/team-mode/team-state-store/store" +import { createTeamMemberErrorHandler } from "./team-member-error-handler" + +const temporaryDirectories: string[] = [] + +async function createTemporaryBaseDir(): Promise { + const baseDir = await mkdtemp(path.join(tmpdir(), "team-member-error-handler-")) + temporaryDirectories.push(baseDir) + return baseDir +} + +function createConfig(baseDir: string): TeamModeConfig { + return TeamModeConfigSchema.parse({ base_dir: baseDir, enabled: true }) +} + +function createRuntimeState(teamRunId: string): RuntimeState { + return { + version: 1, + teamRunId, + teamName: "team-alpha", + specSource: "project", + createdAt: 1, + status: "active", + leadSessionId: "lead-session", + members: [ + { + name: "worker", + sessionId: "member-session", + agentType: "general-purpose", + status: "running", + pendingInjectedMessageIds: [], + }, + ], + shutdownRequests: [], + bounds: { + maxMembers: 8, + maxParallelMembers: 4, + maxMessagesPerRun: 10000, + maxWallClockMinutes: 120, + maxMemberTurns: 500, + }, + } +} + +async function seedRuntimeState(runtimeState: RuntimeState, config: TeamModeConfig): Promise { + await mkdir(path.join(config.base_dir ?? "", "runtime", runtimeState.teamRunId), { recursive: true }) + await saveRuntimeState(runtimeState, config) +} + +afterEach(async () => { + clearTeamSessionRegistry() + await Promise.all(temporaryDirectories.splice(0).map(async (directoryPath) => { + await rm(directoryPath, { recursive: true, force: true }) + })) +}) + +describe("createTeamMemberErrorHandler", () => { + test("marks the matching member errored without changing team status", async () => { + // given + const baseDir = await createTemporaryBaseDir() + const config = createConfig(baseDir) + const teamRunId = randomUUID() + await seedRuntimeState(createRuntimeState(teamRunId), config) + const handler = createTeamMemberErrorHandler(config) + + // when + await handler({ + event: { + type: "session.error", + properties: { sessionID: "member-session", error: new Error("boom") }, + }, + }) + + // then + const runtimeState = await loadRuntimeState(teamRunId, config) + expect(runtimeState.status).toBe("active") + expect(runtimeState.members[0]?.status).toBe("errored") + }) + + test("marks the member errored during the spawn race when the registry tracks the fresh session before disk state persists it", async () => { + // given + const baseDir = await createTemporaryBaseDir() + const config = createConfig(baseDir) + const teamRunId = randomUUID() + await seedRuntimeState({ + ...createRuntimeState(teamRunId), + members: [ + { + name: "worker", + agentType: "general-purpose", + status: "running", + pendingInjectedMessageIds: [], + }, + ], + }, config) + registerTeamSession("member-session", { + teamRunId, + memberName: "worker", + role: "member", + }) + const handler = createTeamMemberErrorHandler(config) + + // when + await handler({ + event: { + type: "session.error", + properties: { sessionID: "member-session", error: new Error("boom") }, + }, + }) + + // then + const runtimeState = await loadRuntimeState(teamRunId, config) + expect(runtimeState.status).toBe("active") + expect(runtimeState.members[0]?.status).toBe("errored") + }) + + test("falls back to disk lookup when the registry points the member session at the wrong teamRunId", async () => { + // given + const baseDir = await createTemporaryBaseDir() + const config = createConfig(baseDir) + const correctTeamRunId = randomUUID() + const wrongTeamRunId = randomUUID() + await seedRuntimeState(createRuntimeState(correctTeamRunId), config) + await seedRuntimeState({ + ...createRuntimeState(wrongTeamRunId), + members: [ + { + name: "worker", + sessionId: "other-session", + agentType: "general-purpose", + status: "running", + pendingInjectedMessageIds: [], + }, + ], + }, config) + registerTeamSession("member-session", { + teamRunId: wrongTeamRunId, + memberName: "worker", + role: "member", + }) + const handler = createTeamMemberErrorHandler(config) + + // when + await handler({ + event: { + type: "session.error", + properties: { sessionID: "member-session", error: new Error("boom") }, + }, + }) + + // then + const correctRuntimeState = await loadRuntimeState(correctTeamRunId, config) + const wrongRuntimeState = await loadRuntimeState(wrongTeamRunId, config) + expect(correctRuntimeState.members[0]?.status).toBe("errored") + expect(wrongRuntimeState.members[0]?.status).toBe("running") + }) +}) diff --git a/src/hooks/team-session-events/team-member-error-handler.ts b/src/hooks/team-session-events/team-member-error-handler.ts new file mode 100644 index 000000000..89dc67601 --- /dev/null +++ b/src/hooks/team-session-events/team-member-error-handler.ts @@ -0,0 +1,53 @@ +import type { TeamModeConfig } from "../../config/schema/team-mode" +import { findResolvedMemberSession } from "../../features/team-mode/member-session-resolution" +import { loadRuntimeState, transitionRuntimeState } from "../../features/team-mode/team-state-store/store" +import { log } from "../../shared/logger" + +type HookInput = { event: { type: string; properties?: unknown } } +export type HookImpl = (input: HookInput) => Promise + +function getErroredSessionID(properties: unknown): string | undefined { + const record = properties as { sessionID?: string } | undefined + return record?.sessionID +} + +export function createTeamMemberErrorHandler(config: TeamModeConfig): HookImpl { + return async ({ event }: HookInput): Promise => { + if (event.type !== "session.error") return + + const erroredSessionID = getErroredSessionID(event.properties) + if (!erroredSessionID) return + + try { + const runtimeMember = await findResolvedMemberSession(erroredSessionID, config, "team member error handler") + if (runtimeMember === null) { + return + } + + const runtimeState = await loadRuntimeState(runtimeMember.teamRunId, config) + await transitionRuntimeState(runtimeState.teamRunId, (currentRuntimeState) => ({ + ...currentRuntimeState, + members: currentRuntimeState.members.map((member) => ( + member.name === runtimeMember.memberName + ? { ...member, status: "errored" } + : member + )), + }), config) + + log("team member session errored", { + event: "team-mode-member-errored", + teamRunId: runtimeState.teamRunId, + teamName: runtimeState.teamName, + memberName: runtimeMember.memberName, + sessionID: erroredSessionID, + runtimeStatus: runtimeState.status, + }) + } catch (error) { + log("team member error handler failed", { + event: "team-mode-member-error-handler-error", + sessionID: erroredSessionID, + error: error instanceof Error ? error.message : String(error), + }) + } + } +}