Files
oh-my-pi/packages/coding-agent/test/agent-session-auto-compaction-queue.test.ts
T
metaphorics 46f4993f65 fix(coding-agent): preserved queued steers/follow-ups during auto-compaction
Two defects dropped the first steering/follow-up message typed as
auto-compaction began:

- The compaction AbortController (which backs isCompacting) was installed
  AFTER auto_compaction_start was emitted. The emit awaits extension
  delivery and yields to the event loop, so a message typed as the loader
  appeared was read while isCompacting was false and mis-routed into the
  core agent queue (which the handoff reset then wiped). Install the
  controller before the emit, and move the emit to the first statement
  inside the existing try so the catch/finally cleanup still runs, in both
  #runAutoCompaction and #runAutoShake. The handoff branch now passes the
  run's local abort signal instead of the mutable controller field, so a
  superseded run bails at the handoff entry check rather than resetting the
  session.

- handoff() calls agent.reset(), which clears the core steering/follow-up
  queues. Capture both queues immediately before reset and restore them
  immediately after (synchronous, no await gap), so queued steers and
  follow-ups -- including in-flight RPC/SDK steer()/followUp() and a hidden
  user companion such as an ultrathink notice -- survive the new-session
  reset instead of being silently dropped.

Adds regression tests for both defects (controller-before-emit for the
context-full and shake paths, queue preservation across the reset for both
pre-enqueued and in-flight messages).
2026-06-18 09:15:01 +09:00

302 lines
11 KiB
TypeScript

import { afterEach, beforeEach, describe, expect, it, vi } from "bun:test";
import * as fs from "node:fs";
import * as path from "node:path";
import { Agent } from "@oh-my-pi/pi-agent-core";
import { getBundledModel } from "@oh-my-pi/pi-catalog/models";
import { ModelRegistry } from "@oh-my-pi/pi-coding-agent/config/model-registry";
import { Settings } from "@oh-my-pi/pi-coding-agent/config/settings";
import { loadExtensions } from "@oh-my-pi/pi-coding-agent/extensibility/extensions/loader";
import { ExtensionRunner } from "@oh-my-pi/pi-coding-agent/extensibility/extensions/runner";
import { AgentSession } from "@oh-my-pi/pi-coding-agent/session/agent-session";
import { AuthStorage } from "@oh-my-pi/pi-coding-agent/session/auth-storage";
import { SessionManager } from "@oh-my-pi/pi-coding-agent/session/session-manager";
import { getProjectAgentDir, TempDir, withTimeout } from "@oh-my-pi/pi-utils";
const runtimeSignalStoreKey = "__ompRuntimeSignals";
type RuntimeSignalGlobal = typeof globalThis & { [runtimeSignalStoreKey]?: string[] };
function getRuntimeSignals(): string[] {
const globalWithSignals = globalThis as RuntimeSignalGlobal;
if (!globalWithSignals[runtimeSignalStoreKey]) {
globalWithSignals[runtimeSignalStoreKey] = [];
}
return globalWithSignals[runtimeSignalStoreKey];
}
/**
* Regression test: auto-compaction completion should resume the agent loop when
* there are queued agent-level messages (follow-up/steering/custom).
*/
describe("AgentSession auto-compaction queue resume", () => {
let tempDir: TempDir;
let session: AgentSession;
let sessionManager: SessionManager;
let authStorage: AuthStorage;
let modelRegistry: ModelRegistry;
beforeEach(async () => {
tempDir = TempDir.createSync("@pi-auto-compaction-queue-");
vi.useFakeTimers();
// Provide an extension that short-circuits compaction so the test doesn't
// make any LLM calls.
const extensionsDir = path.join(getProjectAgentDir(tempDir.path()), "extensions");
fs.mkdirSync(extensionsDir, { recursive: true });
const extensionPath = path.join(extensionsDir, "compaction-short-circuit.ts");
fs.writeFileSync(
extensionPath,
[
"export default function(pi) {",
'\tpi.on("session_before_compact", async (event) => {',
"\t\treturn {",
"\t\t\tcompaction: {",
'\t\t\t\tsummary: "compacted",',
"\t\t\t\tshortSummary: undefined,",
"\t\t\t\tfirstKeptEntryId: event.preparation.firstKeptEntryId,",
"\t\t\t\ttokensBefore: event.preparation.tokensBefore,",
"\t\t\t\tdetails: {},",
"\t\t\t},",
"\t\t};",
"\t});",
'\tpi.on("auto_compaction_start", async (event) => {',
`\t\tconst signals = globalThis.${runtimeSignalStoreKey} ?? (globalThis.${runtimeSignalStoreKey} = []);`,
'\t\tsignals.push("compaction:start:" + event.reason);',
"\t});",
'\tpi.on("auto_compaction_end", async (event) => {',
`\t\tconst signals = globalThis.${runtimeSignalStoreKey} ?? (globalThis.${runtimeSignalStoreKey} = []);`,
'\t\tsignals.push("compaction:end:" + (event.aborted ? "aborted" : "ok"));',
"\t});",
'\tpi.on("todo_reminder", async (event) => {',
`\t\tconst signals = globalThis.${runtimeSignalStoreKey} ?? (globalThis.${runtimeSignalStoreKey} = []);`,
'\t\tsignals.push("todo:" + event.attempt + "/" + event.maxAttempts);',
"\t});",
"}",
].join("\n"),
);
authStorage = await AuthStorage.create(path.join(tempDir.path(), "testauth.db"));
authStorage.setRuntimeApiKey("anthropic", "test-key");
modelRegistry = new ModelRegistry(authStorage);
sessionManager = SessionManager.create(tempDir.path(), tempDir.path());
getRuntimeSignals().length = 0;
const extensionsResult = await loadExtensions([extensionPath], tempDir.path());
const extensionRunner = new ExtensionRunner(
extensionsResult.extensions,
extensionsResult.runtime,
tempDir.path(),
sessionManager,
modelRegistry,
);
const model = getBundledModel("anthropic", "claude-sonnet-4-5");
if (!model) {
throw new Error("Expected built-in anthropic model to exist");
}
const agent = new Agent({
initialState: {
model,
systemPrompt: ["Test"],
tools: [],
messages: [],
},
});
// Seed a minimal session branch so prepareCompaction() returns a preparation.
sessionManager.appendMessage({
role: "user",
content: "hello",
timestamp: Date.now(),
});
session = new AgentSession({
agent,
sessionManager,
settings: Settings.isolated({
"compaction.autoContinue": false,
"todo.reminders": true,
"todo.reminders.max": 3,
}),
modelRegistry,
extensionRunner,
});
});
afterEach(async () => {
await session.dispose();
authStorage.close();
tempDir.removeSync();
vi.useRealTimers();
getRuntimeSignals().length = 0;
vi.restoreAllMocks();
});
it("resumes after threshold compaction when only agent-level queued messages exist", async () => {
session.agent.followUp({
role: "custom",
customType: "test",
content: [{ type: "text", text: "Queued custom" }],
display: false,
timestamp: Date.now(),
});
expect(session.agent.hasQueuedMessages()).toBe(true);
const continueSpy = vi.spyOn(session.agent, "continue").mockImplementation(async () => {
// Real continue() polls and consumes the queued steering/follow-up
// messages. Mirror that here so the stranded-queue drain settles after
// one resume instead of rescheduling itself forever (a no-op mock
// leaves the queue populated, spinning the drain into an OOM loop).
session.agent.clearAllQueues();
});
// Wait for auto_compaction_end event to know when the async handler is done
const { promise: compactionDone, resolve: onCompactionDone } = Promise.withResolvers<void>();
session.subscribe(event => {
if (event.type === "auto_compaction_end") onCompactionDone();
});
// Build a fake AssistantMessage with high token usage to trigger threshold
// compaction (contextWindow=200000, threshold ~80%).
const assistantMsg = {
role: "assistant" as const,
// Non-empty content: an empty `stop` turn would trip the empty-stop guard
// (#handleEmptyAssistantStop) and short-circuit the agent_end handler before
// compaction/todo checks run — hanging this test forever under fake timers.
content: [{ type: "text" as const, text: "Done." }],
api: "anthropic-messages" as const,
provider: "anthropic" as const,
model: "claude-sonnet-4-5",
stopReason: "stop" as const,
usage: {
input: 190000,
output: 1000,
cacheRead: 0,
cacheWrite: 0,
totalTokens: 191000,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
},
timestamp: Date.now(),
};
// Drive auto-compaction through the event flow:
// message_end → stores #lastAssistantMessage
// agent_end → #checkCompaction → shouldCompact → #runAutoCompaction
session.agent.emitExternalEvent({ type: "message_end", message: assistantMsg });
session.agent.emitExternalEvent({ type: "agent_end", messages: [assistantMsg] });
// Wait for compaction completion, then verify waitForIdle blocks on queued continuation.
await compactionDone;
await Promise.resolve();
const idlePromise = session.waitForIdle();
let idleResolved = false;
void idlePromise.then(() => {
idleResolved = true;
});
await Promise.resolve();
expect(idleResolved).toBe(false);
vi.advanceTimersByTime(200);
await idlePromise;
expect(continueSpy).toHaveBeenCalledTimes(1);
const runtimeSignals = getRuntimeSignals();
expect(runtimeSignals).toContain("compaction:start:threshold");
expect(runtimeSignals.some(signal => signal.startsWith("compaction:end:"))).toBe(true);
});
it("has isCompacting true when the auto_compaction_start event fires", async () => {
// Defect 1: the compaction AbortController (which backs isCompacting) must be
// installed before auto_compaction_start is emitted. If it is installed after,
// a message typed the instant the loader appears is read while
// isCompacting === false and mis-routed into the core steering queue (which a
// later handoff reset would wipe) instead of the safe UI compaction queue.
let capturedIsCompacting: boolean | undefined;
const { promise: compactionDone, resolve: onCompactionDone } = Promise.withResolvers<void>();
session.subscribe(event => {
if (event.type === "auto_compaction_start") {
capturedIsCompacting = session.isCompacting;
} else if (event.type === "auto_compaction_end") {
onCompactionDone();
}
});
// Defensive: mirror the resume-drain stub so any queued continuation settles
// instead of spinning the drain (see the threshold test above).
vi.spyOn(session.agent, "continue").mockImplementation(async () => {
session.agent.clearAllQueues();
});
const assistantMsg = {
role: "assistant" as const,
content: [{ type: "text" as const, text: "Done." }],
api: "anthropic-messages" as const,
provider: "anthropic" as const,
model: "claude-sonnet-4-5",
stopReason: "stop" as const,
usage: {
input: 190000,
output: 1000,
cacheRead: 0,
cacheWrite: 0,
totalTokens: 191000,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
},
timestamp: Date.now(),
};
session.agent.emitExternalEvent({ type: "message_end", message: assistantMsg });
session.agent.emitExternalEvent({ type: "agent_end", messages: [assistantMsg] });
await compactionDone;
expect(capturedIsCompacting).toBe(true);
});
it("forwards todo reminder lifecycle signals to extensions", async () => {
const continueSpy = vi.spyOn(session.agent, "continue").mockResolvedValue();
session.setTodoPhases([
{
name: "Execution",
tasks: [{ content: "Finish pending task", status: "in_progress" }],
},
]);
const { promise: reminderDone, resolve: onReminderDone } = Promise.withResolvers<void>();
session.subscribe(event => {
if (event.type === "todo_reminder") onReminderDone();
});
const assistantMsg = {
role: "assistant" as const,
// Non-empty content: see comment on the first test's assistantMsg.
content: [{ type: "text" as const, text: "Done." }],
api: "anthropic-messages" as const,
provider: "anthropic" as const,
model: "claude-sonnet-4-5",
stopReason: "stop" as const,
usage: {
input: 100,
output: 20,
cacheRead: 0,
cacheWrite: 0,
totalTokens: 120,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
},
timestamp: Date.now(),
};
session.agent.emitExternalEvent({ type: "message_end", message: assistantMsg });
session.agent.emitExternalEvent({ type: "agent_end", messages: [assistantMsg] });
await withTimeout(reminderDone, 1000, "Todo reminder timed out");
await Promise.resolve();
expect(getRuntimeSignals()).toContain("todo:1/3");
expect(continueSpy).toHaveBeenCalledTimes(1);
await session.waitForIdle();
});
});