From 6df148f0f69108638ac1e4463af57c576f15d2d8 Mon Sep 17 00:00:00 2001 From: YeonGyu-Kim Date: Wed, 20 May 2026 11:42:38 +0900 Subject: [PATCH] fix(team-mode): close peer message delivery races --- .../team-mode/team-state-store/resume.test.ts | 119 ++++++++++++++++ .../team-mode/team-state-store/resume.ts | 93 ++++++++++++- .../team-mode/tools/messaging.test.ts | 63 +++++++++ src/features/team-mode/tools/messaging.ts | 17 +-- .../team-idle-wake-hint.test.ts | 49 +++++++ .../team-idle-wake-hint.ts | 131 +++++++++++++----- 6 files changed, 426 insertions(+), 46 deletions(-) diff --git a/src/features/team-mode/team-state-store/resume.test.ts b/src/features/team-mode/team-state-store/resume.test.ts index 959261f92..432582498 100644 --- a/src/features/team-mode/team-state-store/resume.test.ts +++ b/src/features/team-mode/team-state-store/resume.test.ts @@ -95,15 +95,18 @@ function createSpecWithTwoWorkers(name = `team-${randomUUID().slice(0, 8)}`): Te } type SessionGetMock = (input: { path: { id: string } }) => Promise +type SessionMessagesMock = (input: { path: { id: string } }) => Promise function createExecutorContext( directory: string, sessionGet: SessionGetMock = mock(async () => ({ data: null })), + sessionMessages?: SessionMessagesMock, ): ExecutorContext { return { client: { session: { get: sessionGet, + ...(sessionMessages ? { messages: sessionMessages } : {}), }, } as ExecutorContext["client"], manager: {} as ExecutorContext["manager"], @@ -291,6 +294,122 @@ describe("resumeAllTeams", () => { expect(entries).not.toContain(`.delivering-${strandedMessageId}.json`) }) + test("#given reclaimed stale reservation is still pending but absent from session history #when active team resumes #then pending state is cleared for mailbox injection", async () => { + // given + const baseDir = await createTemporaryBaseDir() + temporaryDirectories.push(baseDir) + const config = createConfig(baseDir) + const runtimeState = await createRuntimeState(createSpec(), "ses_alive_lead", "user", config) + const workerMessageId = randomUUID() + await transitionRuntimeState(runtimeState.teamRunId, (currentRuntimeState) => ({ + ...currentRuntimeState, + status: "active", + leadSessionId: "ses_alive_lead", + members: currentRuntimeState.members.map((member) => { + if (member.name === "lead") return { ...member, sessionId: "ses_alive_lead", status: "running" as const } + return { + ...member, + sessionId: "ses_worker", + status: "running" as const, + pendingInjectedMessageIds: [workerMessageId], + } + }), + }), config) + const workerInbox = getInboxDir(resolveBaseDir(config), runtimeState.teamRunId, "worker") + await mkdir(workerInbox, { recursive: true, mode: 0o700 }) + const reservedPath = path.join(workerInbox, `.delivering-${workerMessageId}.json`) + await writeFile(reservedPath, JSON.stringify({ + version: 1, + messageId: workerMessageId, + from: "lead", + to: "worker", + kind: "message", + body: "retry after restart", + timestamp: Date.now(), + })) + const ancientMtime = new Date(Date.now() - 60 * 60 * 1000) + await utimes(reservedPath, ancientMtime, ancientMtime) + const sessionGet = mock(async () => ({ data: { id: "alive" } })) + const sessionMessages = mock(async () => ({ data: [] })) + + // when + await resumeAllTeams(createExecutorContext(baseDir, sessionGet, sessionMessages), config) + + // then + const entries = await readdir(workerInbox) + expect(entries).toContain(`${workerMessageId}.json`) + expect(entries).not.toContain(`.delivering-${workerMessageId}.json`) + + const persistedState = await loadRuntimeState(runtimeState.teamRunId, config) + const worker = persistedState.members.find((member) => member.name === "worker") + expect(worker?.pendingInjectedMessageIds).toEqual([]) + }) + + test("#given accepted live delivery lost its pending mark #when stale reservation is reclaimed #then resume processes it instead of exposing duplicate unread", async () => { + // given + const baseDir = await createTemporaryBaseDir() + temporaryDirectories.push(baseDir) + const config = createConfig(baseDir) + const runtimeState = await createRuntimeState(createSpec(), "ses_alive_lead", "user", config) + const workerMessageId = randomUUID() + await transitionRuntimeState(runtimeState.teamRunId, (currentRuntimeState) => ({ + ...currentRuntimeState, + status: "active", + leadSessionId: "ses_alive_lead", + members: currentRuntimeState.members.map((member) => { + if (member.name === "lead") return { ...member, sessionId: "ses_alive_lead", status: "running" as const } + return { + ...member, + sessionId: "ses_worker", + status: "running" as const, + pendingInjectedMessageIds: [], + } + }), + }), config) + const workerInbox = getInboxDir(resolveBaseDir(config), runtimeState.teamRunId, "worker") + await mkdir(workerInbox, { recursive: true, mode: 0o700 }) + const reservedPath = path.join(workerInbox, `.delivering-${workerMessageId}.json`) + await writeFile(reservedPath, JSON.stringify({ + version: 1, + messageId: workerMessageId, + from: "lead", + to: "worker", + kind: "message", + body: "already accepted", + timestamp: Date.now(), + })) + const ancientMtime = new Date(Date.now() - 60 * 60 * 1000) + await utimes(reservedPath, ancientMtime, ancientMtime) + const sessionGet = mock(async () => ({ data: { id: "alive" } })) + const sessionMessages = mock(async ({ path: sessionPath }: { path: { id: string } }) => ({ + data: sessionPath.id === "ses_worker" + ? [ + { + info: { role: "user" }, + parts: [ + { + type: "text", + text: `already accepted`, + }, + ], + }, + ] + : [], + })) + + // when + await resumeAllTeams(createExecutorContext(baseDir, sessionGet, sessionMessages), config) + + // then + const entries = await readdir(workerInbox) + expect(entries).not.toContain(`${workerMessageId}.json`) + expect(entries).not.toContain(`.delivering-${workerMessageId}.json`) + expect(entries).toContain("processed") + + const processedEntries = await readdir(path.join(workerInbox, "processed")) + expect(processedEntries).toContain(`${workerMessageId}.json`) + }) + test("leaves fresh .delivering-* reservations in place on resume", async () => { // given const baseDir = await createTemporaryBaseDir() diff --git a/src/features/team-mode/team-state-store/resume.ts b/src/features/team-mode/team-state-store/resume.ts index 96608bd4c..62ba958aa 100644 --- a/src/features/team-mode/team-state-store/resume.ts +++ b/src/features/team-mode/team-state-store/resume.ts @@ -3,9 +3,10 @@ import { rm, stat } from "node:fs/promises" import type { TeamModeConfig } from "../../../config/schema/team-mode" import { log } from "../../../shared/logger" import type { ExecutorContext } from "../../../tools/delegate-task/executor-types" +import { ackMessages } from "../team-mailbox/ack" import { reclaimStaleReservations } from "../team-mailbox/reservation" import { getRuntimeStateDir, resolveBaseDir } from "../team-registry/paths" -import type { RuntimeState } from "../types" +import type { RuntimeState, RuntimeStateMember } from "../types" import { listActiveTeams, loadRuntimeState, transitionRuntimeState } from "./store" const CREATING_TIMEOUT_MS = 30 * 60 * 1000 @@ -89,6 +90,88 @@ function isCreatingStateStuck(runtimeState: RuntimeState, now: number): boolean return runtimeState.status === "creating" && now - runtimeState.createdAt > CREATING_TIMEOUT_MS } +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 valueContainsMessageId(value: unknown, messageId: string): boolean { + if (typeof value === "string") { + return value.includes(messageId) + } + + if (Array.isArray(value)) { + return value.some((entry) => valueContainsMessageId(entry, messageId)) + } + + if (isRecord(value)) { + return Object.values(value).some((entry) => valueContainsMessageId(entry, messageId)) + } + + return false +} + +async function findAcceptedReclaimedMessageIds( + ctx: ExecutorContext, + member: RuntimeStateMember, + messageIds: readonly string[], +): Promise { + if (messageIds.length === 0 || member.sessionId === undefined) { + return [] + } + + try { + const response = await ctx.client.session.messages({ path: { id: member.sessionId } }) + const messages = getMessagesData(response) + return messageIds.filter((messageId) => messages.some((message) => valueContainsMessageId(message, messageId))) + } catch (historyError) { + log("team mailbox reclaimed reservation history check failed", { + event: "team-mailbox-reclaim-history-check-failed", + member: member.name, + sessionID: member.sessionId, + error: historyError instanceof Error ? historyError.message : String(historyError), + }) + return [] + } +} + +async function reconcileReclaimedReservations( + ctx: ExecutorContext, + teamRunId: string, + member: RuntimeStateMember, + reclaimedMessageIds: readonly string[], + config: TeamModeConfig, +): Promise { + if (reclaimedMessageIds.length === 0) { + return + } + + const acceptedMessageIds = await findAcceptedReclaimedMessageIds(ctx, member, reclaimedMessageIds) + if (acceptedMessageIds.length > 0) { + await ackMessages(teamRunId, member.name, acceptedMessageIds, config) + } + + const reclaimedMessageIdSet = new Set(reclaimedMessageIds) + await transitionRuntimeState(teamRunId, (currentRuntimeState) => ({ + ...currentRuntimeState, + members: currentRuntimeState.members.map((currentMember) => ( + currentMember.name === member.name + ? { + ...currentMember, + pendingInjectedMessageIds: currentMember.pendingInjectedMessageIds.filter((messageId) => !reclaimedMessageIdSet.has(messageId)), + } + : currentMember + )), + }), config) +} + interface WorkerLiveness { readonly name: string readonly wasSpawned: boolean @@ -157,7 +240,13 @@ export async function resumeAllTeams( await Promise.all(runtimeState.members.map(async (member) => { try { - await reclaimStaleReservations(runtimeState.teamRunId, member.name, config, STALE_RESERVATION_TTL_MS) + const reclaimedMessageIds = await reclaimStaleReservations( + runtimeState.teamRunId, + member.name, + config, + STALE_RESERVATION_TTL_MS, + ) + await reconcileReclaimedReservations(ctx, runtimeState.teamRunId, member, reclaimedMessageIds, config) } catch (reclaimError) { log("team mailbox reservation reclaim failed", { event: "team-mailbox-reclaim-failed", diff --git a/src/features/team-mode/tools/messaging.test.ts b/src/features/team-mode/tools/messaging.test.ts index 36ba69ebd..e733ef3ee 100644 --- a/src/features/team-mode/tools/messaging.test.ts +++ b/src/features/team-mode/tools/messaging.test.ts @@ -16,6 +16,7 @@ import { } from "../../../shared/session-prompt-params-state" import { releaseAllPromptAsyncReservationsForTesting } from "../../../hooks/shared/prompt-async-gate" import { listUnreadMessages } from "../team-mailbox/inbox" +import { pollAndBuildInjection } from "../team-mailbox/poll" import { BroadcastNotPermittedError } from "../team-mailbox/send" import { getInboxDir, resolveBaseDir } from "../team-registry/paths" import { createRuntimeState, saveRuntimeState } from "../team-state-store/store" @@ -766,6 +767,68 @@ describe("createTeamSendMessageTool", () => { expect(unreadDuringDelivery).toHaveLength(0) }) + test("#given transform already listed a peer message while recipient becomes idle #when live delivery races it #then only the live prompt receives the message", async () => { + // given + const fixture = await createTeamFixture() + const { loadRuntimeState: loadState, saveRuntimeState: saveState } = await import("../team-state-store/store") + let loadCount = 0 + let transformResult: Awaited> | undefined + const deps = { + loadRuntimeState: async (teamRunId: string) => { + loadCount += 1 + const state = await loadState(teamRunId, fixture.config) + if (loadCount === 2) { + return { + ...state, + members: state.members.map((member) => ( + member.name === "m2" + ? { ...member, status: "running" as const, pendingInjectedMessageIds: [] } + : member + )), + } + } + if (loadCount === 3) { + const staleIdleSnapshot = { + ...state, + members: state.members.map((member) => ( + member.name === "m2" + ? { ...member, status: "idle" as const, pendingInjectedMessageIds: [] } + : member + )), + } + await saveState(staleIdleSnapshot, fixture.config) + transformResult = await pollAndBuildInjection( + fixture.memberTwoSessionId, + "m2", + fixture.teamRunId, + fixture.config, + "turn-race", + ) + return staleIdleSnapshot + } + return state + }, + } + const { client, calls } = createRecordingClient() + const liveTool = createTeamSendMessageTool(fixture.config, client, deps) + + // when + await liveTool.execute({ + teamRunId: fixture.teamRunId, + to: "m2", + body: "race payload", + }, fixture.toolContext(fixture.memberOneSessionId)) + + // then + expect(transformResult).toEqual({ + injected: false, + messageIds: [], + reason: "no unread", + }) + expect(calls).toHaveLength(1) + expect(calls[0]?.parts[0]?.text).toContain("race payload") + }) + test("hides the message from the inbox from the moment it is written for a live recipient", 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 d08d18ca5..338c6713a 100644 --- a/src/features/team-mode/tools/messaging.ts +++ b/src/features/team-mode/tools/messaging.ts @@ -6,7 +6,7 @@ import { z } from "zod" import type { TeamModeConfig } from "../../../config/schema/team-mode" import { dispatchInternalPrompt, isInternalPromptDispatchAccepted } from "../../../hooks/shared/prompt-async-gate" import { log } from "../../../shared/logger" -import { isAmbiguousPromptDispatchFailure } from "../../../shared/prompt-failure-classifier" +import { isAmbiguousPostDispatchPromptFailure } from "../../../shared/prompt-failure-classifier" import { applyMemberSessionRouting, buildMemberPromptBody } from "../member-session-routing" import { buildEnvelope } from "../team-mailbox/poll" import { @@ -70,11 +70,12 @@ const TeamSendMessageArgsSchema = z.object({ type DeliveryReservation = Awaited> type RuntimeMember = RuntimeState["members"][number] -function canPreReserveForLiveDelivery(member: RuntimeMember, senderName: string): boolean { - return member.name !== senderName - && member.sessionId !== undefined - && member.status === "idle" - && member.pendingInjectedMessageIds.length === 0 +function shouldReserveRecipientMailbox(member: RuntimeMember, message: Message, senderName: string): boolean { + if (message.to === "*") { + return member.name !== senderName + } + + return member.name === message.to } async function resolveTeamRuntimeDetails( @@ -273,7 +274,7 @@ async function deliverLive( query: { directory: recipientMember.worktreePath ?? directory }, }, }) - if (promptResult.status === "failed" && isAmbiguousPromptDispatchFailure(promptResult.error)) { + if (promptResult.status === "failed" && isAmbiguousPostDispatchPromptFailure(promptResult)) { try { await markLiveDeliveryPending(teamRunId, recipientName, message.messageId, config) } catch (markError) { @@ -399,7 +400,7 @@ export function createTeamSendMessageTool( const runtimeState = await deps.loadRuntimeState(teamRuntime.teamRunId, config) const reservedRecipients = new Set( runtimeState.members - .filter((member) => canPreReserveForLiveDelivery(member, teamRuntime.senderName)) + .filter((member) => shouldReserveRecipientMailbox(member, message, teamRuntime.senderName)) .map((member) => member.name), ) 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 e3fa374fb..bd91e58cc 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 @@ -26,6 +26,7 @@ import { releaseAllPromptAsyncReservationsForTesting, releasePromptAsyncReservation, } from "../shared/prompt-async-gate" +import { createTeamMemberErrorHandler } from "./team-member-error-handler" import { createTeamIdleWakeHint } from "./team-idle-wake-hint" type WakeHintPromptInput = { @@ -609,6 +610,54 @@ describe("createTeamIdleWakeHint", () => { expect(inboxEntries).not.toContain(`${messageId}.json`) }) + test("#given pending live delivery is requeued by session.error #when stale idle handler continues #then it does not ack the requeued message", 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 errorHandler = createTeamMemberErrorHandler(config) + + const handler = createTeamIdleWakeHint({ + directory: "/tmp/project", + client: { + session: { + promptAsync: mock(async (_input: WakeHintPromptInput) => ({})), + status: async () => { + await errorHandler({ + event: { + type: "session.error", + properties: { sessionID: "member-session", error: new Error("late prompt failure") }, + }, + }) + return { data: { "member-session": { type: "idle" } } } + }, + }, + }, + }, config, { idleSettleMs: 0 }) + + // 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") + expect(runtimeState.members[0]?.pendingInjectedMessageIds).toEqual([]) + + const inboxDir = getInboxDir(resolveBaseDir(config), teamRunId, "worker") + const inboxEntries = await readdir(inboxDir) + expect(inboxEntries).toContain(`${messageId}.json`) + expect(inboxEntries).not.toContain(`.delivering-${messageId}.json`) + expect(inboxEntries).not.toContain("processed") + }) + 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 fcfecde95..6a96ac63d 100644 --- a/src/hooks/team-session-events/team-idle-wake-hint.ts +++ b/src/hooks/team-session-events/team-idle-wake-hint.ts @@ -8,7 +8,7 @@ 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 { isAmbiguousPromptDispatchFailure } from "../../shared/prompt-failure-classifier" +import { isAmbiguousPostDispatchPromptFailure } 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" @@ -51,6 +51,45 @@ function buildWakeHintBatchKey(teamRunId: string, memberName: string, messageIds return `${teamRunId}:${memberName}:${messageIds.toSorted().join(",")}` } +async function claimPendingMessageAcks( + teamRunId: string, + memberName: string, + messageIds: readonly string[], + config: TeamModeConfig, +): Promise { + if (messageIds.length === 0) return [] + + let claimedMessageIds: string[] = [] + const candidateMessageIds = new Set(messageIds) + await transitionRuntimeState(teamRunId, (currentRuntimeState) => { + const currentMember = currentRuntimeState.members.find((member) => member.name === memberName) + if (currentMember === undefined) { + claimedMessageIds = [] + return currentRuntimeState + } + + claimedMessageIds = currentMember.pendingInjectedMessageIds.filter((messageId) => candidateMessageIds.has(messageId)) + if (claimedMessageIds.length === 0) { + return currentRuntimeState + } + + const claimedMessageIdSet = new Set(claimedMessageIds) + return { + ...currentRuntimeState, + members: currentRuntimeState.members.map((member) => ( + member.name === memberName + ? { + ...member, + pendingInjectedMessageIds: member.pendingInjectedMessageIds.filter((messageId) => !claimedMessageIdSet.has(messageId)), + } + : member + )), + } + }, config) + + return claimedMessageIds +} + export function createTeamIdleWakeHint(ctx: TeamIdleWakeHintContext, config: TeamModeConfig, options?: TeamIdleWakeHintOptions): HookImpl { const recentWakeHintBatches = new Map() @@ -88,41 +127,61 @@ export function createTeamIdleWakeHint(ctx: TeamIdleWakeHintContext, config: Tea } } - await ackMessages(runtimeState.teamRunId, memberEntry.name, pendingInjectedMessageIds, config) - await transitionRuntimeState(runtimeState.teamRunId, (currentRuntimeState) => ({ - ...currentRuntimeState, - members: currentRuntimeState.members.map((member) => ( - member.name === memberEntry.name - ? { ...member, pendingInjectedMessageIds: [] } - : member - )), - }), config) - log("team idle handled pending live delivery ack", { - event: "team-mode-idle-pending-ack", - teamRunId: runtimeState.teamRunId, - memberName: memberEntry.name, - sessionID, - ackedCount: pendingInjectedMessageIds.length, - }) + const claimedMessageIds = await claimPendingMessageAcks( + runtimeState.teamRunId, + memberEntry.name, + pendingInjectedMessageIds, + config, + ) + if (claimedMessageIds.length > 0) { + await ackMessages(runtimeState.teamRunId, memberEntry.name, claimedMessageIds, config) + log("team idle handled pending live delivery ack", { + event: "team-mode-idle-pending-ack", + teamRunId: runtimeState.teamRunId, + memberName: memberEntry.name, + sessionID, + ackedCount: claimedMessageIds.length, + }) + } } - const unreadMessages = await listUnreadMessages(runtimeState.teamRunId, memberEntry.name, config) + const latestRuntimeState = await loadRuntimeState(runtimeMember.teamRunId, config) + const latestMemberEntry = latestRuntimeState.members.find((member) => member.name === runtimeMember.memberName) + if (!latestMemberEntry) { + return + } + if ( + latestMemberEntry.status === "errored" + || latestMemberEntry.status === "completed" + || latestMemberEntry.status === "shutdown_approved" + ) { + log("team idle wake hint skipped because member is no longer idle", { + event: "team-mode-idle-member-not-idle", + teamRunId: latestRuntimeState.teamRunId, + memberName: latestMemberEntry.name, + sessionID, + status: latestMemberEntry.status, + }) + return + } + + const unreadMessages = await listUnreadMessages(latestRuntimeState.teamRunId, latestMemberEntry.name, config) if (unreadMessages.length === 0) { log("team idle handled without wake hint", { event: "team-mode-idle-ack-only", - teamRunId: runtimeState.teamRunId, - memberName: memberEntry.name, + teamRunId: latestRuntimeState.teamRunId, + memberName: latestMemberEntry.name, sessionID, ackedCount: pendingInjectedMessageIds.length, }) return } - if (memberEntry.agentType === "leader") { + if (latestMemberEntry.agentType === "leader") { log("team lead idle handled without wake hint", { event: "team-mode-lead-idle-ack-only", - teamRunId: runtimeState.teamRunId, - memberName: memberEntry.name, + teamRunId: latestRuntimeState.teamRunId, + memberName: latestMemberEntry.name, sessionID, ackedCount: pendingInjectedMessageIds.length, }) @@ -132,8 +191,8 @@ export function createTeamIdleWakeHint(ctx: TeamIdleWakeHintContext, config: Tea if (typeof ctx.client.session.promptAsync !== "function") { log("team idle wake hint skipped without promptAsync", { event: "team-mode-idle-wake-hint-skipped", - teamRunId: runtimeState.teamRunId, - memberName: memberEntry.name, + teamRunId: latestRuntimeState.teamRunId, + memberName: latestMemberEntry.name, sessionID, unreadCount: unreadMessages.length, }) @@ -142,16 +201,16 @@ export function createTeamIdleWakeHint(ctx: TeamIdleWakeHintContext, config: Tea const now = Date.now() const wakeHintBatchKey = buildWakeHintBatchKey( - runtimeState.teamRunId, - memberEntry.name, + latestRuntimeState.teamRunId, + latestMemberEntry.name, unreadMessages.map((message) => message.messageId), ) const suppressedUntil = recentWakeHintBatches.get(wakeHintBatchKey) if (suppressedUntil !== undefined && suppressedUntil > now) { log("team idle wake hint skipped for recently hinted unread batch", { event: "team-mode-idle-wake-hint-duplicate-suppressed", - teamRunId: runtimeState.teamRunId, - memberName: memberEntry.name, + teamRunId: latestRuntimeState.teamRunId, + memberName: latestMemberEntry.name, sessionID, unreadCount: unreadMessages.length, }) @@ -161,7 +220,7 @@ export function createTeamIdleWakeHint(ctx: TeamIdleWakeHintContext, config: Tea recentWakeHintBatches.delete(wakeHintBatchKey) } - applyMemberSessionRouting(sessionID, memberEntry) + applyMemberSessionRouting(sessionID, latestMemberEntry) const promptResult = await dispatchInternalPrompt({ mode: "async", client: ctx.client, @@ -171,18 +230,18 @@ export function createTeamIdleWakeHint(ctx: TeamIdleWakeHintContext, config: Tea queueBehavior: "defer", input: { path: { id: sessionID }, - body: buildMemberPromptBody(memberEntry, buildWakeHint(unreadMessages.length)), + body: buildMemberPromptBody(latestMemberEntry, buildWakeHint(unreadMessages.length)), query: { directory: ctx.directory }, }, }) if (!isInternalPromptDispatchAccepted(promptResult)) { - if (promptResult.status === "failed" && isAmbiguousPromptDispatchFailure(promptResult.error)) { + if (promptResult.status === "failed" && isAmbiguousPostDispatchPromptFailure(promptResult)) { recentWakeHintBatches.set(wakeHintBatchKey, Date.now() + WAKE_HINT_DUPLICATE_SUPPRESSION_MS) } log("team idle wake hint skipped by promptAsync gate", { event: "team-mode-idle-wake-hint-gated", - teamRunId: runtimeState.teamRunId, - memberName: memberEntry.name, + teamRunId: latestRuntimeState.teamRunId, + memberName: latestMemberEntry.name, sessionID, unreadCount: unreadMessages.length, status: promptResult.status, @@ -193,8 +252,8 @@ export function createTeamIdleWakeHint(ctx: TeamIdleWakeHintContext, config: Tea log("team idle wake hint sent", { event: "team-mode-idle-wake-hint", - teamRunId: runtimeState.teamRunId, - memberName: memberEntry.name, + teamRunId: latestRuntimeState.teamRunId, + memberName: latestMemberEntry.name, sessionID, unreadCount: unreadMessages.length, ackedCount: pendingInjectedMessageIds.length,