Files
PaperClipAI/patches/acpx@0.13.1.patch
DottaandPaperclip 7d89e304f0 fix(acpx): separate restored history from active turn events
Keep history reconstruction without publishing restored text and tools as live output. Version the changed Cursor contract as profile 16 and retain replay compatibility.

Co-Authored-By: Paperclip <noreply@paperclip.ing>
2026-10-08 12:07:09 -05:00

1317 lines
62 KiB
Diff

diff --git a/dist/client-CxNllqui.d.ts b/dist/client-CxNllqui.d.ts
--- a/dist/client-CxNllqui.d.ts
+++ b/dist/client-CxNllqui.d.ts
@@ -120,6 +120,8 @@
private initializeProtocolConnection;
private handleInitializeFailure;
private createTappedStream;
+ private cursorSubagentStreamOwner;
+ private handleCursorSubagentNotification;
createSession(cwd?: string): Promise<SessionCreateResult>;
loadSession(sessionId: string, cwd?: string): Promise<SessionLoadResult>;
loadSessionWithOptions(sessionId: string, cwd?: string, options?: LoadSessionOptions): Promise<SessionLoadResult>;
@@ -135,6 +137,7 @@
private throwPromptPermissionFailureIfPresent;
setSessionMode(sessionId: string, modeId: string): Promise<void>;
setSessionConfigOption(sessionId: string, configId: string, value: string): Promise<SetSessionConfigOptionResponse>;
+ requestExtension(method: string, params: Record<string, unknown>): Promise<Record<string, unknown>>;
setSessionModel(sessionId: string, modelId: string, controlOverride?: ModelControlOverride): Promise<SetSessionConfigOptionResponse | undefined>;
private setSessionModelThroughConfig;
private setSessionModelThroughLegacyMethod;
diff --git a/dist/live-checkpoint-BSIrfgVo.js b/dist/live-checkpoint-BSIrfgVo.js
--- a/dist/live-checkpoint-BSIrfgVo.js
+++ b/dist/live-checkpoint-BSIrfgVo.js
@@ -6,7 +6,7 @@
import { randomUUID } from "node:crypto";
import { execFile, spawn } from "node:child_process";
import { Readable, Writable } from "node:stream";
-import { PROTOCOL_VERSION, client, methods } from "@agentclientprotocol/sdk";
+import { PROTOCOL_VERSION, RequestError, client, methods } from "@agentclientprotocol/sdk";
import readline from "node:readline/promises";
//#region src/errors.ts
var AcpxOperationalError = class extends Error {
@@ -1068,12 +1068,14 @@
last_agent_disconnect_reason: canonical.lastAgentDisconnectReason,
protocol_version: canonical.protocolVersion,
agent_capabilities: canonical.agentCapabilities,
+ agent_goal_capability: canonical.agentGoalCapability,
title: canonical.title,
messages: canonical.messages,
updated_at: canonical.updated_at,
cumulative_token_usage: canonical.cumulative_token_usage,
cumulative_cost: canonical.cumulative_cost,
request_token_usage: canonical.request_token_usage,
+ cursor_prompt_usage: boundCursorPromptUsage(canonical.cursor_prompt_usage, canonical.messages),
acpx: canonical.acpx,
imported_from: canonical.importedFrom ? {
record_id: canonical.importedFrom.recordId,
@@ -1239,7 +1241,8 @@
for (const [key, value] of Object.entries(record)) {
const parsed = parseTokenUsage(value);
if (parsed == null) return null;
- usage[key] = parsed;
+ const receipt = parsePaperclipPiReceipt(asRecord$4(value)?.paperclip_pi);
+ usage[key] = receipt ? { ...parsed, paperclip_pi: receipt } : parsed;
}
return usage;
}
@@ -1325,6 +1328,93 @@
function isConversationMessage(raw) {
return raw === "Resume" || isUserMessage$1(raw) || isAgentMessage$1(raw);
}
+function parseCursorPromptUsage(raw) {
+ try {
+ const text = JSON.stringify(raw);
+ if (!text || Buffer.byteLength(text, "utf8") > 16384) return null;
+ const v = JSON.parse(text);
+ const object = (x)=>x !== null && typeof x === "object" && !Array.isArray(x);
+ const exact = (x, keys)=>Object.keys(x).length === keys.length && keys.every((k)=>Object.hasOwn(x, k));
+ const id = (x)=>typeof x === "string" && /^[A-Za-z0-9][A-Za-z0-9._:-]{0,239}$/.test(x);
+ if (!object(v) || !exact(v, [
+ "request_id",
+ "prompt_message_id",
+ "receipt"
+ ]) || !id(v.request_id) || !id(v.prompt_message_id)) return null;
+ const r = v.receipt;
+ if (!object(r) || !exact(r, [
+ "schema",
+ "source",
+ "promptId",
+ "completeness",
+ "reasons",
+ "observations",
+ "limits",
+ "truncated"
+ ])) return null;
+ if (r.schema !== "paperclip.cursor.native-usage.v1" || r.source !== "native_turn_ended" || r.completeness !== "partial" || typeof r.promptId !== "string" || !/^[a-f0-9]{8}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{12}$/.test(r.promptId) || typeof r.truncated !== "boolean") return null;
+ const reasons = [
+ "native_counter_semantics_unverified",
+ "native_counters_missing",
+ "invalid_native_counter",
+ "multiple_terminal_observations",
+ "native_terminal_not_observed",
+ "observation_limit_reached",
+ "child_run_attribution_unverified"
+ ];
+ if (!Array.isArray(r.reasons) || r.reasons.length > reasons.length || !r.reasons.includes(reasons[0]) || new Set(r.reasons).size !== r.reasons.length || !r.reasons.every((x)=>reasons.includes(x))) return null;
+ if (!object(r.limits) || !exact(r.limits, [
+ "maxObservations",
+ "maxInvocations",
+ "maxBytes"
+ ]) || r.limits.maxObservations !== 64 || r.limits.maxInvocations !== 64 || r.limits.maxBytes !== 16384 || !Array.isArray(r.observations) || r.observations.length > 64) return null;
+ const fields = [
+ "inputTokens",
+ "outputTokens",
+ "cacheReadTokens",
+ "cacheWriteTokens",
+ "reasoningTokens"
+ ];
+ const invocations = new Map();
+ for (const o of r.observations){
+ if (!object(o) || !exact(o, [
+ "invocationId",
+ "role",
+ "nativeRun",
+ "sequence",
+ "counters"
+ ]) || typeof o.invocationId !== "string" || !/^invocation-([1-9]|[1-5][0-9]|6[0-4])$/.test(o.invocationId) || o.role !== "parent" && o.role !== "child" || !Number.isSafeInteger(o.nativeRun) || o.nativeRun < 1 || !Number.isSafeInteger(o.sequence) || o.sequence < 1 || !object(o.counters)) return null;
+ if (!Object.entries(o.counters).every(([k, n])=>fields.includes(k) && Number.isSafeInteger(n) && n >= 0)) return null;
+ const previous = invocations.get(o.invocationId);
+ if (previous && (previous.role !== o.role || previous.nativeRun !== o.nativeRun || o.sequence <= previous.sequence)) return null;
+ invocations.set(o.invocationId, {
+ role: o.role,
+ nativeRun: o.nativeRun,
+ sequence: o.sequence
+ });
+ }
+ return v;
+ } catch {
+ return null;
+ }
+}
+function boundCursorPromptUsage(raw, messages) {
+ const receipt = parseCursorPromptUsage(raw);
+ if (!receipt || !Array.isArray(messages) || messages.filter(message => message?.User?.id === receipt.prompt_message_id).length !== 1) return undefined;
+ return receipt;
+}
+function recordCursorPromptUsage(conversation, response, requestId, promptMessageId) {
+ // Optional diagnostics cannot affect settlement or standard token/cost accounting.
+ try {
+ if (response.stopReason !== "end_turn" && response.stopReason !== "cancelled") return;
+ const next = boundCursorPromptUsage({request_id: requestId, prompt_message_id: promptMessageId, receipt: response._meta?.paperclipCursorUsage}, conversation.messages);
+ if (!next) return;
+ const previous = parseCursorPromptUsage(conversation.cursor_prompt_usage);
+ if (previous && (previous.request_id === next.request_id || previous.prompt_message_id === next.prompt_message_id || previous.receipt.promptId === next.receipt.promptId)) return;
+ conversation.cursor_prompt_usage = next;
+ } catch {}
+}
+
function parseConversationRecord(record) {
if (!hasValidConversationCore(record)) return;
const title = parseConversationTitle(record.title);
@@ -1339,7 +1429,8 @@
updated_at: record.updated_at,
cumulative_token_usage: cumulativeTokenUsage ?? {},
cumulative_cost: cumulativeCost,
- request_token_usage: requestTokenUsage ?? {}
+ request_token_usage: requestTokenUsage ?? {},
+ cursor_prompt_usage: boundCursorPromptUsage(record.cursor_prompt_usage, record.messages)
};
}
const INVALID_VALUE = Symbol("invalid");
@@ -1541,12 +1632,14 @@
lastAgentDisconnectReason: optionals.lastAgentDisconnectReason,
protocolVersion: typeof record.protocol_version === "number" ? record.protocol_version : void 0,
agentCapabilities: asRecord$4(record.agent_capabilities),
+ agentGoalCapability: asRecord$4(record.agent_goal_capability),
title: conversation.title,
messages: conversation.messages,
updated_at: conversation.updated_at,
cumulative_token_usage: conversation.cumulative_token_usage,
cumulative_cost: conversation.cumulative_cost,
request_token_usage: conversation.request_token_usage,
+ cursor_prompt_usage: conversation.cursor_prompt_usage,
acpx: parseAcpxState(record.acpx),
importedFrom: metadata.importedFrom
};
@@ -1661,8 +1754,9 @@
"RedactedThinking",
"ToolUse"
]);
-const MAP_OBJECT_PATHS = /* @__PURE__ */ new Set(["request_token_usage", "messages.Agent.tool_results"]);
+const MAP_OBJECT_PATHS = /* @__PURE__ */ new Set(["request_token_usage", "messages.Agent.tool_results", "acpx.session_options.env"]);
const OPAQUE_VALUE_PATHS = /* @__PURE__ */ new Set([
+ "cursor_prompt_usage.receipt",
"agent_capabilities",
"messages.Agent.content.ToolUse.input",
"acpx.desired_config_options",
@@ -3104,10 +3198,10 @@
authEnvKeyCache.set(methodId, key);
return key;
}
-function readEnvCredential(methodId) {
+function readEnvCredential(methodId, environment = process.env) {
const key = authEnvKey(methodId);
if (!key) return;
- const value = process.env[key];
+ const value = environment[key];
if (typeof value === "string" && value.trim().length > 0) return value;
}
function protectedEnvKey(key) {
@@ -3135,8 +3229,21 @@
}
return protectedKeys;
}
-function buildAgentEnvironment(authCredentials, sessionEnv) {
- const env = { ...process.env };
+function isPlainStringEnvironment(value) {
+ if (value === null || typeof value !== "object" || Array.isArray(value)) return false;
+ const prototype = Object.getPrototypeOf(value);
+ return (prototype === Object.prototype || prototype === null) && Object.values(value).every((entry) => typeof entry === "string");
+}
+function resolveAgentEnvironment(spawnEnvironment) {
+ let sourceEnvironment = process.env;
+ if (spawnEnvironment !== void 0) {
+ sourceEnvironment = spawnEnvironment();
+ if (!isPlainStringEnvironment(sourceEnvironment)) throw new TypeError("ACPX spawn environment must be a plain record of string values");
+ }
+ return sourceEnvironment;
+}
+function buildAgentEnvironment(authCredentials, sessionEnv, spawnEnvironment) {
+ const env = { ...resolveAgentEnvironment(spawnEnvironment) };
const protectedAuthEnvKeys = promotePrefixedAuthEnvironment(env);
if (authCredentials) for (const [methodId, credential] of Object.entries(authCredentials)) {
addAuthCredentialEnvKeys(protectedAuthEnvKeys, methodId, credential);
@@ -3178,10 +3285,10 @@
const configCredentials = authCredentials ?? {};
return configCredentials[methodId] ?? configCredentials[toEnvToken(methodId)];
}
-function buildAgentSpawnOptions(cwd, authCredentials, sessionEnv) {
+function buildAgentSpawnOptions(cwd, authCredentials, sessionEnv, spawnEnvironment) {
return {
cwd,
- env: buildAgentEnvironment(authCredentials, sessionEnv),
+ env: buildAgentEnvironment(authCredentials, sessionEnv, spawnEnvironment),
stdio: [
"pipe",
"pipe",
@@ -3286,28 +3393,12 @@
const ids = models?.availableModels.map((model) => model.modelId.trim()).filter((modelId) => modelId.length > 0) ?? [];
return ids.length > 0 ? ids.join(", ") : "none advertised";
}
-function resolveRequestedModelId(params) {
- if (!params.models || !isCursorAcpCommandForModelAlias(params.agentCommand)) return params.requestedModel;
- if (params.models.availableModels.some((model) => model.modelId === params.requestedModel)) return params.requestedModel;
- const candidates = params.models.availableModels.map((model) => model.modelId).filter((modelId) => modelId.startsWith(`${params.requestedModel}[`));
- return candidates.length === 1 ? candidates[0] : params.requestedModel;
-}
-function isCursorAcpCommandForModelAlias(agentCommand) {
- if (!agentCommand) return false;
- const { command, args } = splitCommandLine(agentCommand);
- return isCursorAcpCommand(command, args);
-}
function assertRequestedModelSupported(params) {
if (!params.models) {
if (supportsLegacyClaudeCodeModelMetadata(params.agentCommand)) return;
throw new RequestedModelUnsupportedError(`Cannot ${params.context === "replay" ? "replay saved model" : "apply --model"} "${params.requestedModel}": the ACP agent did not advertise model support through a session config option or legacy models metadata, and the adapter does not support a startup model flag.`, "missing-capability");
}
- if (!new Set(params.models.availableModels.map((model) => model.modelId)).has(params.requestedModel)) {
- const resolvedModel = resolveRequestedModelId(params);
- if (resolvedModel !== params.requestedModel) return `Cursor ACP advertised "${resolvedModel}" for requested model "${params.requestedModel}"; using the advertised id.`;
- if (supportsLegacyClaudeCodeModelMetadata(params.agentCommand)) return `requested model "${params.requestedModel}" was not in the Claude ACP advertised model list (${formatAvailableModelIds(params.models)}); forwarding it to Claude Code so the adapter can accept or reject it.`;
- throw new RequestedModelUnsupportedError(`Cannot ${params.context === "replay" ? "replay saved model" : "apply --model"} "${params.requestedModel}": the ACP agent did not advertise that model. Available models: ${formatAvailableModelIds(params.models)}.`, "unadvertised-model");
- }
+ // Discovery catalogs may be incomplete; the provider accepts or rejects the exact ID.
}
//#endregion
//#region src/acp/session-control-errors.ts
@@ -3961,10 +4052,13 @@
...params.elicitationModes.includes("url") ? { url: {} } : {}
} } : {}
};
- if (!params.devinAcp) return baseCapabilities;
+ const typedSessionFailureMeta = {
+ jetbrains: { air: { version: 1, capabilities: ["sessionFailure"] } }
+ };
+ if (!params.devinAcp) return { ...baseCapabilities, _meta: typedSessionFailureMeta };
return {
...baseCapabilities,
- _meta: DEVIN_COMPATIBILITY_CLIENT_CAPABILITIES_META
+ _meta: { ...typedSessionFailureMeta, ...DEVIN_COMPATIBILITY_CLIENT_CAPABILITIES_META }
};
}
function hasResponseField(response, field) {
@@ -4067,6 +4161,27 @@
} })
};
}
+function boundedExtensionRecord(value) {
+ if (value === null || typeof value !== "object" || Array.isArray(value)) throw new TypeError("ACP extension payload must be an object");
+ const prototype = Object.getPrototypeOf(value);
+ if (prototype !== Object.prototype && prototype !== null) throw new TypeError("ACP extension payload must be a plain object");
+ const serialized = JSON.stringify(value);
+ if (Buffer.byteLength(serialized) > 256 * 1024) throw new TypeError("ACP extension payload exceeds its bounded size");
+ return JSON.parse(serialized);
+}
+function closedExtensionMethods(value) {
+ if (value === void 0) return Object.freeze([]);
+ if (!Array.isArray(value) || value.length > 64) throw new TypeError("ACP extension methods must be a bounded allowlist");
+ for (const method of value) {
+ if (typeof method !== "string" || !/^[A-Za-z_][A-Za-z0-9._/-]{1,255}$/.test(method) || !method.includes("/") || /^(?:session|fs|terminal|elicitation|notifications)\//.test(method)) throw new TypeError("ACP extension method cannot override a standard method");
+ }
+ return Object.freeze([...new Set(value)]);
+}
+function mergeRuntimeClientCapabilities(base, extra) {
+ if (extra === void 0) return base;
+ const additional = boundedExtensionRecord(extra);
+ return { ...additional, ...base, _meta: { ...additional._meta, ...base._meta } };
+}
var AcpClient = class {
options;
connection;
@@ -4105,7 +4220,9 @@
cwd: asAbsoluteCwd(options.cwd),
authPolicy: options.authPolicy ?? "skip",
permissionPolicy: snapshotPermissionPolicy(options.permissionPolicy),
- elicitationModes: normalizeElicitationModes(options.elicitationModes)
+ elicitationModes: normalizeElicitationModes(options.elicitationModes),
+ extensionMethods: closedExtensionMethods(options.extensionMethods),
+ clientCapabilities: options.clientCapabilities === void 0 ? void 0 : boundedExtensionRecord(options.clientCapabilities)
};
this.eventHandlers = {
onAcpMessage: this.options.onAcpMessage,
@@ -4216,9 +4333,21 @@
this.lastAgentExit = void 0;
this.lastKnownPid = child.pid ?? void 0;
this.attachAgentLifecycleObservers(child);
+ if (this.options.onAgentSpawn) {
+ try {
+ if (!child.pid) throw new Error("ACPX agent spawn did not expose a valid process id.");
+ await this.options.onAgentSpawn({ pid: child.pid, startedAt: this.agentStartedAt });
+ } catch (error) {
+ try {
+ child.kill("SIGKILL");
+ } catch {}
+ throw error;
+ }
+ }
const startupStderr = [];
child.stderr.on("data", (chunk) => {
this.captureStartupStderr(startupStderr, chunk);
+ this.options.onAgentStderr?.(String(chunk));
if (!this.options.verbose) return;
process.stderr.write(chunk);
});
@@ -4226,6 +4355,7 @@
const output = Readable.toWeb(child.stdout);
const stream = this.createTappedStream(createNdJsonMessageStream(this.options.agentCommand, input, output));
const connection = this.createConnection(stream, launch);
+ connection.signal.addEventListener("abort", () => stream.abortResponseDeliveries(), { once: true });
connection.signal.addEventListener("abort", () => {
this.recordAgentExit("connection_close", child.exitCode ?? null, child.signalCode ?? null);
}, { once: true });
@@ -4253,7 +4383,12 @@
geminiAcp: isGeminiAcpCommand(spawnCommand, args),
copilotAcp: isCopilotAcpCommand(spawnCommand, args),
claudeAcp: isClaudeAcpCommand(spawnCommand, args),
- spawnOptions: buildAgentSpawnOptions(this.options.cwd, this.options.authCredentials, this.options.sessionOptions?.env)
+ spawnOptions: buildAgentSpawnOptions(
+ this.options.spawnCwd ?? this.options.cwd,
+ this.options.authCredentials,
+ this.options.sessionOptions?.env,
+ this.options.spawnEnvironment
+ )
};
}
logAgentLaunch(plan) {
@@ -4279,10 +4414,17 @@
}
async spawnAgentProcess(plan) {
const spawnCommand = buildAgentSpawnCommand(plan.spawnCommand, plan.args, process.platform, plan.spawnOptions.env);
- const spawnedChild = spawn(spawnCommand.command, spawnCommand.args, {
+ const options = {
...plan.spawnOptions,
windowsVerbatimArguments: spawnCommand.windowsVerbatimArguments
- });
+ };
+ const spawnedChild = this.options.spawnAgent
+ ? this.options.spawnAgent({
+ command: spawnCommand.command,
+ args: spawnCommand.args,
+ options
+ })
+ : spawn(spawnCommand.command, spawnCommand.args, options);
try {
await waitForSpawn$1(spawnedChild);
} catch (error) {
@@ -4293,10 +4435,10 @@
createConnection(stream, launch) {
const app = client({ name: "acpx" }).onNotification(methods.client.session.update, async ({ params }) => {
await this.handleSessionUpdate(params);
- }).onNotification(methods.client.elicitation.complete, async () => {}).onRequest(methods.client.session.requestPermission, async ({ params }) => {
- return await this.handlePermissionRequest(params);
+ }).onNotification(methods.client.elicitation.complete, async () => {}).onRequest(methods.client.session.requestPermission, async ({ params, requestId, signal }) => {
+ return await stream.withResponseDelivery(requestId, signal, (responseDelivery) => this.handlePermissionRequest(params, responseDelivery));
}).onRequest(methods.client.elicitation.create, async ({ params, requestId, signal }) => {
- return await this.handleElicitationRequest(params, requestId, signal);
+ return await stream.withResponseDelivery(requestId, signal, (responseDelivery) => this.handleElicitationRequest(params, requestId, signal, responseDelivery));
}).onRequest(methods.client.fs.readTextFile, async ({ params }) => {
return await this.handleReadTextFile(params);
}).onRequest(methods.client.fs.writeTextFile, async ({ params }) => {
@@ -4312,6 +4454,11 @@
}).onRequest(methods.client.terminal.release, async ({ params }) => {
return await this.handleReleaseTerminal(params);
});
+ for (const method of this.options.extensionMethods) {
+ if (this.options.onExtensionRequest) app.onRequest(method, boundedExtensionRecord, async ({ params, requestId, signal }) => {
+ return await stream.withResponseDelivery(requestId, signal, (responseDelivery) => this.handleExtensionRequest(method, params, requestId, signal, responseDelivery));
+ });
+ }
if (launch.devinAcp) app.onRequest("_cognition.ai/request_diagnostics", (params) => {
return params && typeof params === "object" && !Array.isArray(params) ? params : {};
}, async () => ({}));
@@ -4333,12 +4480,12 @@
async initializeProtocolConnection(connection, launch) {
const initializePromise = connection.initialize({
protocolVersion: PROTOCOL_VERSION,
- clientCapabilities: resolveClientCapabilities({
+ clientCapabilities: mergeRuntimeClientCapabilities(resolveClientCapabilities({
devinAcp: launch.devinAcp,
fs: this.options.fs !== false,
terminal: this.options.terminal !== false,
elicitationModes: this.options.elicitationModes ?? []
- }),
+ }), this.options.clientCapabilities),
clientInfo: resolveClientInfo(launch.devinAcp)
});
const initialized = launch.geminiAcp ? await withTimeout(initializePromise, resolveGeminiAcpStartupTimeoutMs()) : await initializePromise;
@@ -4358,6 +4505,49 @@
throw normalizedError;
}
createTappedStream(base) {
+ // Policy admission is distinct from best-effort observation. Each new
+ // stdio connection gets fresh authority and a violation tears it down.
+ const protocolGuard = this.options.protocolGuardFactory?.();
+ const enforceProtocol = (direction, message) => {
+ try { protocolGuard?.(direction, message); }
+ catch (error) {
+ this.rejectPendingConnectionRequests(error);
+ abortResponseDeliveries();
+ this.activePrompt?.elicitationController.abort(error);
+ try { this.agent?.kill("SIGTERM"); } catch {}
+ throw error;
+ }
+ };
+ const deliveries = new Map();
+ const scope = new AbortController();
+ const abortResponseDeliveries = () => scope.abort(new Error("ACP response connection closed before delivery"));
+ const withResponseDelivery = async (id, requestSignal, handler) => {
+ if (!(id === null || typeof id === "string" || typeof id === "number" && Number.isFinite(id))) throw new Error("ACP response request id is invalid");
+ if (deliveries.has(id) || deliveries.size >= 128) {
+ abortResponseDeliveries();
+ throw new Error("ACP response request id is duplicated or pending deliveries exceeded their bound");
+ }
+ const active = this.activePrompt;
+ const signal = AbortSignal.any([scope.signal, requestSignal, ...(active ? [active.elicitationController.signal] : [])]);
+ let settle, reject;
+ const responseDelivery = new Promise((resolve, fail) => { settle = resolve; reject = fail; });
+ // The host may decline the request without awaiting its receipt.
+ void responseDelivery.catch(() => {});
+ const abort = () => reject(new Error("ACP response delivery was cancelled or disconnected"));
+ const receipt = { settle, reject, expected: void 0, cleanup: () => signal.removeEventListener("abort", abort) };
+ deliveries.set(id, receipt);
+ signal.addEventListener("abort", abort, { once: true });
+ if (signal.aborted) abort();
+ try {
+ const result = await handler(responseDelivery);
+ if (signal.aborted || this.activePrompt !== active) throw new Error("ACP response outlived its admitted prompt");
+ receipt.expected = JSON.stringify(result);
+ return result;
+ } catch (error) {
+ reject(new Error("ACP request failed before response delivery"));
+ throw error;
+ }
+ };
const onAcpMessage = () => this.eventHandlers.onAcpMessage;
const onAcpOutputMessage = () => this.eventHandlers.onAcpOutputMessage;
const elicitationRequestIds = /* @__PURE__ */ new Set();
@@ -4367,8 +4557,12 @@
const shouldSuppressInboundReplaySessionUpdate = (message) => {
return this.suppressReplaySessionUpdateMessages && isSessionUpdateNotification(message);
};
+ const cursorSubagentOwner = { active: null, children: new Map() };
+ this.cursorSubagentStreamOwner = cursorSubagentOwner;
+ const onCursorSubagentNotification = (message) => this.handleCursorSubagentNotification(message, cursorSubagentOwner);
+ const onExtensionNotification = (message) => this.handleExtensionNotification(message);
const observeInbound = (message) => {
- const requestId = elicitationRequestId(message);
+ const requestId = elicitationRequestId(message) ?? ("id" in message && this.options.extensionMethods.includes(message.method) ? message.id : void 0);
if (requestId !== void 0) {
elicitationRequestIds.add(requestId);
return;
@@ -4379,6 +4573,7 @@
}
};
return {
+ withResponseDelivery, abortResponseDeliveries,
readable: new ReadableStream({ async start(controller) {
const reader = base.readable.getReader();
try {
@@ -4386,16 +4581,27 @@
const { value, done } = await reader.read();
if (done) break;
if (!value) continue;
+ enforceProtocol("inbound", value);
observeInbound(value);
- if (isExtensionNotification(value)) continue;
+ if (onCursorSubagentNotification(value)) continue;
+ if (isExtensionNotification(value)) {
+ onExtensionNotification(value);
+ continue;
+ }
controller.enqueue(value);
}
+ controller.close();
+ } catch (error) {
+ controller.error(error);
} finally {
+ abortResponseDeliveries();
+ for (const receipt of deliveries.values()) receipt.cleanup();
+ deliveries.clear();
reader.releaseLock();
- controller.close();
}
} }),
writable: new WritableStream({ async write(message) {
+ enforceProtocol("outbound", message);
const promptOwner = promptRequestOwner(message);
if (promptOwner) bindPromptOwner(promptOwner);
const id = responseId(message);
@@ -4403,10 +4609,20 @@
onAcpOutputMessage()?.("outbound", message);
onAcpMessage()?.("outbound", message);
}
+ const receipt = id === void 0 ? void 0 : deliveries.get(id);
const writer = base.writable.getWriter();
try {
await writer.write(message);
+ if (receipt) {
+ if (!("result" in message) || receipt.expected === void 0 || JSON.stringify(message.result) !== receipt.expected) receipt.reject(new Error("ACP wrote an error or changed response instead of the selected resolution"));
+ else receipt.settle();
+ }
+ } catch (error) {
+ receipt?.reject(new Error("ACP response write failed"));
+ abortResponseDeliveries();
+ throw error;
} finally {
+ if (receipt) { receipt.cleanup(); deliveries.delete(id); }
writer.releaseLock();
}
} })
@@ -4422,7 +4638,10 @@
const createPromise = this.runConnectionRequest(() => connection.newSession({
cwd: sessionCwd,
mcpServers: this.options.mcpServers ?? [],
- _meta: buildClaudeCodeOptionsMeta(this.options.sessionOptions, claudeAcp)
+ _meta: this.isGrokBuildAcpCommand()
+ ? { rules: typeof this.options.sessionOptions?.systemPrompt === "string"
+ ? this.options.sessionOptions.systemPrompt : this.options.sessionOptions?.systemPrompt?.append ?? "" }
+ : buildClaudeCodeOptionsMeta(this.options.sessionOptions, claudeAcp)
}));
result = claudeAcp ? await withTimeout(createPromise, resolveClaudeAcpSessionCreateTimeoutMs()) : await createPromise;
} catch (error) {
@@ -4590,15 +4809,17 @@
throw maybeWrapSessionControlError("session/set_config_option", error, `for "${configId}"="${value}"`);
}
}
+ async requestExtension(method, params) {
+ closedExtensionMethods([method]);
+ if (!this.options.extensionMethods.includes(method)) throw new Error("ACP extension method is not enabled");
+ const bounded = boundedExtensionRecord(params);
+ if (!this.loadedSessionId || bounded.sessionId !== this.loadedSessionId) throw new Error("ACP extension session mismatch");
+ return boundedExtensionRecord(await this.runConnectionRequest(() => this.getConnection().extMethod(method, bounded)));
+ }
async setSessionModel(sessionId, modelId, controlOverride) {
const control = this.resolveModelControl(sessionId, controlOverride);
if (!control) throw new RequestedModelUnsupportedError(`Cannot set model "${modelId}": the ACP session did not advertise a model config option or legacy session/set_model support.`, "missing-capability");
- const resolvedModelId = resolveRequestedModelId({
- requestedModel: modelId,
- models: controlOverride?.availableModels ? { availableModels: controlOverride.availableModels } : void 0,
- agentCommand: this.options.agentCommand
- });
- return control.kind === "config_option" ? await this.setSessionModelThroughConfig(sessionId, resolvedModelId, control.configId) : await this.setSessionModelThroughLegacyMethod(sessionId, resolvedModelId);
+ return control.kind === "config_option" ? await this.setSessionModelThroughConfig(sessionId, modelId, control.configId) : await this.setSessionModelThroughLegacyMethod(sessionId, modelId);
}
async setSessionModelThroughConfig(sessionId, modelId, configId) {
const connection = this.getConnection();
@@ -4608,7 +4829,9 @@
configId,
value: modelId
}));
- this.rememberSessionModels(sessionId, modelStateFromConfigOptions(response.configOptions));
+ const selected = modelStateFromConfigOptions(response.configOptions);
+ if (selected?.currentModelId !== modelId) throw Object.assign(new Error(`ACP effective model mismatch: requested ${modelId}, received ${selected?.currentModelId ?? "unverified"}`), { code: "ACPX_EFFECTIVE_MODEL_MISMATCH" });
+ this.rememberSessionModels(sessionId, selected);
return response;
} catch (error) {
return this.throwSessionModelError("session/set_config_option", modelId, error);
@@ -4864,7 +5087,7 @@
}
selectAuthMethod(methods) {
for (const method of methods) {
- const envCredential = readEnvCredential(method.id);
+ const envCredential = readEnvCredential(method.id, resolveAgentEnvironment(this.options.spawnEnvironment));
if (envCredential) return {
methodId: method.id,
credential: envCredential,
@@ -4890,7 +5113,7 @@
}
readAgentSpecificEnvCredential(methodId) {
if (!this.isGrokBuildAcpCommand() || methodId !== "xai.api_key") return;
- const value = process.env.XAI_API_KEY;
+ const value = (resolveAgentEnvironment(this.options.spawnEnvironment)).XAI_API_KEY;
return typeof value === "string" && value.trim().length > 0 ? value : void 0;
}
selectAgentManagedAuthMethod(methodId) {
@@ -4915,9 +5138,9 @@
await connection.authenticate({ methodId: selected.methodId });
this.log(`authenticated with method ${selected.methodId} (${selected.source})`);
}
- async handlePermissionRequest(params) {
+ async handlePermissionRequest(params, responseDelivery) {
if (this.cancellingSessionIds.has(params.sessionId)) return cancelledPermissionResponse();
- const hostResponse = await this.tryHandlePermissionRequestWithHost(params);
+ const hostResponse = await this.tryHandlePermissionRequestWithHost(params, responseDelivery);
if (hostResponse) return hostResponse;
const { response, recorded } = await this.resolvePermissionRequestFromMode(params);
if (!recorded) {
@@ -4926,14 +5149,81 @@
}
return response;
}
- async handleElicitationRequest(request, requestId, requestSignal) {
+ async handleExtensionRequest(method, params, requestId, requestSignal, responseDelivery) {
+ const active = this.activePrompt;
+ const connection = this.connection;
+ if (!this.options.extensionMethods.includes(method) || !this.options.onExtensionRequest) throw new Error("ACP extension method is not enabled");
+ if (!active || !connection || this.closing || (params.sessionId !== void 0 && params.sessionId !== active.sessionId)) throw new Error("ACP extension request has no matching active prompt");
+ if (!(typeof requestId === "string" || typeof requestId === "number" && Number.isFinite(requestId))) throw new Error("ACP extension request id is invalid");
+ const signal = AbortSignal.any([requestSignal, connection.signal, active.elicitationController.signal]);
+ signal.throwIfAborted();
+ let onAbort;
+ const aborted = new Promise((_, reject) => {
+ onAbort = () => reject(new Error("ACP extension request was cancelled"));
+ signal.addEventListener("abort", onAbort, { once: true });
+ });
+ try {
+ const result = await Promise.race([Promise.resolve().then(() => this.options.onExtensionRequest(method, boundedExtensionRecord(params), { requestId, signal, responseDelivery })), aborted]);
+ signal.throwIfAborted();
+ if (this.activePrompt !== active || this.connection !== connection) throw new Error("ACP extension request outlived its prompt");
+ return boundedExtensionRecord(result);
+ } finally {
+ signal.removeEventListener("abort", onAbort);
+ }
+ }
+ handleCursorSubagentNotification(message, owner) {
+ // Keep the SDK's standard union closed. Cursor child transcripts use
+ // their own session IDs and must never be flattened into the parent.
+ const method = "cursor/subagent_update";
+ if (!this.options.extensionMethods.includes(method) || this.options.clientCapabilities?._meta?.subagents !== true) return false;
+ if (message?.method !== "session/update" || "id" in message) return false;
+ const lifecycle = ["subagent_spawned", "subagent_state_update"].includes(message?.params?.update?.sessionUpdate);
+ const active = this.activePrompt;
+ const connection = this.connection;
+ if (this.cursorSubagentStreamOwner !== owner || this.closing || connection?.signal.aborted) return true;
+ if (!active || !connection || active.elicitationController.signal.aborted) return lifecycle || owner.children.has(message?.params?.sessionId);
+ if (owner.active !== active) { owner.active = active; owner.children.clear(); }
+ if (!lifecycle && message?.params?.sessionId === active.sessionId) return false;
+ try {
+ const params = boundedExtensionRecord(message.params);
+ if (lifecycle) {
+ if (params.sessionId !== active.sessionId && !owner.children.has(params.sessionId)) return true;
+ const child = params.update?.subagentSessionId;
+ if (typeof child !== "string" || !child || child.length > 1000 || child === active.sessionId) return true;
+ if (params.update.sessionUpdate === "subagent_spawned") {
+ if (owner.children.has(child) || owner.children.size >= 256) return true;
+ owner.children.set(child, params.sessionId);
+ } else if (owner.children.get(child) !== params.sessionId) return true;
+ this.options.onExtensionNotification?.(method, { sessionId: active.sessionId, parentSessionId: params.sessionId, update: params.update });
+ } else if (owner.children.has(params.sessionId)) {
+ this.options.onExtensionNotification?.(method, { sessionId: active.sessionId, parentSessionId: owner.children.get(params.sessionId), childSessionId: params.sessionId, update: params.update });
+ }
+ } catch {
+ this.log("Cursor child activity notification was rejected");
+ }
+ return true;
+ }
+ handleExtensionNotification(message) {
+ if (!this.options.extensionMethods.includes(message.method) || !this.options.onExtensionNotification || this.closing) return;
+ const active = this.activePrompt;
+ if (!active || active.elicitationController.signal.aborted) return;
+ try {
+ const params = boundedExtensionRecord(message.params);
+ if (params.sessionId !== void 0 && params.sessionId !== active.sessionId) return;
+ this.options.onExtensionNotification(message.method, params);
+ } catch {
+ this.log("ACP extension notification was rejected");
+ }
+ }
+ async handleElicitationRequest(request, requestId, requestSignal, responseDelivery) {
const resolved = this.resolveElicitationOwner(request);
if ("response" in resolved) return resolved.response;
const { active, handler } = resolved.owner;
const signal = AbortSignal.any([requestSignal, active.elicitationController.signal]);
const handlerAttempt = Promise.resolve().then(async () => await handler(request, {
requestId,
- signal
+ signal,
+ responseDelivery
})).then((response) => ({
kind: "response",
response
@@ -4970,7 +5260,7 @@
isElicitationSessionCancelling(sessionId) {
return this.closing || this.cancellingSessionIds.has(sessionId);
}
- async tryHandlePermissionRequestWithHost(params) {
+ async tryHandlePermissionRequestWithHost(params, responseDelivery) {
if (!this.options.onPermissionRequest) return;
const signal = this.cancellationSignalForSession(params.sessionId);
try {
@@ -4978,7 +5268,7 @@
sessionId: params.sessionId,
raw: params,
inferredKind: inferToolKind(params)
- }, { signal });
+ }, { signal, responseDelivery });
return this.hostPermissionDecisionResponse(params, signal, decision);
} catch (error) {
return this.hostPermissionErrorResponse(params, signal, error);
@@ -5028,6 +5318,12 @@
attachAgentLifecycleObservers(child) {
child.once("exit", (exitCode, signal) => {
this.recordAgentExit("process_exit", exitCode, signal);
+ this.options.onAgentExit?.({
+ pid: child.pid ?? this.lastKnownPid,
+ exitCode,
+ signal,
+ exitedAt: isoNow$1()
+ });
});
child.once("close", (exitCode, signal) => {
this.recordAgentExit("process_close", exitCode, signal);
@@ -5100,6 +5396,9 @@
return await this.filesystem.readTextFile(params);
} catch (error) {
this.recordPermissionError(params.sessionId, error);
+ // Keep ENOENT distinct from permission and other filesystem failures.
+ // The ACP peer must receive a resource error before it creates a file.
+ if (error && typeof error === "object" && error.code === "ENOENT") throw RequestError.resourceNotFound(params.path);
throw error;
}
}
@@ -5239,6 +5538,7 @@
record.cumulative_token_usage = conversation.cumulative_token_usage;
record.cumulative_cost = conversation.cumulative_cost;
record.request_token_usage = conversation.request_token_usage;
+ record.cursor_prompt_usage = boundCursorPromptUsage(conversation.cursor_prompt_usage, conversation.messages);
}
//#endregion
//#region src/runtime/engine/session-options.ts
@@ -5701,7 +6001,8 @@
updated_at: conversation.updated_at,
cumulative_token_usage: deepClone(conversation.cumulative_token_usage ?? {}),
cumulative_cost: cloneUsageCost(conversation.cumulative_cost),
- request_token_usage: deepClone(conversation.request_token_usage ?? {})
+ request_token_usage: deepClone(conversation.request_token_usage ?? {}),
+ cursor_prompt_usage: boundCursorPromptUsage(conversation.cursor_prompt_usage, conversation.messages)
};
}
function cloneUsageCost(cost) {
@@ -5771,10 +6072,20 @@
trimConversationForRuntime(conversation);
return acpx;
}
+function parsePaperclipPiReceipt(raw, costField = "cost_usd") {
+ const receipt = asRecord$4(raw);
+ if (receipt?.provenance !== "assistant_message_receipts" && receipt?.provenance !== "assistant_message_and_compaction_receipts") return;
+ const cost = receipt[costField];
+ if (cost !== void 0 && !isNonNegativeFiniteNumber(cost)) return;
+ return { provenance: receipt.provenance, ...cost !== void 0 ? { cost_usd: cost } : {} };
+}
function recordPromptResponseUsage(conversation, usage, promptMessageId, timestamp = isoNow()) {
const tokenUsage = sourceToTokenUsage(usage);
- if (!tokenUsage) return false;
- applyTokenUsage(conversation, tokenUsage, promptMessageId);
+ const receipt = parsePaperclipPiReceipt(asRecord$4(asRecord$4(usage)?._meta)?.paperclipPi, "costUsd");
+ if (!tokenUsage && !receipt) return false;
+ if (tokenUsage) applyTokenUsage(conversation, tokenUsage, promptMessageId);
+ const userId = promptMessageId ?? lastUserMessageId(conversation);
+ if (receipt && userId) conversation.request_token_usage[userId] = { ...conversation.request_token_usage[userId], paperclip_pi: receipt };
updateConversationTimestamp(conversation, timestamp);
trimConversationForRuntime(conversation);
return true;
@@ -6136,6 +6447,7 @@
pendingAgentSessionId = loadState.pendingAgentSessionId;
sessionModels = loadState.sessionModels;
const preferenceReplay = await replayFreshSessionPreferences({
+ reusingLoadedSession,
client,
record,
createdFreshSession,
@@ -6193,14 +6505,16 @@
if (shouldReconnect) process.stderr.write(`[acpx] saved session pid ${record.pid} is dead; respawning agent and attempting session reconnect\n`);
}
async function replayFreshSessionPreferences(params) {
- if (!params.createdFreshSession) return {
+ // A load acknowledgement can reset the model even when the conversation survives.
+ // Retained connections have no new model acknowledgement to replay.
+ if (params.reusingLoadedSession) return {
modelReplay: { replayed: false },
configReplay: { replayed: false }
};
let modelReplay = { replayed: false };
let configReplay = { replayed: false };
try {
- await replayDesiredMode({
+ if (params.createdFreshSession) await replayDesiredMode({
client: params.client,
sessionId: params.sessionId,
desiredModeId: params.desiredModeId,
@@ -6219,7 +6533,7 @@
verbose: params.verbose,
suppressWarnings: params.suppressWarnings
});
- configReplay = await replayDesiredConfigOptions({
+ if (params.createdFreshSession) configReplay = await replayDesiredConfigOptions({
client: params.client,
record: params.record,
sessionId: params.sessionId,
@@ -6435,6 +6749,27 @@
//#region src/runtime/engine/prompt-turn.ts
const SESSION_REPLY_IDLE_MS = 1e3;
const SESSION_REPLY_DRAIN_TIMEOUT_MS = 5e3;
+const TYPED_SESSION_FAILURE_CATEGORIES = /* @__PURE__ */ new Set([
+ "connection",
+ "access",
+ "limit",
+ "service",
+ "request",
+ "unknown"
+]);
+function typedTerminalSessionFailureCategory(response) {
+ if (response === null || typeof response !== "object" || Array.isArray(response)) return null;
+ const meta = response._meta;
+ if (meta === null || typeof meta !== "object" || Array.isArray(meta)) return null;
+ const jetbrains = meta.jetbrains;
+ if (jetbrains === null || typeof jetbrains !== "object" || Array.isArray(jetbrains)) return null;
+ const air = jetbrains.air;
+ if (air === null || typeof air !== "object" || Array.isArray(air)) return null;
+ if (!Number.isInteger(air.version) || air.version < 1) return null;
+ const failure = air.sessionFailure;
+ if (failure === null || typeof failure !== "object" || Array.isArray(failure) || failure.severity !== "error") return null;
+ return typeof failure.category === "string" && TYPED_SESSION_FAILURE_CATEGORIES.has(failure.category) ? failure.category : "unknown";
+}
async function runPromptTurn(params) {
try {
const promptPromise = params.client.prompt(params.sessionId, params.prompt, params.onPromptRequestStarted, params.onElicitation);
@@ -6444,7 +6779,23 @@
idleMs: SESSION_REPLY_IDLE_MS,
timeoutMs: SESSION_REPLY_DRAIN_TIMEOUT_MS
}).catch(() => {});
+ // A failed terminal result may still include authoritative billable usage.
recordPromptResponseUsage(params.conversation, response.usage, params.promptMessageId);
+ recordCursorPromptUsage(params.conversation, response, params.requestId, params.promptMessageId);
+ const terminalFailureCategory = typedTerminalSessionFailureCategory(response);
+ if (terminalFailureCategory !== null) {
+ // Pass complete provider text to the diagnostic callback. Consumers redact
+ // before bounding it; truncating here can split and expose a credential.
+ const failure = response._meta.jetbrains.air.sessionFailure;
+ try {
+ params.onTerminalSessionFailure?.({
+ category: terminalFailureCategory,
+ ...(typeof failure.title === "string" ? { title: failure.title } : {}),
+ ...(typeof failure.details === "string" ? { details: failure.details } : {})
+ });
+ } catch {}
+ throw new Error(`ACP agent reported a terminal ${terminalFailureCategory} failure.`);
+ }
return {
stopReason: response.stopReason,
source: "rpc"
diff --git a/dist/runtime.d.ts b/dist/runtime.d.ts
--- a/dist/runtime.d.ts
+++ b/dist/runtime.d.ts
@@ -1,7 +1,8 @@
import { _ as SessionRecord, a as AcpElicitationHandler, c as AcpElicitationResponse, f as McpServer$1, h as PermissionPolicy, i as AcpElicitationContext, l as AcpPermissionDecision, m as PermissionMode, n as SystemPromptOption, o as AcpElicitationMode, p as NonInteractivePermissionPolicy, s as AcpElicitationRequest, t as SessionAgentOptions, u as AcpPermissionRequest } from "./session-options-DwRDODlr.js";
import { a as RequestedModelUnsupportedErrorCode, i as RequestedModelUnsupportedError, n as REQUESTED_MODEL_UNSUPPORTED_ERROR_CODE, o as RequestedModelUnsupportedReason, r as REQUESTED_MODEL_UNSUPPORTED_REASONS, s as isRequestedModelUnsupportedError, t as AcpClient } from "./client-CxNllqui.js";
+import { ChildProcess, SpawnOptionsWithoutStdio } from "node:child_process";
import fs from "node:fs";
-import { ToolCallContent, ToolCallLocation, ToolKind } from "@agentclientprotocol/sdk";
+import { SessionNotification, ToolCallContent, ToolCallLocation, ToolKind } from "@agentclientprotocol/sdk";
//#region src/agent-registry.d.ts
declare const DEFAULT_AGENT_NAME = "codex";
//#endregion
@@ -44,6 +45,7 @@
data: string;
};
type AcpRuntimeTurnInput = {
+ onTerminalSessionFailure?: (failure: { category: string; title?: string; details?: string }) => void;
handle: AcpRuntimeHandle;
text: string;
attachments?: AcpRuntimeTurnAttachment[];
@@ -141,6 +143,11 @@
kind?: string;
source?: string;
};
+type AcpRuntimePlanEntry = {
+ content: string;
+ priority?: string;
+ status: "pending" | "in_progress" | "completed";
+};
type AcpRuntimeEvent = {
type: "text_delta";
text: string;
@@ -186,6 +193,10 @@
* non-null `input` schema.
*/
availableCommands?: AcpRuntimeAvailableCommand[];
+} | {
+ type: "plan";
+ tag: "plan";
+ entries: AcpRuntimePlanEntry[];
} | {
type: "tool_call";
text: string;
@@ -312,7 +323,51 @@
elicitationModes?: readonly AcpElicitationMode[];
onPermissionRequest?: (req: AcpPermissionRequest, ctx: {
signal: AbortSignal;
+ /** Resolves after this exact response is written to the current provider pipe. */
+ responseDelivery?: Promise<void>;
}) => Promise<AcpPermissionDecision | undefined>;
+ /** Ephemeral initialize capabilities; mandatory base capabilities take precedence. */
+ clientCapabilities?: Record<string, unknown>;
+ /** Closed list of inbound extension request and notification method names. */
+ extensionMethods?: readonly string[];
+ /** Bound to the current connection and active prompt; never persisted. */
+ onExtensionRequest?: (method: string, params: Record<string, unknown>, context: {
+ requestId: string | number;
+ signal: AbortSignal;
+ responseDelivery?: Promise<void>;
+ }) => Promise<Record<string, unknown>>;
+ onExtensionNotification?: (method: string, params: Record<string, unknown>) => void;
+ /** Ephemeral allowlisted environment evaluated immediately before child spawn. */
+ spawnEnvironment?: () => Record<string, string>;
+ /** Host-only spawn cwd; does not change the cwd advertised in session/new. */
+ spawnCwd?: string;
+ /** Host-owned verified executable launch. */
+ spawnAgent?: (input: {
+ command: string;
+ args: readonly string[];
+ options: SpawnOptionsWithoutStdio;
+ }) => ChildProcess;
+ onAgentSpawn?: (meta: { pid: number; startedAt: string }) => Promise<void> | void;
+ onAgentStderr?: (chunk: string) => void;
+ onAgentExit?: (meta: {
+ pid?: number;
+ exitCode: number | null;
+ signal: NodeJS.Signals | null;
+ exitedAt: string;
+ }) => void;
+ /** Synchronous admission guard. Exceptions fail the connection, never an observer. */
+ protocolGuardFactory?: () => (direction: "inbound" | "outbound", message: unknown) => void;
+ onAcpMessage?: (direction: "inbound" | "outbound", message: unknown) => void;
+ /** Ephemeral full ACP notification callback; never written to session records. */
+ onSessionNotification?: (notification: SessionNotification) => void;
+ /** Ephemeral normalized ACP client filesystem/terminal operation callback. */
+ onClientOperation?: (operation: {
+ method: string;
+ status: string;
+ summary: string;
+ details?: string;
+ timestamp: string;
+ }) => void;
};
type AcpFileSessionStoreOptions = {
stateDir: string;
@@ -369,6 +424,7 @@
private createAndSaveRuntimeRecord;
private retainInitializedSessionOwner;
startTurn(input: {
+ onTerminalSessionFailure?: (failure: { category: string; title?: string; details?: string }) => void;
handle: AcpRuntimeHandle;
text: string;
attachments?: AcpRuntimeTurnAttachment[];
@@ -414,6 +470,7 @@
private cleanupRuntimeTurn;
private finalizeRuntimeTurnRecord;
runTurn(input: {
+ onTerminalSessionFailure?: (failure: { category: string; title?: string; details?: string }) => void;
handle: AcpRuntimeHandle;
text: string;
attachments?: AcpRuntimeTurnAttachment[];
diff --git a/dist/runtime.js b/dist/runtime.js
--- a/dist/runtime.js
+++ b/dist/runtime.js
@@ -153,7 +153,7 @@
const contentType = asTrimmedString(contentRaw.type);
if (contentType && contentType !== "text") return null;
const text = asString(contentRaw.text);
- if (text && text.length > 0) return {
+ if (typeof text === "string" && (text.length > 0 || origin.messageId)) return {
type: "text_delta",
text,
stream: params.stream,
@@ -162,7 +162,7 @@
};
}
const text = asString(params.payload.text);
- if (!text || text.length === 0) return null;
+ if (typeof text !== "string" || (text.length === 0 && !origin.messageId)) return null;
return {
type: "text_delta",
text,
@@ -371,7 +371,7 @@
current_mode_update: (payload) => statusUpdateEvent("current_mode_update", payload),
config_option_update: (payload) => statusUpdateEvent("config_option_update", payload),
session_info_update: (payload) => statusUpdateEvent("session_info_update", payload),
- plan: (payload) => statusUpdateEvent("plan", payload),
+ plan: planUpdateEvent,
client_operation: clientOperationEvent,
update: updateStatusEvent,
done: () => null,
@@ -423,6 +423,34 @@
tag: "available_commands_update",
availableCommands
};
+}
+function persistedGoalCapability(goal) {
+ if (!isRecord(goal) || goal.version !== 1 || goal.controlMethod !== "_session/goal" || !Array.isArray(goal.actions)) return;
+ const actions = goal.actions.filter((action) => ["set", "pause", "resume", "clear"].includes(action));
+ if (!actions.includes("set") || !actions.includes("clear")) return;
+ return { version: 1, control_method: goal.controlMethod, actions };
+}
+function restoredGoalCapability(goal) {
+ if (!isRecord(goal)) return;
+ const canonical = persistedGoalCapability({ ...goal, controlMethod: goal.control_method ?? goal.controlMethod });
+ if (!canonical) return;
+ return { version: canonical.version, controlMethod: canonical.control_method, actions: canonical.actions };
+}
+
+
+function planUpdateEvent(payload) {
+ const raw = Array.isArray(payload.entries) ? payload.entries : [];
+ const entries = [];
+ for (const entry of raw) {
+ if (!isRecord(entry)) continue;
+ const content = asTrimmedString(entry.content);
+ if (!content) continue;
+ const status = asTrimmedString(entry.status);
+ if (status !== "pending" && status !== "in_progress" && status !== "completed") continue;
+ const priority = asTrimmedString(entry.priority);
+ entries.push({ content, status, ...priority ? { priority } : {} });
+ }
+ return { type: "plan", tag: "plan", entries };
}
function normalizeUsageCost(value) {
if (!isRecord(value)) return;
@@ -812,7 +840,22 @@
this.deps = deps;
}
createClient(options) {
- return this.deps.clientFactory?.(options) ?? new AcpClient(options);
+ const patchedOptions = {
+ ...options,
+ spawnCwd: this.options.spawnCwd,
+ spawnEnvironment: this.options.spawnEnvironment,
+ spawnAgent: this.options.spawnAgent,
+ onAgentSpawn: this.options.onAgentSpawn,
+ onAgentStderr: this.options.onAgentStderr,
+ onAgentExit: this.options.onAgentExit,
+ onAcpMessage: this.options.onAcpMessage,
+ clientCapabilities: this.options.clientCapabilities,
+ protocolGuardFactory: this.options.protocolGuardFactory,
+ extensionMethods: this.options.extensionMethods,
+ onExtensionRequest: this.options.onExtensionRequest,
+ onExtensionNotification: this.options.onExtensionNotification
+ };
+ return this.deps.clientFactory?.(patchedOptions) ?? new AcpClient(patchedOptions);
}
createSessionOwner(input) {
const owner = {
@@ -821,12 +864,16 @@
pendingSessionUpdates: []
};
input.client.setEventHandlers({
+ onAcpMessage: this.options.onAcpMessage,
onSessionUpdate: (notification) => this.routeOwnedSessionUpdate(owner, notification),
onClientOperation: (operation) => this.routeOwnedClientOperation(owner, operation)
});
return owner;
}
routeOwnedSessionUpdate(owner, notification) {
+ try {
+ this.options.onSessionNotification?.(notification);
+ } catch {}
const active = owner.activeTurn;
if (active) {
const { task, turn } = active;
@@ -833,9 +880,7 @@
if (turn.connectionUpdates) {
turn.connectionUpdates.push(notification);
- this.emitRuntimeTurnEvent(task, {
- jsonrpc: "2.0",
- method: "session/update",
- params: notification
- });
+ // ACP load/resume replays saved history before accepting a prompt.
+ // Retain it in the session projection without publishing it as new
+ // turn output, including tool calls and interruption metadata.
return;
}
@@ -863,6 +910,9 @@
projection.checkpoint.request();
}
routeOwnedClientOperation(owner, operation) {
+ try {
+ this.options.onClientOperation?.(operation);
+ } catch {}
const active = owner.activeTurn;
if (!active) return;
const { task, turn } = active;
@@ -1008,6 +1058,8 @@
record.closedAt = void 0;
record.protocolVersion = owner.client.initializeResult?.protocolVersion;
record.agentCapabilities = owner.client.initializeResult?.agentCapabilities;
+ record.agentGoalCapability = persistedGoalCapability(owner.client.initializeResult?._meta?.goal);
+ this.options.onAgentInitialize?.(owner.client.initializeResult);
applyLifecycleSnapshotToRecord(record, owner.client.getAgentLifecycleSnapshot());
}
async finishBufferedOwnerControl(owner, record) {
@@ -1086,6 +1138,11 @@
record.closed = false;
record.closedAt = void 0;
this.closingActiveRecords.delete(record.acpxRecordId);
+ this.options.onAgentInitialize?.({
+ protocolVersion: record.protocolVersion,
+ agentCapabilities: record.agentCapabilities,
+ _meta: record.agentGoalCapability ? { goal: restoredGoalCapability(record.agentGoalCapability) } : void 0
+ });
await this.options.sessionStore.save(record);
return record;
}
@@ -1149,6 +1206,8 @@
this.closingActiveRecords.delete(record.acpxRecordId);
record.protocolVersion = client.initializeResult?.protocolVersion;
record.agentCapabilities = client.initializeResult?.agentCapabilities;
+ record.agentGoalCapability = persistedGoalCapability(client.initializeResult?._meta?.goal);
+ this.options.onAgentInitialize?.(client.initializeResult);
applyConfigOptionsToRecord(record, session.sessionResult);
const modelApplication = await applyRequestedModelIfAdvertised({
client,
@@ -1301,9 +1360,11 @@
client: turn.client,
sessionId,
prompt: task.promptInput,
+ onTerminalSessionFailure: task.input.onTerminalSessionFailure,
timeoutMs: task.input.timeoutMs ?? this.options.timeoutMs,
conversation: turn.conversation,
promptMessageId: turn.promptMessageId,
+ requestId: task.input.requestId,
onPromptRequestStarted: () => task.promptStarted.resolve(),
onElicitation: task.input.onElicitation
});
@@ -1469,6 +1530,10 @@
setSessionConfigOption: async (configId, value) => {
return (await task.state.activeController.setResolvedSessionConfigOption(configId, value)).response;
},
+ requestExtension: async (method, params) => {
+ await this.waitForRuntimeControlSession(task, turn);
+ return await turn.client.requestExtension(method, params);
+ },
setResolvedSessionConfigOption: async (configId, value) => await this.setRuntimeResolvedSessionConfigOption(task, turn, configId, value)
};
}
@@ -1575,6 +1640,8 @@
reconcileAgentSessionId(turn.record, turn.record.agentSessionId);
turn.record.protocolVersion = turn.client.initializeResult?.protocolVersion;
turn.record.agentCapabilities = turn.client.initializeResult?.agentCapabilities;
+ turn.record.agentGoalCapability = persistedGoalCapability(turn.client.initializeResult?._meta?.goal);
+ this.options.onAgentInitialize?.(turn.client.initializeResult);
turn.record.acpx = turn.acpxState;
applyConversation(turn.record, turn.conversation);
applyLifecycleSnapshotToRecord(turn.record, turn.client.getAgentLifecycleSnapshot());
@@ -1706,6 +1773,17 @@
applyDesiredConfigOptionToRecord(connectedRecord, configId, value);
});
await this.options.sessionStore.save(result.record);
+ }
+ async requestExtension(input) {
+ const recordId = input.handle.acpxRecordId ?? input.handle.sessionKey;
+ return await this.withManagerLock(this.runtimeOperationLocks, recordId, async () => {
+ const record = await this.requireRecord(recordId);
+ const controller = this.activeControllers.get(record.acpxRecordId);
+ if (controller) return await controller.requestExtension(input.method, input.params);
+ return (await this.withRuntimeControlSession(record, input.sessionMode ?? "persistent", async ({ client }) => {
+ return await client.requestExtension(input.method, input.params);
+ })).value;
+ });
}
async cancel(handle) {
await this.activeControllers.get(handle.acpxRecordId ?? handle.sessionKey)?.requestCancelActivePrompt();
@@ -2055,6 +2133,7 @@
const turnPromise = this.getManager().then((manager) => manager.startTurn({
handle,
text: input.text,
+ onTerminalSessionFailure: input.onTerminalSessionFailure,
attachments: input.attachments,
mode: input.mode,
sessionMode: state.mode,
@@ -2087,6 +2166,7 @@
yield* (await this.getManager()).runTurn({
handle,
text: input.text,
+ onTerminalSessionFailure: input.onTerminalSessionFailure,
attachments: input.attachments,
mode: input.mode,
sessionMode: state.mode,
@@ -2118,6 +2198,14 @@
async setConfigOption(input) {
const { handle, state } = this.resolveManagerHandle(input.handle);
await (await this.getManager()).setConfigOption(handle, input.key, input.value, state.mode);
+ }
+ async requestExtension(input) {
+ const { handle, state } = this.resolveManagerHandle(input.handle);
+ return await (await this.getManager()).requestExtension({
+ ...input,
+ handle,
+ sessionMode: input.sessionMode ?? state.mode
+ });
}
async cancel(input) {
const { handle } = this.resolveManagerHandle(input.handle);
diff --git a/dist/session-options-DwRDODlr.d.ts b/dist/session-options-DwRDODlr.d.ts
--- a/dist/session-options-DwRDODlr.d.ts
+++ b/dist/session-options-DwRDODlr.d.ts
@@ -1,4 +1,5 @@
import { AgentCapabilities, AnyMessage, ContentBlock, CreateElicitationRequest, ElicitationContentValue, JsonRpcId, McpServer, McpServer as McpServer$1, RequestPermissionRequest, SessionConfigOption, SessionNotification, SetSessionConfigOptionResponse, ToolKind } from "@agentclientprotocol/sdk";
+import { ChildProcess, SpawnOptionsWithoutStdio } from "node:child_process";
//#region src/prompt-content.d.ts
type PromptInput = ContentBlock[];
//#endregion
@@ -27,6 +28,8 @@
requestId: JsonRpcId;
/** Aborts with the request itself or its owning prompt turn/session. */
signal: AbortSignal;
+ /** Await only after returning the selected response from this callback. */
+ responseDelivery?: Promise<void>;
};
type AcpElicitationResponseMeta = {
_meta?: Record<string, unknown> | null;
@@ -116,6 +119,37 @@
};
env?: Record<string, string>;
};
+ /** Ephemeral initialize capabilities; mandatory base capabilities take precedence. */
+ clientCapabilities?: Record<string, unknown>;
+ /** Closed list of inbound extension request and notification method names. */
+ extensionMethods?: readonly string[];
+ /** Bound to the current connection and active prompt; never persisted. */
+ onExtensionRequest?: (method: string, params: Record<string, unknown>, context: {
+ requestId: string | number;
+ signal: AbortSignal;
+ responseDelivery?: Promise<void>;
+ }) => Promise<Record<string, unknown>>;
+ onExtensionNotification?: (method: string, params: Record<string, unknown>) => void;
+ /** Ephemeral child environment factory; its return value is never persisted. */
+ spawnEnvironment?: () => Record<string, string>;
+ /** Host-only child cwd, separate from the cwd advertised to ACP. */
+ spawnCwd?: string;
+ /** Host-owned verified executable launch. */
+ spawnAgent?: (input: {
+ command: string;
+ args: readonly string[];
+ options: SpawnOptionsWithoutStdio;
+ }) => ChildProcess;
+ onAgentSpawn?: (meta: { pid: number; startedAt: string }) => Promise<void> | void;
+ onAgentStderr?: (chunk: string) => void;
+ onAgentExit?: (meta: {
+ pid?: number;
+ exitCode: number | null;
+ signal: NodeJS.Signals | null;
+ exitedAt: string;
+ }) => void;
+ /** Synchronous admission guard. Exceptions fail the connection, never an observer. */
+ protocolGuardFactory?: () => (direction: "inbound" | "outbound", message: unknown) => void;
onAcpMessage?: (direction: AcpMessageDirection, message: AcpJsonRpcMessage) => void;
onAcpOutputMessage?: (direction: AcpMessageDirection, message: AcpJsonRpcMessage) => void;
onSessionUpdate?: (notification: SessionNotification) => void;
@@ -123,6 +157,8 @@
onPermissionEscalation?: (event: PermissionEscalationEvent) => void;
onPermissionRequest?: (req: AcpPermissionRequest, ctx: {
signal: AbortSignal;
+ /** Resolves after this exact response is written to the current provider pipe. */
+ responseDelivery?: Promise<void>;
}) => Promise<AcpPermissionDecision | undefined>;
};
declare const SESSION_RECORD_SCHEMA: "acpx.session.v1";
@@ -263,12 +299,16 @@
lastAgentDisconnectReason?: string;
protocolVersion?: number;
agentCapabilities?: AgentCapabilities;
+ agentGoalCapability?: Record<string, unknown>;
title?: string | null;
messages: SessionMessage[];
updated_at: string;
cumulative_token_usage: SessionTokenUsage;
cumulative_cost?: SessionUsageCost;
- request_token_usage: Record<string, SessionTokenUsage>;
+ request_token_usage: Record<string, SessionTokenUsage & {
+ paperclip_pi?: { provenance: "assistant_message_receipts" | "assistant_message_and_compaction_receipts"; cost_usd?: number };
+ }>;
+ cursor_prompt_usage?: { request_id: string; prompt_message_id: string; receipt: Record<string, unknown> };
acpx?: SessionAcpxState;
importedFrom?: SessionImportedFrom;
};