mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-10 12:07:09 +02:00
## Thinking Path > - Paperclip is the open source app people use to manage AI agents for work. > - Users connect accounts and choose an agent harness and model. > - The runtime change in #14970 supports custom providers on those connections. > - Normal setup must stay simple while advanced users can choose a compatible gateway. > - Shared connector rows and access controls keep these choices consistent. > - This pull request refines the agent setup UI and adds review stories and repeatable browser qualification. > - The qualification checks real tools and downloaded outputs, not only a successful run status. ## Linked Issues or Issue Description Refs #14970, #37, #13083, #14104, #14565, #12692. The core implementation in #14970 is merged. This branch incorporates its squash commit and targets `master`. Both PRs contain our implementation. #14016 is a reference only and is not a dependency. This PR has 96 changed files. ## What Changed - Complete model-provider connector presentation beside other connectors. Each row uses the existing Connect action and connection list. Tags are stored without category UI. The base PR includes the provider forms and routes. - Show persistent Subscription, API Key, and Advanced choices. Label Advanced as Custom Gateway. Reuse provider logos, connection lists, and permissions controls. Default access to the organization and all agents when permitted; keep narrowing controls under Advanced. - Keep Configure reachable before subscription sign-in, so users can select a supported environment when the default cannot sign in. Testing and saving still require a connection. Show the execution environment in Configure. Preserve the confirmed Connect choice. Editing a method, credential, saved account, or advanced choice requires that current choice to connect before testing or saving. Use matching model and thinking-effort dropdowns and retain connection icons in selected values. - Preserve the new harness model default when switching an existing OpenCode agent to Codex or Claude, and resolve user-selected model names with the effective harness. - Load popular OpenRouter models through the shared connection-model discovery path. Keep explicit model lists and manual model entry available. - Group onboarding, connection setup, agent runtime, management, recovery, and production-component stories under AI Connections / Provider routing. - Add an explicit-only provider-connections browser suite for managed local or existing local/staging targets. Use private browser profiles and credential handoffs. Support human-assisted subscription sign-in without sharing passwords or tokens in reports. - Verify persisted connection identity, runtime probes, tool execution, exact artifact bytes, completion, and context-dependent follow-up. Retain source/model provenance, cost bounds, closed error diagnostics, original failures, and cleanup evidence. - Add Gemini startup-model and skill-root fixes, Grok private-history detection, ACP filesystem regression fixtures, selected-workspace handling for local Hermes, and artifact-helper workspace fallback. - Keep managed Grok runtime homes disposable. Remove host-side transcript retention/restoration because private file modes do not isolate same-user agent processes. Ignore earlier development archives and use a fresh task handoff when history is unavailable. Verify the absence of restored transcripts with a separate same-user process. - Capture stopped-run diagnostics before deleting an attached-company fixture agent. Track creation and owned sign-in receipts; revoke only this attempt's accounts and never adopt a concurrent campaign's newly created account. Preserve failure signals and final status through cleanup. - Require the requested environment in the saved agent and every run, including follow-ups. Reject a forced incompatible target. Keep one cancellation state through startup, every cell, reporting, and teardown for SIGINT, SIGTERM, and SIGHUP. Stop further paid cells after interruption. Document qualification limits. ## Verification - Current head `b3bb3e94d577d43d9965a6b9daba039f599b2e49` includes master `d9f600043`. The security fix in `a758fde31` passes full workspace typecheck, production build, and 119 connection/Grok regressions. The unchanged UI passes all 126 configuration/model-discovery tests and token gates. The final published-guide correction passes Grok adapter typecheck. Earlier head `eebd8225c` passed the complete deterministic runner suite (1,404 Vitest tests and 128 Node tests) and all CI jobs. Current-head CI run `37520147514` passed all 47 jobs, including the full sharded Vitest and browser matrix, production build, and canary dry run. All 55 checks completed: 53 successes and two expected skips. The current-head security scan passed, Greptile is 5/5, and no review threads remain open. - A separate same-user process reproduced reading a restored Grok transcript before the security fix. The regression now finds no transcript. Existing fresh-session fallback and ordinary session metadata behavior pass. - The final account-choice and cleanup fixes pass 85 setup tests and 26 qualification-harness tests. Regressions verify that editing a connection invalidates confirmation, Configure remains reachable before sign-in, diagnostics are captured before fixture deletion, and concurrent campaigns cannot adopt or revoke each other's accounts. UI and E2E typechecks pass. - The Storybook build and actual Chromium production-component stories passed during this change. Review the neighboring AI Connections / Provider routing stories, regular connector rows, three connection modes, model discovery, and the single execution-environment control in Configure. - Cancellation smoke verified authenticated cleanup before browser close for SIGINT, SIGTERM, and SIGHUP. Regressions cover interruption during startup and reporting, missing-file ACP resource errors, and preserved permission denials. Both ACP runtime versions and 54 ACPX/Grok regressions passed. The deterministic connection-intent browser suite passed two tests. - Historical local qualification retained 43 passing API/gateway cells out of 46, with downloaded outputs and follow-up receipts. These attempts span earlier builds; they do not qualify this exact commit or staging. Subscription combinations, Gemini overloads, and the unresolved follow-up failure remain recorded rather than counted as passing. - Use `pnpm test:e2e:runner -- --list --suite provider-connections` to inspect the matrix. Follow `tests/runner-e2e/PROVIDER-CONNECTIONS.md` for credentials, target URL, sign-in assistance, budget, evidence, and cleanup. Paid live tests remain opt-in. ## Risks - The core implementation in #14970 is merged. This PR adds no database migration of its own. - Subscription login needs an interactive provider session. Dedicated accounts and staging qualification remain follow-up work; this PR does not certify every login combination for production. - Managed Grok transcript resume is deferred until provider history has an OS isolation or authorized broker solution. Follow-ups start fresh with Paperclip task context; earlier live Grok results do not qualify this behavior. - Gemini CLI 0.58.0 has an upstream ACP new-file error conversion defect. Live overloads and one unresolved follow-up timeout remain recorded. The stock CLI is unchanged, and those cases are not marked as passing. - Real-provider tests spend credits and use private credential/evidence directories. The launcher requires explicit selection and checks target ownership. It must not attach to a developer's database by accident. - OpenClaw Gateway, Hermes Gateway, Claude Managed, AWS AgentCore, Process, HTTP, and legacy ACPX local remain outside custom provider setup. ## Model Used OpenAI GPT-6 through Codex, with reasoning, repository tools, code execution, and browser testing. The exact deployment model ID and context window size were not exposed in this session. ## Checklist - [x] I have included a thinking path that traces from project context to this change - [x] I have specified the model used (with version and capability details) - [x] I have checked ROADMAP.md and confirmed this PR does not duplicate planned core work - [x] I have searched GitHub for duplicate or related PRs and linked them above - [x] I have either (a) linked existing issues with `Fixes: #` / `Closes #` / `Refs #` OR (b) described the issue in-PR following the relevant issue template - [x] I have not referenced internal/instance-local Paperclip issues or links (only public GitHub `#NNN` / `github.com/paperclipai/paperclip` URLs) - [x] My branch name describes the change (e.g. `docs/...`, `fix/...`) and contains no internal Paperclip ticket id or instance-derived details - [x] I have run tests locally and they pass - [x] I have added or updated tests where applicable - [x] I have updated relevant documentation to reflect my changes - [x] I have considered and documented any risks above - [x] All Paperclip CI gates are green - [x] Greptile is 5/5 with no open P2s, recommendations, or follow-ups - [x] I will address all Greptile and reviewer comments before requesting merge --------- Co-authored-by: Paperclip <noreply@paperclip.ing>
971 lines
45 KiB
Diff
971 lines
45 KiB
Diff
diff --git a/dist/client-CxNllqui.d.ts b/dist/client-CxNllqui.d.ts
|
|
index 5e2113ac0a1a92c99322cf01e5c106760a0b0dbf..b7b5151f3da45b5e0d12aea55e9f3b071ab54a8d 100644
|
|
--- a/dist/client-CxNllqui.d.ts
|
|
+++ b/dist/client-CxNllqui.d.ts
|
|
@@ -135,6 +135,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
|
|
--- 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,6 +1068,7 @@
|
|
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,
|
|
@@ -1239,7 +1240,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;
|
|
}
|
|
@@ -1541,6 +1543,7 @@
|
|
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,
|
|
@@ -1661,7 +1664,7 @@
|
|
"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([
|
|
"agent_capabilities",
|
|
"messages.Agent.content.ToolUse.input",
|
|
@@ -3104,10 +3107,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 +3138,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 +3194,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",
|
|
@@ -3961,10 +3977,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 +4086,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 +4145,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 +4258,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 +4280,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 +4308,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 +4339,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 +4360,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 +4379,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 +4405,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 +4430,36 @@
|
|
throw normalizedError;
|
|
}
|
|
createTappedStream(base) {
|
|
+ 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 +4469,9 @@
|
|
const shouldSuppressInboundReplaySessionUpdate = (message) => {
|
|
return this.suppressReplaySessionUpdateMessages && isSessionUpdateNotification(message);
|
|
};
|
|
+ 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 +4482,7 @@
|
|
}
|
|
};
|
|
return {
|
|
+ withResponseDelivery, abortResponseDeliveries,
|
|
readable: new ReadableStream({ async start(controller) {
|
|
const reader = base.readable.getReader();
|
|
try {
|
|
@@ -4387,10 +4491,16 @@
|
|
if (done) break;
|
|
if (!value) continue;
|
|
observeInbound(value);
|
|
- if (isExtensionNotification(value)) continue;
|
|
+ if (isExtensionNotification(value)) {
|
|
+ onExtensionNotification(value);
|
|
+ continue;
|
|
+ }
|
|
controller.enqueue(value);
|
|
}
|
|
} finally {
|
|
+ abortResponseDeliveries();
|
|
+ for (const receipt of deliveries.values()) receipt.cleanup();
|
|
+ deliveries.clear();
|
|
reader.releaseLock();
|
|
controller.close();
|
|
}
|
|
@@ -4403,10 +4513,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 +4542,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) {
|
|
@@ -4864,7 +4987,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 +5013,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 +5038,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 +5049,49 @@
|
|
}
|
|
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);
|
|
+ }
|
|
+ }
|
|
+ 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 +5128,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 +5136,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 +5186,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 +5264,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;
|
|
}
|
|
}
|
|
@@ -5771,10 +5938,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;
|
|
@@ -6435,6 +6612,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 +6642,22 @@
|
|
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);
|
|
+ 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 e8102acb03c4c38830ad5ec22f356125eb0423b7..d3a9ba6266a29926ed53302410534415d1fba4bf 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,49 @@ 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;
|
|
+ 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 +422,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 +468,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 a1f4a70a003792c6eacf68b6b038f37bfec1db53..fe2484de6978fd16bfdd55ce69066fd43b1a0f77 100644
|
|
--- a/dist/runtime.js
|
|
+++ b/dist/runtime.js
|
|
@@ -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,21 @@ 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,
|
|
+ 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 +863,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 +909,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 +1057,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 +1137,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 +1205,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,6 +1359,7 @@ 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,
|
|
@@ -1469,6 +1528,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 +1638,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 +1772,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 +2131,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 +2164,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 +2197,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 c3da1645235bbea22de3f8484149051cd7dca56b..5ccc8dda101e7c90f61483c86fd2731855ce2e51 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,35 @@ 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;
|
|
onAcpMessage?: (direction: AcpMessageDirection, message: AcpJsonRpcMessage) => void;
|
|
onAcpOutputMessage?: (direction: AcpMessageDirection, message: AcpJsonRpcMessage) => void;
|
|
onSessionUpdate?: (notification: SessionNotification) => void;
|
|
@@ -123,6 +155,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 +297,15 @@ 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 };
|
|
+ }>;
|
|
acpx?: SessionAcpxState;
|
|
importedFrom?: SessionImportedFrom;
|
|
};
|