fix(team-mode): defer live mailbox acks

This commit is contained in:
YeonGyu-Kim
2026-05-15 22:49:10 +09:00
parent ae7ff3bb7e
commit 196f6512ae
7 changed files with 301 additions and 27 deletions
+15 -9
View File
@@ -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
}
}
}
}
@@ -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 () => {
+25 -8
View File
@@ -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<void> {
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,
@@ -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<void> {
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<void> {
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<void> {
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()
@@ -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",
@@ -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<void> {
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<void> {
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")
})
})
@@ -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<void> {
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<void> => {
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)