mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-07 16:11:46 +02:00
1189 lines
57 KiB
Diff
1189 lines
57 KiB
Diff
diff --git a/dist/client-CxNllqui.d.ts b/dist/client-CxNllqui.d.ts
|
|
index 5e2113ac0..45aff03c3 100644
|
|
--- a/dist/client-CxNllqui.d.ts
|
|
+++ b/dist/client-CxNllqui.d.ts
|
|
@@ -120,6 +120,8 @@ declare class AcpClient {
|
|
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 @@ declare class AcpClient {
|
|
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
|
|
index d454fd7c5..6eba74345 100644
|
|
--- a/dist/live-checkpoint-BSIrfgVo.js
|
|
+++ b/dist/live-checkpoint-BSIrfgVo.js
|
|
@@ -1068,12 +1068,14 @@ function serializeSessionRecordForDisk(record) {
|
|
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 @@ function parseRequestTokenUsage(raw) {
|
|
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 isAgentMessage$1(raw) {
|
|
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 @@ function parseConversationRecord(record) {
|
|
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 @@ function parseSessionRecord(raw) {
|
|
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 @@ const ZED_TAG_KEYS = /* @__PURE__ */ new Set([
|
|
"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 @@ function authEnvKey(methodId) {
|
|
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 @@ function promotePrefixedAuthEnvironment(env) {
|
|
}
|
|
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 @@ function resolveConfiguredAuthCredential(methodId, authCredentials) {
|
|
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",
|
|
@@ -3961,10 +4068,13 @@ function resolveClientCapabilities(params) {
|
|
...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 +4177,27 @@ function createNdJsonMessageStream(agentCommand, output, input) {
|
|
} })
|
|
};
|
|
}
|
|
+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 +4236,9 @@ var AcpClient = class {
|
|
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 +4349,21 @@ var AcpClient = class {
|
|
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 +4371,7 @@ var AcpClient = class {
|
|
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 +4399,12 @@ var AcpClient = class {
|
|
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 +4430,17 @@ var AcpClient = class {
|
|
}
|
|
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 +4451,10 @@ var AcpClient = class {
|
|
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 +4470,11 @@ var AcpClient = class {
|
|
}).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 +4496,12 @@ var AcpClient = class {
|
|
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 +4521,49 @@ var AcpClient = class {
|
|
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 +4573,12 @@ var AcpClient = class {
|
|
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 +4589,7 @@ var AcpClient = class {
|
|
}
|
|
};
|
|
return {
|
|
+ withResponseDelivery, abortResponseDeliveries,
|
|
readable: new ReadableStream({ async start(controller) {
|
|
const reader = base.readable.getReader();
|
|
try {
|
|
@@ -4386,16 +4597,27 @@ var AcpClient = class {
|
|
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 +4625,20 @@ var AcpClient = class {
|
|
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 +4654,10 @@ var AcpClient = class {
|
|
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) {
|
|
@@ -4864,7 +5099,7 @@ var AcpClient = class {
|
|
}
|
|
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 +5125,7 @@ var AcpClient = class {
|
|
}
|
|
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 +5150,9 @@ var AcpClient = class {
|
|
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 +5161,81 @@ var AcpClient = class {
|
|
}
|
|
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 +5272,7 @@ var AcpClient = class {
|
|
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 +5280,7 @@ var AcpClient = class {
|
|
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 +5330,12 @@ var AcpClient = class {
|
|
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);
|
|
@@ -5239,6 +5547,7 @@ function applyConversation(record, conversation) {
|
|
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 +6010,8 @@ function cloneSessionConversation(conversation) {
|
|
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 +6081,20 @@ function recordSessionUpdate(conversation, state, notification, timestamp = isoN
|
|
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;
|
|
@@ -6435,6 +6755,27 @@ async function withConnectedSession(options) {
|
|
//#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 +6785,23 @@ async function runPromptTurn(params) {
|
|
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
|
|
index e8102acb0..377c5d838 100644
|
|
--- 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 @@ type AcpRuntimeTurnAttachment = {
|
|
data: string;
|
|
};
|
|
type AcpRuntimeTurnInput = {
|
|
+ onTerminalSessionFailure?: (failure: { category: string; title?: string; details?: string }) => void;
|
|
handle: AcpRuntimeHandle;
|
|
text: string;
|
|
attachments?: AcpRuntimeTurnAttachment[];
|
|
@@ -141,6 +143,11 @@ type AcpTextDeltaOriginMeta = {
|
|
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 @@ type AcpRuntimeEvent = {
|
|
* non-null `input` schema.
|
|
*/
|
|
availableCommands?: AcpRuntimeAvailableCommand[];
|
|
+} | {
|
|
+ type: "plan";
|
|
+ tag: "plan";
|
|
+ entries: AcpRuntimePlanEntry[];
|
|
} | {
|
|
type: "tool_call";
|
|
text: string;
|
|
@@ -312,7 +323,51 @@ type AcpRuntimeOptions = {
|
|
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 @@ declare class AcpRuntimeManager {
|
|
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 @@ declare class AcpRuntimeManager {
|
|
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
|
|
index a1f4a70a0..a17fce3b3 100644
|
|
--- a/dist/runtime.js
|
|
+++ b/dist/runtime.js
|
|
@@ -153,7 +153,7 @@ function resolveTextChunk(params) {
|
|
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 @@ function resolveTextChunk(params) {
|
|
};
|
|
}
|
|
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 @@ const PROMPT_EVENT_PARSERS = {
|
|
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,
|
|
@@ -424,6 +424,34 @@ function availableCommandsUpdateEvent(payload) {
|
|
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;
|
|
const amount = asOptionalFiniteNumber(value.amount);
|
|
@@ -812,7 +840,22 @@ var AcpRuntimeManager = class {
|
|
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 @@ var AcpRuntimeManager = class {
|
|
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;
|
|
@@ -863,6 +910,9 @@ var AcpRuntimeManager = class {
|
|
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 @@ var AcpRuntimeManager = class {
|
|
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 @@ var AcpRuntimeManager = class {
|
|
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 @@ var AcpRuntimeManager = class {
|
|
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 @@ var AcpRuntimeManager = class {
|
|
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 @@ var AcpRuntimeManager = class {
|
|
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 @@ var AcpRuntimeManager = class {
|
|
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());
|
|
@@ -1707,6 +1774,17 @@ var AcpRuntimeManager = class {
|
|
});
|
|
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 @@ var AcpxRuntime = class {
|
|
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 @@ var AcpxRuntime = class {
|
|
yield* (await this.getManager()).runTurn({
|
|
handle,
|
|
text: input.text,
|
|
+ onTerminalSessionFailure: input.onTerminalSessionFailure,
|
|
attachments: input.attachments,
|
|
mode: input.mode,
|
|
sessionMode: state.mode,
|
|
@@ -2119,6 +2199,14 @@ var AcpxRuntime = class {
|
|
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);
|
|
await (await this.getManager()).cancel(handle);
|
|
diff --git a/dist/session-options-DwRDODlr.d.ts b/dist/session-options-DwRDODlr.d.ts
|
|
index c3da16452..40ae05653 100644
|
|
--- 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 @@ type AcpElicitationContext = {
|
|
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 @@ type AcpClientOptions = {
|
|
};
|
|
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 @@ type AcpClientOptions = {
|
|
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 @@ type SessionRecord = {
|
|
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;
|
|
};
|