fix(team-mode): close peer message delivery races

This commit is contained in:
YeonGyu-Kim
2026-05-20 11:42:38 +09:00
parent d3e218f912
commit 6df148f0f6
6 changed files with 426 additions and 46 deletions
@@ -95,15 +95,18 @@ function createSpecWithTwoWorkers(name = `team-${randomUUID().slice(0, 8)}`): Te
}
type SessionGetMock = (input: { path: { id: string } }) => Promise<unknown>
type SessionMessagesMock = (input: { path: { id: string } }) => Promise<unknown>
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: `<peer_message from="lead" messageId="${workerMessageId}" kind="message">already accepted</peer_message>`,
},
],
},
]
: [],
}))
// 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()
@@ -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<string, unknown> {
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<string[]> {
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<void> {
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",
@@ -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<ReturnType<typeof pollAndBuildInjection>> | 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()
+9 -8
View File
@@ -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<ReturnType<typeof reserveMessageForDelivery>>
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<string>(
runtimeState.members
.filter((member) => canPreReserveForLiveDelivery(member, teamRuntime.senderName))
.filter((member) => shouldReserveRecipientMailbox(member, message, teamRuntime.senderName))
.map((member) => member.name),
)