fix(extensions): preserve mutation queue ownership

This commit is contained in:
Duncan Ogilvie
2026-08-10 01:40:18 +02:00
parent bbca0153eb
commit 0cf54f428c
5 changed files with 93 additions and 9 deletions
@@ -1053,13 +1053,6 @@ export class ExtensionRunner {
} catch (error) {
handlerFailure ??= { error };
}
const registrationBarrier = this.#toolRegistrationBarrier;
if (registrationBarrier) {
// Also wait for detached registrations that were scheduled before this
// handler completed. Their failures are reported at registration time,
// not attributed to this unrelated handler.
await registrationBarrier;
}
return result;
},
timeoutMs,
@@ -386,7 +386,7 @@ export class SessionTools {
return this.#toolRegistryMutationScope.run(true, mutation);
});
const operation = untilAborted(signal, serialized);
this.#toolRegistryMutationTail = operation.then(
this.#toolRegistryMutationTail = serialized.then(
() => undefined,
() => undefined,
);
@@ -294,6 +294,37 @@ describe("AgentSession refreshMCPTools rebuild skipping", () => {
expect(session.systemPrompt).toEqual(["tools:read,mcp__nucleus_search,late_prompt_tool"]);
});
it("keeps queued mutations serialized when a waiting caller aborts", async () => {
const firstMutationEntered = Promise.withResolvers<void>();
const releaseFirstMutation = Promise.withResolvers<void>();
const { session } = newSession(async toolNames => `tools:${toolNames.join(",")}`);
const firstMutation = session.runToolRegistryMutation(async () => {
firstMutationEntered.resolve();
await releaseFirstMutation.promise;
});
await firstMutationEntered.promise;
const controller = new AbortController();
let abortedMutationRan = false;
const abortedMutation = session.runToolRegistryMutation(async () => {
abortedMutationRan = true;
}, controller.signal);
controller.abort(new Error("cancel queued mutation"));
await expect(abortedMutation).rejects.toThrow("cancel queued mutation");
let thirdMutationRan = false;
const thirdMutation = session.runToolRegistryMutation(async () => {
thirdMutationRan = true;
});
await Promise.resolve();
expect(thirdMutationRan).toBe(false);
releaseFirstMutation.resolve();
await Promise.all([firstMutation, thirdMutation]);
expect(abortedMutationRan).toBe(false);
expect(thirdMutationRan).toBe(true);
});
it("drops queued and in-flight MCP prompt commits when disposal begins", async () => {
const firstRebuildStarted = Promise.withResolvers<void>();
const releaseFirstRebuild = Promise.withResolvers<void>();
@@ -1448,6 +1448,56 @@ describe("ExtensionRunner", () => {
});
});
it("does not charge detached registrations to unrelated tool-call handlers", async () => {
const extensionPath = path.join(tempDir.path(), "detached-registration-barrier.ts");
fs.writeFileSync(
extensionPath,
`
export default function(pi) {
const { Type } = pi.typebox;
pi.registerTool({
name: "detached_source_tool",
label: "Detached Source Tool",
description: "Provides a registration event for the detached barrier test.",
parameters: Type.Object({}),
execute: async () => ({ content: [{ type: "text", text: "ok" }], details: {} }),
});
pi.on("tool_call", () => undefined);
}
`,
);
const loaded = await loadTestExtensions([extensionPath]);
const runner = new ExtensionRunner(
loaded.extensions,
loaded.runtime,
tempDir.path(),
sessionManager,
modelRegistry,
);
runner.onToolRegistered(() => Promise.withResolvers<void>().promise);
const extension = loaded.extensions[0];
const registrationListener = extension?.toolRegistrationListeners?.values().next().value;
if (!registrationListener) throw new Error("expected registration listener");
registrationListener("detached_source_tool");
const errors: ExtensionError[] = [];
runner.onError(error => {
errors.push(error);
});
testSetExtensionHandlerTimeoutMs(10);
const result = await runner.emitToolCall({
type: "tool_call",
toolName: "unrelated",
toolCallId: "unrelated-call",
input: {},
});
expect(result).toBeUndefined();
expect(errors).toEqual([]);
});
it("aborts a tool_call handler's confirmation before returning its timeout block", async () => {
const extensionPath = path.join(tempDir.path(), "confirm-tool-call.ts");
const markerPath = path.join(tempDir.path(), "confirm-settled.txt");
@@ -178,6 +178,7 @@ describe("createAgentSession defaultInactive tool activation", () => {
});
pi.on("input", async () => {
await startupPromise;
await pi.setActiveTools([...pi.getActiveTools(), "late_active_tool"]);
});
};
@@ -190,10 +191,19 @@ describe("createAgentSession defaultInactive tool activation", () => {
expect(session.getAllToolNames()).not.toContain("late_active_tool");
const runner = session.extensionRunner;
if (!runner) throw new Error("expected extension runner");
await runner.emit({ type: "session_start" });
const errors: string[] = [];
const unsubscribe = runner.onError(error => {
errors.push(error.error);
});
await initializeExtensions(session, {
reportSendError: vi.fn(),
reportRuntimeError: vi.fn(),
});
expect(session.getAllToolNames()).not.toContain("late_active_tool");
startupGate.resolve();
await runner.emitInput("probe", undefined, "interactive");
unsubscribe();
expect(errors).toEqual([]);
expect(session.getAllToolNames()).toEqual(expect.arrayContaining(["late_active_tool", "late_inactive_tool"]));
expect(session.getEnabledToolNames()).toContain("late_active_tool");