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!
This commit is contained in:
fsdwen
2026-07-02 03:55:15 +08:00
committed by GitHub
parent ab2f6f5642
commit 44cb3f7b98
5 changed files with 336 additions and 17 deletions
@@ -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<string, unknown>,
>(): Promise<T> => {
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);
});
});
+82 -13
View File
@@ -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<Record<AnnounceErrorChainKey, unknown>> & {
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<object> = 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<void> {
@@ -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 };
+55 -3
View File
@@ -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<string, unknown> };
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<AgentCallResponse> => ({
@@ -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,
});
});
});
+1 -1
View File
@@ -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}`,
@@ -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({