diff --git a/src/features/team-mode/team-mailbox/ack.ts b/src/features/team-mode/team-mailbox/ack.ts index 9428f7ae5..d2d8c09cd 100644 --- a/src/features/team-mode/team-mailbox/ack.ts +++ b/src/features/team-mode/team-mailbox/ack.ts @@ -17,18 +17,24 @@ export async function ackMessages( for (const messageId of messageIds) { const messageFileName = `${messageId}.json` - const sourcePath = path.join(inboxDir, messageFileName) + const sourcePaths = [ + path.join(inboxDir, messageFileName), + path.join(inboxDir, `.delivering-${messageFileName}`), + ] const targetPath = path.join(processedDir, messageFileName) - try { - await rename(sourcePath, targetPath) - } catch (error) { - const err = error as NodeJS.ErrnoException - if (err.code === "ENOENT") { - continue - } + for (const sourcePath of sourcePaths) { + try { + await rename(sourcePath, targetPath) + break + } catch (error) { + const err = error as NodeJS.ErrnoException + if (err.code === "ENOENT") { + continue + } - throw error + throw error + } } } } diff --git a/src/features/team-mode/tools/messaging.test.ts b/src/features/team-mode/tools/messaging.test.ts index 3b31e0621..bb1474ad2 100644 --- a/src/features/team-mode/tools/messaging.test.ts +++ b/src/features/team-mode/tools/messaging.test.ts @@ -530,7 +530,7 @@ describe("createTeamSendMessageTool", () => { expect(message.from).toBe("m1") }) - test("acks the message after live delivery so the transform hook does not redeliver", async () => { + test("keeps live-delivered messages reserved until the recipient idles", async () => { // given const fixture = await createTeamFixture() const { client } = createRecordingClient() @@ -546,9 +546,13 @@ describe("createTeamSendMessageTool", () => { // then const inboxDir = getInboxDir(resolveBaseDir(fixture.config), fixture.teamRunId, "m2") const inboxEntries = (await readdir(inboxDir)).filter((entry) => entry.endsWith(".json")) - const processedEntries = (await readdir(path.join(inboxDir, "processed"))).filter((entry) => entry.endsWith(".json")) - expect(inboxEntries).toHaveLength(0) - expect(processedEntries).toHaveLength(1) + expect(inboxEntries).toHaveLength(1) + expect(inboxEntries[0]?.startsWith(".delivering-")).toBe(true) + + const { loadRuntimeState: loadState } = await import("../team-state-store/store") + const runtimeState = await loadState(fixture.teamRunId, fixture.config) + const recipient = runtimeState.members.find((member) => member.name === "m2") + expect(recipient?.pendingInjectedMessageIds).toHaveLength(1) }) test("broadcast fans out live delivery to every member except the sender", async () => { diff --git a/src/features/team-mode/tools/messaging.ts b/src/features/team-mode/tools/messaging.ts index d3bcb5c5a..d18791a14 100644 --- a/src/features/team-mode/tools/messaging.ts +++ b/src/features/team-mode/tools/messaging.ts @@ -1,22 +1,20 @@ import { randomUUID } from "node:crypto" -import { tool, type ToolDefinition } from "@opencode-ai/plugin/tool" +import { type ToolDefinition, tool } from "@opencode-ai/plugin/tool" import { z } from "zod" import type { TeamModeConfig } from "../../../config/schema/team-mode" +import { promptAsyncAfterSessionIdle } from "../../../hooks/shared/prompt-async-gate" import { log } from "../../../shared/logger" import { applyMemberSessionRouting, buildMemberPromptBody } from "../member-session-routing" -import { lookupTeamSession } from "../team-session-registry" -import { loadRuntimeState } from "../team-state-store/store" import { buildEnvelope } from "../team-mailbox/poll" import { - commitDeliveryReservation, releaseDeliveryReservation, reserveMessageForDelivery, } from "../team-mailbox/reservation" import { BroadcastNotPermittedError, sendMessage } from "../team-mailbox/send" -import { promptAsyncAfterSessionIdle } from "../../../hooks/shared/prompt-async-gate" - +import { lookupTeamSession } from "../team-session-registry" +import { loadRuntimeState, transitionRuntimeState } from "../team-state-store/store" import type { Message } from "../types" import { MessageSchema } from "../types" @@ -135,6 +133,25 @@ async function releaseReservationSafely( } } +async function markLiveDeliveryPending( + teamRunId: string, + recipientName: string, + messageId: string, + config: TeamModeConfig, +): Promise { + await transitionRuntimeState(teamRunId, (currentRuntimeState) => ({ + ...currentRuntimeState, + members: currentRuntimeState.members.map((member) => ( + member.name === recipientName + ? { + ...member, + pendingInjectedMessageIds: Array.from(new Set([...member.pendingInjectedMessageIds, messageId])), + } + : member + )), + }), config) +} + async function deliverLive( client: LiveDeliveryClient, message: Message, @@ -207,8 +224,8 @@ async function deliverLive( }) continue } - await commitDeliveryReservation(reservation) - log("[team-mailbox] live delivery committed", { + await markLiveDeliveryPending(teamRunId, recipientName, message.messageId, config) + log("[team-mailbox] live delivery reserved until recipient idle", { teamRunId, recipient: recipientName, recipientSessionId, diff --git a/src/hooks/team-session-events/team-idle-wake-hint.test.ts b/src/hooks/team-session-events/team-idle-wake-hint.test.ts index 67b81473d..905ea2cfa 100644 --- a/src/hooks/team-session-events/team-idle-wake-hint.test.ts +++ b/src/hooks/team-session-events/team-idle-wake-hint.test.ts @@ -80,6 +80,42 @@ function createRuntimeState(teamRunId: string, pendingInjectedMessageIds: string } } +function createLeaderRuntimeState(teamRunId: string, pendingInjectedMessageIds: string[]): RuntimeState { + return { + version: 1, + teamRunId, + teamName: "team-alpha", + specSource: "project", + createdAt: 1, + status: "active", + leadSessionId: "lead-session", + members: [ + { + name: "lead", + sessionId: "lead-session", + agentType: "leader", + status: "idle", + pendingInjectedMessageIds, + }, + { + name: "worker", + sessionId: "member-session", + agentType: "general-purpose", + status: "idle", + 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) @@ -103,6 +139,46 @@ async function seedUnreadMessage( }, teamRunId, config, { isLead: true, activeMembers: ["worker"] }) } +async function seedReservedUnreadMessage( + teamRunId: string, + config: TeamModeConfig, + messageId: string, + body: string, + timestamp: number, +): Promise { + await sendMessage({ + version: 1, + messageId, + from: "lead", + to: "worker", + kind: "message", + body, + timestamp, + }, teamRunId, config, { + isLead: true, + activeMembers: ["worker"], + reservedRecipients: new Set(["worker"]), + }) +} + +async function seedLeadUnreadMessage( + teamRunId: string, + config: TeamModeConfig, + messageId: string, + body: string, + timestamp: number, +): Promise { + await sendMessage({ + version: 1, + messageId, + from: "worker", + to: "lead", + kind: "message", + body, + timestamp, + }, teamRunId, config, { isLead: false, activeMembers: ["lead"] }) +} + afterEach(async () => { clearTeamSessionRegistry() SessionCategoryRegistry.clear() @@ -414,6 +490,81 @@ describe("createTeamIdleWakeHint", () => { expect(processedEntries.sort()).toEqual(messageIds.map((messageId) => `${messageId}.json`).sort()) }) + test("acks pending reserved live-delivery messages on idle", async () => { + // given + const baseDir = await createTemporaryBaseDir() + const config = createConfig(baseDir) + const teamRunId = randomUUID() + const messageId = randomUUID() + await seedRuntimeState(createRuntimeState(teamRunId, [messageId]), config) + await seedReservedUnreadMessage(teamRunId, config, messageId, "live delivery body", 100) + + const handler = createTeamIdleWakeHint({ + directory: "/tmp/project", + client: { session: { promptAsync: mock(async (_input: WakeHintPromptInput) => ({})) } }, + }, config) + + // when + await handler({ + event: { + type: "session.idle", + properties: { sessionID: "member-session" }, + }, + }) + + // then + const runtimeState = await loadRuntimeState(teamRunId, config) + expect(runtimeState.members[0]?.pendingInjectedMessageIds).toEqual([]) + + const inboxDir = getInboxDir(resolveBaseDir(config), teamRunId, "worker") + const inboxEntries = await readdir(inboxDir) + expect(inboxEntries).not.toContain(`.delivering-${messageId}.json`) + + const processedEntries = await readdir(path.join(inboxDir, "processed")) + expect(processedEntries).toContain(`${messageId}.json`) + }) + + test("acks pending lead messages on idle without sending a wake hint", async () => { + // given + const baseDir = await createTemporaryBaseDir() + const config = createConfig(baseDir) + const teamRunId = randomUUID() + const messageIds = [randomUUID(), randomUUID()] + await seedRuntimeState(createLeaderRuntimeState(teamRunId, messageIds), config) + await seedLeadUnreadMessage(teamRunId, config, messageIds[0], "one", 100) + await seedLeadUnreadMessage(teamRunId, config, messageIds[1], "two", 200) + + const ackSpy = spyOn(ackModule, "ackMessages") + const promptAsyncSpy = mock(async (_input: WakeHintPromptInput) => ({})) + const handler = createTeamIdleWakeHint({ + directory: "/tmp/project", + client: { session: { promptAsync: promptAsyncSpy } }, + }, config) + + // when + await handler({ + event: { + type: "session.idle", + properties: { sessionID: "lead-session" }, + }, + }) + + // then + expect(ackSpy).toHaveBeenCalledTimes(1) + expect(ackSpy).toHaveBeenCalledWith(teamRunId, "lead", messageIds, config) + expect(promptAsyncSpy).not.toHaveBeenCalled() + + const runtimeState = await loadRuntimeState(teamRunId, config) + expect(runtimeState.members[0]?.pendingInjectedMessageIds).toEqual([]) + + const inboxDir = getInboxDir(resolveBaseDir(config), teamRunId, "lead") + const inboxEntries = await readdir(inboxDir) + expect(inboxEntries).toContain("processed") + + const processedEntries = await readdir(path.join(inboxDir, "processed")) + expect(processedEntries.sort()).toEqual(messageIds.map((messageId) => `${messageId}.json`).sort()) + }) + test("sends a wake hint during the spawn race when the registry tracks the fresh member session before disk state persists it", async () => { // given const baseDir = await createTemporaryBaseDir() diff --git a/src/hooks/team-session-events/team-idle-wake-hint.ts b/src/hooks/team-session-events/team-idle-wake-hint.ts index 347da7f73..c9f8fbd95 100644 --- a/src/hooks/team-session-events/team-idle-wake-hint.ts +++ b/src/hooks/team-session-events/team-idle-wake-hint.ts @@ -1,12 +1,12 @@ import type { TeamModeConfig } from "../../config/schema/team-mode" -import { ackMessages } from "../../features/team-mode/team-mailbox/ack" -import { listUnreadMessages } from "../../features/team-mode/team-mailbox/inbox" -import { loadRuntimeState, listActiveTeams, transitionRuntimeState } from "../../features/team-mode/team-state-store/store" import { findResolvedMemberSession } from "../../features/team-mode/member-session-resolution" import { applyMemberSessionRouting, buildMemberPromptBody, } from "../../features/team-mode/member-session-routing" +import { ackMessages } from "../../features/team-mode/team-mailbox/ack" +import { listUnreadMessages } from "../../features/team-mode/team-mailbox/inbox" +import { loadRuntimeState, transitionRuntimeState } from "../../features/team-mode/team-state-store/store" import { resolveSessionEventID } from "../../shared/event-session-id" import { log } from "../../shared/logger" import { promptAsyncAfterSessionIdle } from "../shared/prompt-async-gate" @@ -59,7 +59,7 @@ export function createTeamIdleWakeHint(ctx: TeamIdleWakeHintContext, config: Tea const runtimeState = await loadRuntimeState(runtimeMember.teamRunId, config) const memberEntry = runtimeState.members.find((member) => member.name === runtimeMember.memberName) - if (!memberEntry || memberEntry.agentType === "leader") { + if (!memberEntry) { return } @@ -88,6 +88,17 @@ export function createTeamIdleWakeHint(ctx: TeamIdleWakeHintContext, config: Tea return } + if (memberEntry.agentType === "leader") { + log("team lead idle handled without wake hint", { + event: "team-mode-lead-idle-ack-only", + teamRunId: runtimeState.teamRunId, + memberName: memberEntry.name, + sessionID, + ackedCount: pendingInjectedMessageIds.length, + }) + return + } + if (typeof ctx.client.session.promptAsync !== "function") { log("team idle wake hint skipped without promptAsync", { event: "team-mode-idle-wake-hint-skipped", 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 index 8d1418929..39fbfe452 100644 --- a/src/hooks/team-session-events/team-member-error-handler.test.ts +++ b/src/hooks/team-session-events/team-member-error-handler.test.ts @@ -2,12 +2,14 @@ import { afterEach, describe, expect, test } from "bun:test" import { randomUUID } from "node:crypto" -import { mkdtemp, mkdir, rm } from "node:fs/promises" +import { mkdtemp, mkdir, readdir, 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 { sendMessage } from "../../features/team-mode/team-mailbox/send" +import { getInboxDir, resolveBaseDir } from "../../features/team-mode/team-registry/paths" import { clearTeamSessionRegistry, registerTeamSession, @@ -57,11 +59,37 @@ function createRuntimeState(teamRunId: string): RuntimeState { } } +function createRuntimeStateWithPendingMessage(teamRunId: string, messageId: string): RuntimeState { + const runtimeState = createRuntimeState(teamRunId) + const worker = runtimeState.members[0] + if (worker === undefined) { + throw new Error("worker member missing from fixture") + } + worker.pendingInjectedMessageIds = [messageId] + return runtimeState +} + async function seedRuntimeState(runtimeState: RuntimeState, config: TeamModeConfig): Promise { await mkdir(path.join(config.base_dir ?? "", "runtime", runtimeState.teamRunId), { recursive: true }) await saveRuntimeState(runtimeState, config) } +async function seedReservedMessage(teamRunId: string, config: TeamModeConfig, messageId: string): Promise { + await sendMessage({ + version: 1, + messageId, + from: "lead", + to: "worker", + kind: "message", + body: "pending live delivery", + timestamp: 1, + }, teamRunId, config, { + isLead: true, + activeMembers: ["worker"], + reservedRecipients: new Set(["worker"]), + }) +} + afterEach(async () => { clearTeamSessionRegistry() await Promise.all(temporaryDirectories.splice(0).map(async (directoryPath) => { @@ -169,4 +197,33 @@ describe("createTeamMemberErrorHandler", () => { expect(correctRuntimeState.members[0]?.status).toBe("errored") expect(wrongRuntimeState.members[0]?.status).toBe("running") }) + + test("requeues pending live-delivery messages when the recipient session errors before idle ack", async () => { + // given + const baseDir = await createTemporaryBaseDir() + const config = createConfig(baseDir) + const teamRunId = randomUUID() + const messageId = randomUUID() + await seedRuntimeState(createRuntimeStateWithPendingMessage(teamRunId, messageId), config) + await seedReservedMessage(teamRunId, config, messageId) + const handler = createTeamMemberErrorHandler(config) + + // when + await handler({ + event: { + type: "session.error", + properties: { sessionID: "member-session", error: new Error("late prompt failure") }, + }, + }) + + // then + const runtimeState = await loadRuntimeState(teamRunId, config) + expect(runtimeState.members[0]?.status).toBe("errored") + expect(runtimeState.members[0]?.pendingInjectedMessageIds).toEqual([]) + + const inboxEntries = await readdir(getInboxDir(resolveBaseDir(config), teamRunId, "worker")) + expect(inboxEntries).toContain(`${messageId}.json`) + expect(inboxEntries).not.toContain(`.delivering-${messageId}.json`) + expect(inboxEntries).not.toContain("processed") + }) }) diff --git a/src/hooks/team-session-events/team-member-error-handler.ts b/src/hooks/team-session-events/team-member-error-handler.ts index c7f669c8c..0fe927b49 100644 --- a/src/hooks/team-session-events/team-member-error-handler.ts +++ b/src/hooks/team-session-events/team-member-error-handler.ts @@ -1,5 +1,9 @@ import type { TeamModeConfig } from "../../config/schema/team-mode" import { findResolvedMemberSession } from "../../features/team-mode/member-session-resolution" +import { + releaseDeliveryReservation, + reserveMessageForDelivery, +} from "../../features/team-mode/team-mailbox/reservation" import { loadRuntimeState, transitionRuntimeState } from "../../features/team-mode/team-state-store/store" import { resolveSessionEventID } from "../../shared/event-session-id" import { log } from "../../shared/logger" @@ -11,6 +15,22 @@ function getErroredSessionID(properties: unknown): string | undefined { return resolveSessionEventID(properties) } +async function requeuePendingLiveDeliveries( + teamRunId: string, + memberName: string, + messageIds: readonly string[], + config: TeamModeConfig, +): Promise { + for (const messageId of messageIds) { + const reservation = await reserveMessageForDelivery(teamRunId, memberName, messageId, config) + if (reservation === null) { + continue + } + + await releaseDeliveryReservation(reservation) + } +} + export function createTeamMemberErrorHandler(config: TeamModeConfig): HookImpl { return async ({ event }: HookInput): Promise => { if (event.type !== "session.error") return @@ -25,11 +45,19 @@ export function createTeamMemberErrorHandler(config: TeamModeConfig): HookImpl { } const runtimeState = await loadRuntimeState(runtimeMember.teamRunId, config) + const memberEntry = runtimeState.members.find((member) => member.name === runtimeMember.memberName) + const pendingInjectedMessageIds = memberEntry?.pendingInjectedMessageIds ?? [] + await requeuePendingLiveDeliveries( + runtimeState.teamRunId, + runtimeMember.memberName, + pendingInjectedMessageIds, + config, + ) await transitionRuntimeState(runtimeState.teamRunId, (currentRuntimeState) => ({ ...currentRuntimeState, members: currentRuntimeState.members.map((member) => ( member.name === runtimeMember.memberName - ? { ...member, status: "errored" } + ? { ...member, status: "errored", pendingInjectedMessageIds: [] } : member )), }), config)