From b2918fd4dbf2d5e65e302633cbed1a659ce1113e Mon Sep 17 00:00:00 2001 From: YeonGyu-Kim Date: Tue, 19 May 2026 18:23:57 +0900 Subject: [PATCH] fix(team-mode): close peer message delivery races --- .../background-agent/parent-wake-notifier.ts | 1 - .../team-mode/team-mailbox/poll.test.ts | 24 ++++ src/features/team-mode/team-mailbox/poll.ts | 83 +++++++------- .../team-mode/tools/messaging.test.ts | 70 +++++++++++- src/features/team-mode/tools/messaging.ts | 56 +++++++++- src/hooks/shared/prompt-async-gate.test.ts | 40 +++++++ .../team-idle-wake-hint.test.ts | 42 +++++++ .../team-idle-wake-hint.ts | 15 +++ .../team-member-error-handler.test.ts | 85 ++++++++++++++ .../team-member-error-handler.ts | 105 +++++++++++++++++- src/plugin/event.ts | 2 +- src/shared/prompt-async-gate.ts | 4 +- src/shared/prompt-async-route-audit.test.ts | 12 ++ 13 files changed, 492 insertions(+), 47 deletions(-) diff --git a/src/features/background-agent/parent-wake-notifier.ts b/src/features/background-agent/parent-wake-notifier.ts index 283791399..d150f9ee4 100644 --- a/src/features/background-agent/parent-wake-notifier.ts +++ b/src/features/background-agent/parent-wake-notifier.ts @@ -196,7 +196,6 @@ export class ParentWakeNotifier { sessionID, source: "background-agent-parent-wake", settleMs: 0, - postDispatchHoldMs: 250, queueBehavior: "defer", checkToolState: !toolWaitDecision.skipPromptGateToolStateCheck, input: { diff --git a/src/features/team-mode/team-mailbox/poll.test.ts b/src/features/team-mode/team-mailbox/poll.test.ts index 424c60eaa..a92b926d9 100644 --- a/src/features/team-mode/team-mailbox/poll.test.ts +++ b/src/features/team-mode/team-mailbox/poll.test.ts @@ -79,6 +79,30 @@ describe("pollAndBuildInjection", () => { }) }) + test("#given concurrent transforms for one turn #when mailbox injection is claimed #then only one call injects the peer message", async () => { + // given + const { teamRunId, config } = await setupRuntime(["m1"]) + + await sendMessage({ + version: 1, + messageId: randomUUID(), + from: "lead", + to: "m1", + kind: "message", + body: "race", + timestamp: 100, + }, teamRunId, config, { isLead: true, activeMembers: ["m1"] }) + + // when + const results = await Promise.all(Array.from({ length: 8 }, () => + pollAndBuildInjection("session-1", "m1", teamRunId, config, "turn-race") + )) + + // then + expect(results.filter((result) => result.injected)).toHaveLength(1) + expect(results.filter((result) => !result.injected)).toHaveLength(7) + }) + test("wraps hostile message bodies in a literal peer_message envelope", async () => { // given const { teamRunId, config } = await setupRuntime(["m1"]) diff --git a/src/features/team-mode/team-mailbox/poll.ts b/src/features/team-mode/team-mailbox/poll.ts index 2e8dcd8f4..27cbd5929 100644 --- a/src/features/team-mode/team-mailbox/poll.ts +++ b/src/features/team-mode/team-mailbox/poll.ts @@ -48,47 +48,54 @@ export async function pollAndBuildInjection( config: TeamModeConfig, turnMarker: string, ): Promise { - const runtimeState = await loadRuntimeState(teamRunId, config) - const runtimeMember = runtimeState.members.find((member) => member.name === memberName) - if (runtimeMember === undefined) { - throw new Error(`runtime member not found for session ${sessionID}: ${memberName}`) - } + const unreadMessages = await listUnreadMessages(teamRunId, memberName, config) + let result: InjectionResult | undefined - if (runtimeMember.lastInjectedTurnMarker === turnMarker) { - return { injected: false, messageIds: [], reason: "already injected this turn" } - } - - const pendingMessageIds = new Set(runtimeMember.pendingInjectedMessageIds) - const unreadMessages = (await listUnreadMessages(teamRunId, memberName, config)) - .filter((message) => !pendingMessageIds.has(message.messageId)) - if (unreadMessages.length === 0) { - if (pendingMessageIds.size > 0) { - return { injected: false, messageIds: [], reason: "pending ack" } + await transitionRuntimeState(teamRunId, (currentRuntimeState) => { + const runtimeMember = currentRuntimeState.members.find((member) => member.name === memberName) + if (runtimeMember === undefined) { + throw new Error(`runtime member not found for session ${sessionID}: ${memberName}`) } - return { injected: false, messageIds: [], reason: "no unread" } + if (runtimeMember.lastInjectedTurnMarker === turnMarker) { + result = { injected: false, messageIds: [], reason: "already injected this turn" } + return currentRuntimeState + } + + const pendingMessageIds = new Set(runtimeMember.pendingInjectedMessageIds) + const injectableMessages = unreadMessages.filter((message) => !pendingMessageIds.has(message.messageId)) + if (injectableMessages.length === 0) { + result = pendingMessageIds.size > 0 + ? { injected: false, messageIds: [], reason: "pending ack" } + : { injected: false, messageIds: [], reason: "no unread" } + return currentRuntimeState + } + + const messageIds: string[] = [] + const envelopes: string[] = [] + for (const unreadMessage of injectableMessages) { + messageIds.push(unreadMessage.messageId) + envelopes.push(buildEnvelope(unreadMessage)) + } + result = { injected: true, content: envelopes.join("\n"), messageIds } + + return { + ...currentRuntimeState, + members: currentRuntimeState.members.map((member) => ( + member.name === memberName + ? { + ...member, + lastInjectedTurnMarker: turnMarker, + pendingInjectedMessageIds: Array.from(new Set([...member.pendingInjectedMessageIds, ...messageIds])), + } + : member + )), + } + }, config) + + if (result === undefined) { + throw new Error(`mailbox injection claim failed for session ${sessionID}: ${memberName}`) } - const messageIds: string[] = [] - const envelopes: string[] = [] - for (const unreadMessage of unreadMessages) { - messageIds.push(unreadMessage.messageId) - envelopes.push(buildEnvelope(unreadMessage)) - } - const content = envelopes.join("\n") - - await transitionRuntimeState(teamRunId, (currentRuntimeState) => ({ - ...currentRuntimeState, - members: currentRuntimeState.members.map((member) => ( - member.name === memberName - ? { - ...member, - lastInjectedTurnMarker: turnMarker, - pendingInjectedMessageIds: Array.from(new Set([...member.pendingInjectedMessageIds, ...messageIds])), - } - : member - )), - }), config) - - return { injected: true, content, messageIds } + return result } diff --git a/src/features/team-mode/tools/messaging.test.ts b/src/features/team-mode/tools/messaging.test.ts index 0c12b29b8..36ba69ebd 100644 --- a/src/features/team-mode/tools/messaging.test.ts +++ b/src/features/team-mode/tools/messaging.test.ts @@ -1,7 +1,7 @@ /// import { afterEach, describe, expect, test } from "bun:test" -import { mkdtemp, readdir, readFile } from "node:fs/promises" +import { mkdtemp, readdir, readFile, rm } from "node:fs/promises" import { randomUUID } from "node:crypto" import { tmpdir } from "node:os" import path from "node:path" @@ -673,6 +673,74 @@ describe("createTeamSendMessageTool", () => { expect(recipient?.pendingInjectedMessageIds).toHaveLength(1) }) + test("#given live delivery prompt dispatches but pending mark fails #when delivery finishes #then the message is not re-exposed as unread", async () => { + // given + const fixture = await createTeamFixture() + let promptCalls = 0 + const client = { + session: { + promptAsync: async () => { + promptCalls += 1 + await rm(path.join(resolveBaseDir(fixture.config), "runtime", fixture.teamRunId, "state.json")) + }, + }, + } satisfies LiveDeliveryClient + const liveTool = createTeamSendMessageTool(fixture.config, client) + + // when + await liveTool.execute({ + teamRunId: fixture.teamRunId, + to: "m2", + body: "accepted before state vanished", + }, fixture.toolContext(fixture.memberOneSessionId)) + + // then + expect(promptCalls).toBe(1) + const unread = await listUnreadMessages(fixture.teamRunId, "m2", fixture.config) + expect(unread).toHaveLength(0) + + const inboxDir = getInboxDir(resolveBaseDir(fixture.config), fixture.teamRunId, "m2") + const inboxEntries = (await readdir(inboxDir)).filter((entry) => entry.endsWith(".json")) + expect(inboxEntries).toHaveLength(1) + expect(inboxEntries[0]?.startsWith(".delivering-")).toBe(true) + }) + + test("#given live delivery cannot reload runtime after pre-reserve #when delivery aborts #then the message is released for mailbox injection", async () => { + // given + const fixture = await createTeamFixture() + const { loadRuntimeState: loadState } = await import("../team-state-store/store") + const runtimeState = await loadState(fixture.teamRunId, fixture.config) + let loadCount = 0 + const deps = { + loadRuntimeState: async () => { + loadCount += 1 + if (loadCount === 3) { + throw new Error("runtime reload failed") + } + return runtimeState + }, + } + const { client, calls } = createRecordingClient() + const liveTool = createTeamSendMessageTool(fixture.config, client, deps) + + // when + await liveTool.execute({ + teamRunId: fixture.teamRunId, + to: "m2", + body: "fallback unread", + }, fixture.toolContext(fixture.memberOneSessionId)) + + // then + expect(calls).toHaveLength(0) + const unread = await listUnreadMessages(fixture.teamRunId, "m2", fixture.config) + expect(unread).toHaveLength(1) + expect(unread[0]?.body).toBe("fallback unread") + + const inboxEntries = await readdir(getInboxDir(resolveBaseDir(fixture.config), fixture.teamRunId, "m2")) + expect(inboxEntries.filter((entry) => entry.endsWith(".json") && !entry.startsWith("."))).toHaveLength(1) + expect(inboxEntries.some((entry) => entry.startsWith(".delivering-"))).toBe(false) + }) + test("reserves the message during live delivery so concurrent listings cannot surface it", async () => { // given const fixture = await createTeamFixture() diff --git a/src/features/team-mode/tools/messaging.ts b/src/features/team-mode/tools/messaging.ts index c0c0027c3..d08d18ca5 100644 --- a/src/features/team-mode/tools/messaging.ts +++ b/src/features/team-mode/tools/messaging.ts @@ -161,6 +161,22 @@ async function markLiveDeliveryPending( }), config) } +async function releaseReservationsForRecipients( + teamRunId: string, + recipientNames: readonly string[], + messageId: string, + config: TeamModeConfig, +): Promise { + for (const recipientName of recipientNames) { + const reservation = await reserveMessageForDelivery(teamRunId, recipientName, messageId, config) + await releaseReservationSafely(reservation, { + teamRunId, + recipient: recipientName, + messageId, + }) + } +} + async function deliverLive( client: LiveDeliveryClient, message: Message, @@ -170,7 +186,19 @@ async function deliverLive( directory: string, deps: TeamSendMessageToolDeps, ): Promise { - const runtimeState = await deps.loadRuntimeState(teamRunId, config) + let runtimeState: RuntimeState + try { + runtimeState = await deps.loadRuntimeState(teamRunId, config) + } catch (error) { + await releaseReservationsForRecipients(teamRunId, deliveredTo, message.messageId, config) + log("[team-mailbox] live delivery unavailable after pre-reserve, released recipients to inbox", { + teamRunId, + messageId: message.messageId, + deliveredTo, + error: error instanceof Error ? error.message : String(error), + }) + return + } const envelope = buildEnvelope(message) for (const recipientName of deliveredTo) { @@ -246,7 +274,18 @@ async function deliverLive( }, }) if (promptResult.status === "failed" && isAmbiguousPromptDispatchFailure(promptResult.error)) { - await markLiveDeliveryPending(teamRunId, recipientName, message.messageId, config) + try { + await markLiveDeliveryPending(teamRunId, recipientName, message.messageId, config) + } catch (markError) { + log("[team-mailbox] live delivery prompt may be accepted but pending mark failed, keeping reservation hidden", { + teamRunId, + recipient: recipientName, + recipientSessionId, + messageId: message.messageId, + error: markError instanceof Error ? markError.message : String(markError), + }) + continue + } log("[team-mailbox] live delivery prompt failed after dispatch attempt, keeping reservation pending", { teamRunId, recipient: recipientName, @@ -271,7 +310,18 @@ async function deliverLive( }) continue } - await markLiveDeliveryPending(teamRunId, recipientName, message.messageId, config) + try { + await markLiveDeliveryPending(teamRunId, recipientName, message.messageId, config) + } catch (markError) { + log("[team-mailbox] live delivery prompt dispatched but pending mark failed, keeping reservation hidden", { + teamRunId, + recipient: recipientName, + recipientSessionId, + messageId: message.messageId, + error: markError instanceof Error ? markError.message : String(markError), + }) + continue + } log("[team-mailbox] live delivery reserved until recipient idle", { teamRunId, recipient: recipientName, diff --git a/src/hooks/shared/prompt-async-gate.test.ts b/src/hooks/shared/prompt-async-gate.test.ts index 159b8b0f3..4a7bdeb9a 100644 --- a/src/hooks/shared/prompt-async-gate.test.ts +++ b/src/hooks/shared/prompt-async-gate.test.ts @@ -711,6 +711,46 @@ describe("dispatchInternalPrompt shared gate behavior", () => { expect(promptCalls).toBe(0) }) + test("#given latest assistant turn has unknown finish #when an internal promptAsync is requested #then no prompt is sent", async () => { + // given + let promptCalls = 0 + const client = { + session: { + status: async () => ({ data: { ses_unknown_finish: { type: "idle" } } }), + messages: async () => ({ + data: [ + { + info: { id: "msg_user", role: "user" }, + parts: [{ type: "text", text: "run work" }], + }, + { + info: { id: "msg_assistant", role: "assistant", finish: "unknown" }, + parts: [{ type: "reasoning", text: "still resolving" }], + }, + ], + }), + promptAsync: async () => { + promptCalls += 1 + }, + }, + } + + // when + const result = await dispatchInternalPrompt({ + mode: "async", + client, + sessionID: "ses_unknown_finish", + input: { path: { id: "ses_unknown_finish" }, body: { parts: [] } }, + source: "test:unknown-finish", + settleMs: 0, + postDispatchHoldMs: 0, + }) + + // then + expect(result.status).toBe("queued") + expect(promptCalls).toBe(0) + }) + test("#given internal user tail follows an assistant waiting on tools #when an internal promptAsync is requested #then no prompt is sent", async () => { // given let promptCalls = 0 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 cdad22a01..e3fa374fb 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 @@ -567,6 +567,48 @@ describe("createTeamIdleWakeHint", () => { expect(processedEntries).toContain(`${messageId}.json`) }) + test("#given stale idle while pending live delivery is still busy #when idle wake runs #then it keeps the reservation pending", 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 ackSpy = spyOn(ackModule, "ackMessages") + const promptAsyncSpy = mock(async (_input: WakeHintPromptInput) => ({})) + const handler = createTeamIdleWakeHint({ + directory: "/tmp/project", + client: { + session: { + promptAsync: promptAsyncSpy, + status: async () => ({ data: { "member-session": { type: "busy" } } }), + }, + }, + }, config, { idleSettleMs: 0 }) + + // when + await handler({ + event: { + type: "session.idle", + properties: { sessionID: "member-session" }, + }, + }) + + // then + expect(ackSpy).not.toHaveBeenCalled() + expect(promptAsyncSpy).not.toHaveBeenCalled() + + const runtimeState = await loadRuntimeState(teamRunId, config) + expect(runtimeState.members[0]?.pendingInjectedMessageIds).toEqual([messageId]) + + const inboxDir = getInboxDir(resolveBaseDir(config), teamRunId, "worker") + const inboxEntries = await readdir(inboxDir) + expect(inboxEntries).toContain(`.delivering-${messageId}.json`) + expect(inboxEntries).not.toContain(`${messageId}.json`) + }) + test("#given a pending live-delivery ack and later unread message #when member idles after the live reply #then it wakes the member for the unread message", 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 16a1a1289..fcfecde95 100644 --- a/src/hooks/team-session-events/team-idle-wake-hint.ts +++ b/src/hooks/team-session-events/team-idle-wake-hint.ts @@ -10,6 +10,7 @@ import { loadRuntimeState, transitionRuntimeState } from "../../features/team-mo import { resolveSessionEventID } from "../../shared/event-session-id" import { isAmbiguousPromptDispatchFailure } from "../../shared/prompt-failure-classifier" import { log } from "../../shared/logger" +import { isSessionActive, settleAfterSessionIdle } from "../../shared/session-idle-settle" import { dispatchInternalPrompt, isInternalPromptDispatchAccepted } from "../shared/prompt-async-gate" type PromptAsyncInput = { @@ -73,6 +74,20 @@ export function createTeamIdleWakeHint(ctx: TeamIdleWakeHintContext, config: Tea const pendingInjectedMessageIds = [...memberEntry.pendingInjectedMessageIds] if (pendingInjectedMessageIds.length > 0) { + if (typeof ctx.client.session.status === "function") { + await settleAfterSessionIdle(options?.idleSettleMs ?? 0) + if (await isSessionActive(ctx.client, sessionID)) { + log("team idle pending ack skipped while session remains active", { + event: "team-mode-idle-pending-ack-active", + teamRunId: runtimeState.teamRunId, + memberName: memberEntry.name, + sessionID, + pendingCount: pendingInjectedMessageIds.length, + }) + return + } + } + await ackMessages(runtimeState.teamRunId, memberEntry.name, pendingInjectedMessageIds, config) await transitionRuntimeState(runtimeState.teamRunId, (currentRuntimeState) => ({ ...currentRuntimeState, 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 39fbfe452..cd9cbcfff 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 @@ -226,4 +226,89 @@ describe("createTeamMemberErrorHandler", () => { expect(inboxEntries).not.toContain(`.delivering-${messageId}.json`) expect(inboxEntries).not.toContain("processed") }) + + test("#given session.error arrives while OpenCode still reports busy #when pending live delivery exists #then it does not requeue a duplicate peer message", 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, { + settleMs: 0, + client: { + session: { + status: async () => ({ data: { "member-session": { type: "busy" } } }), + }, + }, + }) + + // when + await handler({ + event: { + type: "session.error", + properties: { sessionID: "member-session", error: new Error("transient provider error") }, + }, + }) + + // then + const runtimeState = await loadRuntimeState(teamRunId, config) + expect(runtimeState.members[0]?.status).toBe("running") + expect(runtimeState.members[0]?.pendingInjectedMessageIds).toEqual([messageId]) + + const inboxEntries = await readdir(getInboxDir(resolveBaseDir(config), teamRunId, "worker")) + expect(inboxEntries).toContain(`.delivering-${messageId}.json`) + expect(inboxEntries).not.toContain(`${messageId}.json`) + expect(inboxEntries).not.toContain("processed") + }) + + test("#given session.error arrives after peer message reached history #when pending live delivery exists #then it does not requeue a duplicate peer message", 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, { + settleMs: 0, + client: { + session: { + status: async () => ({ data: { "member-session": { type: "idle" } } }), + messages: async () => ({ + data: [ + { + info: { role: "user" }, + parts: [ + { + type: "text", + text: `pending live delivery`, + }, + ], + }, + ], + }), + }, + }, + }) + + // when + await handler({ + event: { + type: "session.error", + properties: { sessionID: "member-session", error: new Error("late session.error after accepted prompt") }, + }, + }) + + // then + const runtimeState = await loadRuntimeState(teamRunId, config) + expect(runtimeState.members[0]?.status).toBe("running") + expect(runtimeState.members[0]?.pendingInjectedMessageIds).toEqual([messageId]) + + const inboxEntries = await readdir(getInboxDir(resolveBaseDir(config), teamRunId, "worker")) + expect(inboxEntries).toContain(`.delivering-${messageId}.json`) + expect(inboxEntries).not.toContain(`${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 0fe927b49..9fd6f4860 100644 --- a/src/hooks/team-session-events/team-member-error-handler.ts +++ b/src/hooks/team-session-events/team-member-error-handler.ts @@ -7,9 +7,23 @@ import { import { loadRuntimeState, transitionRuntimeState } from "../../features/team-mode/team-state-store/store" import { resolveSessionEventID } from "../../shared/event-session-id" import { log } from "../../shared/logger" +import { + DEFAULT_SESSION_IDLE_SETTLE_MS, + isSessionActive, + settleAfterSessionIdle, +} from "../../shared/session-idle-settle" type HookInput = { event: { type: string; properties?: unknown } } export type HookImpl = (input: HookInput) => Promise +type TeamMemberErrorHandlerDeps = { + client?: { + session?: { + status?: () => Promise + messages?: (input: { path: { id: string } }) => Promise + } + } + settleMs?: number +} function getErroredSessionID(properties: unknown): string | undefined { return resolveSessionEventID(properties) @@ -31,7 +45,73 @@ async function requeuePendingLiveDeliveries( } } -export function createTeamMemberErrorHandler(config: TeamModeConfig): HookImpl { +async function shouldKeepPendingLiveDeliveries( + deps: TeamMemberErrorHandlerDeps, + sessionID: string, +): Promise { + if (typeof deps.client?.session?.status !== "function") { + return false + } + + await settleAfterSessionIdle(deps.settleMs ?? DEFAULT_SESSION_IDLE_SETTLE_MS) + return await isSessionActive(deps.client, sessionID) +} + +function isRecord(value: unknown): value is Record { + return typeof value === "object" && value !== null +} + +function getMessagesData(response: unknown): unknown[] { + if (isRecord(response) && Array.isArray(response.data)) { + return response.data + } + + return Array.isArray(response) ? response : [] +} + +function valueContainsAnyMessageId(value: unknown, messageIds: ReadonlySet): boolean { + if (typeof value === "string") { + return [...messageIds].some((messageId) => value.includes(messageId)) + } + + if (Array.isArray(value)) { + return value.some((entry) => valueContainsAnyMessageId(entry, messageIds)) + } + + if (isRecord(value)) { + return Object.values(value).some((entry) => valueContainsAnyMessageId(entry, messageIds)) + } + + return false +} + +async function sessionHistoryContainsPendingMessage( + deps: TeamMemberErrorHandlerDeps, + sessionID: string, + messageIds: readonly string[], +): Promise { + if (messageIds.length === 0 || typeof deps.client?.session?.messages !== "function") { + return false + } + + try { + const response = await deps.client.session.messages({ path: { id: sessionID } }) + const pendingMessageIds = new Set(messageIds) + return getMessagesData(response).some((message) => valueContainsAnyMessageId(message, pendingMessageIds)) + } catch (error) { + log("team member session history check failed", { + event: "team-mode-member-error-history-check-failed", + sessionID, + error: error instanceof Error ? error.message : String(error), + }) + return false + } +} + +export function createTeamMemberErrorHandler( + config: TeamModeConfig, + deps: TeamMemberErrorHandlerDeps = {}, +): HookImpl { return async ({ event }: HookInput): Promise => { if (event.type !== "session.error") return @@ -47,6 +127,29 @@ 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 ?? [] + if (await shouldKeepPendingLiveDeliveries(deps, erroredSessionID)) { + log("team member session error ignored while session remains active", { + event: "team-mode-member-error-active", + teamRunId: runtimeState.teamRunId, + teamName: runtimeState.teamName, + memberName: runtimeMember.memberName, + sessionID: erroredSessionID, + pendingCount: pendingInjectedMessageIds.length, + }) + return + } + if (await sessionHistoryContainsPendingMessage(deps, erroredSessionID, pendingInjectedMessageIds)) { + log("team member session error ignored after pending peer message reached history", { + event: "team-mode-member-error-peer-message-accepted", + teamRunId: runtimeState.teamRunId, + teamName: runtimeState.teamName, + memberName: runtimeMember.memberName, + sessionID: erroredSessionID, + pendingCount: pendingInjectedMessageIds.length, + }) + return + } + await requeuePendingLiveDeliveries( runtimeState.teamRunId, runtimeMember.memberName, diff --git a/src/plugin/event.ts b/src/plugin/event.ts index 937309722..44c6c2ec4 100644 --- a/src/plugin/event.ts +++ b/src/plugin/event.ts @@ -342,7 +342,7 @@ export function createEventHandler(args: { ? createTeamLeadOrphanHandler(teamModeConfig, managers.tmuxSessionManager, managers.backgroundManager) : undefined; const teamMemberErrorHandler = teamModeConfig - ? createTeamMemberErrorHandler(teamModeConfig) + ? createTeamMemberErrorHandler(teamModeConfig, { client: pluginContext.client }) : undefined; const teamMemberStatusHandler = teamModeConfig ? createTeamMemberStatusHandler(teamModeConfig) diff --git a/src/shared/prompt-async-gate.ts b/src/shared/prompt-async-gate.ts index 25e35aa83..107f81a56 100644 --- a/src/shared/prompt-async-gate.ts +++ b/src/shared/prompt-async-gate.ts @@ -10,7 +10,7 @@ import { settleAfterSessionIdle, } from "./session-idle-settle" -export const DEFAULT_PROMPT_ASYNC_POST_DISPATCH_HOLD_MS = 250 +export const DEFAULT_PROMPT_ASYNC_POST_DISPATCH_HOLD_MS = 2_000 export const DEFAULT_PROMPT_DISPATCH_TIMEOUT_MS = 30_000 export const DEFAULT_PROMPT_GATE_MESSAGES_FETCH_TIMEOUT_MS = 5_000 export const DEFAULT_PROMPT_QUEUE_RETRY_MS = 250 @@ -420,7 +420,7 @@ function latestAssistantTurnBlocksInternalPrompt(messages: unknown[]): boolean { const role = messageRole(message) if (role === "assistant") { const finish = messageFinish(message) - if (finish === undefined) { + if (finish === undefined || finish === "unknown") { return true } if (!isRecord(message) || !Array.isArray(message.parts)) { diff --git a/src/shared/prompt-async-route-audit.test.ts b/src/shared/prompt-async-route-audit.test.ts index 30841edab..fc6f371fc 100644 --- a/src/shared/prompt-async-route-audit.test.ts +++ b/src/shared/prompt-async-route-audit.test.ts @@ -344,6 +344,9 @@ await dispatchInternalPrompt(options) } const contents = await readFile(filePath, "utf8") + if (!contents.includes("prompt")) { + continue + } if (detectRawPromptInSnippet(contents)) { offenders.push(relativeSourcePath(filePath)) } @@ -395,6 +398,9 @@ await dispatchInternalPrompt(options) // when for (const filePath of files) { const contents = await readFile(filePath, "utf8") + if (!contents.includes("dispatchInternalPrompt")) { + continue + } const missingLines = findPromptGateCallsWithoutQueueBehavior(filePath, contents) for (const line of missingLines) { offenders.push(`${relativeSourcePath(filePath)}:${line}`) @@ -413,6 +419,12 @@ await dispatchInternalPrompt(options) // when for (const filePath of files) { const contents = await readFile(filePath, "utf8") + if ( + !contents.includes("promptWithModelSuggestionRetry") + && !contents.includes("promptSyncWithModelSuggestionRetry") + ) { + continue + } const missingLines = findPromptRetryCallsWithoutQueueBehavior(filePath, contents) for (const line of missingLines) { offenders.push(`${relativeSourcePath(filePath)}:${line}`)