diff --git a/src/hooks/team-session-events/team-member-status-handler.test.ts b/src/hooks/team-session-events/team-member-status-handler.test.ts new file mode 100644 index 000000000..85f1d87cc --- /dev/null +++ b/src/hooks/team-session-events/team-member-status-handler.test.ts @@ -0,0 +1,221 @@ +/// + +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, RuntimeStateMember } from "../../features/team-mode/types" +import { loadRuntimeState, saveRuntimeState } from "../../features/team-mode/team-state-store/store" +import { createTeamMemberStatusHandler } from "./team-member-status-handler" + +const temporaryDirectories: string[] = [] + +async function createTemporaryBaseDir(): Promise { + const baseDir = await mkdtemp(path.join(tmpdir(), "team-member-status-handler-")) + temporaryDirectories.push(baseDir) + return baseDir +} + +function createConfig(baseDir: string): TeamModeConfig { + return TeamModeConfigSchema.parse({ base_dir: baseDir, enabled: true }) +} + +function buildMember(overrides?: Partial): RuntimeStateMember { + return { + name: "worker", + sessionId: "member-session", + agentType: "general-purpose", + status: "running", + pendingInjectedMessageIds: [], + ...overrides, + } +} + +function createRuntimeState(teamRunId: string, member: RuntimeStateMember = buildMember()): RuntimeState { + return { + version: 1, + teamRunId, + teamName: "team-alpha", + specSource: "project", + createdAt: 1, + status: "active", + leadSessionId: "lead-session", + members: [member], + 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("createTeamMemberStatusHandler", () => { + test("transitions a running member to idle when its session becomes idle", async () => { + // given + const baseDir = await createTemporaryBaseDir() + const config = createConfig(baseDir) + const teamRunId = randomUUID() + await seedRuntimeState(createRuntimeState(teamRunId, buildMember({ status: "running" })), config) + const handler = createTeamMemberStatusHandler(config) + + // when + await handler({ event: { type: "session.idle", properties: { sessionID: "member-session" } } }) + + // then + const runtimeState = await loadRuntimeState(teamRunId, config) + expect(runtimeState.members[0]?.status).toBe("idle") + }) + + test("leaves an already-idle member untouched on a subsequent session.idle", async () => { + // given + const baseDir = await createTemporaryBaseDir() + const config = createConfig(baseDir) + const teamRunId = randomUUID() + await seedRuntimeState(createRuntimeState(teamRunId, buildMember({ status: "idle" })), config) + const handler = createTeamMemberStatusHandler(config) + + // when + await handler({ event: { type: "session.idle", properties: { sessionID: "member-session" } } }) + + // then + const runtimeState = await loadRuntimeState(teamRunId, config) + expect(runtimeState.members[0]?.status).toBe("idle") + }) + + test("never overrides a terminal errored status on session.idle", async () => { + // given + const baseDir = await createTemporaryBaseDir() + const config = createConfig(baseDir) + const teamRunId = randomUUID() + await seedRuntimeState(createRuntimeState(teamRunId, buildMember({ status: "errored" })), config) + const handler = createTeamMemberStatusHandler(config) + + // when + await handler({ event: { type: "session.idle", properties: { sessionID: "member-session" } } }) + + // then + const runtimeState = await loadRuntimeState(teamRunId, config) + expect(runtimeState.members[0]?.status).toBe("errored") + }) + + test("marks a running member completed when its session is deleted", async () => { + // given + const baseDir = await createTemporaryBaseDir() + const config = createConfig(baseDir) + const teamRunId = randomUUID() + await seedRuntimeState(createRuntimeState(teamRunId, buildMember({ status: "running" })), config) + const handler = createTeamMemberStatusHandler(config) + + // when + await handler({ event: { type: "session.deleted", properties: { info: { id: "member-session" } } } }) + + // then + const runtimeState = await loadRuntimeState(teamRunId, config) + expect(runtimeState.members[0]?.status).toBe("completed") + }) + + test("marks an idle member completed when its session is deleted", async () => { + // given + const baseDir = await createTemporaryBaseDir() + const config = createConfig(baseDir) + const teamRunId = randomUUID() + await seedRuntimeState(createRuntimeState(teamRunId, buildMember({ status: "idle" })), config) + const handler = createTeamMemberStatusHandler(config) + + // when + await handler({ event: { type: "session.deleted", properties: { info: { id: "member-session" } } } }) + + // then + const runtimeState = await loadRuntimeState(teamRunId, config) + expect(runtimeState.members[0]?.status).toBe("completed") + }) + + test("preserves a terminal errored status even when the session is deleted", async () => { + // given + const baseDir = await createTemporaryBaseDir() + const config = createConfig(baseDir) + const teamRunId = randomUUID() + await seedRuntimeState(createRuntimeState(teamRunId, buildMember({ status: "errored" })), config) + const handler = createTeamMemberStatusHandler(config) + + // when + await handler({ event: { type: "session.deleted", properties: { info: { id: "member-session" } } } }) + + // then + const runtimeState = await loadRuntimeState(teamRunId, config) + expect(runtimeState.members[0]?.status).toBe("errored") + }) + + test("ignores session.idle events for sessions that are not team members", async () => { + // given + const baseDir = await createTemporaryBaseDir() + const config = createConfig(baseDir) + const teamRunId = randomUUID() + await seedRuntimeState(createRuntimeState(teamRunId), config) + const handler = createTeamMemberStatusHandler(config) + + // when + await handler({ event: { type: "session.idle", properties: { sessionID: "unknown-session" } } }) + + // then + const runtimeState = await loadRuntimeState(teamRunId, config) + expect(runtimeState.members[0]?.status).toBe("running") + }) + + test("ignores session.deleted events when the deleted session is the team lead", async () => { + // given + const baseDir = await createTemporaryBaseDir() + const config = createConfig(baseDir) + const teamRunId = randomUUID() + await seedRuntimeState(createRuntimeState(teamRunId), config) + registerTeamSession("lead-session", { teamRunId, memberName: "lead", role: "lead" }) + const handler = createTeamMemberStatusHandler(config) + + // when + await handler({ event: { type: "session.deleted", properties: { info: { id: "lead-session" } } } }) + + // then + const runtimeState = await loadRuntimeState(teamRunId, config) + expect(runtimeState.members[0]?.status).toBe("running") + }) + + test("uses the in-memory registry to recognize a fresh session during the spawn race", async () => { + // given + const baseDir = await createTemporaryBaseDir() + const config = createConfig(baseDir) + const teamRunId = randomUUID() + await seedRuntimeState(createRuntimeState(teamRunId, buildMember({ sessionId: undefined, status: "running" })), config) + registerTeamSession("member-session", { teamRunId, memberName: "worker", role: "member" }) + const handler = createTeamMemberStatusHandler(config) + + // when + await handler({ event: { type: "session.idle", properties: { sessionID: "member-session" } } }) + + // then + const runtimeState = await loadRuntimeState(teamRunId, config) + expect(runtimeState.members[0]?.status).toBe("idle") + }) +}) diff --git a/src/hooks/team-session-events/team-member-status-handler.ts b/src/hooks/team-session-events/team-member-status-handler.ts new file mode 100644 index 000000000..3e31173c7 --- /dev/null +++ b/src/hooks/team-session-events/team-member-status-handler.ts @@ -0,0 +1,93 @@ +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 type { RuntimeStateMember } from "../../features/team-mode/types" +import { log } from "../../shared/logger" + +type HookInput = { event: { type: string; properties?: unknown } } +export type HookImpl = (input: HookInput) => Promise + +type MemberStatus = RuntimeStateMember["status"] + +const IDLE_TRANSITION_SOURCE_STATUSES: ReadonlySet = new Set(["running"]) +const COMPLETED_TRANSITION_SOURCE_STATUSES: ReadonlySet = new Set(["running", "idle", "pending"]) + +function getSessionIDFromIdleEvent(properties: unknown): string | undefined { + const record = properties as { sessionID?: string } | undefined + return record?.sessionID +} + +function getSessionIDFromDeletedEvent(properties: unknown): string | undefined { + const record = properties as { info?: { id?: string } } | undefined + return record?.info?.id +} + +async function transitionMemberStatus( + runtimeMember: { teamRunId: string; memberName: string }, + allowedSources: ReadonlySet, + nextStatus: MemberStatus, + config: TeamModeConfig, + sessionID: string, + eventLabel: string, +): Promise { + const runtimeState = await loadRuntimeState(runtimeMember.teamRunId, config) + const currentEntry = runtimeState.members.find((member) => member.name === runtimeMember.memberName) + if (currentEntry === undefined) return + if (!allowedSources.has(currentEntry.status)) return + + await transitionRuntimeState(runtimeState.teamRunId, (currentRuntimeState) => ({ + ...currentRuntimeState, + members: currentRuntimeState.members.map((member) => ( + member.name === runtimeMember.memberName + ? { ...member, status: nextStatus } + : member + )), + }), config) + + log(`team member ${eventLabel}`, { + event: `team-mode-member-${eventLabel}`, + teamRunId: runtimeState.teamRunId, + teamName: runtimeState.teamName, + memberName: runtimeMember.memberName, + sessionID, + previousStatus: currentEntry.status, + nextStatus, + }) +} + +export function createTeamMemberStatusHandler(config: TeamModeConfig): HookImpl { + return async ({ event }: HookInput): Promise => { + if (event.type === "session.idle") { + const sessionID = getSessionIDFromIdleEvent(event.properties) + if (!sessionID) return + try { + const runtimeMember = await findResolvedMemberSession(sessionID, config, "team member status handler") + if (runtimeMember === null) return + await transitionMemberStatus(runtimeMember, IDLE_TRANSITION_SOURCE_STATUSES, "idle", config, sessionID, "idled") + } catch (error) { + log("team member status handler failed on session.idle", { + event: "team-mode-member-status-handler-error", + sessionID, + error: error instanceof Error ? error.message : String(error), + }) + } + return + } + + if (event.type === "session.deleted") { + const sessionID = getSessionIDFromDeletedEvent(event.properties) + if (!sessionID) return + try { + const runtimeMember = await findResolvedMemberSession(sessionID, config, "team member status handler") + if (runtimeMember === null) return + await transitionMemberStatus(runtimeMember, COMPLETED_TRANSITION_SOURCE_STATUSES, "completed", config, sessionID, "completed") + } catch (error) { + log("team member status handler failed on session.deleted", { + event: "team-mode-member-status-handler-error", + sessionID, + error: error instanceof Error ? error.message : String(error), + }) + } + } + } +}