Skip to content

Commit 0607489

Browse files
committed
fix(reply): preserve upstream preflight cancellation
1 parent e3e6bed commit 0607489

2 files changed

Lines changed: 196 additions & 42 deletions

File tree

src/auto-reply/reply/agent-runner-memory.test.ts

Lines changed: 110 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1078,11 +1078,22 @@ describe("runMemoryFlushIfNeeded", () => {
10781078
};
10791079
const replyAbortController = new AbortController();
10801080
const explicitAbortController = new AbortController();
1081-
replyAbortController.abort(new Error("reply lifecycle timeout"));
1081+
const timeoutError = new Error("reply lifecycle timeout");
1082+
timeoutError.name = "TimeoutError";
10821083
const replyOperation = createReplyOperation({
10831084
abortSignal: replyAbortController.signal,
10841085
explicitAbortSignal: explicitAbortController.signal,
10851086
});
1087+
let compactAbortSignal: AbortSignal | undefined;
1088+
compactEmbeddedAgentSessionMock.mockImplementationOnce(async (call) => {
1089+
compactAbortSignal = (call as CompactEmbeddedAgentSessionParams).abortSignal;
1090+
expect(compactAbortSignal?.aborted).toBe(false);
1091+
1092+
replyAbortController.abort(timeoutError);
1093+
1094+
expect(compactAbortSignal?.aborted).toBe(false);
1095+
return { ok: true, compacted: true, result: { tokensAfter: 42 } };
1096+
});
10861097

10871098
await runPreflightCompactionIfNeeded({
10881099
cfg: { agents: { defaults: { compaction: { memoryFlush: {} } } } },
@@ -1102,9 +1113,105 @@ describe("runMemoryFlushIfNeeded", () => {
11021113

11031114
expect(compactEmbeddedAgentSessionMock).toHaveBeenCalledTimes(1);
11041115
const compactCall = requireCompactEmbeddedAgentSessionCall();
1105-
expect(compactCall.abortSignal).toBe(explicitAbortController.signal);
1116+
expect(compactCall.abortSignal).not.toBe(explicitAbortController.signal);
11061117
expect(compactCall.abortSignal).not.toBe(replyOperation.abortSignal);
1107-
expect(compactCall.abortSignal?.aborted).toBe(false);
1118+
expect(compactCall.abortSignal).toBe(compactAbortSignal);
1119+
expect(compactAbortSignal?.aborted).toBe(false);
1120+
});
1121+
1122+
it("cancels required preflight compaction on upstream non-timeout abort", async () => {
1123+
const sessionEntry: SessionEntry = {
1124+
sessionId: "session",
1125+
updatedAt: Date.now(),
1126+
totalTokens: 180_499,
1127+
totalTokensFresh: true,
1128+
};
1129+
const replyAbortController = new AbortController();
1130+
const replyOperation = createReplyOperation({
1131+
abortSignal: replyAbortController.signal,
1132+
});
1133+
const abortError = new Error("chat run aborted");
1134+
abortError.name = "AbortError";
1135+
let compactAbortSignal: AbortSignal | undefined;
1136+
compactEmbeddedAgentSessionMock.mockImplementationOnce(async (call) => {
1137+
compactAbortSignal = (call as CompactEmbeddedAgentSessionParams).abortSignal;
1138+
expect(compactAbortSignal?.aborted).toBe(false);
1139+
1140+
replyAbortController.abort(abortError);
1141+
1142+
expect(compactAbortSignal?.aborted).toBe(true);
1143+
expect(compactAbortSignal?.reason).toBe(abortError);
1144+
return { ok: true, compacted: true, result: { tokensAfter: 42 } };
1145+
});
1146+
1147+
await runPreflightCompactionIfNeeded({
1148+
cfg: { agents: { defaults: { compaction: { memoryFlush: {} } } } },
1149+
followupRun: createTestFollowupRun({
1150+
sessionId: "session",
1151+
sessionKey: "agent:main:main",
1152+
}),
1153+
defaultModel: "anthropic/claude-opus-4-6",
1154+
agentCfgContextTokens: 200_000,
1155+
sessionEntry,
1156+
sessionStore: { "agent:main:main": sessionEntry },
1157+
sessionKey: "agent:main:main",
1158+
storePath: path.join(rootDir, "sessions.json"),
1159+
isHeartbeat: false,
1160+
replyOperation,
1161+
});
1162+
1163+
expect(compactEmbeddedAgentSessionMock).toHaveBeenCalledTimes(1);
1164+
expect(compactAbortSignal?.aborted).toBe(true);
1165+
});
1166+
1167+
it("still cancels required preflight compaction on explicit abort after lifecycle timeout", async () => {
1168+
const sessionEntry: SessionEntry = {
1169+
sessionId: "session",
1170+
updatedAt: Date.now(),
1171+
totalTokens: 180_499,
1172+
totalTokensFresh: true,
1173+
};
1174+
const replyAbortController = new AbortController();
1175+
const explicitAbortController = new AbortController();
1176+
const timeoutError = new Error("reply lifecycle timeout");
1177+
timeoutError.name = "TimeoutError";
1178+
replyAbortController.abort(timeoutError);
1179+
const replyOperation = createReplyOperation({
1180+
abortSignal: replyAbortController.signal,
1181+
explicitAbortSignal: explicitAbortController.signal,
1182+
});
1183+
const explicitAbortError = new Error("reply operation aborted by user");
1184+
explicitAbortError.name = "AbortError";
1185+
let compactAbortSignal: AbortSignal | undefined;
1186+
compactEmbeddedAgentSessionMock.mockImplementationOnce(async (call) => {
1187+
compactAbortSignal = (call as CompactEmbeddedAgentSessionParams).abortSignal;
1188+
expect(compactAbortSignal?.aborted).toBe(false);
1189+
1190+
explicitAbortController.abort(explicitAbortError);
1191+
1192+
expect(compactAbortSignal?.aborted).toBe(true);
1193+
expect(compactAbortSignal?.reason).toBe(explicitAbortError);
1194+
return { ok: true, compacted: true, result: { tokensAfter: 42 } };
1195+
});
1196+
1197+
await runPreflightCompactionIfNeeded({
1198+
cfg: { agents: { defaults: { compaction: { memoryFlush: {} } } } },
1199+
followupRun: createTestFollowupRun({
1200+
sessionId: "session",
1201+
sessionKey: "agent:main:main",
1202+
}),
1203+
defaultModel: "anthropic/claude-opus-4-6",
1204+
agentCfgContextTokens: 200_000,
1205+
sessionEntry,
1206+
sessionStore: { "agent:main:main": sessionEntry },
1207+
sessionKey: "agent:main:main",
1208+
storePath: path.join(rootDir, "sessions.json"),
1209+
isHeartbeat: false,
1210+
replyOperation,
1211+
});
1212+
1213+
expect(compactEmbeddedAgentSessionMock).toHaveBeenCalledTimes(1);
1214+
expect(compactAbortSignal?.aborted).toBe(true);
11081215
});
11091216

11101217
it("fails when required preflight context-engine compaction is deferred to background maintenance", async () => {

src/auto-reply/reply/agent-runner-memory.ts

Lines changed: 86 additions & 39 deletions
Original file line numberDiff line numberDiff line change
@@ -9,6 +9,7 @@ import {
99
import { resolveBootstrapWarningSignaturesSeen } from "../../agents/bootstrap-budget.js";
1010
import { estimateMessagesTokens } from "../../agents/compaction.js";
1111
import { classifyCompactionReason } from "../../agents/embedded-agent-runner/compact-reasons.js";
12+
import { isSignalTimeoutReason } from "../../agents/failover-error.js";
1213
import { resolveAgentHarnessPolicy } from "../../agents/harness/policy.js";
1314
import { ensureSelectedAgentHarnessPlugin } from "../../agents/harness/runtime-plugin.js";
1415
import { runWithModelFallback } from "../../agents/model-fallback.js";
@@ -63,6 +64,10 @@ import type { ReplyOperation } from "./reply-run-registry.js";
6364
import { incrementCompactionCount } from "./session-updates.js";
6465

6566
type EmbeddedAgentRuntime = typeof import("../../agents/embedded-agent.js");
67+
type ScopedAbortSignal = {
68+
abortSignal: AbortSignal;
69+
cleanup: () => void;
70+
};
6671
type UpdateSessionEntryParams = {
6772
storePath: string;
6873
sessionKey: string;
@@ -85,6 +90,44 @@ function loadEmbeddedAgentRuntime(): Promise<EmbeddedAgentRuntime> {
8590
return embeddedAgentRuntimeLoader.load();
8691
}
8792

93+
function createPreflightCompactionAbortSignal(operation: ReplyOperation): ScopedAbortSignal {
94+
const controller = new AbortController();
95+
const cleanupHandlers: Array<() => void> = [];
96+
const abortFrom = (signal: AbortSignal, opts?: { ignoreTimeoutReason?: boolean }) => {
97+
if (controller.signal.aborted) {
98+
return;
99+
}
100+
if (opts?.ignoreTimeoutReason && isSignalTimeoutReason(signal.reason)) {
101+
return;
102+
}
103+
controller.abort(signal.reason);
104+
};
105+
const watch = (signal: AbortSignal | undefined, opts?: { ignoreTimeoutReason?: boolean }) => {
106+
if (!signal) {
107+
return;
108+
}
109+
if (signal.aborted) {
110+
abortFrom(signal, opts);
111+
return;
112+
}
113+
const onAbort = () => abortFrom(signal, opts);
114+
signal.addEventListener("abort", onAbort, { once: true });
115+
cleanupHandlers.push(() => signal.removeEventListener("abort", onAbort));
116+
};
117+
118+
watch(operation.abortSignal, { ignoreTimeoutReason: true });
119+
watch(operation.explicitAbortSignal);
120+
121+
return {
122+
abortSignal: controller.signal,
123+
cleanup: () => {
124+
for (const cleanup of cleanupHandlers.splice(0)) {
125+
cleanup();
126+
}
127+
},
128+
};
129+
}
130+
88131
async function compactEmbeddedAgentSessionDefault(
89132
...args: Parameters<typeof import("../../agents/embedded-agent.js").compactEmbeddedAgentSession>
90133
): Promise<
@@ -943,45 +986,49 @@ export async function runPreflightCompactionIfNeeded(params: {
943986
params.sessionKey ?? params.followupRun.run.sessionKey,
944987
{ storePath: params.storePath },
945988
);
946-
const result = await deps.compactEmbeddedAgentSession({
947-
sessionId: entry.sessionId,
948-
sessionKey: params.sessionKey,
949-
sandboxSessionKey: params.runtimePolicySessionKey,
950-
allowGatewaySubagentBinding: true,
951-
messageChannel: params.followupRun.run.messageProvider,
952-
groupId: entry.groupId ?? params.followupRun.run.groupId,
953-
groupChannel: entry.groupChannel ?? params.followupRun.run.groupChannel,
954-
groupSpace: entry.space ?? params.followupRun.run.groupSpace,
955-
senderId: params.followupRun.run.senderId,
956-
senderName: params.followupRun.run.senderName,
957-
senderUsername: params.followupRun.run.senderUsername,
958-
senderE164: params.followupRun.run.senderE164,
959-
sessionFile: sessionFile ?? params.followupRun.run.sessionFile,
960-
workspaceDir: params.followupRun.run.workspaceDir,
961-
cwd: params.followupRun.run.cwd,
962-
agentDir: params.followupRun.run.agentDir,
963-
config: params.cfg,
964-
skillsSnapshot: entry.skillsSnapshot ?? params.followupRun.run.skillsSnapshot,
965-
provider: params.followupRun.run.provider,
966-
model: params.followupRun.run.model,
967-
authProfileId: params.followupRun.run.authProfileId,
968-
agentHarnessId:
969-
entry.sessionId === params.followupRun.run.sessionId ? entry.agentHarnessId : undefined,
970-
thinkLevel: params.followupRun.run.thinkLevel,
971-
bashElevated: params.followupRun.run.bashElevated,
972-
...(params.replyOperation.explicitAbortSignal
973-
? { abortSignal: params.replyOperation.explicitAbortSignal }
974-
: {}),
975-
trigger: "budget",
976-
force: true,
977-
forcePreflight: true,
978-
preflightRequired: true,
979-
preflightCompactionTrigger: compactionTrigger,
980-
deferOwningContextEngineCompaction: false,
981-
contextTokenBudget: contextWindowTokens,
982-
currentTokenCount: tokenCountForCompaction ?? freshPersistedTokens,
983-
ownerNumbers: params.followupRun.run.ownerNumbers,
984-
});
989+
const preflightAbort = createPreflightCompactionAbortSignal(params.replyOperation);
990+
let result: Awaited<ReturnType<typeof deps.compactEmbeddedAgentSession>> | undefined;
991+
try {
992+
result = await deps.compactEmbeddedAgentSession({
993+
sessionId: entry.sessionId,
994+
sessionKey: params.sessionKey,
995+
sandboxSessionKey: params.runtimePolicySessionKey,
996+
allowGatewaySubagentBinding: true,
997+
messageChannel: params.followupRun.run.messageProvider,
998+
groupId: entry.groupId ?? params.followupRun.run.groupId,
999+
groupChannel: entry.groupChannel ?? params.followupRun.run.groupChannel,
1000+
groupSpace: entry.space ?? params.followupRun.run.groupSpace,
1001+
senderId: params.followupRun.run.senderId,
1002+
senderName: params.followupRun.run.senderName,
1003+
senderUsername: params.followupRun.run.senderUsername,
1004+
senderE164: params.followupRun.run.senderE164,
1005+
sessionFile: sessionFile ?? params.followupRun.run.sessionFile,
1006+
workspaceDir: params.followupRun.run.workspaceDir,
1007+
cwd: params.followupRun.run.cwd,
1008+
agentDir: params.followupRun.run.agentDir,
1009+
config: params.cfg,
1010+
skillsSnapshot: entry.skillsSnapshot ?? params.followupRun.run.skillsSnapshot,
1011+
provider: params.followupRun.run.provider,
1012+
model: params.followupRun.run.model,
1013+
authProfileId: params.followupRun.run.authProfileId,
1014+
agentHarnessId:
1015+
entry.sessionId === params.followupRun.run.sessionId ? entry.agentHarnessId : undefined,
1016+
thinkLevel: params.followupRun.run.thinkLevel,
1017+
bashElevated: params.followupRun.run.bashElevated,
1018+
abortSignal: preflightAbort.abortSignal,
1019+
trigger: "budget",
1020+
force: true,
1021+
forcePreflight: true,
1022+
preflightRequired: true,
1023+
preflightCompactionTrigger: compactionTrigger,
1024+
deferOwningContextEngineCompaction: false,
1025+
contextTokenBudget: contextWindowTokens,
1026+
currentTokenCount: tokenCountForCompaction ?? freshPersistedTokens,
1027+
ownerNumbers: params.followupRun.run.ownerNumbers,
1028+
});
1029+
} finally {
1030+
preflightAbort.cleanup();
1031+
}
9851032

9861033
if (!result?.ok) {
9871034
const reason = result?.reason ?? "not_compacted";

0 commit comments

Comments
 (0)