fix: revoke CLI tool access when a run is canceled (#149154)

Stop canceled CLI sources and execution attempts from retaining Gateway tool access through delayed process cleanup, while preserving accepted tool results during normal completion.
This commit is contained in:
Shakker 2026-09-15 16:53:51 +01:00 • committed by GitHub
parent 7bdca7d8d1
commit 684e301de4
No known key found for this signature in database
GPG key ID: B5690EEEBB952194
21 changed files with 516 additions and 124 deletions

View file

@ -0,0 +1,151 @@
import { afterEach, describe, expect, it, vi } from "vitest";
import { createDeferred } from "../../../test/helpers/promise.js";
import { createReplyOperation } from "../../auto-reply/reply/reply-run-registry.operation.js";
import {
activateMcpLoopbackClientGrantCapture,
deactivateMcpLoopbackClientGrantCapture,
mintMcpLoopbackClientGrant,
resolveMcpLoopbackClientGrant,
revokeMcpLoopbackClientGrant,
transferMcpLoopbackClientGrant,
} from "../../gateway/mcp-grant-store.js";
import type { CliBackendLiveSessionHandle } from "../../plugins/cli-backend.types.js";
import { getAdmittedRunDelegatedAuthority } from "../admitted-run-context.js";
import {
closePluginTestAdmissions,
createExecution,
runPlugin,
SUCCESS_RESULT,
waitUntilAborted,
} from "./execute-plugin.test-support.js";
import { createCliToolTracking } from "./execute-tool-tracking.js";
vi.mock("../tools/gateway.js", () => ({ callGatewayTool: vi.fn() }));
const activeSessions = new Set<CliBackendLiveSessionHandle>();
afterEach(() => {
for (const session of activeSessions) {
session.close("restart");
}
activeSessions.clear();
closePluginTestAdmissions();
vi.restoreAllMocks();
vi.useRealTimers();
});
describe("CLI MCP capture authority", () => {
it.each([
{ name: "one-shot timeout", liveSession: false },
{ name: "live backend cancellation", liveSession: true },
])("revokes MCP authority before iterator cleanup after $name", async ({ liveSession }) => {
vi.useFakeTimers();
const source = new AbortController();
const { context } = await createExecution({ abortSignal: source.signal, timeoutMs: 100 });
const operation = createReplyOperation({
sessionKey: context.params.sessionKey!,
sessionId: context.params.sessionId,
resetTriggered: false,
});
context.params.replyOperation = operation;
const runtimeOwnerToken = `runtime-${context.params.runId}`;
const grant = mintMcpLoopbackClientGrant({
context: { sessionKey: context.params.sessionKey!, senderIsOwner: false },
runtimeOwnerToken,
admittedRunContext: context.params.admittedRunContext,
abortSignal: source.signal,
});
context.preparedBackend.mcpClientGrantCapture = {
transportToken: grant.token,
adoptProcessToken: (targetToken) => {
transferMcpLoopbackClientGrant({
sourceToken: grant.token,
targetToken,
runtimeOwnerToken,
});
},
revokeProcessToken: () => {
revokeMcpLoopbackClientGrant(grant.token);
},
activate: (captureKey, assertCurrent) => {
activateMcpLoopbackClientGrantCapture({
token: grant.token,
runtimeOwnerToken,
captureKey,
assertCurrent,
});
},
deactivate: (captureKey) => {
deactivateMcpLoopbackClientGrantCapture({
token: grant.token,
runtimeOwnerToken,
captureKey,
});
},
};
const tracking = createCliToolTracking(context);
const captureKey = `capture-${context.params.runId}`;
const capture = { token: grant.token, runtimeOwnerToken, captureKey };
const streamStarted = createDeferred();
const streamClosing = createDeferred();
const releaseCleanup = createDeferred();
const run = runPlugin(
context,
async function* (execution) {
if (liveSession) {
const capability = execution.liveSession;
if (!capability) {
throw new Error("expected live CLI session capability");
}
const handle: CliBackendLiveSessionHandle = {
generation: context.params.runId,
fingerprint: capability.fingerprint,
isIdle: () => true,
close: () => capability.remove(handle),
waitForExit: async () => {},
};
capability.register(handle);
activeSessions.add(handle);
capability.activate(handle);
}
const aborted = waitUntilAborted(execution);
streamStarted.resolve();
try {
await aborted;
yield SUCCESS_RESULT;
} finally {
streamClosing.resolve();
await releaseCleanup.promise;
}
},
{ liveSession, mcpCapture: { captureKey, beginCapture: tracking.beginGatewayCapture } },
);
const observedRun = run.catch((error: unknown) => error);
try {
await streamStarted.promise;
const retained = resolveMcpLoopbackClientGrant(capture);
expect(retained?.isCurrent()).toBe(true);
if (liveSession) {
expect(operation.abortByUser()).toBe(true);
} else {
await vi.advanceTimersByTimeAsync(100);
}
await streamClosing.promise;
expect(source.signal.aborted).toBe(false);
expect(getAdmittedRunDelegatedAuthority(context.params.admittedRunContext)).toBeDefined();
expect(retained?.isCurrent()).toBe(false);
expect(resolveMcpLoopbackClientGrant(capture)).toBeUndefined();
} finally {
releaseCleanup.resolve();
await observedRun;
tracking.finalizeCapture(() => {});
revokeMcpLoopbackClientGrant(grant.token);
operation.complete();
}
expect(await observedRun).toMatchObject(
liveSession ? { name: "AbortError" } : { reason: "overall-timeout", timedOut: true },
);
});
});

View file

@ -74,6 +74,7 @@ export function runPlugin(
useResume?: boolean;
forceNewSession?: boolean;
liveSession?: boolean;
mcpCapture?: Parameters<typeof executePluginOwnedProcess>[0]["mcpCapture"];
requiredGeneration?: string;
onNoOutputTimeout?: NonNullable<
Parameters<typeof executePluginOwnedProcess>[0]["onNoOutputTimeout"]
@ -95,15 +96,11 @@ export function runPlugin(
promptContext: context.promptContext,
useResume: options.useResume ?? Boolean(options.requiredGeneration),
sessionId: options.sessionId ?? "sdk-session",
mcpCapture: options.mcpCapture,
...(options.forceNewSession ? { forceNewSession: true } : {}),
...(options.liveSession || options.requiredGeneration
? {
liveSession: {
beginCapture: () => {},
...(options.requiredGeneration
? { requiredGeneration: options.requiredGeneration }
: {}),
},
liveSession: { requiredGeneration: options.requiredGeneration },
}
: {}),
...(options.onNoOutputTimeout ? { onNoOutputTimeout: options.onNoOutputTimeout } : {}),
@ -135,3 +132,20 @@ export function closePluginTestAdmissions(): void {
admission.close();
}
}
export function waitUntilAborted(execution: CliBackendExecuteContext): Promise<void> {
const signal = execution.abortSignal;
if (!signal) {
throw new Error("Host execution did not expose its abort signal.");
}
return new Promise((_, reject) => {
signal.addEventListener(
"abort",
() =>
reject(
signal.reason instanceof Error ? signal.reason : new Error("CLI test run was aborted."),
),
{ once: true },
);
});
}

View file

@ -19,6 +19,7 @@ import {
requestNativeTool,
runPlugin,
SUCCESS_RESULT,
waitUntilAborted,
} from "./execute-plugin.test-support.js";
import type { PreparedCliRunContext } from "./types.js";
@ -50,23 +51,6 @@ function registerOwnerSession(context: PreparedCliRunContext, generation: string
return { handle: session, close };
}
function waitUntilAborted(execution: CliBackendExecuteContext): Promise<void> {
const signal = execution.abortSignal;
if (!signal) {
throw new Error("Host execution did not expose its abort signal.");
}
return new Promise((_, reject) => {
signal.addEventListener(
"abort",
() =>
reject(
signal.reason instanceof Error ? signal.reason : new Error("CLI test run was aborted."),
),
{ once: true },
);
});
}
afterEach(() => {
for (const session of activeSessions) {
session.close("restart");

View file

@ -426,9 +426,11 @@ export async function executePluginOwnedProcess(params: {
onActiveLoopbackAskUserDeadlineChange?: (listener: () => void) => () => void;
onNoOutputTimeout?: (error: FailoverError) => void;
onInterrupted?: (reason: CliTerminalInterruption["reason"]) => boolean;
liveSession?: {
mcpCapture?: {
captureKey?: string;
beginCapture: (captureKey: string | undefined) => void;
beginCapture: (captureKey: string | undefined, assertCurrent: () => void) => void;
};
liveSession?: {
requiredGeneration?: string;
};
}): Promise<RunExit> {
@ -445,7 +447,12 @@ export async function executePluginOwnedProcess(params: {
? AbortSignal.any([controller.signal, run.abortSignal])
: controller.signal;
const assertCurrent = createCliRunCurrentAssertion(run, signal);
const assertRunCurrent = createCliRunCurrentAssertion(run);
const termination: { reason: TerminationReason } = { reason: "exit" };
// Normal cleanup closes native callbacks; MCP results retain their admitted
// caller through the capture drain. Cancellation and timeouts still close both.
const assertCaptureCurrent = () =>
(termination.reason === "exit" ? assertRunCurrent : assertCurrent)();
const outstanding = {
approvals: 0,
background: 0,
@ -568,9 +575,14 @@ export async function executePluginOwnedProcess(params: {
argv0: params.executionArgv0,
env: params.env,
...params.liveSession,
captureKey: params.mcpCapture?.captureKey,
beginCapture: (captureKey) =>
params.mcpCapture?.beginCapture(captureKey, assertCaptureCurrent),
abortSignal: signal,
claimResources: params.context.preparedBackend.claimLiveSessionResources,
});
} else {
params.mcpCapture?.beginCapture(params.mcpCapture.captureKey, assertCaptureCurrent);
}
assertCurrent();
const execution = params.execute({

View file

@ -240,11 +240,13 @@ export async function executeCliProcess(params: {
terminalInterruption = { reason };
return true;
},
mcpCapture: {
captureKey: params.initialGatewayCaptureKey,
beginCapture: params.toolTracking.beginGatewayCapture,
},
...(params.useManagedClaudeLiveSession
? {
liveSession: {
captureKey: params.initialGatewayCaptureKey,
beginCapture: params.toolTracking.beginGatewayCapture,
requiredGeneration: params.cliSessionIdToUse
? context.requiredClaudeLiveSessionGeneration
: undefined,
@ -270,9 +272,23 @@ export async function executeCliProcess(params: {
}
// Startup can wait behind another scoped run. Reserve cancellation under
// the caller's run id before awaiting the child or replacement fence.
const abortManagedRun = () => supervisor.cancel(runParams.runId, "manual-cancel");
let processCancelled = false;
const assertProcessCurrent = () => {
params.assertCurrent();
if (processCancelled) {
throw new Error("CLI process authority is no longer active");
}
};
const abortManagedRun = () => {
processCancelled = true;
supervisor.cancel(runParams.runId, "manual-cancel");
};
runParams.abortSignal?.addEventListener("abort", abortManagedRun, { once: true });
try {
params.toolTracking.beginGatewayCapture(
params.initialGatewayCaptureKey,
assertProcessCurrent,
);
const managedRun = await supervisor.spawn({
assertCurrent: params.assertCurrent,
runId: runParams.runId,
@ -284,6 +300,9 @@ export async function executeCliProcess(params: {
argv0: params.executionArgv0,
timeoutMs: runParams.timeoutMs,
noOutputTimeoutMs: params.noOutputTimeoutMs,
onCancel: () => {
processCancelled = true;
},
cwd: context.cwd ?? context.workspaceDir,
env: params.env,
input: params.stdin ?? "",
@ -298,7 +317,10 @@ export async function executeCliProcess(params: {
kind: "cli" as const,
runId: runParams.runId,
toolAuthorityFingerprint: runParams.toolAuthorityFingerprint,
cancel: () => managedRun.cancel("manual-cancel"),
cancel: () => {
processCancelled = true;
managedRun.cancel("manual-cancel");
},
}
: undefined;
if (replyBackendHandle) {
@ -306,6 +328,7 @@ export async function executeCliProcess(params: {
}
try {
result = await managedRun.wait();
processCancelled ||= result.reason !== "exit";
} finally {
if (replyBackendHandle) {
runParams.replyOperation?.detachBackend(replyBackendHandle);

View file

@ -19,7 +19,7 @@ function createTracking() {
deactivate: vi.fn(),
};
const tracking = createCliToolTracking(context);
tracking.beginGatewayCapture("deadline-test");
tracking.beginGatewayCapture("deadline-test", () => {});
return tracking;
}

View file

@ -342,14 +342,14 @@ export function createCliToolTracking(context: PreparedCliRunContext) {
call.args,
);
};
const beginGatewayCapture = (captureKey: string | undefined) => {
const beginGatewayCapture = (captureKey: string | undefined, assertCurrent: () => void) => {
if (!captureKey || gatewayCaptureKey === captureKey) {
return;
}
if (gatewayCaptureKey) {
throw new Error("CLI MCP capture key changed during an active attempt");
}
context.preparedBackend.mcpClientGrantCapture?.activate(captureKey);
context.preparedBackend.mcpClientGrantCapture?.activate(captureKey, assertCurrent);
gatewayCaptureKey = captureKey;
const isPotentialDelivery = (toolName: string) =>
isMessagingTool(normalizeCliMessagingToolName(toolName));

View file

@ -31,7 +31,10 @@ import type { CliBackendParseJsonlEvent } from "../../plugins/cli-backend.types.
import { getPluginModuleLoaderStats } from "../../plugins/plugin-module-loader-cache.js";
import { createEmptyPluginRegistry } from "../../plugins/registry-empty.js";
import { setActivePluginRegistry } from "../../plugins/runtime.js";
import { createChildAdapter } from "../../process/supervisor/adapters/child.js";
import type { getProcessSupervisor } from "../../process/supervisor/index.js";
import { createProcessSupervisor } from "../../process/supervisor/supervisor.js";
import { createStubChildAdapter } from "../../process/supervisor/supervisor.test-support.js";
import { createUserTurnTranscriptRecorder } from "../../sessions/user-turn-transcript.js";
import { createTestUserTurnTranscriptTarget } from "../../sessions/user-turn-transcript.test-support.js";
import { prepareSystemAgentRunAdmission } from "../admitted-run-context.js";
@ -52,6 +55,10 @@ import type { PreparedCliRunContext } from "./types.js";
const executePreparedCliRun = wrapPreparedCliRunWithTestAdmission(executePreparedCliRunImpl);
vi.mock("../../process/supervisor/adapters/child.js", () => ({
createChildAdapter: vi.fn(),
}));
// Gateway unit coverage owns quiet-admission timing. These integration cases only
// need to drain calls already in flight, so skip the repeated 250 ms quiet window.
vi.mock("../../gateway/mcp-http.loopback-runtime.js", async (importOriginal) => {
@ -3051,8 +3058,12 @@ describe("executePreparedCliRun supervisor output capture", () => {
it("captures non-Claude JSONL sends and fences every attempt with a unique key", async () => {
const context = buildPreparedCliRunContext({ output: "jsonl", provider: "local-cli" });
context.mcpDeliveryCapture = true;
const activateCapture = vi.fn<(captureKey: string) => void>();
const deactivateCapture = vi.fn<(captureKey: string) => void>();
const activateCapture = vi.fn<(captureKey: string, assertCurrent: () => void) => void>();
const deactivateCapture = vi.fn((_captureKey: string) => {
const assertion = activateCapture.mock.calls.at(-1)?.[1];
expect(assertion).toBeTypeOf("function");
expect(assertion).not.toThrow();
});
context.preparedBackend.mcpClientGrantCapture = {
transportToken: "capture-test-token",
adoptProcessToken: vi.fn(),
@ -3064,6 +3075,7 @@ describe("executePreparedCliRun supervisor output capture", () => {
supervisorSpawnMock.mockImplementation(async (...args: unknown[]) => {
const input = args[0] as SupervisorSpawnInput;
const captureKey = input.env?.OPENCLAW_MCP_CLI_CAPTURE_KEY ?? "";
expect(activateCapture.mock.calls.at(-1)?.[1]).not.toThrow();
captureKeys.push(captureKey);
recordMcpLoopbackToolCallResult({
captureKey,
@ -3094,5 +3106,72 @@ describe("executePreparedCliRun supervisor output capture", () => {
activateCapture.mock.invocationCallOrder[1] ?? Number.POSITIVE_INFINITY,
);
});
it("revokes MCP capture at the supervisor deadline before native exit", async () => {
vi.useFakeTimers();
const context = buildPreparedCliRunContext({ output: "text", provider: "local-cli" });
context.mcpDeliveryCapture = true;
const activateCapture = vi.fn<(captureKey: string, assertCurrent: () => void) => void>();
const deactivateCapture = vi.fn();
context.preparedBackend.mcpClientGrantCapture = {
transportToken: "capture-test-token",
adoptProcessToken: vi.fn(),
revokeProcessToken: vi.fn(),
activate: activateCapture,
deactivate: deactivateCapture,
};
const adapter = createStubChildAdapter();
vi.mocked(createChildAdapter).mockResolvedValueOnce({
...adapter,
onExit: vi.fn(),
onError: vi.fn(),
});
const supervisor = createProcessSupervisor();
const spawned = createDeferred<Awaited<ReturnType<ProcessSupervisor["spawn"]>>>();
supervisorSpawnMock.mockImplementationOnce(async (...args: unknown[]) => {
const input = args[0] as SupervisorSpawnInput;
if (input.mode !== "child") {
throw new Error("Expected the CLI child transport");
}
// The shared execution fixture already expands the deferred arguments.
const { resolveArgs: _resolveArgs, ...spawnInput } = input;
const managed = await supervisor.spawn(spawnInput);
spawned.resolve(managed);
return managed;
});
const execution = executePreparedCliRun(context);
const settled = execution.catch(() => undefined);
try {
const managed = await Promise.race([
spawned.promise,
execution.then(() => {
throw new Error("CLI completed before the supervisor started");
}),
]);
const assertCaptureCurrent = activateCapture.mock.calls[0]?.[1];
expect(assertCaptureCurrent).toBeTypeOf("function");
expect(assertCaptureCurrent).not.toThrow();
adapter.killMock.mockImplementation(() => {
expect(assertCaptureCurrent).toThrow("CLI process authority is no longer active");
});
await vi.advanceTimersByTimeAsync(context.params.timeoutMs);
expect(adapter.killMock).toHaveBeenCalledOnce();
expect(managed.activity.resultSettled).toBe(false);
expect(deactivateCapture).not.toHaveBeenCalled();
expect(assertCaptureCurrent).toThrow("CLI process authority is no longer active");
adapter.settle(null, "SIGTERM");
await expect(execution).rejects.toMatchObject({ reason: "timeout" });
expect(deactivateCapture).toHaveBeenCalledOnce();
} finally {
adapter.settle(0);
await settled;
await supervisor.shutdown();
vi.mocked(createChildAdapter).mockReset();
vi.clearAllTimers();
vi.useRealTimers();
}
});
});
/* oxlint-disable max-lines -- TODO: split this grandfathered oversized file. */

View file

@ -580,9 +580,6 @@ export async function executePreparedCliRun(
useResume,
trigger: params.trigger,
});
if (!useManagedClaudeLiveSession) {
toolTracking.beginGatewayCapture(initialGatewayCaptureKey);
}
runOutput = await executeCliProcess({
context,
assertCurrent,

View file

@ -4340,7 +4340,8 @@ describe("prepareCliRunContext", () => {
});
expect(context.preparedBackend.mcpClientGrantCapture?.transportToken).toBe("loopback-token");
context.preparedBackend.mcpClientGrantCapture?.adoptProcessToken("stable-loopback-token");
context.preparedBackend.mcpClientGrantCapture?.activate("capture-test");
const assertCaptureCurrent = () => {};
context.preparedBackend.mcpClientGrantCapture?.activate("capture-test", assertCaptureCurrent);
context.preparedBackend.mcpClientGrantCapture?.deactivate("capture-test");
expect(transferMcpLoopbackClientGrant).toHaveBeenCalledExactlyOnceWith({
sourceToken: "loopback-token",
@ -4351,6 +4352,7 @@ describe("prepareCliRunContext", () => {
token: "stable-loopback-token",
runtimeOwnerToken: "loopback-owner-token",
captureKey: "capture-test",
assertCurrent: assertCaptureCurrent,
});
expect(deactivateMcpLoopbackClientGrantCapture).toHaveBeenCalledExactlyOnceWith({
token: "stable-loopback-token",
@ -4459,7 +4461,7 @@ describe("prepareCliRunContext", () => {
).toBeNull();
expect(projectNativeToolAuthority).not.toHaveBeenCalled();
expect(captureNativeToolAuthority).not.toHaveBeenCalled();
capture.activate("native-capture");
capture.activate("native-capture", () => {});
observe(["Read", "Bash"]);
expect(projectNativeToolAuthority).toHaveBeenCalledExactlyOnceWith(["Read", "Bash"]);
@ -4505,7 +4507,7 @@ describe("prepareCliRunContext", () => {
? {}
: { cliToolAvailability: { native: selected, openClaw: ["message"] } },
);
capture.activate("native-capture");
capture.activate("native-capture", () => {});
observe(observed);
expect(projectNativeToolAuthority).toHaveBeenCalledExactlyOnceWith(projected);
@ -4553,7 +4555,7 @@ describe("prepareCliRunContext", () => {
["read", "web_fetch", "web_search"],
{ toolOverrides: { webSearch: false } },
);
capture.activate("native-capture");
capture.activate("native-capture", () => {});
observe(["Read", "WebFetch", "WebSearch"]);
expect(captureNativeToolAuthority).toHaveBeenLastCalledWith(["read", "web_fetch"]);
@ -4567,7 +4569,7 @@ describe("prepareCliRunContext", () => {
])("clears native authority before rejecting a $name runtime snapshot", async ({ tools }) => {
const { capture, observe, projectNativeToolAuthority, captureNativeToolAuthority } =
await prepareNativeAuthority(["read"]);
capture.activate("native-capture");
capture.activate("native-capture", () => {});
observe(["Read"]);
projectNativeToolAuthority.mockClear();
@ -4579,7 +4581,7 @@ describe("prepareCliRunContext", () => {
it("clears native authority before rejecting a non-canonical backend projection", async () => {
const { capture, observe, projectNativeToolAuthority, captureNativeToolAuthority } =
await prepareNativeAuthority(["read"]);
capture.activate("native-capture");
capture.activate("native-capture", () => {});
observe(["Read"]);
projectNativeToolAuthority.mockReturnValue(["Bash"]);
@ -4593,7 +4595,7 @@ describe("prepareCliRunContext", () => {
const { capture, observe, projectNativeToolAuthority, captureNativeToolAuthority } =
await prepareNativeAuthority(["read"]);
if (state === "stale") {
capture.activate("native-capture");
capture.activate("native-capture", () => {});
observe(["Read"]);
captureNativeToolAuthority.mockReturnValue(false);
projectNativeToolAuthority.mockClear();

View file

@ -1469,6 +1469,8 @@ async function prepareCliRunContextWithinReadFence(
context: mcpGrantContext,
runtimeOwnerToken: mcpLoopbackRuntime.ownerToken,
admittedRunContext: params.admittedRunContext,
abortSignal: params.abortSignal,
assertCurrent: params.assertCurrent,
// MCP owns a canonical main target even when the native callback is sessionless.
bindQuestionAnswerAuthority: (assertActive) =>
bindQuestionAnswerAuthorityForSession(mcpGrantContext.sessionKey, assertActive),
@ -1519,11 +1521,12 @@ async function prepareCliRunContextWithinReadFence(
revokeProcessToken: () => {
prepareDeps.revokeMcpLoopbackClientGrant(activeToken);
},
activate: (captureKey: string) => {
activate: (captureKey: string, assertCurrent: () => void) => {
const activated = prepareDeps.activateMcpLoopbackClientGrantCapture({
token: activeToken,
runtimeOwnerToken: mcpLoopbackRuntime.ownerToken,
captureKey,
assertCurrent,
});
if (!activated) {
throw new Error(

View file

@ -347,7 +347,7 @@ type CliPreparedBackend = {
adoptProcessToken: (processToken: string) => void;
/** Revoke the bearer when the child process that holds it exits. */
revokeProcessToken: () => void;
activate: (captureKey: string) => void;
activate: (captureKey: string, assertCurrent: () => void) => void;
deactivate: (captureKey: string) => void;
captureNativeTools?: (tools: unknown) => void;
};

View file

@ -321,57 +321,91 @@ describe("mcp-grant-store", () => {
expect(capture.captureNativeToolAuthority(["exec"])).toBe(false);
});
it.each(["deactivate", "rebind", "transfer", "close", "revoke"] as const)(
"rejects a retained native capture after %s",
async (invalidation) => {
const admittedRunContext = await admitted("run-stale-native");
const grant = mintMcpLoopbackClientGrant({
context: {
sessionKey: "agent:main:first",
senderIsOwner: false,
nativeCronCreatorToolAllowlist: [],
},
runtimeOwnerToken: "runtime-one",
admittedRunContext,
it.each([
"deactivate",
"rebind",
"transfer",
"close",
"revoke",
"source-abort",
"caller-revocation",
"capture-abort",
"revoke-during-assertion",
] as const)("rejects a retained native capture after %s", async (invalidation) => {
const admittedRunContext = await admitted("run-stale-native");
const sourceController = new AbortController();
const captureController = new AbortController();
let callerCurrent = true;
let revokeDuringAssertion = false;
const grant = mintMcpLoopbackClientGrant({
context: {
sessionKey: "agent:main:first",
senderIsOwner: false,
nativeCronCreatorToolAllowlist: [],
},
runtimeOwnerToken: "runtime-one",
admittedRunContext,
abortSignal: sourceController.signal,
assertCurrent: () => {
if (!callerCurrent) {
throw new Error("caller revoked");
}
if (revokeDuringAssertion) {
revokeMcpLoopbackClientGrant(grant.token);
}
},
});
const params = {
token: grant.token,
runtimeOwnerToken: "runtime-one",
captureKey: "capture-stale-native",
};
const capture = activateMcpLoopbackClientGrantCapture({
...params,
assertCurrent: () => captureController.signal.throwIfAborted(),
});
if (!capture) {
throw new Error("expected an active native capture");
}
expect(capture.captureNativeToolAuthority(["read"])).toBe(true);
const retained = resolveMcpLoopbackClientGrant(params);
expect(retained?.isCurrent()).toBe(true);
if (invalidation === "deactivate") {
deactivateMcpLoopbackClientGrantCapture(params);
} else if (invalidation === "rebind") {
bindMcpLoopbackClientGrantAdmission({ ...params, admittedRunContext });
} else if (invalidation === "transfer") {
const next = mintMcpLoopbackClientGrant({
context: grant.context,
runtimeOwnerToken: params.runtimeOwnerToken,
admittedRunContext: await admitted("run-next-native"),
});
const params = {
token: grant.token,
runtimeOwnerToken: "runtime-one",
captureKey: "capture-stale-native",
};
const capture = activateMcpLoopbackClientGrantCapture(params);
if (!capture) {
throw new Error("expected an active native capture");
}
expect(capture.captureNativeToolAuthority(["read"])).toBe(true);
if (invalidation === "deactivate") {
deactivateMcpLoopbackClientGrantCapture(params);
} else if (invalidation === "rebind") {
bindMcpLoopbackClientGrantAdmission({ ...params, admittedRunContext });
} else if (invalidation === "transfer") {
const next = mintMcpLoopbackClientGrant({
context: grant.context,
runtimeOwnerToken: params.runtimeOwnerToken,
admittedRunContext: await admitted("run-next-native"),
});
transferMcpLoopbackClientGrant({
sourceToken: next.token,
targetToken: grant.token,
runtimeOwnerToken: params.runtimeOwnerToken,
});
activateMcpLoopbackClientGrantCapture(params);
} else if (invalidation === "close") {
admissions.at(-1)?.close();
} else {
revokeMcpLoopbackClientGrant(grant.token);
}
expect(capture.captureNativeToolAuthority(["exec"])).toBe(false);
expect(capture.captureNativeToolAuthority(null)).toBe(false);
expect(resolveMcpLoopbackClientGrant(params)?.context.nativeCronCreatorToolAllowlist).toEqual(
invalidation === "rebind" ? ["read"] : invalidation === "transfer" ? null : undefined,
);
},
);
transferMcpLoopbackClientGrant({
sourceToken: next.token,
targetToken: grant.token,
runtimeOwnerToken: params.runtimeOwnerToken,
});
activateMcpLoopbackClientGrantCapture(params);
} else if (invalidation === "close") {
admissions.at(-1)?.close();
} else if (invalidation === "source-abort") {
sourceController.abort();
} else if (invalidation === "caller-revocation") {
callerCurrent = false;
} else if (invalidation === "capture-abort") {
captureController.abort();
} else if (invalidation === "revoke-during-assertion") {
revokeDuringAssertion = true;
} else {
revokeMcpLoopbackClientGrant(grant.token);
}
expect(retained?.isCurrent()).toBe(false);
expect(capture.captureNativeToolAuthority(["exec"])).toBe(false);
expect(capture.captureNativeToolAuthority(null)).toBe(false);
expect(resolveMcpLoopbackClientGrant(params)?.context.nativeCronCreatorToolAllowlist).toEqual(
invalidation === "rebind" ? ["read"] : invalidation === "transfer" ? null : undefined,
);
});
it("rejects an active bearer and capture after its admitted authority closes", async () => {
const admittedRunContext = await admitted("run-closed-grant");
@ -432,6 +466,8 @@ describe("mcp-grant-store", () => {
it("transfers fresh turn authority onto a process-stable bearer", async () => {
const firstAdmission = await admitted("run-first-turn");
const nextAdmission = await admitted("run-next-turn");
const firstController = new AbortController();
const nextController = new AbortController();
const skillLibraryAuthoring: SkillLibraryAuthoringCapability = {
target: "personal",
defaultTarget: "personal",
@ -445,11 +481,13 @@ describe("mcp-grant-store", () => {
context: { sessionKey: "agent:main:first", runId: "run-first-turn", senderIsOwner: false },
runtimeOwnerToken: "runtime-one",
admittedRunContext: firstAdmission,
abortSignal: firstController.signal,
});
const next = mintMcpLoopbackClientGrant({
context: { sessionKey: "agent:main:next", runId: "run-next-turn", senderIsOwner: true },
runtimeOwnerToken: "runtime-one",
admittedRunContext: nextAdmission,
abortSignal: nextController.signal,
skillLibraryAuthoring,
toolAuth: {
agentDir: "/tmp/next-agent",
@ -471,6 +509,7 @@ describe("mcp-grant-store", () => {
expect(first?.isCurrent()).toBe(true);
// Turn cleanup revokes the process bearer while the warm child still holds its token.
// The next admitted turn must be able to restore that exact inactive bearer.
firstController.abort();
expect(revokeMcpLoopbackClientGrant(stable.token)).toBe(true);
expect(first?.isCurrent()).toBe(false);
const revocations: Array<{ token: string; runtimeOwnerToken: string }> = [];
@ -549,6 +588,15 @@ describe("mcp-grant-store", () => {
{ token: stable.token, runtimeOwnerToken: "runtime-one" },
{ token: next.token, runtimeOwnerToken: "runtime-one" },
]);
nextController.abort();
expect(transferred?.isCurrent()).toBe(false);
expect(
resolveMcpLoopbackClientGrant({
token: stable.token,
runtimeOwnerToken: "runtime-one",
captureKey: "next-capture",
}),
).toBeUndefined();
});
it("revokes client grants by token or exact Gateway runtime", () => {

View file

@ -16,6 +16,7 @@ import type {
} from "../auto-reply/get-reply-options.types.js";
import type { InboundEventKind } from "../channels/inbound-event/kind.js";
import type { CronScheduledToolCallerOrigin } from "../cron/scheduled-tool-policy.js";
import type { AgentRunDelegatedAuthority } from "../infra/agent-run-registry.js";
import type { ExecMode } from "../infra/exec-approvals.js";
import type { PluginHookChannelContext } from "../plugins/hook-types.js";
import { resolveGlobalMap } from "../shared/global-singleton.js";
@ -122,11 +123,15 @@ type StoredMcpLoopbackClientGrant = McpLoopbackClientGrant & {
runtimeOwnerToken: string;
/** Exact host admission retained outside the child-visible request context. */
admittedRunContext?: AdmittedRunContext;
abortSignal?: AbortSignal;
assertCurrent?: () => void;
/** Original CLI policy, rebound only to this stored row's exact lifetime. */
bindQuestionAnswerAuthority?: (assertActive: () => void) => PreparedQuestionAnswerAuthority;
skillLibraryAuthoring?: SkillLibraryAuthoringCapability;
rootedExecution?: PreparedRootedExecutionCapability;
activeCaptureKey?: string;
/** Effective attempt authority, including plugin-owned timeout and cancellation. */
assertCaptureCurrent?: () => void;
toolAuth?: McpLoopbackToolAuth;
};
@ -234,6 +239,8 @@ export function mintMcpLoopbackClientGrant(params: {
context: McpLoopbackRequestContext;
runtimeOwnerToken: string;
admittedRunContext?: AdmittedRunContext;
abortSignal?: AbortSignal;
assertCurrent?: () => void;
bindQuestionAnswerAuthority?: StoredMcpLoopbackClientGrant["bindQuestionAnswerAuthority"];
skillLibraryAuthoring?: SkillLibraryAuthoringCapability;
rootedExecution?: PreparedRootedExecutionCapability;
@ -252,6 +259,8 @@ export function mintMcpLoopbackClientGrant(params: {
context: structuredClone({ ...params.context, sessionKey }),
runtimeOwnerToken,
...(params.admittedRunContext ? { admittedRunContext: params.admittedRunContext } : {}),
abortSignal: params.abortSignal,
assertCurrent: params.assertCurrent,
bindQuestionAnswerAuthority: params.bindQuestionAnswerAuthority,
...(params.skillLibraryAuthoring
? { skillLibraryAuthoring: params.skillLibraryAuthoring }
@ -275,6 +284,27 @@ function replaceMcpLoopbackClientGrant(grant: StoredMcpLoopbackClientGrant): voi
});
}
function isMcpLoopbackClientGrantCurrent(
grant: StoredMcpLoopbackClientGrant,
authority: AgentRunDelegatedAuthority | undefined,
): boolean {
if (!grant.admittedRunContext || !authority || grant.abortSignal?.aborted) {
return false;
}
try {
grant.assertCurrent?.();
grant.assertCaptureCurrent?.();
} catch {
return false;
}
// Caller assertions can revoke or replace the row while checking their own owner.
return (
getAdmittedRunDelegatedAuthority(grant.admittedRunContext) === authority &&
!grant.abortSignal?.aborted &&
clientGrantsByToken.get(grant.token) === grant
);
}
/** Attaches the exact late CLI admission before the grant can execute tools. */
export function bindMcpLoopbackClientGrantAdmission(params: {
token: string;
@ -298,6 +328,7 @@ export function activateMcpLoopbackClientGrantCapture(params: {
token: string;
runtimeOwnerToken: string;
captureKey: string;
assertCurrent?: () => void;
}): false | { captureNativeToolAuthority: (toolNames: readonly string[] | null) => boolean } {
const captureKey = params.captureKey.trim();
if (!captureKey) {
@ -310,6 +341,7 @@ export function activateMcpLoopbackClientGrantCapture(params: {
let activeGrant = {
...grant,
activeCaptureKey: captureKey,
assertCaptureCurrent: params.assertCurrent,
context: {
...grant.context,
...(grant.context.nativeCronCreatorToolAllowlist !== undefined
@ -327,8 +359,7 @@ export function activateMcpLoopbackClientGrantCapture(params: {
if (
!authority ||
!admission ||
clientGrantsByToken.get(params.token) !== activeGrant ||
getAdmittedRunDelegatedAuthority(admission) !== authority ||
!isMcpLoopbackClientGrantCurrent(activeGrant, authority) ||
activeGrant.context.nativeCronCreatorToolAllowlist === undefined
) {
return false;
@ -361,7 +392,11 @@ export function deactivateMcpLoopbackClientGrantCapture(params: {
) {
return false;
}
const { activeCaptureKey: _activeCaptureKey, ...inactiveGrant } = grant;
const {
activeCaptureKey: _activeCaptureKey,
assertCaptureCurrent: _assertCaptureCurrent,
...inactiveGrant
} = grant;
replaceMcpLoopbackClientGrant(inactiveGrant);
return true;
}
@ -387,7 +422,11 @@ export function transferMcpLoopbackClientGrant(params: {
// The child cannot replace its bearer after launch. Turn cleanup may already
// have revoked that bearer, so recreate it only from this fresh admitted grant.
// An existing bearer owned by another runtime is never replaceable.
const { activeCaptureKey: _activeCaptureKey, ...inactiveSource } = source;
const {
activeCaptureKey: _activeCaptureKey,
assertCaptureCurrent: _assertCaptureCurrent,
...inactiveSource
} = source;
clientGrantsByToken.set(params.targetToken, {
...inactiveSource,
token: params.targetToken,
@ -439,9 +478,10 @@ export function resolveMcpLoopbackClientGrant(params: {
return undefined;
}
// Every bind, capture change, and transfer replaces the row, fencing even same-reference reuse.
const isCurrent = () =>
clientGrantsByToken.get(token) === grant &&
getAdmittedRunDelegatedAuthority(admittedRunContext) === delegatedAuthority;
const isCurrent = () => isMcpLoopbackClientGrantCurrent(grant, delegatedAuthority);
if (!isCurrent()) {
return undefined;
}
const questionAnswerAuthority = grant.bindQuestionAnswerAuthority?.(() => {
if (!isCurrent()) {
throw new Error("question creator MCP grant is no longer active");

View file

@ -12,6 +12,7 @@ import {
buildDefaultTestCliBackend,
createCliRunnerPrepareFixture,
} from "../agents/cli-runner.test-helpers.js";
import { createCliRunCurrentAssertion } from "../agents/cli-runner/execution-target.js";
import { prepareCliRunContext } from "../agents/cli-runner/prepare.js";
import {
resetCliRunnerPrepareTestDeps,
@ -237,7 +238,10 @@ async function withCliQuestionLoopback(
context.preparedBackend.env?.OPENCLAW_MCP_TOKEN,
"prepared CLI grant",
);
context.preparedBackend.mcpClientGrantCapture?.activate(captureKey);
context.preparedBackend.mcpClientGrantCapture?.activate(
captureKey,
createCliRunCurrentAssertion(context.params),
);
expect(
resolveMcpLoopbackClientGrant({
token,
@ -395,8 +399,8 @@ describe("CLI loopback question creator authority", () => {
token: owner.token,
runtimeOwnerToken: fixture.runtimeOwnerToken,
captureKey,
})?.isCurrent(),
).toBe(true);
}),
).toBeUndefined();
} else {
owner.admission.close();
}

View file

@ -6,12 +6,16 @@ import { createDeferred } from "../../../test/helpers/promise.js";
import { runQaGatewayFixture } from "../../../test/helpers/qa-gateway-cleanup.js";
import "../../agents/test-helpers/fast-coding-tools.js";
import "../../agents/test-helpers/fast-openclaw-tools.js";
import { prepareSystemAgentRunAdmission } from "../../agents/admitted-run-context.js";
import {
getAdmittedRunDelegatedAuthority,
prepareSystemAgentRunAdmission,
} from "../../agents/admitted-run-context.js";
import { testing as cliBackendsTesting } from "../../agents/cli-backends.test-support.js";
import {
buildDefaultTestCliBackend,
createCliRunnerPrepareFixture,
} from "../../agents/cli-runner.test-helpers.js";
import { createCliRunCurrentAssertion } from "../../agents/cli-runner/execution-target.js";
import { prepareCliRunContext } from "../../agents/cli-runner/prepare.js";
import {
resetCliRunnerPrepareTestDeps,
@ -80,6 +84,9 @@ async function withRootedCli(
write: (filePath: string, content: string) => Promise<McpResponse>;
revoke: () => boolean;
replace: () => boolean;
source: AbortController;
requestSignal: AbortSignal;
admittedRunContext: PreparedCliRunContext["params"]["admittedRunContext"];
}) => Promise<void>,
) {
const cli = createCliRunnerPrepareFixture(prepareCliRunContext);
@ -94,6 +101,7 @@ async function withRootedCli(
const admission = prepareSystemAgentRunAdmission(config, "rooted-mcp-run", "main", "rooted-test");
const requests: Promise<McpResponse>[] = [];
const controller = new AbortController();
const source = new AbortController();
let prepared: PreparedCliRunContext | undefined;
await runQaGatewayFixture(
async () => {
@ -105,6 +113,7 @@ async function withRootedCli(
runId: "rooted-mcp-run",
sessionKey: "agent:main:main",
preparedRunAdmission: admission,
abortSignal: source.signal,
rootedExecution: { root },
skillsSnapshot: { prompt: "", skills: [] },
cliToolAvailability: { native: [], openClaw: ["read", "write"] },
@ -119,6 +128,7 @@ async function withRootedCli(
};
expectDefined(prepared.preparedBackend.mcpClientGrantCapture, "CLI capture").activate(
capture.captureKey,
createCliRunCurrentAssertion(prepared.params),
);
const request = (method: string, params?: Record<string, unknown>) => {
const response = (async () => {
@ -147,8 +157,12 @@ async function withRootedCli(
request("tools/call", { name: "write", arguments: { path: filePath, content } }),
revoke: () => revokeMcpLoopbackClientGrant(token),
replace: () => activateMcpLoopbackClientGrantCapture(capture) !== false,
source,
requestSignal: controller.signal,
admittedRunContext: prepared.params.admittedRunContext,
});
},
() => source.abort(),
() => controller.abort(),
() => Promise.allSettled(requests),
() => closeMcpLoopbackServer(),
@ -179,7 +193,7 @@ describe("rooted CLI grants through MCP HTTP dispatch", () => {
});
});
it.each(["revoke", "replace"] as const)(
it.each(["revoke", "replace", "source-cancel"] as const)(
"prevents a pending file write when the grant is changed by %s",
async (change) => {
await withRootedCli(async (fixture) => {
@ -205,14 +219,20 @@ describe("rooted CLI grants through MCP HTTP dispatch", () => {
throw new Error("write completed before filesystem preparation was paused");
}),
]);
expect(fixture[change]()).toBe(true);
if (change === "source-cancel") {
fixture.source.abort();
expect(fixture.requestSignal.aborted).toBe(false);
expect(getAdmittedRunDelegatedAuthority(fixture.admittedRunContext)).toBeDefined();
} else {
expect(fixture[change]()).toBe(true);
}
} finally {
release.resolve();
}
const result = await response;
await expect(fs.readFile(report, "utf8")).resolves.toBe("original report");
expect(result.result.isError).toBe(true);
expect(JSON.stringify(result.result.content)).toContain("authority is no longer active");
await expect(fs.readFile(report, "utf8")).resolves.toBe("original report");
if (change === "replace") {
expect((await fixture.write("report.md", "fresh write")).result.isError).toBe(false);
await expect(fs.readFile(report, "utf8")).resolves.toBe("fresh write");

View file

@ -123,6 +123,7 @@ describe("process supervisor queued cancellation", () => {
);
const replacementRunId = `cancel-queued-${mode}-replacement`;
const resolveArgs = vi.fn(() => ["must-not-resolve"]);
let captureCurrent = true;
const replacementPromise = supervisor.spawn(
Object.assign(
createSpawnInput({
@ -132,6 +133,11 @@ describe("process supervisor queued cancellation", () => {
replaceExistingScope: true,
}),
mode === "child" ? { resolveArgs } : {},
{
onCancel: () => {
captureCurrent = false;
},
},
),
);
@ -139,6 +145,7 @@ describe("process supervisor queued cancellation", () => {
expect(createPtyAdapterMock).not.toHaveBeenCalled();
supervisor.cancel(replacementRunId, "manual-cancel");
expect(captureCurrent).toBe(false);
firstStartup.resolve(first);
const [firstRun, replacementRun] = await Promise.all([firstRunPromise, replacementPromise]);

View file

@ -400,15 +400,20 @@ describe("process supervisor", () => {
scopeKey: "scope:cancel-fenced",
argv: createSilentIdleArgv(),
});
let replacementCurrent = true;
const replacementPromise = spawnChild(supervisor, {
runId: "cancel-fenced-replacement",
scopeKey: "scope:cancel-fenced",
replaceExistingScope: true,
argv: createSilentIdleArgv(),
onCancel: () => {
replacementCurrent = false;
},
});
expect(createChildAdapterMock).toHaveBeenCalledTimes(1);
supervisor.cancelScope("scope:cancel-fenced", "manual-cancel");
expect(replacementCurrent).toBe(false);
const laterPromise = spawnChild(supervisor, {
runId: "cancel-fenced-later",

View file

@ -24,7 +24,7 @@ type OwnedRun = {
runId: string;
scopeKey?: string;
terminationReason?: TerminationReason;
cancel?: (reason: TerminationReason) => void;
cancel: (reason: TerminationReason) => void;
pending?: Promise<ManagedRun>;
waitForExtinction?: () => Promise<void>;
cleanupOwners: ScopeCleanupOwner[];
@ -118,18 +118,10 @@ export function createProcessSupervisor(): ProcessSupervisor & {
let shutdownPromise: Promise<void> | null = null;
let cleanupFailure: { error: unknown } | undefined;
const cancelOwner = (current: OwnedRun, reason: TerminationReason) => {
if (current.cancel) {
current.cancel(reason);
return;
}
current.terminationReason ??= reason;
};
const cancel = (runId: string, reason: TerminationReason = "manual-cancel") => {
for (const current of ownedRuns) {
if (current.runId === runId) {
cancelOwner(current, reason);
current.cancel(reason);
}
}
};
@ -137,7 +129,7 @@ export function createProcessSupervisor(): ProcessSupervisor & {
const cancelActiveScope = (scopeKey: string, reason: TerminationReason) => {
for (const current of ownedRuns) {
if (current.waitForExtinction && current.scopeKey === scopeKey) {
cancelOwner(current, reason);
current.cancel(reason);
}
}
};
@ -148,7 +140,7 @@ export function createProcessSupervisor(): ProcessSupervisor & {
}
for (const current of ownedRuns) {
if (current.scopeKey === scopeKey) {
cancelOwner(current, reason);
current.cancel(reason);
}
}
};
@ -317,6 +309,7 @@ export function createProcessSupervisor(): ProcessSupervisor & {
const requestCancel = (reason: TerminationReason) => {
setForcedReason(reason);
input.onCancel?.(reason);
cancelAdapter?.(reason);
// Any cancel must abort construction: the relay may already be spawned
// and waiting for ready, and a later deadline must not replace this reason.
@ -645,6 +638,10 @@ export function createProcessSupervisor(): ProcessSupervisor & {
const owner: OwnedRun = {
runId,
scopeKey,
cancel: (reason) => {
owner.terminationReason ??= reason;
input.onCancel?.(reason);
},
cleanupOwners: scopeKey ? [...(scopeCleanupOwners.get(scopeKey) ?? [])] : [],
};
// Reserve cancellation before either adapter startup or a replacement
@ -700,7 +697,7 @@ export function createProcessSupervisor(): ProcessSupervisor & {
return (shutdownPromise ??= Promise.resolve().then(async () => {
while (ownedRuns.size) {
for (const owner of ownedRuns) {
cancelOwner(owner, "manual-cancel");
owner.cancel("manual-cancel");
}
// A failed startup owns no live process; only failed owner extinction
// must keep the process-wide supervisor fenced for operator recovery.

View file

@ -109,6 +109,8 @@ type SpawnBaseInput = {
maxCapturedOutputChars?: number;
onStdout?: (chunk: string) => void;
onStderr?: (chunk: string) => void;
/** Revoke caller-owned capabilities when cancellation starts, before native termination. */
onCancel?: (reason: TerminationReason) => void;
};
type SpawnChildInput = SpawnBaseInput & {

View file

@ -24,6 +24,7 @@ import {
buildDefaultTestCliBackend,
createCliRunnerPrepareFixture,
} from "../src/agents/cli-runner.test-helpers.js";
import { createCliRunCurrentAssertion } from "../src/agents/cli-runner/execution-target.js";
import { prepareCliRunContext } from "../src/agents/cli-runner/prepare.js";
import {
resetCliRunnerPrepareTestDeps,
@ -288,7 +289,10 @@ describe("loopback ask_user Telegram channel transport", () => {
context.preparedBackend.env?.OPENCLAW_MCP_TOKEN,
"prepared CLI grant",
);
context.preparedBackend.mcpClientGrantCapture?.activate(captureKey);
context.preparedBackend.mcpClientGrantCapture?.activate(
captureKey,
createCliRunCurrentAssertion(context.params),
);
expect(
resolveMcpLoopbackClientGrant({
token,