mirror of
https://github.com/paperclipai/paperclip.git
synced 2026-10-10 12:07:09 +02:00
## Thinking Path > - Paperclip manages AI agents and their work. > - The experimental Runner owns provider processes and durable sessions. > - Pi needs working task execution and human controls. > - The five-PR stack must preserve changes already on master. > - Each layer now carries the complete integrated source for a safe sequential fallback. > - This PR belongs to native GitHub stack #15602, ending at #14956. ## Linked Issues or Issue Description Refs #14436, #14631, #14743 and #14956. Ship Pi 1.0 through the experimental Paperclip Runner. The five PRs are #14921, #14922, #14923, #14924 and #14956. The user authorized the complete merge after checks pass. Existing `pi_local` execution is unchanged. Accounting and wider provider/platform qualification remain deferred. ## What Changed - Recover missing final replies after workspace finalization changes owners, using accepted-turn evidence without rerunning work or granting external-chat publication. - Preserve the admitted Pi instruction root across warm runs, while retaining changed-root rejection. - Give Pi a bounded 15-second default shutdown grace so stop, drain acknowledgement and durable suspension can complete. Explicit deadlines and other providers retain their existing behavior. - Integrate the Pi 1.0 runtime and master contracts. - Use Pi profile 22. Preserve explicit caller-selected models and exact native thinking levels. Keep Pi's wrapper, helper, extension and question/control behavior unchanged from the qualified profile-19 runtime. - Preserve master's Dot lifecycle and consent fields, configured task environment, status guards and current Codex/Claude dependency versions. Cursor stays qualified. Copilot stays pending; profile 17 binds the changed shared protocol validation sources. - Exclude general AWS IAM credentials from Pi static/custom provider bindings and selected task projections; preserve the provider-scoped Bedrock bearer key. Profile 21 is retained as historical provenance. Rust and cloud install probes use the current declaration. - Patch bundled brace-expansion 5.0.9 to the exact official 5.0.12 payload. Pin the patch and complete runtime closures. Include the patch in normal installed setup tooling. Keep the upstream Pi shrinkwrap as provenance and permit only this exact security correction. - Include current attestation files in the Docker build context. Keep the repository lockfile unchanged from master. CI and private image builds resolve manifest changes before their frozen installation. ## Verification - Full local `pnpm -r typecheck` passes, including Runner Rust, server and UI. Focused integration checks pass: 194 Runner admission/environment tests, 63 profile/credential tests with one expected skip, 152 Dot/UI configuration tests, and Pi transcript/notice tests. - Full local `pnpm build` passes on the final source. - Fresh final-source checks pass: all 698 Rust workspace tests (32 binaries), 156 credential/profile/controller tests with one expected skip, Runner TypeScript typecheck, and 20 package/setup/sandbox tests. - The profile-21 Pi materializer passes on the native host with the official pinned Node 24.21.0 and its npm. It verifies all 150 locked packages, the patched dependency and the exact closure. Setup/package bundle tests and UI token gates pass. - The old hashes were reproduced for all three supported targets before calculating the patched graph. New closure hashes are darwin-arm64 `282022db10150c6632b3444df421342e7d534bdf5d5fb1097a2e79d0625a2bcf`, darwin-x64 `64e251e19009f755c0b04f73ce2138246faab71a961b0f13d75ebfcc34bef12e`, and linux-x64 `713b1fdff42fb56a1518bdc084f181d70bee8ebadc3e4b1d76321ed9108c8410`. Independent native platform execution is separate from graph identity reproduction. - Historical cloud qualification remains unchanged: all seven core cases pass on shipping source `10dc43c9ec65d88c2f782d62afb296d09494f215`, harness `1a4408a48cfb5a1f094a311141c257c92cd7a893`, image `sha256:5b3a775b383591bda1b0c1889e509acc70ce7f37c53f09733c81d59037f02280`, and accepted Sonnet 4.6/low fixture. All 215 canonical files and all seven cleanup checks pass independent verification. These are profile-19 results and are not relabeled as fresh profile-22 runs. - Current Pi digest: `sha256:e92078bee3c23bec4100aa589013a44613d054cd686826534025d8019e9f39a9`. [The readiness plan](https://github.com/paperclipai/paperclip/blob/codex/pi-production-readiness/doc/plans/2026-10-02-pi-production-readiness.md) preserves campaign and failed-attempt provenance. - Merge only after every PR's current-head CI and fresh review pass. Linux CI covers the full suites, build and browser tests. The local embedded Postgres API-authority suite cannot start on this macOS/Node 26 host, so Linux CI must confirm that suite. ### Fresh profile-22 core qualification — 2026-10-08 All seven accepted core cases pass canonically on Pi profile 22, with `openrouter/anthropic/claude-sonnet-4.6` and native-confirmed low thinking. This model is a fixture; production accepts the caller's explicit Pi provider/model. Runtime/install source: `3241a992f2a7703e59e97ed0fd3e5d6405de4401`. Frozen accepted harness: `1a4408a48cfb5a1f094a311141c257c92cd7a893`. Immutable cloud image: `ghcr.io/paperclipai/paperclip-daytona-runner@sha256:506f22db7edd78f37c0c40bec1cc084af1850455026dbf467194bfbb8fcef141`. Pi digest: `sha256:e92078bee3c23bec4100aa589013a44613d054cd686826534025d8019e9f39a9`. [Hosted Linux image and clean-install verification](https://github.com/paperclipai/paperclip/actions/runs/37868328023) passes, including all 20 source-bound archives, normal CLI/Pi setup, companion import and the production pack reader. This exact installation source includes the latest master integration and the corrected Pi warm instruction-root fence. Full local typecheck/build and current-head hosted CI verify the final stack. All 13 focused real-root regressions pass. The full local executor suite passed 662 tests; 15 database tests could not start the Mac embedded PostgreSQL service. Hosted Linux CI passes the full required verification and E2E checks. These fresh results keep their own source identity; profile-19 results remain historical. | Core path | Canonical campaign | Retained archive SHA-256 | | --- | --- | --- | | File edit, validation, download and Done | `pi-core22-replyfix-0-1791511228` | 23 files; `a473e8603a3dd4737863291f8d3d1e392391f0b16d433c3e0e0e9d8baf7a97b0` | | Pending question and controller restart | `pi-core22-replyfix-1-1791511376` | 33 files; `6b829c4eb74e1f32a89c692a4ae7130dbfc1c6d3cf13915effe2103d9e242c8e` | | Three-turn session/process/workspace continuity | `pi-core22-replyfix-2-1791511587` | 23 files; `7a87021f8f9a3fdd3c58bb4467f8d82c635e3ea4795d6e75f144d9aa14818df8` | | Four typed questions and browser reconnects | `pi-core22-replyfix-3-1791511881` | 42 files; `9e31755252be1f4f9cb0626c984c142d4d1ae5f5bee3a7af08444db8d12c280a` | | Plan approval and completion | `pi-core22-replyfix-4-1791512031` | 22 files; `a0383ce1aab38e7b5a25ce0e9dd3bebea5c037ebd96ae6b29dae19015da2ae2c` | | Same-turn steering and permission denial | `pi-core22-replyfix-5-1791512261` | 39 files; `c929b8c7070f0b66aedc17e65ca46e6beab1e363926ac9f7e2a75fb250f05949` | | Stop during pending permission | `pi-core22-replyfix-6-1791512390` | 33 files; `7f0a58ae0f4d5bfc76149435f4e322537089c5bd16e7ffe9b5ad71f10a621a07` | All 215 canonical files (28714587 bytes) are independently hash-verified. All seven cleanup grades pass, with no owned runtime process or temporary root after each case. Automatic retries are zero. The owned cloud host stopped normally after retention. The prior profile-22 warm attempt remains failed and separately retained: archive SHA-256 `1e54eba5ec72b50cee1534b23d1d1d4f21a090006b8a64501ba70db972abfde5`. Its original canonical classification is preserved. Diagnosis reproduced a product bug comparing an agent-files root against an unset checkpoint-only field. The fix stores the admitted physical root separately from the adopted per-run collection capability. The real-root regression fails before the fix and passes afterward, including rejection of a changed physical root. Fixture, grader, model and all seven accepted case IDs are unchanged; this fresh campaign tests final-reply publication after file registration first. The intermediate restart attempt also remains failed and retained: archive SHA-256 `5dcaefdf1d17cf4cd54fd4cf810f45e736667392339b8ce7caf08bb4e225277f`. Its original canonical classification is preserved. Pi resumed, wrote the verified answer and completed its task; exact runner suspension was proven, but idle stop consumed about 5.2s and left under 3s for the drain acknowledgement. The Pi-only default shutdown grace is now 15s, preserving a full 5s drain round trip and a finite suspension reserve. Explicit caller deadlines, other provider defaults, literal drain receipts and exact suspension identity checks remain unchanged. The timing regression fails before this correction and passes afterward; all 18 focused settlement tests and Runner typecheck pass. The final-source file attempt is also preserved as failed (`candidate_failure`), archive SHA-256 `db6767b6773ea618997927ac77bdb005a5ac81492c7b9c0ffbc900449f829bc9`. Native edit, validation, exact downloadable artifact and Done/succeeded all passed, and the exact final reply was durably recorded. A workspace recovery owner completed before the live heartbeat reached presentation, leaving that reply absent from task chat. Recovery now materializes only a completed final reply from the accepted turn of an ordinary internal Done task, preserving issue/run/contract binding, suppression, external-chat authorization and same-run deduplication. The database regression covers the generated file-preparation receipt, suppression, unapproved external continuation and replay. Server typecheck and all 49 response-selection tests pass; hosted Linux verifies the database regression because embedded PostgreSQL cannot start on this Mac. The delayed-final-answer database regression passes on [the final root-source Linux server shard](https://github.com/paperclipai/paperclip/actions/runs/37868262553/job/113628594152), alongside 1,108 passing tests. The first root Runner shard had one unchanged durable-resume test exceed its 5-second timeout; the identical top-source shard and the isolated exact test passed. One rerun of that failed job and its required aggregate passed without source or test changes. The original failed job log and the single-rerun receipt remain retained. ### October 9 merge verification Current merge head: `5a8fe63512a7166aaef5cf50065a25008aa8b44b`. All current-head checks pass, including `ci / verify` and `ci / e2e`; exact-head Greptile review is 5/5 with no unresolved threads. Current master conflicts are resolved. The user authorized the maintainer override of the code-owner review gate after these checks. The seven retained live core cases remain bound to source `3241a992f2a7703e59e97ed0fd3e5d6405de4401` and its recorded cloud image. ## Risks - The security correction changes the dependency closure and profile identity. Old sessions must reopen on the new profile. Exact identities and credential bindings fail closed. - The runner remains experimental and requires explicit selection. Legacy Pi Local is unchanged. Caller model IDs pass through; the E2E model is a fixture. - Accounting and the broad platform/provider matrix remain deferred. This merge does not publish a release or deploy a service. ## Model Used OpenAI GPT-6 through Codex assisted with reasoning, repository inspection, editing and tool use. The exact serving ID and context window are not exposed in this session. Final live qualification uses Pi 1.0.0 with `openrouter/anthropic/claude-sonnet-4.6` and native-confirmed low thinking. ## 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 #` 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>
1312 lines
63 KiB
Diff
1312 lines
63 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
|
|
--- a/dist/live-checkpoint-BSIrfgVo.js
|
|
+++ b/dist/live-checkpoint-BSIrfgVo.js
|
|
@@ -6,7 +6,7 @@
|
|
import { randomUUID } from "node:crypto";
|
|
import { execFile, spawn } from "node:child_process";
|
|
import { Readable, Writable } from "node:stream";
|
|
-import { PROTOCOL_VERSION, client, methods } from "@agentclientprotocol/sdk";
|
|
+import { PROTOCOL_VERSION, RequestError, client, methods } from "@agentclientprotocol/sdk";
|
|
import readline from "node:readline/promises";
|
|
//#region src/errors.ts
|
|
var AcpxOperationalError = class extends Error {
|
|
@@ -1068,12 +1068,14 @@
|
|
last_agent_disconnect_reason: canonical.lastAgentDisconnectReason,
|
|
protocol_version: canonical.protocolVersion,
|
|
agent_capabilities: canonical.agentCapabilities,
|
|
+ agent_goal_capability: canonical.agentGoalCapability,
|
|
title: canonical.title,
|
|
messages: canonical.messages,
|
|
updated_at: canonical.updated_at,
|
|
cumulative_token_usage: canonical.cumulative_token_usage,
|
|
cumulative_cost: canonical.cumulative_cost,
|
|
request_token_usage: canonical.request_token_usage,
|
|
+ cursor_prompt_usage: boundCursorPromptUsage(canonical.cursor_prompt_usage, canonical.messages),
|
|
acpx: canonical.acpx,
|
|
imported_from: canonical.importedFrom ? {
|
|
record_id: canonical.importedFrom.recordId,
|
|
@@ -1239,7 +1241,8 @@
|
|
for (const [key, value] of Object.entries(record)) {
|
|
const parsed = parseTokenUsage(value);
|
|
if (parsed == null) return null;
|
|
- usage[key] = parsed;
|
|
+ const receipt = parsePaperclipPiReceipt(asRecord$4(value)?.paperclip_pi);
|
|
+ usage[key] = receipt ? { ...parsed, paperclip_pi: receipt } : parsed;
|
|
}
|
|
return usage;
|
|
}
|
|
@@ -1325,6 +1328,93 @@
|
|
function isConversationMessage(raw) {
|
|
return raw === "Resume" || isUserMessage$1(raw) || isAgentMessage$1(raw);
|
|
}
|
|
+function parseCursorPromptUsage(raw) {
|
|
+ try {
|
|
+ const text = JSON.stringify(raw);
|
|
+ if (!text || Buffer.byteLength(text, "utf8") > 16384) return null;
|
|
+ const v = JSON.parse(text);
|
|
+ const object = (x)=>x !== null && typeof x === "object" && !Array.isArray(x);
|
|
+ const exact = (x, keys)=>Object.keys(x).length === keys.length && keys.every((k)=>Object.hasOwn(x, k));
|
|
+ const id = (x)=>typeof x === "string" && /^[A-Za-z0-9][A-Za-z0-9._:-]{0,239}$/.test(x);
|
|
+ if (!object(v) || !exact(v, [
|
|
+ "request_id",
|
|
+ "prompt_message_id",
|
|
+ "receipt"
|
|
+ ]) || !id(v.request_id) || !id(v.prompt_message_id)) return null;
|
|
+ const r = v.receipt;
|
|
+ if (!object(r) || !exact(r, [
|
|
+ "schema",
|
|
+ "source",
|
|
+ "promptId",
|
|
+ "completeness",
|
|
+ "reasons",
|
|
+ "observations",
|
|
+ "limits",
|
|
+ "truncated"
|
|
+ ])) return null;
|
|
+ if (r.schema !== "paperclip.cursor.native-usage.v1" || r.source !== "native_turn_ended" || r.completeness !== "partial" || typeof r.promptId !== "string" || !/^[a-f0-9]{8}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{4}-[a-f0-9]{12}$/.test(r.promptId) || typeof r.truncated !== "boolean") return null;
|
|
+ const reasons = [
|
|
+ "native_counter_semantics_unverified",
|
|
+ "native_counters_missing",
|
|
+ "invalid_native_counter",
|
|
+ "multiple_terminal_observations",
|
|
+ "native_terminal_not_observed",
|
|
+ "observation_limit_reached",
|
|
+ "child_run_attribution_unverified"
|
|
+ ];
|
|
+ if (!Array.isArray(r.reasons) || r.reasons.length > reasons.length || !r.reasons.includes(reasons[0]) || new Set(r.reasons).size !== r.reasons.length || !r.reasons.every((x)=>reasons.includes(x))) return null;
|
|
+ if (!object(r.limits) || !exact(r.limits, [
|
|
+ "maxObservations",
|
|
+ "maxInvocations",
|
|
+ "maxBytes"
|
|
+ ]) || r.limits.maxObservations !== 64 || r.limits.maxInvocations !== 64 || r.limits.maxBytes !== 16384 || !Array.isArray(r.observations) || r.observations.length > 64) return null;
|
|
+ const fields = [
|
|
+ "inputTokens",
|
|
+ "outputTokens",
|
|
+ "cacheReadTokens",
|
|
+ "cacheWriteTokens",
|
|
+ "reasoningTokens"
|
|
+ ];
|
|
+ const invocations = new Map();
|
|
+ for (const o of r.observations){
|
|
+ if (!object(o) || !exact(o, [
|
|
+ "invocationId",
|
|
+ "role",
|
|
+ "nativeRun",
|
|
+ "sequence",
|
|
+ "counters"
|
|
+ ]) || typeof o.invocationId !== "string" || !/^invocation-([1-9]|[1-5][0-9]|6[0-4])$/.test(o.invocationId) || o.role !== "parent" && o.role !== "child" || !Number.isSafeInteger(o.nativeRun) || o.nativeRun < 1 || !Number.isSafeInteger(o.sequence) || o.sequence < 1 || !object(o.counters)) return null;
|
|
+ if (!Object.entries(o.counters).every(([k, n])=>fields.includes(k) && Number.isSafeInteger(n) && n >= 0)) return null;
|
|
+ const previous = invocations.get(o.invocationId);
|
|
+ if (previous && (previous.role !== o.role || previous.nativeRun !== o.nativeRun || o.sequence <= previous.sequence)) return null;
|
|
+ invocations.set(o.invocationId, {
|
|
+ role: o.role,
|
|
+ nativeRun: o.nativeRun,
|
|
+ sequence: o.sequence
|
|
+ });
|
|
+ }
|
|
+ return v;
|
|
+ } catch {
|
|
+ return null;
|
|
+ }
|
|
+}
|
|
+function boundCursorPromptUsage(raw, messages) {
|
|
+ const receipt = parseCursorPromptUsage(raw);
|
|
+ if (!receipt || !Array.isArray(messages) || messages.filter(message => message?.User?.id === receipt.prompt_message_id).length !== 1) return undefined;
|
|
+ return receipt;
|
|
+}
|
|
+function recordCursorPromptUsage(conversation, response, requestId, promptMessageId) {
|
|
+ // Optional diagnostics cannot affect settlement or standard token/cost accounting.
|
|
+ try {
|
|
+ if (response.stopReason !== "end_turn" && response.stopReason !== "cancelled") return;
|
|
+ const next = boundCursorPromptUsage({request_id: requestId, prompt_message_id: promptMessageId, receipt: response._meta?.paperclipCursorUsage}, conversation.messages);
|
|
+ if (!next) return;
|
|
+ const previous = parseCursorPromptUsage(conversation.cursor_prompt_usage);
|
|
+ if (previous && (previous.request_id === next.request_id || previous.prompt_message_id === next.prompt_message_id || previous.receipt.promptId === next.receipt.promptId)) return;
|
|
+ conversation.cursor_prompt_usage = next;
|
|
+ } catch {}
|
|
+}
|
|
+
|
|
function parseConversationRecord(record) {
|
|
if (!hasValidConversationCore(record)) return;
|
|
const title = parseConversationTitle(record.title);
|
|
@@ -1339,7 +1429,8 @@
|
|
updated_at: record.updated_at,
|
|
cumulative_token_usage: cumulativeTokenUsage ?? {},
|
|
cumulative_cost: cumulativeCost,
|
|
- request_token_usage: requestTokenUsage ?? {}
|
|
+ request_token_usage: requestTokenUsage ?? {},
|
|
+ cursor_prompt_usage: boundCursorPromptUsage(record.cursor_prompt_usage, record.messages)
|
|
};
|
|
}
|
|
const INVALID_VALUE = Symbol("invalid");
|
|
@@ -1541,12 +1632,14 @@
|
|
lastAgentDisconnectReason: optionals.lastAgentDisconnectReason,
|
|
protocolVersion: typeof record.protocol_version === "number" ? record.protocol_version : void 0,
|
|
agentCapabilities: asRecord$4(record.agent_capabilities),
|
|
+ agentGoalCapability: asRecord$4(record.agent_goal_capability),
|
|
title: conversation.title,
|
|
messages: conversation.messages,
|
|
updated_at: conversation.updated_at,
|
|
cumulative_token_usage: conversation.cumulative_token_usage,
|
|
cumulative_cost: conversation.cumulative_cost,
|
|
request_token_usage: conversation.request_token_usage,
|
|
+ cursor_prompt_usage: conversation.cursor_prompt_usage,
|
|
acpx: parseAcpxState(record.acpx),
|
|
importedFrom: metadata.importedFrom
|
|
};
|
|
@@ -1661,8 +1754,9 @@
|
|
"RedactedThinking",
|
|
"ToolUse"
|
|
]);
|
|
-const MAP_OBJECT_PATHS = /* @__PURE__ */ new Set(["request_token_usage", "messages.Agent.tool_results"]);
|
|
+const MAP_OBJECT_PATHS = /* @__PURE__ */ new Set(["request_token_usage", "messages.Agent.tool_results", "acpx.session_options.env"]);
|
|
const OPAQUE_VALUE_PATHS = /* @__PURE__ */ new Set([
|
|
+ "cursor_prompt_usage.receipt",
|
|
"agent_capabilities",
|
|
"messages.Agent.content.ToolUse.input",
|
|
"acpx.desired_config_options",
|
|
@@ -3104,10 +3198,10 @@
|
|
authEnvKeyCache.set(methodId, key);
|
|
return key;
|
|
}
|
|
-function readEnvCredential(methodId) {
|
|
+function readEnvCredential(methodId, environment = process.env) {
|
|
const key = authEnvKey(methodId);
|
|
if (!key) return;
|
|
- const value = process.env[key];
|
|
+ const value = environment[key];
|
|
if (typeof value === "string" && value.trim().length > 0) return value;
|
|
}
|
|
function protectedEnvKey(key) {
|
|
@@ -3135,8 +3229,21 @@
|
|
}
|
|
return protectedKeys;
|
|
}
|
|
-function buildAgentEnvironment(authCredentials, sessionEnv) {
|
|
- const env = { ...process.env };
|
|
+function isPlainStringEnvironment(value) {
|
|
+ if (value === null || typeof value !== "object" || Array.isArray(value)) return false;
|
|
+ const prototype = Object.getPrototypeOf(value);
|
|
+ return (prototype === Object.prototype || prototype === null) && Object.values(value).every((entry) => typeof entry === "string");
|
|
+}
|
|
+function resolveAgentEnvironment(spawnEnvironment) {
|
|
+ let sourceEnvironment = process.env;
|
|
+ if (spawnEnvironment !== void 0) {
|
|
+ sourceEnvironment = spawnEnvironment();
|
|
+ if (!isPlainStringEnvironment(sourceEnvironment)) throw new TypeError("ACPX spawn environment must be a plain record of string values");
|
|
+ }
|
|
+ return sourceEnvironment;
|
|
+}
|
|
+function buildAgentEnvironment(authCredentials, sessionEnv, spawnEnvironment) {
|
|
+ const env = { ...resolveAgentEnvironment(spawnEnvironment) };
|
|
const protectedAuthEnvKeys = promotePrefixedAuthEnvironment(env);
|
|
if (authCredentials) for (const [methodId, credential] of Object.entries(authCredentials)) {
|
|
addAuthCredentialEnvKeys(protectedAuthEnvKeys, methodId, credential);
|
|
@@ -3178,10 +3285,10 @@
|
|
const configCredentials = authCredentials ?? {};
|
|
return configCredentials[methodId] ?? configCredentials[toEnvToken(methodId)];
|
|
}
|
|
-function buildAgentSpawnOptions(cwd, authCredentials, sessionEnv) {
|
|
+function buildAgentSpawnOptions(cwd, authCredentials, sessionEnv, spawnEnvironment) {
|
|
return {
|
|
cwd,
|
|
- env: buildAgentEnvironment(authCredentials, sessionEnv),
|
|
+ env: buildAgentEnvironment(authCredentials, sessionEnv, spawnEnvironment),
|
|
stdio: [
|
|
"pipe",
|
|
"pipe",
|
|
@@ -3286,28 +3393,12 @@
|
|
const ids = models?.availableModels.map((model) => model.modelId.trim()).filter((modelId) => modelId.length > 0) ?? [];
|
|
return ids.length > 0 ? ids.join(", ") : "none advertised";
|
|
}
|
|
-function resolveRequestedModelId(params) {
|
|
- if (!params.models || !isCursorAcpCommandForModelAlias(params.agentCommand)) return params.requestedModel;
|
|
- if (params.models.availableModels.some((model) => model.modelId === params.requestedModel)) return params.requestedModel;
|
|
- const candidates = params.models.availableModels.map((model) => model.modelId).filter((modelId) => modelId.startsWith(`${params.requestedModel}[`));
|
|
- return candidates.length === 1 ? candidates[0] : params.requestedModel;
|
|
-}
|
|
-function isCursorAcpCommandForModelAlias(agentCommand) {
|
|
- if (!agentCommand) return false;
|
|
- const { command, args } = splitCommandLine(agentCommand);
|
|
- return isCursorAcpCommand(command, args);
|
|
-}
|
|
function assertRequestedModelSupported(params) {
|
|
if (!params.models) {
|
|
if (supportsLegacyClaudeCodeModelMetadata(params.agentCommand)) return;
|
|
throw new RequestedModelUnsupportedError(`Cannot ${params.context === "replay" ? "replay saved model" : "apply --model"} "${params.requestedModel}": the ACP agent did not advertise model support through a session config option or legacy models metadata, and the adapter does not support a startup model flag.`, "missing-capability");
|
|
}
|
|
- if (!new Set(params.models.availableModels.map((model) => model.modelId)).has(params.requestedModel)) {
|
|
- const resolvedModel = resolveRequestedModelId(params);
|
|
- if (resolvedModel !== params.requestedModel) return `Cursor ACP advertised "${resolvedModel}" for requested model "${params.requestedModel}"; using the advertised id.`;
|
|
- if (supportsLegacyClaudeCodeModelMetadata(params.agentCommand)) return `requested model "${params.requestedModel}" was not in the Claude ACP advertised model list (${formatAvailableModelIds(params.models)}); forwarding it to Claude Code so the adapter can accept or reject it.`;
|
|
- throw new RequestedModelUnsupportedError(`Cannot ${params.context === "replay" ? "replay saved model" : "apply --model"} "${params.requestedModel}": the ACP agent did not advertise that model. Available models: ${formatAvailableModelIds(params.models)}.`, "unadvertised-model");
|
|
- }
|
|
+ // Discovery catalogs may be incomplete; the provider accepts or rejects the exact ID.
|
|
}
|
|
//#endregion
|
|
//#region src/acp/session-control-errors.ts
|
|
@@ -3961,10 +4052,13 @@
|
|
...params.elicitationModes.includes("url") ? { url: {} } : {}
|
|
} } : {}
|
|
};
|
|
- if (!params.devinAcp) return baseCapabilities;
|
|
+ const typedSessionFailureMeta = {
|
|
+ jetbrains: { air: { version: 1, capabilities: ["sessionFailure"] } }
|
|
+ };
|
|
+ if (!params.devinAcp) return { ...baseCapabilities, _meta: typedSessionFailureMeta };
|
|
return {
|
|
...baseCapabilities,
|
|
- _meta: DEVIN_COMPATIBILITY_CLIENT_CAPABILITIES_META
|
|
+ _meta: { ...typedSessionFailureMeta, ...DEVIN_COMPATIBILITY_CLIENT_CAPABILITIES_META }
|
|
};
|
|
}
|
|
function hasResponseField(response, field) {
|
|
@@ -4067,6 +4161,27 @@
|
|
} })
|
|
};
|
|
}
|
|
+function boundedExtensionRecord(value) {
|
|
+ if (value === null || typeof value !== "object" || Array.isArray(value)) throw new TypeError("ACP extension payload must be an object");
|
|
+ const prototype = Object.getPrototypeOf(value);
|
|
+ if (prototype !== Object.prototype && prototype !== null) throw new TypeError("ACP extension payload must be a plain object");
|
|
+ const serialized = JSON.stringify(value);
|
|
+ if (Buffer.byteLength(serialized) > 256 * 1024) throw new TypeError("ACP extension payload exceeds its bounded size");
|
|
+ return JSON.parse(serialized);
|
|
+}
|
|
+function closedExtensionMethods(value) {
|
|
+ if (value === void 0) return Object.freeze([]);
|
|
+ if (!Array.isArray(value) || value.length > 64) throw new TypeError("ACP extension methods must be a bounded allowlist");
|
|
+ for (const method of value) {
|
|
+ if (typeof method !== "string" || !/^[A-Za-z_][A-Za-z0-9._/-]{1,255}$/.test(method) || !method.includes("/") || /^(?:session|fs|terminal|elicitation|notifications)\//.test(method)) throw new TypeError("ACP extension method cannot override a standard method");
|
|
+ }
|
|
+ return Object.freeze([...new Set(value)]);
|
|
+}
|
|
+function mergeRuntimeClientCapabilities(base, extra) {
|
|
+ if (extra === void 0) return base;
|
|
+ const additional = boundedExtensionRecord(extra);
|
|
+ return { ...additional, ...base, _meta: { ...additional._meta, ...base._meta } };
|
|
+}
|
|
var AcpClient = class {
|
|
options;
|
|
connection;
|
|
@@ -4105,7 +4220,9 @@
|
|
cwd: asAbsoluteCwd(options.cwd),
|
|
authPolicy: options.authPolicy ?? "skip",
|
|
permissionPolicy: snapshotPermissionPolicy(options.permissionPolicy),
|
|
- elicitationModes: normalizeElicitationModes(options.elicitationModes)
|
|
+ elicitationModes: normalizeElicitationModes(options.elicitationModes),
|
|
+ extensionMethods: closedExtensionMethods(options.extensionMethods),
|
|
+ clientCapabilities: options.clientCapabilities === void 0 ? void 0 : boundedExtensionRecord(options.clientCapabilities)
|
|
};
|
|
this.eventHandlers = {
|
|
onAcpMessage: this.options.onAcpMessage,
|
|
@@ -4216,9 +4333,21 @@
|
|
this.lastAgentExit = void 0;
|
|
this.lastKnownPid = child.pid ?? void 0;
|
|
this.attachAgentLifecycleObservers(child);
|
|
+ if (this.options.onAgentSpawn) {
|
|
+ try {
|
|
+ if (!child.pid) throw new Error("ACPX agent spawn did not expose a valid process id.");
|
|
+ await this.options.onAgentSpawn({ pid: child.pid, startedAt: this.agentStartedAt });
|
|
+ } catch (error) {
|
|
+ try {
|
|
+ child.kill("SIGKILL");
|
|
+ } catch {}
|
|
+ throw error;
|
|
+ }
|
|
+ }
|
|
const startupStderr = [];
|
|
child.stderr.on("data", (chunk) => {
|
|
this.captureStartupStderr(startupStderr, chunk);
|
|
+ this.options.onAgentStderr?.(String(chunk));
|
|
if (!this.options.verbose) return;
|
|
process.stderr.write(chunk);
|
|
});
|
|
@@ -4226,6 +4355,7 @@
|
|
const output = Readable.toWeb(child.stdout);
|
|
const stream = this.createTappedStream(createNdJsonMessageStream(this.options.agentCommand, input, output));
|
|
const connection = this.createConnection(stream, launch);
|
|
+ connection.signal.addEventListener("abort", () => stream.abortResponseDeliveries(), { once: true });
|
|
connection.signal.addEventListener("abort", () => {
|
|
this.recordAgentExit("connection_close", child.exitCode ?? null, child.signalCode ?? null);
|
|
}, { once: true });
|
|
@@ -4253,7 +4383,12 @@
|
|
geminiAcp: isGeminiAcpCommand(spawnCommand, args),
|
|
copilotAcp: isCopilotAcpCommand(spawnCommand, args),
|
|
claudeAcp: isClaudeAcpCommand(spawnCommand, args),
|
|
- spawnOptions: buildAgentSpawnOptions(this.options.cwd, this.options.authCredentials, this.options.sessionOptions?.env)
|
|
+ spawnOptions: buildAgentSpawnOptions(
|
|
+ this.options.spawnCwd ?? this.options.cwd,
|
|
+ this.options.authCredentials,
|
|
+ this.options.sessionOptions?.env,
|
|
+ this.options.spawnEnvironment
|
|
+ )
|
|
};
|
|
}
|
|
logAgentLaunch(plan) {
|
|
@@ -4279,10 +4414,17 @@
|
|
}
|
|
async spawnAgentProcess(plan) {
|
|
const spawnCommand = buildAgentSpawnCommand(plan.spawnCommand, plan.args, process.platform, plan.spawnOptions.env);
|
|
- const spawnedChild = spawn(spawnCommand.command, spawnCommand.args, {
|
|
+ const options = {
|
|
...plan.spawnOptions,
|
|
windowsVerbatimArguments: spawnCommand.windowsVerbatimArguments
|
|
- });
|
|
+ };
|
|
+ const spawnedChild = this.options.spawnAgent
|
|
+ ? this.options.spawnAgent({
|
|
+ command: spawnCommand.command,
|
|
+ args: spawnCommand.args,
|
|
+ options
|
|
+ })
|
|
+ : spawn(spawnCommand.command, spawnCommand.args, options);
|
|
try {
|
|
await waitForSpawn$1(spawnedChild);
|
|
} catch (error) {
|
|
@@ -4293,10 +4435,10 @@
|
|
createConnection(stream, launch) {
|
|
const app = client({ name: "acpx" }).onNotification(methods.client.session.update, async ({ params }) => {
|
|
await this.handleSessionUpdate(params);
|
|
- }).onNotification(methods.client.elicitation.complete, async () => {}).onRequest(methods.client.session.requestPermission, async ({ params }) => {
|
|
- return await this.handlePermissionRequest(params);
|
|
+ }).onNotification(methods.client.elicitation.complete, async () => {}).onRequest(methods.client.session.requestPermission, async ({ params, requestId, signal }) => {
|
|
+ return await stream.withResponseDelivery(requestId, signal, (responseDelivery) => this.handlePermissionRequest(params, responseDelivery));
|
|
}).onRequest(methods.client.elicitation.create, async ({ params, requestId, signal }) => {
|
|
- return await this.handleElicitationRequest(params, requestId, signal);
|
|
+ return await stream.withResponseDelivery(requestId, signal, (responseDelivery) => this.handleElicitationRequest(params, requestId, signal, responseDelivery));
|
|
}).onRequest(methods.client.fs.readTextFile, async ({ params }) => {
|
|
return await this.handleReadTextFile(params);
|
|
}).onRequest(methods.client.fs.writeTextFile, async ({ params }) => {
|
|
@@ -4312,6 +4454,11 @@
|
|
}).onRequest(methods.client.terminal.release, async ({ params }) => {
|
|
return await this.handleReleaseTerminal(params);
|
|
});
|
|
+ for (const method of this.options.extensionMethods) {
|
|
+ if (this.options.onExtensionRequest) app.onRequest(method, boundedExtensionRecord, async ({ params, requestId, signal }) => {
|
|
+ return await stream.withResponseDelivery(requestId, signal, (responseDelivery) => this.handleExtensionRequest(method, params, requestId, signal, responseDelivery));
|
|
+ });
|
|
+ }
|
|
if (launch.devinAcp) app.onRequest("_cognition.ai/request_diagnostics", (params) => {
|
|
return params && typeof params === "object" && !Array.isArray(params) ? params : {};
|
|
}, async () => ({}));
|
|
@@ -4333,12 +4480,12 @@
|
|
async initializeProtocolConnection(connection, launch) {
|
|
const initializePromise = connection.initialize({
|
|
protocolVersion: PROTOCOL_VERSION,
|
|
- clientCapabilities: resolveClientCapabilities({
|
|
+ clientCapabilities: mergeRuntimeClientCapabilities(resolveClientCapabilities({
|
|
devinAcp: launch.devinAcp,
|
|
fs: this.options.fs !== false,
|
|
terminal: this.options.terminal !== false,
|
|
elicitationModes: this.options.elicitationModes ?? []
|
|
- }),
|
|
+ }), this.options.clientCapabilities),
|
|
clientInfo: resolveClientInfo(launch.devinAcp)
|
|
});
|
|
const initialized = launch.geminiAcp ? await withTimeout(initializePromise, resolveGeminiAcpStartupTimeoutMs()) : await initializePromise;
|
|
@@ -4358,6 +4505,49 @@
|
|
throw normalizedError;
|
|
}
|
|
createTappedStream(base) {
|
|
+ // Policy admission is distinct from best-effort observation. Each new
|
|
+ // stdio connection gets fresh authority and a violation tears it down.
|
|
+ const protocolGuard = this.options.protocolGuardFactory?.();
|
|
+ const enforceProtocol = (direction, message) => {
|
|
+ try { protocolGuard?.(direction, message); }
|
|
+ catch (error) {
|
|
+ this.rejectPendingConnectionRequests(error);
|
|
+ abortResponseDeliveries();
|
|
+ this.activePrompt?.elicitationController.abort(error);
|
|
+ try { this.agent?.kill("SIGTERM"); } catch {}
|
|
+ throw error;
|
|
+ }
|
|
+ };
|
|
+ const deliveries = new Map();
|
|
+ const scope = new AbortController();
|
|
+ const abortResponseDeliveries = () => scope.abort(new Error("ACP response connection closed before delivery"));
|
|
+ const withResponseDelivery = async (id, requestSignal, handler) => {
|
|
+ if (!(id === null || typeof id === "string" || typeof id === "number" && Number.isFinite(id))) throw new Error("ACP response request id is invalid");
|
|
+ if (deliveries.has(id) || deliveries.size >= 128) {
|
|
+ abortResponseDeliveries();
|
|
+ throw new Error("ACP response request id is duplicated or pending deliveries exceeded their bound");
|
|
+ }
|
|
+ const active = this.activePrompt;
|
|
+ const signal = AbortSignal.any([scope.signal, requestSignal, ...(active ? [active.elicitationController.signal] : [])]);
|
|
+ let settle, reject;
|
|
+ const responseDelivery = new Promise((resolve, fail) => { settle = resolve; reject = fail; });
|
|
+ // The host may decline the request without awaiting its receipt.
|
|
+ void responseDelivery.catch(() => {});
|
|
+ const abort = () => reject(new Error("ACP response delivery was cancelled or disconnected"));
|
|
+ const receipt = { settle, reject, expected: void 0, cleanup: () => signal.removeEventListener("abort", abort) };
|
|
+ deliveries.set(id, receipt);
|
|
+ signal.addEventListener("abort", abort, { once: true });
|
|
+ if (signal.aborted) abort();
|
|
+ try {
|
|
+ const result = await handler(responseDelivery);
|
|
+ if (signal.aborted || this.activePrompt !== active) throw new Error("ACP response outlived its admitted prompt");
|
|
+ receipt.expected = JSON.stringify(result);
|
|
+ return result;
|
|
+ } catch (error) {
|
|
+ reject(new Error("ACP request failed before response delivery"));
|
|
+ throw error;
|
|
+ }
|
|
+ };
|
|
const onAcpMessage = () => this.eventHandlers.onAcpMessage;
|
|
const onAcpOutputMessage = () => this.eventHandlers.onAcpOutputMessage;
|
|
const elicitationRequestIds = /* @__PURE__ */ new Set();
|
|
@@ -4367,8 +4557,12 @@
|
|
const shouldSuppressInboundReplaySessionUpdate = (message) => {
|
|
return this.suppressReplaySessionUpdateMessages && isSessionUpdateNotification(message);
|
|
};
|
|
+ const cursorSubagentOwner = { active: null, children: new Map() };
|
|
+ this.cursorSubagentStreamOwner = cursorSubagentOwner;
|
|
+ const onCursorSubagentNotification = (message) => this.handleCursorSubagentNotification(message, cursorSubagentOwner);
|
|
+ const onExtensionNotification = (message) => this.handleExtensionNotification(message);
|
|
const observeInbound = (message) => {
|
|
- const requestId = elicitationRequestId(message);
|
|
+ const requestId = elicitationRequestId(message) ?? ("id" in message && this.options.extensionMethods.includes(message.method) ? message.id : void 0);
|
|
if (requestId !== void 0) {
|
|
elicitationRequestIds.add(requestId);
|
|
return;
|
|
@@ -4379,6 +4573,7 @@
|
|
}
|
|
};
|
|
return {
|
|
+ withResponseDelivery, abortResponseDeliveries,
|
|
readable: new ReadableStream({ async start(controller) {
|
|
const reader = base.readable.getReader();
|
|
try {
|
|
@@ -4386,16 +4581,27 @@
|
|
const { value, done } = await reader.read();
|
|
if (done) break;
|
|
if (!value) continue;
|
|
+ enforceProtocol("inbound", value);
|
|
observeInbound(value);
|
|
- if (isExtensionNotification(value)) continue;
|
|
+ if (onCursorSubagentNotification(value)) continue;
|
|
+ if (isExtensionNotification(value)) {
|
|
+ onExtensionNotification(value);
|
|
+ continue;
|
|
+ }
|
|
controller.enqueue(value);
|
|
}
|
|
+ controller.close();
|
|
+ } catch (error) {
|
|
+ controller.error(error);
|
|
} finally {
|
|
+ abortResponseDeliveries();
|
|
+ for (const receipt of deliveries.values()) receipt.cleanup();
|
|
+ deliveries.clear();
|
|
reader.releaseLock();
|
|
- controller.close();
|
|
}
|
|
} }),
|
|
writable: new WritableStream({ async write(message) {
|
|
+ enforceProtocol("outbound", message);
|
|
const promptOwner = promptRequestOwner(message);
|
|
if (promptOwner) bindPromptOwner(promptOwner);
|
|
const id = responseId(message);
|
|
@@ -4403,10 +4609,20 @@
|
|
onAcpOutputMessage()?.("outbound", message);
|
|
onAcpMessage()?.("outbound", message);
|
|
}
|
|
+ const receipt = id === void 0 ? void 0 : deliveries.get(id);
|
|
const writer = base.writable.getWriter();
|
|
try {
|
|
await writer.write(message);
|
|
+ if (receipt) {
|
|
+ if (!("result" in message) || receipt.expected === void 0 || JSON.stringify(message.result) !== receipt.expected) receipt.reject(new Error("ACP wrote an error or changed response instead of the selected resolution"));
|
|
+ else receipt.settle();
|
|
+ }
|
|
+ } catch (error) {
|
|
+ receipt?.reject(new Error("ACP response write failed"));
|
|
+ abortResponseDeliveries();
|
|
+ throw error;
|
|
} finally {
|
|
+ if (receipt) { receipt.cleanup(); deliveries.delete(id); }
|
|
writer.releaseLock();
|
|
}
|
|
} })
|
|
@@ -4422,7 +4638,10 @@
|
|
const createPromise = this.runConnectionRequest(() => connection.newSession({
|
|
cwd: sessionCwd,
|
|
mcpServers: this.options.mcpServers ?? [],
|
|
- _meta: buildClaudeCodeOptionsMeta(this.options.sessionOptions, claudeAcp)
|
|
+ _meta: this.isGrokBuildAcpCommand()
|
|
+ ? { rules: typeof this.options.sessionOptions?.systemPrompt === "string"
|
|
+ ? this.options.sessionOptions.systemPrompt : this.options.sessionOptions?.systemPrompt?.append ?? "" }
|
|
+ : buildClaudeCodeOptionsMeta(this.options.sessionOptions, claudeAcp)
|
|
}));
|
|
result = claudeAcp ? await withTimeout(createPromise, resolveClaudeAcpSessionCreateTimeoutMs()) : await createPromise;
|
|
} catch (error) {
|
|
@@ -4589,6 +4808,13 @@
|
|
} catch (error) {
|
|
throw maybeWrapSessionControlError("session/set_config_option", error, `for "${configId}"="${value}"`);
|
|
}
|
|
+ }
|
|
+ async requestExtension(method, params) {
|
|
+ closedExtensionMethods([method]);
|
|
+ const payload = boundedExtensionRecord(params);
|
|
+ const connection = this.getConnection();
|
|
+ const result = await this.runConnectionRequest(() => connection.extMethod(method, payload));
|
|
+ return boundedExtensionRecord(result);
|
|
}
|
|
async setSessionModel(sessionId, modelId, controlOverride) {
|
|
const control = this.resolveModelControl(sessionId, controlOverride);
|
|
@@ -4593,12 +4819,7 @@
|
|
async setSessionModel(sessionId, modelId, controlOverride) {
|
|
const control = this.resolveModelControl(sessionId, controlOverride);
|
|
if (!control) throw new RequestedModelUnsupportedError(`Cannot set model "${modelId}": the ACP session did not advertise a model config option or legacy session/set_model support.`, "missing-capability");
|
|
- const resolvedModelId = resolveRequestedModelId({
|
|
- requestedModel: modelId,
|
|
- models: controlOverride?.availableModels ? { availableModels: controlOverride.availableModels } : void 0,
|
|
- agentCommand: this.options.agentCommand
|
|
- });
|
|
- return control.kind === "config_option" ? await this.setSessionModelThroughConfig(sessionId, resolvedModelId, control.configId) : await this.setSessionModelThroughLegacyMethod(sessionId, resolvedModelId);
|
|
+ return control.kind === "config_option" ? await this.setSessionModelThroughConfig(sessionId, modelId, control.configId) : await this.setSessionModelThroughLegacyMethod(sessionId, modelId);
|
|
}
|
|
async setSessionModelThroughConfig(sessionId, modelId, configId) {
|
|
const connection = this.getConnection();
|
|
@@ -4608,7 +4829,9 @@
|
|
configId,
|
|
value: modelId
|
|
}));
|
|
- this.rememberSessionModels(sessionId, modelStateFromConfigOptions(response.configOptions));
|
|
+ const selected = modelStateFromConfigOptions(response.configOptions);
|
|
+ if (selected?.currentModelId !== modelId) throw Object.assign(new Error(`ACP effective model mismatch: requested ${modelId}, received ${selected?.currentModelId ?? "unverified"}`), { code: "ACPX_EFFECTIVE_MODEL_MISMATCH" });
|
|
+ this.rememberSessionModels(sessionId, selected);
|
|
return response;
|
|
} catch (error) {
|
|
return this.throwSessionModelError("session/set_config_option", modelId, error);
|
|
@@ -4864,7 +5087,7 @@
|
|
}
|
|
selectAuthMethod(methods) {
|
|
for (const method of methods) {
|
|
- const envCredential = readEnvCredential(method.id);
|
|
+ const envCredential = readEnvCredential(method.id, resolveAgentEnvironment(this.options.spawnEnvironment));
|
|
if (envCredential) return {
|
|
methodId: method.id,
|
|
credential: envCredential,
|
|
@@ -4890,7 +5113,7 @@
|
|
}
|
|
readAgentSpecificEnvCredential(methodId) {
|
|
if (!this.isGrokBuildAcpCommand() || methodId !== "xai.api_key") return;
|
|
- const value = process.env.XAI_API_KEY;
|
|
+ const value = (resolveAgentEnvironment(this.options.spawnEnvironment)).XAI_API_KEY;
|
|
return typeof value === "string" && value.trim().length > 0 ? value : void 0;
|
|
}
|
|
selectAgentManagedAuthMethod(methodId) {
|
|
@@ -4915,9 +5138,9 @@
|
|
await connection.authenticate({ methodId: selected.methodId });
|
|
this.log(`authenticated with method ${selected.methodId} (${selected.source})`);
|
|
}
|
|
- async handlePermissionRequest(params) {
|
|
+ async handlePermissionRequest(params, responseDelivery) {
|
|
if (this.cancellingSessionIds.has(params.sessionId)) return cancelledPermissionResponse();
|
|
- const hostResponse = await this.tryHandlePermissionRequestWithHost(params);
|
|
+ const hostResponse = await this.tryHandlePermissionRequestWithHost(params, responseDelivery);
|
|
if (hostResponse) return hostResponse;
|
|
const { response, recorded } = await this.resolvePermissionRequestFromMode(params);
|
|
if (!recorded) {
|
|
@@ -4926,14 +5149,81 @@
|
|
}
|
|
return response;
|
|
}
|
|
- async handleElicitationRequest(request, requestId, requestSignal) {
|
|
+ async handleExtensionRequest(method, params, requestId, requestSignal, responseDelivery) {
|
|
+ const active = this.activePrompt;
|
|
+ const connection = this.connection;
|
|
+ if (!this.options.extensionMethods.includes(method) || !this.options.onExtensionRequest) throw new Error("ACP extension method is not enabled");
|
|
+ if (!active || !connection || this.closing || (params.sessionId !== void 0 && params.sessionId !== active.sessionId)) throw new Error("ACP extension request has no matching active prompt");
|
|
+ if (!(typeof requestId === "string" || typeof requestId === "number" && Number.isFinite(requestId))) throw new Error("ACP extension request id is invalid");
|
|
+ const signal = AbortSignal.any([requestSignal, connection.signal, active.elicitationController.signal]);
|
|
+ signal.throwIfAborted();
|
|
+ let onAbort;
|
|
+ const aborted = new Promise((_, reject) => {
|
|
+ onAbort = () => reject(new Error("ACP extension request was cancelled"));
|
|
+ signal.addEventListener("abort", onAbort, { once: true });
|
|
+ });
|
|
+ try {
|
|
+ const result = await Promise.race([Promise.resolve().then(() => this.options.onExtensionRequest(method, boundedExtensionRecord(params), { requestId, signal, responseDelivery })), aborted]);
|
|
+ signal.throwIfAborted();
|
|
+ if (this.activePrompt !== active || this.connection !== connection) throw new Error("ACP extension request outlived its prompt");
|
|
+ return boundedExtensionRecord(result);
|
|
+ } finally {
|
|
+ signal.removeEventListener("abort", onAbort);
|
|
+ }
|
|
+ }
|
|
+ handleCursorSubagentNotification(message, owner) {
|
|
+ // Keep the SDK's standard union closed. Cursor child transcripts use
|
|
+ // their own session IDs and must never be flattened into the parent.
|
|
+ const method = "cursor/subagent_update";
|
|
+ if (!this.options.extensionMethods.includes(method) || this.options.clientCapabilities?._meta?.subagents !== true) return false;
|
|
+ if (message?.method !== "session/update" || "id" in message) return false;
|
|
+ const lifecycle = ["subagent_spawned", "subagent_state_update"].includes(message?.params?.update?.sessionUpdate);
|
|
+ const active = this.activePrompt;
|
|
+ const connection = this.connection;
|
|
+ if (this.cursorSubagentStreamOwner !== owner || this.closing || connection?.signal.aborted) return true;
|
|
+ if (!active || !connection || active.elicitationController.signal.aborted) return lifecycle || owner.children.has(message?.params?.sessionId);
|
|
+ if (owner.active !== active) { owner.active = active; owner.children.clear(); }
|
|
+ if (!lifecycle && message?.params?.sessionId === active.sessionId) return false;
|
|
+ try {
|
|
+ const params = boundedExtensionRecord(message.params);
|
|
+ if (lifecycle) {
|
|
+ if (params.sessionId !== active.sessionId && !owner.children.has(params.sessionId)) return true;
|
|
+ const child = params.update?.subagentSessionId;
|
|
+ if (typeof child !== "string" || !child || child.length > 1000 || child === active.sessionId) return true;
|
|
+ if (params.update.sessionUpdate === "subagent_spawned") {
|
|
+ if (owner.children.has(child) || owner.children.size >= 256) return true;
|
|
+ owner.children.set(child, params.sessionId);
|
|
+ } else if (owner.children.get(child) !== params.sessionId) return true;
|
|
+ this.options.onExtensionNotification?.(method, { sessionId: active.sessionId, parentSessionId: params.sessionId, update: params.update });
|
|
+ } else if (owner.children.has(params.sessionId)) {
|
|
+ this.options.onExtensionNotification?.(method, { sessionId: active.sessionId, parentSessionId: owner.children.get(params.sessionId), childSessionId: params.sessionId, update: params.update });
|
|
+ }
|
|
+ } catch {
|
|
+ this.log("Cursor child activity notification was rejected");
|
|
+ }
|
|
+ return true;
|
|
+ }
|
|
+ handleExtensionNotification(message) {
|
|
+ if (!this.options.extensionMethods.includes(message.method) || !this.options.onExtensionNotification || this.closing) return;
|
|
+ const active = this.activePrompt;
|
|
+ if (!active || active.elicitationController.signal.aborted) return;
|
|
+ try {
|
|
+ const params = boundedExtensionRecord(message.params);
|
|
+ if (params.sessionId !== void 0 && params.sessionId !== active.sessionId) return;
|
|
+ this.options.onExtensionNotification(message.method, params);
|
|
+ } catch {
|
|
+ this.log("ACP extension notification was rejected");
|
|
+ }
|
|
+ }
|
|
+ async handleElicitationRequest(request, requestId, requestSignal, responseDelivery) {
|
|
const resolved = this.resolveElicitationOwner(request);
|
|
if ("response" in resolved) return resolved.response;
|
|
const { active, handler } = resolved.owner;
|
|
const signal = AbortSignal.any([requestSignal, active.elicitationController.signal]);
|
|
const handlerAttempt = Promise.resolve().then(async () => await handler(request, {
|
|
requestId,
|
|
- signal
|
|
+ signal,
|
|
+ responseDelivery
|
|
})).then((response) => ({
|
|
kind: "response",
|
|
response
|
|
@@ -4970,7 +5260,7 @@
|
|
isElicitationSessionCancelling(sessionId) {
|
|
return this.closing || this.cancellingSessionIds.has(sessionId);
|
|
}
|
|
- async tryHandlePermissionRequestWithHost(params) {
|
|
+ async tryHandlePermissionRequestWithHost(params, responseDelivery) {
|
|
if (!this.options.onPermissionRequest) return;
|
|
const signal = this.cancellationSignalForSession(params.sessionId);
|
|
try {
|
|
@@ -4978,7 +5268,7 @@
|
|
sessionId: params.sessionId,
|
|
raw: params,
|
|
inferredKind: inferToolKind(params)
|
|
- }, { signal });
|
|
+ }, { signal, responseDelivery });
|
|
return this.hostPermissionDecisionResponse(params, signal, decision);
|
|
} catch (error) {
|
|
return this.hostPermissionErrorResponse(params, signal, error);
|
|
@@ -5028,6 +5318,12 @@
|
|
attachAgentLifecycleObservers(child) {
|
|
child.once("exit", (exitCode, signal) => {
|
|
this.recordAgentExit("process_exit", exitCode, signal);
|
|
+ this.options.onAgentExit?.({
|
|
+ pid: child.pid ?? this.lastKnownPid,
|
|
+ exitCode,
|
|
+ signal,
|
|
+ exitedAt: isoNow$1()
|
|
+ });
|
|
});
|
|
child.once("close", (exitCode, signal) => {
|
|
this.recordAgentExit("process_close", exitCode, signal);
|
|
@@ -5100,6 +5396,9 @@
|
|
return await this.filesystem.readTextFile(params);
|
|
} catch (error) {
|
|
this.recordPermissionError(params.sessionId, error);
|
|
+ // Keep ENOENT distinct from permission and other filesystem failures.
|
|
+ // The ACP peer must receive a resource error before it creates a file.
|
|
+ if (error && typeof error === "object" && error.code === "ENOENT") throw RequestError.resourceNotFound(params.path);
|
|
throw error;
|
|
}
|
|
}
|
|
@@ -5239,6 +5538,7 @@
|
|
record.cumulative_token_usage = conversation.cumulative_token_usage;
|
|
record.cumulative_cost = conversation.cumulative_cost;
|
|
record.request_token_usage = conversation.request_token_usage;
|
|
+ record.cursor_prompt_usage = boundCursorPromptUsage(conversation.cursor_prompt_usage, conversation.messages);
|
|
}
|
|
//#endregion
|
|
//#region src/runtime/engine/session-options.ts
|
|
@@ -5701,7 +6001,8 @@
|
|
updated_at: conversation.updated_at,
|
|
cumulative_token_usage: deepClone(conversation.cumulative_token_usage ?? {}),
|
|
cumulative_cost: cloneUsageCost(conversation.cumulative_cost),
|
|
- request_token_usage: deepClone(conversation.request_token_usage ?? {})
|
|
+ request_token_usage: deepClone(conversation.request_token_usage ?? {}),
|
|
+ cursor_prompt_usage: boundCursorPromptUsage(conversation.cursor_prompt_usage, conversation.messages)
|
|
};
|
|
}
|
|
function cloneUsageCost(cost) {
|
|
@@ -5771,10 +6072,20 @@
|
|
trimConversationForRuntime(conversation);
|
|
return acpx;
|
|
}
|
|
+function parsePaperclipPiReceipt(raw, costField = "cost_usd") {
|
|
+ const receipt = asRecord$4(raw);
|
|
+ if (receipt?.provenance !== "assistant_message_receipts" && receipt?.provenance !== "assistant_message_and_compaction_receipts") return;
|
|
+ const cost = receipt[costField];
|
|
+ if (cost !== void 0 && !isNonNegativeFiniteNumber(cost)) return;
|
|
+ return { provenance: receipt.provenance, ...cost !== void 0 ? { cost_usd: cost } : {} };
|
|
+}
|
|
function recordPromptResponseUsage(conversation, usage, promptMessageId, timestamp = isoNow()) {
|
|
const tokenUsage = sourceToTokenUsage(usage);
|
|
- if (!tokenUsage) return false;
|
|
- applyTokenUsage(conversation, tokenUsage, promptMessageId);
|
|
+ const receipt = parsePaperclipPiReceipt(asRecord$4(asRecord$4(usage)?._meta)?.paperclipPi, "costUsd");
|
|
+ if (!tokenUsage && !receipt) return false;
|
|
+ if (tokenUsage) applyTokenUsage(conversation, tokenUsage, promptMessageId);
|
|
+ const userId = promptMessageId ?? lastUserMessageId(conversation);
|
|
+ if (receipt && userId) conversation.request_token_usage[userId] = { ...conversation.request_token_usage[userId], paperclip_pi: receipt };
|
|
updateConversationTimestamp(conversation, timestamp);
|
|
trimConversationForRuntime(conversation);
|
|
return true;
|
|
@@ -6136,6 +6447,7 @@
|
|
pendingAgentSessionId = loadState.pendingAgentSessionId;
|
|
sessionModels = loadState.sessionModels;
|
|
const preferenceReplay = await replayFreshSessionPreferences({
|
|
+ reusingLoadedSession,
|
|
client,
|
|
record,
|
|
createdFreshSession,
|
|
@@ -6193,14 +6505,16 @@
|
|
if (shouldReconnect) process.stderr.write(`[acpx] saved session pid ${record.pid} is dead; respawning agent and attempting session reconnect\n`);
|
|
}
|
|
async function replayFreshSessionPreferences(params) {
|
|
- if (!params.createdFreshSession) return {
|
|
+ // A load acknowledgement can reset the model even when the conversation survives.
|
|
+ // Retained connections have no new model acknowledgement to replay.
|
|
+ if (params.reusingLoadedSession) return {
|
|
modelReplay: { replayed: false },
|
|
configReplay: { replayed: false }
|
|
};
|
|
let modelReplay = { replayed: false };
|
|
let configReplay = { replayed: false };
|
|
try {
|
|
- await replayDesiredMode({
|
|
+ if (params.createdFreshSession) await replayDesiredMode({
|
|
client: params.client,
|
|
sessionId: params.sessionId,
|
|
desiredModeId: params.desiredModeId,
|
|
@@ -6219,7 +6533,7 @@
|
|
verbose: params.verbose,
|
|
suppressWarnings: params.suppressWarnings
|
|
});
|
|
- configReplay = await replayDesiredConfigOptions({
|
|
+ if (params.createdFreshSession) configReplay = await replayDesiredConfigOptions({
|
|
client: params.client,
|
|
record: params.record,
|
|
sessionId: params.sessionId,
|
|
@@ -6435,6 +6749,27 @@
|
|
//#region src/runtime/engine/prompt-turn.ts
|
|
const SESSION_REPLY_IDLE_MS = 1e3;
|
|
const SESSION_REPLY_DRAIN_TIMEOUT_MS = 5e3;
|
|
+const TYPED_SESSION_FAILURE_CATEGORIES = /* @__PURE__ */ new Set([
|
|
+ "connection",
|
|
+ "access",
|
|
+ "limit",
|
|
+ "service",
|
|
+ "request",
|
|
+ "unknown"
|
|
+]);
|
|
+function typedTerminalSessionFailureCategory(response) {
|
|
+ if (response === null || typeof response !== "object" || Array.isArray(response)) return null;
|
|
+ const meta = response._meta;
|
|
+ if (meta === null || typeof meta !== "object" || Array.isArray(meta)) return null;
|
|
+ const jetbrains = meta.jetbrains;
|
|
+ if (jetbrains === null || typeof jetbrains !== "object" || Array.isArray(jetbrains)) return null;
|
|
+ const air = jetbrains.air;
|
|
+ if (air === null || typeof air !== "object" || Array.isArray(air)) return null;
|
|
+ if (!Number.isInteger(air.version) || air.version < 1) return null;
|
|
+ const failure = air.sessionFailure;
|
|
+ if (failure === null || typeof failure !== "object" || Array.isArray(failure) || failure.severity !== "error") return null;
|
|
+ return typeof failure.category === "string" && TYPED_SESSION_FAILURE_CATEGORIES.has(failure.category) ? failure.category : "unknown";
|
|
+}
|
|
async function runPromptTurn(params) {
|
|
try {
|
|
const promptPromise = params.client.prompt(params.sessionId, params.prompt, params.onPromptRequestStarted, params.onElicitation);
|
|
@@ -6444,7 +6779,23 @@
|
|
idleMs: SESSION_REPLY_IDLE_MS,
|
|
timeoutMs: SESSION_REPLY_DRAIN_TIMEOUT_MS
|
|
}).catch(() => {});
|
|
+ // A failed terminal result may still include authoritative billable usage.
|
|
recordPromptResponseUsage(params.conversation, response.usage, params.promptMessageId);
|
|
+ recordCursorPromptUsage(params.conversation, response, params.requestId, params.promptMessageId);
|
|
+ const terminalFailureCategory = typedTerminalSessionFailureCategory(response);
|
|
+ if (terminalFailureCategory !== null) {
|
|
+ // Pass complete provider text to the diagnostic callback. Consumers redact
|
|
+ // before bounding it; truncating here can split and expose a credential.
|
|
+ const failure = response._meta.jetbrains.air.sessionFailure;
|
|
+ try {
|
|
+ params.onTerminalSessionFailure?.({
|
|
+ category: terminalFailureCategory,
|
|
+ ...(typeof failure.title === "string" ? { title: failure.title } : {}),
|
|
+ ...(typeof failure.details === "string" ? { details: failure.details } : {})
|
|
+ });
|
|
+ } catch {}
|
|
+ throw new Error(`ACP agent reported a terminal ${terminalFailureCategory} failure.`);
|
|
+ }
|
|
return {
|
|
stopReason: response.stopReason,
|
|
source: "rpc"
|
|
diff --git a/dist/runtime.d.ts b/dist/runtime.d.ts
|
|
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;
|
|
};
|