Merge PR #7488: fix(task): bound abort cleanup and quarantine late jobs (@metaphorics)

This commit is contained in:
can1357
2026-08-05 01:12:02 +02:00
11 changed files with 432 additions and 35 deletions
@@ -8,6 +8,7 @@
*/
import { afterEach, describe, expect, it, vi } from "bun:test";
import type { AssistantMessage } from "@oh-my-pi/pi-ai";
import { AsyncJobManager } from "@oh-my-pi/pi-coding-agent/async/job-manager";
import type { LoadExtensionsResult } from "@oh-my-pi/pi-coding-agent/extensibility/extensions/types";
import type { CreateAgentSessionResult } from "@oh-my-pi/pi-coding-agent/sdk";
import * as sdkModule from "@oh-my-pi/pi-coding-agent/sdk";
@@ -18,7 +19,7 @@ import { EventBus } from "@oh-my-pi/pi-coding-agent/utils/event-bus";
const baseAgent: AgentDefinition = { name: "task", description: "test", systemPrompt: "test", source: "bundled" };
function assistantStopMessage(text: string): AssistantMessage {
function assistantStopMessage(text: string, totalTokens = 0): AssistantMessage {
return {
role: "assistant",
content: [{ type: "text", text }],
@@ -27,10 +28,10 @@ function assistantStopMessage(text: string): AssistantMessage {
model: "mock",
usage: {
input: 0,
output: 0,
output: totalTokens,
cacheRead: 0,
cacheWrite: 0,
totalTokens: 0,
totalTokens,
cost: { input: 0, output: 0, cacheRead: 0, cacheWrite: 0, total: 0 },
},
stopReason: "stop",
@@ -44,9 +45,15 @@ interface AsyncQuiescenceHarness {
abortCalls: () => number;
settleCalls: () => number;
emitTerminalYield: (data: unknown) => void;
emitAssistant: (text: string, totalTokens?: number) => void;
finishJob: () => void;
}
interface AsyncSessionOptions {
abort?: () => Promise<void>;
dispose?: () => Promise<void>;
}
/**
* Mock session with the owner-async surface the barrier drives:
* `hasPendingAsyncWork` / `getAsyncJobSnapshot` / `settleAsyncWork`. The job
@@ -56,6 +63,7 @@ interface AsyncQuiescenceHarness {
*/
function createAsyncSession(
onPrompt: (params: { text: string; promptIndex: number; harness: AsyncQuiescenceHarness }) => void,
options: AsyncSessionOptions = {},
): AsyncQuiescenceHarness {
const listeners: Array<(event: AgentSessionEvent) => void> = [];
const state = { messages: [] as AssistantMessage[] };
@@ -103,6 +111,11 @@ function createAsyncSession(
state.messages.push(reaction);
emit({ type: "message_end", message: reaction } as AgentSessionEvent);
};
const emitAssistant = (text: string, totalTokens = 0) => {
const message = assistantStopMessage(text, totalTokens);
state.messages.push(message);
emit({ type: "message_end", message } as AgentSessionEvent);
};
const harness: AsyncQuiescenceHarness = {
session: undefined as unknown as AgentSession,
@@ -110,6 +123,7 @@ function createAsyncSession(
abortCalls: () => abortCount,
settleCalls: () => settleCount,
emitTerminalYield,
emitAssistant,
finishJob,
};
@@ -143,8 +157,9 @@ function createAsyncSession(
},
abort: async () => {
abortCount += 1;
await options.abort?.();
},
dispose: async () => {},
dispose: options.dispose ?? (async () => {}),
setIrcWakeTurnObserver: () => {},
};
harness.session = session as unknown as AgentSession;
@@ -163,6 +178,7 @@ function mockCreateAgentSession(session: AgentSession) {
describe("runSubprocess async quiescence fresh-yield contract", () => {
afterEach(() => {
vi.restoreAllMocks();
AsyncJobManager.resetForTests();
});
it("parks a pending yield, injects the result, and completes on the fresh yield", async () => {
@@ -250,4 +266,125 @@ describe("runSubprocess async quiescence fresh-yield contract", () => {
expect(result.exitCode).toBe(0);
expect(result.output).toContain("done");
});
it("does not wait on a second idle barrier after a terminal yield", async () => {
const harness = createAsyncSession(({ promptIndex, harness: h }) => {
if (promptIndex === 1) {
h.finishJob();
h.emitTerminalYield({ report: "done" });
}
});
const idleStarted = Promise.withResolvers<void>();
const releaseIdle = Promise.withResolvers<void>();
let idleCalls = 0;
harness.session.waitForIdle = async () => {
idleCalls += 1;
idleStarted.resolve();
await releaseIdle.promise;
};
mockCreateAgentSession(harness.session);
const run = runSubprocess({
cwd: "/tmp",
agent: baseAgent,
task: "do the work",
index: 0,
id: "quiescence-no-second-idle",
});
const outcome = await Promise.race([
run.then(() => "completed" as const),
idleStarted.promise.then(() => "blocked" as const),
]);
releaseIdle.resolve();
const result = await run;
expect(outcome).toBe("completed");
expect(idleCalls).toBe(0);
expect(result.exitCode).toBe(0);
expect(result.output).toContain("done");
});
it("returns an aborted result after cleanup grace and waits for every late resource", async () => {
const abortStarted = Promise.withResolvers<void>();
const abortGate = Promise.withResolvers<void>();
const disposeGate = Promise.withResolvers<void>();
const lateJobGate = Promise.withResolvers<void>();
const manager = new AsyncJobManager({});
AsyncJobManager.setInstance(manager);
let lateJobId: string | undefined;
let deferredCleanup: Promise<void> | undefined;
const harness = createAsyncSession(
({ promptIndex, harness: h }) => {
if (promptIndex !== 1) return;
h.finishJob();
h.emitAssistant("captured before cleanup", 7);
h.emitTerminalYield({ report: "yielded output" });
},
{
abort: async () => {
abortStarted.resolve();
await abortGate.promise;
},
dispose: async () => {
lateJobId = manager.register(
"task",
"shutdown-time job",
async () => {
await lateJobGate.promise;
return "late result";
},
{ ownerId: "cleanup-timeout" },
);
await disposeGate.promise;
},
},
);
mockCreateAgentSession(harness.session);
const run = runSubprocess({
cwd: "/tmp",
agent: baseAgent,
task: "do the work",
index: 0,
id: "cleanup-timeout",
keepAlive: false,
onCleanupDeferred: completion => {
deferredCleanup = completion;
},
});
await abortStarted.promise;
const result = await run;
expect(result.exitCode).toBe(1);
expect(result.aborted).toBe(true);
expect(result.abortReason).toBe("cleanup exceeded 10000 ms");
expect(result.error).toBe(
"Task aborted. Cleanup did not finish within 10000 ms. This task was not isolated, so its changes may remain in the working directory.",
);
expect(result.output).toContain("yielded output");
expect(result.usage?.totalTokens).toBe(7);
expect(lateJobId).toBeDefined();
expect(deferredCleanup).toBeDefined();
let cleanupSettled = false;
const cleanupOutcome = deferredCleanup?.then(
() => {
cleanupSettled = true;
},
() => {
cleanupSettled = true;
},
);
abortGate.resolve();
disposeGate.reject(new Error("dispose failed"));
for (let attempt = 0; attempt < 10 && manager.getJob(lateJobId ?? "")?.status === "running"; attempt += 1) {
await Promise.resolve();
}
expect(cleanupSettled).toBe(false);
expect(manager.getJob(lateJobId ?? "")?.status).toBe("cancelled");
lateJobGate.resolve();
await cleanupOutcome;
expect(cleanupSettled).toBe(true);
}, 15_000);
});
@@ -139,6 +139,63 @@ describe("runIsolatedSubprocess", () => {
expect(deleteSpy).toHaveBeenCalledWith(repoRoot, "omp/task/PreserveBranchFailure");
expect(cleanupSpy).toHaveBeenCalledTimes(1);
});
it("keeps an isolated worktree until deferred child cleanup settles", async () => {
const cleanupGate = Promise.withResolvers<void>();
vi.spyOn(worktreeModule, "ensureIsolation").mockResolvedValue({
mergedDir: "/repo/isolated",
backend: natives.IsoBackendKind.Rcopy,
fellBack: false,
fallbackReason: null,
});
vi.spyOn(executorModule, "runSubprocess").mockImplementation(async options => {
options.onCleanupDeferred?.(cleanupGate.promise);
return result({ exitCode: 1, aborted: true, error: "cleanup exceeded its deadline" });
});
const cleanupSpy = vi.spyOn(worktreeModule, "cleanupIsolation").mockResolvedValue();
const outcome = await runIsolatedSubprocess({
baseOptions: {
cwd: "/repo",
agent: {
name: "task",
description: "Task agent",
systemPrompt: "test",
source: "bundled",
},
task: "Do work",
index: 0,
id: "DeferredCleanup",
},
context: {
repoRoot: "/repo",
baseline: {
root: {
repoRoot: "/repo",
headCommit: "base",
staged: "",
unstaged: "",
untracked: [],
untrackedPatch: "",
},
nested: [],
},
},
preferredBackend: undefined,
agentId: "DeferredCleanup",
mergeMode: "patch",
artifactsDir: "/artifacts",
buildFailureResult: error => result({ exitCode: 1, error: String(error) }),
});
expect(outcome.exitCode).toBe(1);
expect(cleanupSpy).not.toHaveBeenCalled();
cleanupGate.resolve();
await cleanupGate.promise;
await Promise.resolve();
await Promise.resolve();
expect(cleanupSpy).toHaveBeenCalledTimes(1);
});
});
describe("mergeIsolatedChanges", () => {