From 44cb3f7b980f96314238f3d11c8bb7e3f5d9922c Mon Sep 17 00:00:00 2001 From: fsdwen Date: Thu, 2 Jul 2026 03:55:15 +0800 Subject: [PATCH] fix(agents): prevent duplicate subagent announce sends Fixes the subagent announce prompt-lock race by treating lock-change failures with visible-send evidence as terminal, while keeping pre-send failures retryable. Scoped proof: - node scripts/run-vitest.mjs src/agents/subagent-announce-delivery.test.ts src/agents/subagent-announce.test.ts src/agents/subagent-registry-lifecycle.test.ts - node scripts/run-tsgo.mjs -p tsconfig.core.json - node scripts/run-tsgo.mjs -p test/tsconfig/tsconfig.core.test.agents.json - oxfmt --check on touched agent files - git diff --check - .agents/skills/autoreview/scripts/autoreview --mode local Thanks @fsdwen! --- src/agents/subagent-announce-delivery.test.ts | 154 ++++++++++++++++++ src/agents/subagent-announce-delivery.ts | 95 +++++++++-- src/agents/subagent-announce.test.ts | 58 ++++++- src/agents/subagent-announce.ts | 2 +- .../subagent-registry-lifecycle.test.ts | 44 +++++ 5 files changed, 336 insertions(+), 17 deletions(-) diff --git a/src/agents/subagent-announce-delivery.test.ts b/src/agents/subagent-announce-delivery.test.ts index 8002893be521..4ba394256138 100644 --- a/src/agents/subagent-announce-delivery.test.ts +++ b/src/agents/subagent-announce-delivery.test.ts @@ -4968,4 +4968,158 @@ describe("deliverSubagentAnnouncement completion delivery", () => { bestEffortDeliver: true, }); }); + + it("does not retry session-file-changed failures with send evidence", async () => { + const sendErr = new OutboundDeliveryError("outbound delivery failed", { + cause: new Error("outbound delivery failed"), + results: [{ channel: "telegram", messageId: "msg-1" }], + }); + const callGateway: typeof runtimeCallGateway = vi.fn(async () => { + throw new Error("session file changed while embedded prompt lock was released", { + cause: sendErr, + }); + }); + const queueEmbeddedAgentMessageWithOutcome = createQueueOutcomeSequenceMock(["no_active_run"]); + const result = await deliverSlackChannelAnnouncement({ + callGateway, + queueEmbeddedAgentMessageWithOutcome, + sessionId: "requester-session-lock-race-evidence", + isActive: true, + expectsCompletionMessage: true, + directIdempotencyKey: "announce-permanent-lock-error-evidence", + }); + + expect(result.delivered).toBe(false); + expect(result.path).toBe("direct"); + expect(result.terminal).toBe(true); + expect(result.phases?.map((phase) => phase.phase)).toEqual(["direct-primary"]); + expect(callGateway).toHaveBeenCalledTimes(1); + expect(queueEmbeddedAgentMessageWithOutcome).toHaveBeenCalledTimes(1); + }); + + it("does not fallback-steer after wrapped prompt-lock takeover with send evidence", async () => { + const takeoverErr = Object.assign( + new Error("session file changed while embedded prompt lock was released: /tmp/session.jsonl"), + { name: "EmbeddedAttemptSessionTakeoverError" }, + ); + + const promptErr = Object.assign(new Error("some model error"), { visibleReplySent: true }); + const wrapperErr = Object.assign(new Error("some model error", { cause: takeoverErr }), { + name: "EmbeddedAttemptSessionTakeoverError", + cleanupError: takeoverErr, + promptError: promptErr, + }); + + const callGateway: typeof runtimeCallGateway = vi.fn(async () => { + throw wrapperErr; + }); + const queueEmbeddedAgentMessageWithOutcome = createQueueOutcomeSequenceMock(["no_active_run"]); + const result = await deliverSlackChannelAnnouncement({ + callGateway, + queueEmbeddedAgentMessageWithOutcome, + sessionId: "requester-session-lock-race-wrapped-evidence", + isActive: true, + expectsCompletionMessage: true, + directIdempotencyKey: "announce-permanent-wrapped-lock-error-evidence", + }); + + expect(result.delivered).toBe(false); + expect(result.path).toBe("direct"); + expect(result.error).toBe("some model error"); + expect(result.terminal).toBe(true); + expect(result.phases?.map((phase) => phase.phase)).toEqual(["direct-primary"]); + expect(callGateway).toHaveBeenCalledTimes(1); + expect(queueEmbeddedAgentMessageWithOutcome).toHaveBeenCalledTimes(1); + }); + + it("retries session-file-changed failures without send evidence", async () => { + let attempts = 0; + const callGatewaySpy = vi.fn(); + const callGateway: typeof runtimeCallGateway = async < + T = Record, + >(): Promise => { + callGatewaySpy(); + attempts++; + if (attempts <= 1) { + throw new Error("session file changed while embedded prompt lock was released"); + } + return { + result: { + payloads: [{ text: "recovered after retry" }], + }, + } as T; + }; + const queueEmbeddedAgentMessageWithOutcome = createQueueOutcomeSequenceMock(["no_active_run"]); + const result = await deliverSlackChannelAnnouncement({ + callGateway, + queueEmbeddedAgentMessageWithOutcome, + sessionId: "requester-session-lock-race-no-evidence", + isActive: true, + expectsCompletionMessage: true, + directIdempotencyKey: "announce-retry-lock-error-no-evidence", + }); + + expect(result.delivered).toBe(true); + expect(result.path).toBe("direct"); + expect(callGatewaySpy).toHaveBeenCalledTimes(2); + }); + + it("detects send evidence from OutboundDeliveryError in the error chain", () => { + const err = new Error( + "session file changed while embedded prompt lock was released: /tmp/session.jsonl", + { + cause: new OutboundDeliveryError("outbound delivery failed", { + cause: new Error("outbound delivery failed"), + results: [{ channel: "telegram", messageId: "msg-1" }], + }), + }, + ); + + expect(testing.isSessionFileChangedAnnounceError(err.message)).toBe(true); + expect(testing.hasAnnounceSendEvidence(err)).toBe(true); + }); + + it("classifies session-file-changed error as no-send-evidence when the error chain has no send markers", () => { + const err = new Error( + "session file changed while embedded prompt lock was released: /tmp/session.jsonl", + ); + + expect(testing.isSessionFileChangedAnnounceError(err.message)).toBe(true); + expect(testing.hasAnnounceSendEvidence(err)).toBe(false); + }); + + it("detects send evidence from visibleReplySent flag on session-file-changed error", () => { + const err = Object.assign( + new Error("session file changed while embedded prompt lock was released: /tmp/session.jsonl"), + { visibleReplySent: true }, + ); + + expect(testing.hasAnnounceSendEvidence(err)).toBe(true); + }); + + it("detects send evidence from sentBeforeError flag on session-file-changed error", () => { + const err = Object.assign( + new Error("session file changed while embedded prompt lock was released: /tmp/session.jsonl"), + { sentBeforeError: true }, + ); + + expect(testing.hasAnnounceSendEvidence(err)).toBe(true); + }); + + it("detects send evidence recursively through promptError", () => { + const takeoverErr = Object.assign( + new Error("session file changed while embedded prompt lock was released: /tmp/session.jsonl"), + { name: "EmbeddedAttemptSessionTakeoverError" }, + ); + + const promptErr = Object.assign(new Error("some model error"), { visibleReplySent: true }); + + const wrapperErr = Object.assign(new Error("some model error", { cause: takeoverErr }), { + name: "EmbeddedAttemptSessionTakeoverError", + promptError: promptErr, + }); + + expect(testing.hasAnnounceSendEvidence(wrapperErr)).toBe(true); + expect(testing.hasSessionFileChangedAnnounceError(wrapperErr)).toBe(true); + }); }); diff --git a/src/agents/subagent-announce-delivery.ts b/src/agents/subagent-announce-delivery.ts index 9eec384fd060..3b6e3925b6d7 100644 --- a/src/agents/subagent-announce-delivery.ts +++ b/src/agents/subagent-announce-delivery.ts @@ -380,6 +380,9 @@ const TRANSIENT_ANNOUNCE_DELIVERY_ERROR_PATTERNS: readonly RegExp[] = [ /\b(econnreset|econnrefused|etimedout|enotfound|ehostunreach|network error)\b/i, ]; +const SESSION_FILE_CHANGED_ANNOUNCE_RE = + /session file changed while embedded prompt lock was released/i; + const PERMANENT_ANNOUNCE_DELIVERY_ERROR_PATTERNS: readonly RegExp[] = [ /unsupported channel/i, /unknown channel/i, @@ -390,14 +393,73 @@ const PERMANENT_ANNOUNCE_DELIVERY_ERROR_PATTERNS: readonly RegExp[] = [ /forbidden: bot was kicked/i, /recipient is not a valid/i, /outbound not configured for channel/i, + SESSION_FILE_CHANGED_ANNOUNCE_RE, ]; +function isSessionFileChangedAnnounceError(message: string): boolean { + return SESSION_FILE_CHANGED_ANNOUNCE_RE.test(message); +} + +const ANNOUNCE_ERROR_CHAIN_KEYS = [ + "cause", + "cleanupError", + "error", + "promptError", + "reason", +] as const; +type AnnounceErrorChainKey = (typeof ANNOUNCE_ERROR_CHAIN_KEYS)[number]; +type AnnounceErrorRecord = Partial> & { + sentBeforeError?: unknown; + visibleReplySent?: unknown; +}; + +function isAnnounceErrorRecord(error: unknown): error is AnnounceErrorRecord { + return Boolean(error && typeof error === "object"); +} + +function hasAnnounceErrorMatch( + error: unknown, + matches: (candidate: unknown) => boolean, + seen: Set = new Set(), +): boolean { + if (matches(error)) { + return true; + } + if (!isAnnounceErrorRecord(error)) { + return false; + } + if (seen.has(error)) { + return false; + } + seen.add(error); + + return ANNOUNCE_ERROR_CHAIN_KEYS.some((key) => hasAnnounceErrorMatch(error[key], matches, seen)); +} + +function hasSessionFileChangedAnnounceError(error: unknown): boolean { + return hasAnnounceErrorMatch(error, (candidate) => + isSessionFileChangedAnnounceError(summarizeDeliveryError(candidate)), + ); +} + function isTransientAnnounceDeliveryError(error: unknown): boolean { const message = summarizeDeliveryError(error); + const topLevelPermanent = Boolean( + message && PERMANENT_ANNOUNCE_DELIVERY_ERROR_PATTERNS.some((re) => re.test(message)), + ); + if (topLevelPermanent && !isSessionFileChangedAnnounceError(message)) { + return false; + } + + const sessionFileChanged = hasSessionFileChangedAnnounceError(error); + if (sessionFileChanged) { + return !hasAnnounceSendEvidence(error); + } + if (!message) { return false; } - if (PERMANENT_ANNOUNCE_DELIVERY_ERROR_PATTERNS.some((re) => re.test(message))) { + if (topLevelPermanent) { return false; } return TRANSIENT_ANNOUNCE_DELIVERY_ERROR_PATTERNS.some((re) => re.test(message)); @@ -405,8 +467,9 @@ function isTransientAnnounceDeliveryError(error: unknown): boolean { function isPermanentAnnounceDeliveryError(error: unknown): boolean { const message = summarizeDeliveryError(error); - return Boolean( - message && PERMANENT_ANNOUNCE_DELIVERY_ERROR_PATTERNS.some((re) => re.test(message)), + return ( + (message && PERMANENT_ANNOUNCE_DELIVERY_ERROR_PATTERNS.some((re) => re.test(message))) || + hasSessionFileChangedAnnounceError(error) ); } @@ -426,17 +489,18 @@ function isSessionWriteLockAnnounceAgentError(error: unknown): boolean { ); } -function didVisibleSendFailAfterPartialDelivery(error: unknown): boolean { +function hasDirectAnnounceSendEvidence(error: unknown): boolean { if (isOutboundDeliveryError(error) && error.sentBeforeError) { return true; } - const maybeDeliveryError = error as { - sentBeforeError?: unknown; - visibleReplySent?: unknown; - }; - return ( - maybeDeliveryError.sentBeforeError === true || maybeDeliveryError.visibleReplySent === true - ); + if (!isAnnounceErrorRecord(error)) { + return false; + } + return error.sentBeforeError === true || error.visibleReplySent === true; +} + +function hasAnnounceSendEvidence(error: unknown): boolean { + return hasAnnounceErrorMatch(error, hasDirectAnnounceSendEvidence); } async function waitForAnnounceRetryDelay(ms: number, signal?: AbortSignal): Promise { @@ -891,7 +955,7 @@ async function deliverGeneratedMediaCompletionDirect(params: { path: "direct", }; } catch (err) { - const terminal = didVisibleSendFailAfterPartialDelivery(err); + const terminal = hasAnnounceSendEvidence(err); return { delivered: false, path: "direct", @@ -1508,7 +1572,7 @@ async function sendSubagentAnnounceDirectly(params: { }), }); } catch (err) { - if (isPermanentAnnounceDeliveryError(err)) { + if (isPermanentAnnounceDeliveryError(err) && hasAnnounceSendEvidence(err)) { throw err; } if ( @@ -1695,10 +1759,12 @@ async function sendSubagentAnnounceDirectly(params: { path: "direct", }; } catch (err) { + const terminal = isPermanentAnnounceDeliveryError(err) && hasAnnounceSendEvidence(err); return { delivered: false, path: "direct", error: summarizeDeliveryError(err), + ...(terminal ? { terminal: true } : {}), }; } } @@ -1785,5 +1851,8 @@ export const testing = { } : defaultSubagentAnnounceDeliveryDeps; }, + hasAnnounceSendEvidence, + hasSessionFileChangedAnnounceError, + isSessionFileChangedAnnounceError, }; export { testing as __testing }; diff --git a/src/agents/subagent-announce.test.ts b/src/agents/subagent-announce.test.ts index 480c7266c83d..36cb89d05ff5 100644 --- a/src/agents/subagent-announce.test.ts +++ b/src/agents/subagent-announce.test.ts @@ -5,7 +5,7 @@ import type { EmbeddedAgentQueueMessageOutcome } from "./embedded-agent-runner/r import { createSubagentAnnounceDeliveryRuntimeMock } from "./subagent-announce.test-support.js"; type AgentCallRequest = { method?: string; params?: Record }; -type AgentCallResponse = { runId?: string; status: string; error?: string }; +type AgentCallResponse = { runId?: string; status: string; error?: string; terminal?: boolean }; const agentSpy = vi.fn( async (_req: AgentCallRequest): Promise => ({ @@ -149,10 +149,15 @@ vi.mock("./subagent-announce-delivery.js", () => ({ threadId: effectiveOrigin?.threadId, }), }, - })) as { status?: string; error?: string }; + })) as { status?: string; error?: string; terminal?: boolean }; if (response.status === "error") { - return { delivered: false, path: "direct", error: response.error ?? "agent delivery failed" }; + return { + delivered: false, + path: "direct", + error: response.error ?? "agent delivery failed", + ...(response.terminal === true ? { terminal: true } : {}), + }; } return { delivered: true, path: "direct" }; @@ -599,4 +604,51 @@ describe("subagent announce seam flow", () => { ); logSpy.mockRestore(); }); + + it("treats terminal direct completion failures as announced for cleanup", async () => { + let deliveryResult: + | { + delivered: boolean; + path: string; + error?: string; + terminal?: boolean; + } + | undefined; + agentSpy.mockResolvedValueOnce({ + status: "error", + error: "prompt lock failed after visible send", + terminal: true, + }); + + const didAnnounce = await runSubagentAnnounceFlow({ + childSessionKey: "agent:main:subagent:slack", + childRunId: "run-terminal-direct-failure", + requesterSessionKey: "agent:main:main", + requesterDisplayKey: "main", + requesterOrigin: { + channel: "slack", + to: "C123", + }, + task: "deliver completion", + timeoutMs: 10, + cleanup: "keep", + waitForCompletion: false, + startedAt: 10, + endedAt: 20, + outcome: { status: "ok" }, + roundOneReply: "done", + expectsCompletionMessage: true, + onDeliveryResult: (delivery) => { + deliveryResult = delivery; + }, + }); + + expect(didAnnounce).toBe(true); + expect(deliveryResult).toMatchObject({ + delivered: false, + path: "direct", + error: "prompt lock failed after visible send", + terminal: true, + }); + }); }); diff --git a/src/agents/subagent-announce.ts b/src/agents/subagent-announce.ts index 02ef26e82348..a5ec24fa13bf 100644 --- a/src/agents/subagent-announce.ts +++ b/src/agents/subagent-announce.ts @@ -597,7 +597,7 @@ export async function runSubagentAnnounceFlow(params: { signal: params.signal, }); params.onDeliveryResult?.(delivery); - didAnnounce = delivery.delivered; + didAnnounce = delivery.delivered || delivery.terminal === true; if (!delivery.delivered && delivery.path === "direct" && delivery.error) { defaultRuntime.log( `[warn] Subagent completion direct announce failed for run ${params.childRunId}: ${delivery.error}`, diff --git a/src/agents/subagent-registry-lifecycle.test.ts b/src/agents/subagent-registry-lifecycle.test.ts index 47af8d35cc3b..c92821015293 100644 --- a/src/agents/subagent-registry-lifecycle.test.ts +++ b/src/agents/subagent-registry-lifecycle.test.ts @@ -587,6 +587,50 @@ describe("subagent registry lifecycle hardening", () => { }); }); + it("finalizes terminal visible-send failures without scheduling completion retry", async () => { + const persist = vi.fn(); + const entry = createRunEntry({ + endedAt: 4_000, + expectsCompletionMessage: true, + retainAttachmentsOnKeep: true, + }); + const runSubagentAnnounceFlow: LifecycleControllerParams["runSubagentAnnounceFlow"] = vi.fn( + async (announceParams) => { + announceParams.onDeliveryResult?.({ + delivered: false, + path: "direct", + error: "prompt lock failed after visible send", + terminal: true, + }); + return true; + }, + ); + + const controller = createLifecycleController({ entry, persist, runSubagentAnnounceFlow }); + + await expect( + controller.completeSubagentRun({ + runId: entry.runId, + endedAt: 4_000, + outcome: { status: "ok" }, + reason: SUBAGENT_ENDED_REASON_COMPLETE, + triggerCleanup: true, + }), + ).resolves.toBeUndefined(); + + await vi.waitFor(() => expect(entry.cleanupCompletedAt).toBeTypeOf("number")); + expect(entry.delivery?.status).toBe("delivered"); + expect(entry.delivery?.lastError).toBeUndefined(); + expect(entry.delivery?.payload).toBeUndefined(); + expect(entry.delivery?.suspendedAt).toBeUndefined(); + expect(entry.delivery?.suspendedReason).toBeUndefined(); + expect(runSubagentAnnounceFlow).toHaveBeenCalledTimes(1); + expectFields(firstCallArg(taskExecutorMocks.setDetachedTaskDeliveryStatusByRunId), { + runId: entry.runId, + deliveryStatus: "delivered", + }); + }); + it("skips announce delivery when completion messages are disabled", async () => { const persist = vi.fn(); const entry = createRunEntry({