fix(team-mode): close peer message delivery races
This commit is contained in:
@@ -196,7 +196,6 @@ export class ParentWakeNotifier {
|
|||||||
sessionID,
|
sessionID,
|
||||||
source: "background-agent-parent-wake",
|
source: "background-agent-parent-wake",
|
||||||
settleMs: 0,
|
settleMs: 0,
|
||||||
postDispatchHoldMs: 250,
|
|
||||||
queueBehavior: "defer",
|
queueBehavior: "defer",
|
||||||
checkToolState: !toolWaitDecision.skipPromptGateToolStateCheck,
|
checkToolState: !toolWaitDecision.skipPromptGateToolStateCheck,
|
||||||
input: {
|
input: {
|
||||||
|
|||||||
@@ -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 () => {
|
test("wraps hostile message bodies in a literal peer_message envelope", async () => {
|
||||||
// given
|
// given
|
||||||
const { teamRunId, config } = await setupRuntime(["m1"])
|
const { teamRunId, config } = await setupRuntime(["m1"])
|
||||||
|
|||||||
@@ -48,47 +48,54 @@ export async function pollAndBuildInjection(
|
|||||||
config: TeamModeConfig,
|
config: TeamModeConfig,
|
||||||
turnMarker: string,
|
turnMarker: string,
|
||||||
): Promise<InjectionResult> {
|
): Promise<InjectionResult> {
|
||||||
const runtimeState = await loadRuntimeState(teamRunId, config)
|
const unreadMessages = await listUnreadMessages(teamRunId, memberName, config)
|
||||||
const runtimeMember = runtimeState.members.find((member) => member.name === memberName)
|
let result: InjectionResult | undefined
|
||||||
if (runtimeMember === undefined) {
|
|
||||||
throw new Error(`runtime member not found for session ${sessionID}: ${memberName}`)
|
|
||||||
}
|
|
||||||
|
|
||||||
if (runtimeMember.lastInjectedTurnMarker === turnMarker) {
|
await transitionRuntimeState(teamRunId, (currentRuntimeState) => {
|
||||||
return { injected: false, messageIds: [], reason: "already injected this turn" }
|
const runtimeMember = currentRuntimeState.members.find((member) => member.name === memberName)
|
||||||
}
|
if (runtimeMember === undefined) {
|
||||||
|
throw new Error(`runtime member not found for session ${sessionID}: ${memberName}`)
|
||||||
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" }
|
|
||||||
}
|
}
|
||||||
|
|
||||||
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[] = []
|
return result
|
||||||
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 }
|
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,7 +1,7 @@
|
|||||||
/// <reference types="bun-types" />
|
/// <reference types="bun-types" />
|
||||||
|
|
||||||
import { afterEach, describe, expect, test } from "bun:test"
|
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 { randomUUID } from "node:crypto"
|
||||||
import { tmpdir } from "node:os"
|
import { tmpdir } from "node:os"
|
||||||
import path from "node:path"
|
import path from "node:path"
|
||||||
@@ -673,6 +673,74 @@ describe("createTeamSendMessageTool", () => {
|
|||||||
expect(recipient?.pendingInjectedMessageIds).toHaveLength(1)
|
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 () => {
|
test("reserves the message during live delivery so concurrent listings cannot surface it", async () => {
|
||||||
// given
|
// given
|
||||||
const fixture = await createTeamFixture()
|
const fixture = await createTeamFixture()
|
||||||
|
|||||||
@@ -161,6 +161,22 @@ async function markLiveDeliveryPending(
|
|||||||
}), config)
|
}), config)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
async function releaseReservationsForRecipients(
|
||||||
|
teamRunId: string,
|
||||||
|
recipientNames: readonly string[],
|
||||||
|
messageId: string,
|
||||||
|
config: TeamModeConfig,
|
||||||
|
): Promise<void> {
|
||||||
|
for (const recipientName of recipientNames) {
|
||||||
|
const reservation = await reserveMessageForDelivery(teamRunId, recipientName, messageId, config)
|
||||||
|
await releaseReservationSafely(reservation, {
|
||||||
|
teamRunId,
|
||||||
|
recipient: recipientName,
|
||||||
|
messageId,
|
||||||
|
})
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
async function deliverLive(
|
async function deliverLive(
|
||||||
client: LiveDeliveryClient,
|
client: LiveDeliveryClient,
|
||||||
message: Message,
|
message: Message,
|
||||||
@@ -170,7 +186,19 @@ async function deliverLive(
|
|||||||
directory: string,
|
directory: string,
|
||||||
deps: TeamSendMessageToolDeps,
|
deps: TeamSendMessageToolDeps,
|
||||||
): Promise<void> {
|
): Promise<void> {
|
||||||
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)
|
const envelope = buildEnvelope(message)
|
||||||
|
|
||||||
for (const recipientName of deliveredTo) {
|
for (const recipientName of deliveredTo) {
|
||||||
@@ -246,7 +274,18 @@ async function deliverLive(
|
|||||||
},
|
},
|
||||||
})
|
})
|
||||||
if (promptResult.status === "failed" && isAmbiguousPromptDispatchFailure(promptResult.error)) {
|
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", {
|
log("[team-mailbox] live delivery prompt failed after dispatch attempt, keeping reservation pending", {
|
||||||
teamRunId,
|
teamRunId,
|
||||||
recipient: recipientName,
|
recipient: recipientName,
|
||||||
@@ -271,7 +310,18 @@ async function deliverLive(
|
|||||||
})
|
})
|
||||||
continue
|
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", {
|
log("[team-mailbox] live delivery reserved until recipient idle", {
|
||||||
teamRunId,
|
teamRunId,
|
||||||
recipient: recipientName,
|
recipient: recipientName,
|
||||||
|
|||||||
@@ -711,6 +711,46 @@ describe("dispatchInternalPrompt shared gate behavior", () => {
|
|||||||
expect(promptCalls).toBe(0)
|
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 () => {
|
test("#given internal user tail follows an assistant waiting on tools #when an internal promptAsync is requested #then no prompt is sent", async () => {
|
||||||
// given
|
// given
|
||||||
let promptCalls = 0
|
let promptCalls = 0
|
||||||
|
|||||||
@@ -567,6 +567,48 @@ describe("createTeamIdleWakeHint", () => {
|
|||||||
expect(processedEntries).toContain(`${messageId}.json`)
|
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 () => {
|
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
|
// given
|
||||||
const baseDir = await createTemporaryBaseDir()
|
const baseDir = await createTemporaryBaseDir()
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ import { loadRuntimeState, transitionRuntimeState } from "../../features/team-mo
|
|||||||
import { resolveSessionEventID } from "../../shared/event-session-id"
|
import { resolveSessionEventID } from "../../shared/event-session-id"
|
||||||
import { isAmbiguousPromptDispatchFailure } from "../../shared/prompt-failure-classifier"
|
import { isAmbiguousPromptDispatchFailure } from "../../shared/prompt-failure-classifier"
|
||||||
import { log } from "../../shared/logger"
|
import { log } from "../../shared/logger"
|
||||||
|
import { isSessionActive, settleAfterSessionIdle } from "../../shared/session-idle-settle"
|
||||||
import { dispatchInternalPrompt, isInternalPromptDispatchAccepted } from "../shared/prompt-async-gate"
|
import { dispatchInternalPrompt, isInternalPromptDispatchAccepted } from "../shared/prompt-async-gate"
|
||||||
|
|
||||||
type PromptAsyncInput = {
|
type PromptAsyncInput = {
|
||||||
@@ -73,6 +74,20 @@ export function createTeamIdleWakeHint(ctx: TeamIdleWakeHintContext, config: Tea
|
|||||||
|
|
||||||
const pendingInjectedMessageIds = [...memberEntry.pendingInjectedMessageIds]
|
const pendingInjectedMessageIds = [...memberEntry.pendingInjectedMessageIds]
|
||||||
if (pendingInjectedMessageIds.length > 0) {
|
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 ackMessages(runtimeState.teamRunId, memberEntry.name, pendingInjectedMessageIds, config)
|
||||||
await transitionRuntimeState(runtimeState.teamRunId, (currentRuntimeState) => ({
|
await transitionRuntimeState(runtimeState.teamRunId, (currentRuntimeState) => ({
|
||||||
...currentRuntimeState,
|
...currentRuntimeState,
|
||||||
|
|||||||
@@ -226,4 +226,89 @@ describe("createTeamMemberErrorHandler", () => {
|
|||||||
expect(inboxEntries).not.toContain(`.delivering-${messageId}.json`)
|
expect(inboxEntries).not.toContain(`.delivering-${messageId}.json`)
|
||||||
expect(inboxEntries).not.toContain("processed")
|
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: `<peer_message from="lead" messageId="${messageId}" kind="message">pending live delivery</peer_message>`,
|
||||||
|
},
|
||||||
|
],
|
||||||
|
},
|
||||||
|
],
|
||||||
|
}),
|
||||||
|
},
|
||||||
|
},
|
||||||
|
})
|
||||||
|
|
||||||
|
// 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")
|
||||||
|
})
|
||||||
})
|
})
|
||||||
|
|||||||
@@ -7,9 +7,23 @@ import {
|
|||||||
import { loadRuntimeState, transitionRuntimeState } from "../../features/team-mode/team-state-store/store"
|
import { loadRuntimeState, transitionRuntimeState } from "../../features/team-mode/team-state-store/store"
|
||||||
import { resolveSessionEventID } from "../../shared/event-session-id"
|
import { resolveSessionEventID } from "../../shared/event-session-id"
|
||||||
import { log } from "../../shared/logger"
|
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 } }
|
type HookInput = { event: { type: string; properties?: unknown } }
|
||||||
export type HookImpl = (input: HookInput) => Promise<void>
|
export type HookImpl = (input: HookInput) => Promise<void>
|
||||||
|
type TeamMemberErrorHandlerDeps = {
|
||||||
|
client?: {
|
||||||
|
session?: {
|
||||||
|
status?: () => Promise<unknown>
|
||||||
|
messages?: (input: { path: { id: string } }) => Promise<unknown>
|
||||||
|
}
|
||||||
|
}
|
||||||
|
settleMs?: number
|
||||||
|
}
|
||||||
|
|
||||||
function getErroredSessionID(properties: unknown): string | undefined {
|
function getErroredSessionID(properties: unknown): string | undefined {
|
||||||
return resolveSessionEventID(properties)
|
return resolveSessionEventID(properties)
|
||||||
@@ -31,7 +45,73 @@ async function requeuePendingLiveDeliveries(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
export function createTeamMemberErrorHandler(config: TeamModeConfig): HookImpl {
|
async function shouldKeepPendingLiveDeliveries(
|
||||||
|
deps: TeamMemberErrorHandlerDeps,
|
||||||
|
sessionID: string,
|
||||||
|
): Promise<boolean> {
|
||||||
|
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<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 valueContainsAnyMessageId(value: unknown, messageIds: ReadonlySet<string>): 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<boolean> {
|
||||||
|
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<void> => {
|
return async ({ event }: HookInput): Promise<void> => {
|
||||||
if (event.type !== "session.error") return
|
if (event.type !== "session.error") return
|
||||||
|
|
||||||
@@ -47,6 +127,29 @@ export function createTeamMemberErrorHandler(config: TeamModeConfig): HookImpl {
|
|||||||
const runtimeState = await loadRuntimeState(runtimeMember.teamRunId, config)
|
const runtimeState = await loadRuntimeState(runtimeMember.teamRunId, config)
|
||||||
const memberEntry = runtimeState.members.find((member) => member.name === runtimeMember.memberName)
|
const memberEntry = runtimeState.members.find((member) => member.name === runtimeMember.memberName)
|
||||||
const pendingInjectedMessageIds = memberEntry?.pendingInjectedMessageIds ?? []
|
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(
|
await requeuePendingLiveDeliveries(
|
||||||
runtimeState.teamRunId,
|
runtimeState.teamRunId,
|
||||||
runtimeMember.memberName,
|
runtimeMember.memberName,
|
||||||
|
|||||||
+1
-1
@@ -342,7 +342,7 @@ export function createEventHandler(args: {
|
|||||||
? createTeamLeadOrphanHandler(teamModeConfig, managers.tmuxSessionManager, managers.backgroundManager)
|
? createTeamLeadOrphanHandler(teamModeConfig, managers.tmuxSessionManager, managers.backgroundManager)
|
||||||
: undefined;
|
: undefined;
|
||||||
const teamMemberErrorHandler = teamModeConfig
|
const teamMemberErrorHandler = teamModeConfig
|
||||||
? createTeamMemberErrorHandler(teamModeConfig)
|
? createTeamMemberErrorHandler(teamModeConfig, { client: pluginContext.client })
|
||||||
: undefined;
|
: undefined;
|
||||||
const teamMemberStatusHandler = teamModeConfig
|
const teamMemberStatusHandler = teamModeConfig
|
||||||
? createTeamMemberStatusHandler(teamModeConfig)
|
? createTeamMemberStatusHandler(teamModeConfig)
|
||||||
|
|||||||
@@ -10,7 +10,7 @@ import {
|
|||||||
settleAfterSessionIdle,
|
settleAfterSessionIdle,
|
||||||
} from "./session-idle-settle"
|
} 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_DISPATCH_TIMEOUT_MS = 30_000
|
||||||
export const DEFAULT_PROMPT_GATE_MESSAGES_FETCH_TIMEOUT_MS = 5_000
|
export const DEFAULT_PROMPT_GATE_MESSAGES_FETCH_TIMEOUT_MS = 5_000
|
||||||
export const DEFAULT_PROMPT_QUEUE_RETRY_MS = 250
|
export const DEFAULT_PROMPT_QUEUE_RETRY_MS = 250
|
||||||
@@ -420,7 +420,7 @@ function latestAssistantTurnBlocksInternalPrompt(messages: unknown[]): boolean {
|
|||||||
const role = messageRole(message)
|
const role = messageRole(message)
|
||||||
if (role === "assistant") {
|
if (role === "assistant") {
|
||||||
const finish = messageFinish(message)
|
const finish = messageFinish(message)
|
||||||
if (finish === undefined) {
|
if (finish === undefined || finish === "unknown") {
|
||||||
return true
|
return true
|
||||||
}
|
}
|
||||||
if (!isRecord(message) || !Array.isArray(message.parts)) {
|
if (!isRecord(message) || !Array.isArray(message.parts)) {
|
||||||
|
|||||||
@@ -344,6 +344,9 @@ await dispatchInternalPrompt(options)
|
|||||||
}
|
}
|
||||||
|
|
||||||
const contents = await readFile(filePath, "utf8")
|
const contents = await readFile(filePath, "utf8")
|
||||||
|
if (!contents.includes("prompt")) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
if (detectRawPromptInSnippet(contents)) {
|
if (detectRawPromptInSnippet(contents)) {
|
||||||
offenders.push(relativeSourcePath(filePath))
|
offenders.push(relativeSourcePath(filePath))
|
||||||
}
|
}
|
||||||
@@ -395,6 +398,9 @@ await dispatchInternalPrompt(options)
|
|||||||
// when
|
// when
|
||||||
for (const filePath of files) {
|
for (const filePath of files) {
|
||||||
const contents = await readFile(filePath, "utf8")
|
const contents = await readFile(filePath, "utf8")
|
||||||
|
if (!contents.includes("dispatchInternalPrompt")) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
const missingLines = findPromptGateCallsWithoutQueueBehavior(filePath, contents)
|
const missingLines = findPromptGateCallsWithoutQueueBehavior(filePath, contents)
|
||||||
for (const line of missingLines) {
|
for (const line of missingLines) {
|
||||||
offenders.push(`${relativeSourcePath(filePath)}:${line}`)
|
offenders.push(`${relativeSourcePath(filePath)}:${line}`)
|
||||||
@@ -413,6 +419,12 @@ await dispatchInternalPrompt(options)
|
|||||||
// when
|
// when
|
||||||
for (const filePath of files) {
|
for (const filePath of files) {
|
||||||
const contents = await readFile(filePath, "utf8")
|
const contents = await readFile(filePath, "utf8")
|
||||||
|
if (
|
||||||
|
!contents.includes("promptWithModelSuggestionRetry")
|
||||||
|
&& !contents.includes("promptSyncWithModelSuggestionRetry")
|
||||||
|
) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
const missingLines = findPromptRetryCallsWithoutQueueBehavior(filePath, contents)
|
const missingLines = findPromptRetryCallsWithoutQueueBehavior(filePath, contents)
|
||||||
for (const line of missingLines) {
|
for (const line of missingLines) {
|
||||||
offenders.push(`${relativeSourcePath(filePath)}:${line}`)
|
offenders.push(`${relativeSourcePath(filePath)}:${line}`)
|
||||||
|
|||||||
Reference in New Issue
Block a user