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; setSessionConfigOption(sessionId: string, configId: string, value: string): Promise; + requestExtension(method: string, params: Record): Promise>; setSessionModel(sessionId: string, modelId: string, controlOverride?: ModelControlOverride): Promise; private setSessionModelThroughConfig; private setSessionModelThroughLegacyMethod; diff --git a/dist/live-checkpoint-BSIrfgVo.js b/dist/live-checkpoint-BSIrfgVo.js index d454fd7c5bf742b469be75eb8c3988b694a0ffaa..e3a4989473bba20e6dcdccbac086bb6d2e777e9e 100644 --- a/dist/live-checkpoint-BSIrfgVo.js +++ b/dist/live-checkpoint-BSIrfgVo.js @@ -1068,6 +1068,7 @@ function serializeSessionRecordForDisk(record) { last_agent_disconnect_reason: canonical.lastAgentDisconnectReason, protocol_version: canonical.protocolVersion, agent_capabilities: canonical.agentCapabilities, + agent_goal_capability: canonical.agentGoalCapability, title: canonical.title, messages: canonical.messages, updated_at: canonical.updated_at, @@ -1239,7 +1240,8 @@ function parseRequestTokenUsage(raw) { for (const [key, value] of Object.entries(record)) { const parsed = parseTokenUsage(value); if (parsed == null) return null; - usage[key] = parsed; + const receipt = parsePaperclipPiReceipt(asRecord$4(value)?.paperclip_pi); + usage[key] = receipt ? { ...parsed, paperclip_pi: receipt } : parsed; } return usage; } @@ -1541,6 +1543,7 @@ function parseSessionRecord(raw) { lastAgentDisconnectReason: optionals.lastAgentDisconnectReason, protocolVersion: typeof record.protocol_version === "number" ? record.protocol_version : void 0, agentCapabilities: asRecord$4(record.agent_capabilities), + agentGoalCapability: asRecord$4(record.agent_goal_capability), title: conversation.title, messages: conversation.messages, updated_at: conversation.updated_at, @@ -1661,7 +1664,7 @@ const ZED_TAG_KEYS = /* @__PURE__ */ new Set([ "RedactedThinking", "ToolUse" ]); -const MAP_OBJECT_PATHS = /* @__PURE__ */ new Set(["request_token_usage", "messages.Agent.tool_results"]); +const MAP_OBJECT_PATHS = /* @__PURE__ */ new Set(["request_token_usage", "messages.Agent.tool_results", "acpx.session_options.env"]); const OPAQUE_VALUE_PATHS = /* @__PURE__ */ new Set([ "agent_capabilities", "messages.Agent.content.ToolUse.input", @@ -3104,10 +3107,10 @@ function authEnvKey(methodId) { authEnvKeyCache.set(methodId, key); return key; } -function readEnvCredential(methodId) { +function readEnvCredential(methodId, environment = process.env) { const key = authEnvKey(methodId); if (!key) return; - const value = process.env[key]; + const value = environment[key]; if (typeof value === "string" && value.trim().length > 0) return value; } function protectedEnvKey(key) { @@ -3135,8 +3138,21 @@ function promotePrefixedAuthEnvironment(env) { } return protectedKeys; } -function buildAgentEnvironment(authCredentials, sessionEnv) { - const env = { ...process.env }; +function isPlainStringEnvironment(value) { + if (value === null || typeof value !== "object" || Array.isArray(value)) return false; + const prototype = Object.getPrototypeOf(value); + return (prototype === Object.prototype || prototype === null) && Object.values(value).every((entry) => typeof entry === "string"); +} +function resolveAgentEnvironment(spawnEnvironment) { + let sourceEnvironment = process.env; + if (spawnEnvironment !== void 0) { + sourceEnvironment = spawnEnvironment(); + if (!isPlainStringEnvironment(sourceEnvironment)) throw new TypeError("ACPX spawn environment must be a plain record of string values"); + } + return sourceEnvironment; +} +function buildAgentEnvironment(authCredentials, sessionEnv, spawnEnvironment) { + const env = { ...resolveAgentEnvironment(spawnEnvironment) }; const protectedAuthEnvKeys = promotePrefixedAuthEnvironment(env); if (authCredentials) for (const [methodId, credential] of Object.entries(authCredentials)) { addAuthCredentialEnvKeys(protectedAuthEnvKeys, methodId, credential); @@ -3178,10 +3194,10 @@ function resolveConfiguredAuthCredential(methodId, authCredentials) { const configCredentials = authCredentials ?? {}; return configCredentials[methodId] ?? configCredentials[toEnvToken(methodId)]; } -function buildAgentSpawnOptions(cwd, authCredentials, sessionEnv) { +function buildAgentSpawnOptions(cwd, authCredentials, sessionEnv, spawnEnvironment) { return { cwd, - env: buildAgentEnvironment(authCredentials, sessionEnv), + env: buildAgentEnvironment(authCredentials, sessionEnv, spawnEnvironment), stdio: [ "pipe", "pipe", @@ -3961,10 +3977,13 @@ function resolveClientCapabilities(params) { ...params.elicitationModes.includes("url") ? { url: {} } : {} } } : {} }; - if (!params.devinAcp) return baseCapabilities; + const typedSessionFailureMeta = { + jetbrains: { air: { version: 1, capabilities: ["sessionFailure"] } } + }; + if (!params.devinAcp) return { ...baseCapabilities, _meta: typedSessionFailureMeta }; return { ...baseCapabilities, - _meta: DEVIN_COMPATIBILITY_CLIENT_CAPABILITIES_META + _meta: { ...typedSessionFailureMeta, ...DEVIN_COMPATIBILITY_CLIENT_CAPABILITIES_META } }; } function hasResponseField(response, field) { @@ -4067,6 +4086,27 @@ function createNdJsonMessageStream(agentCommand, output, input) { } }) }; } +function boundedExtensionRecord(value) { + if (value === null || typeof value !== "object" || Array.isArray(value)) throw new TypeError("ACP extension payload must be an object"); + const prototype = Object.getPrototypeOf(value); + if (prototype !== Object.prototype && prototype !== null) throw new TypeError("ACP extension payload must be a plain object"); + const serialized = JSON.stringify(value); + if (Buffer.byteLength(serialized) > 256 * 1024) throw new TypeError("ACP extension payload exceeds its bounded size"); + return JSON.parse(serialized); +} +function closedExtensionMethods(value) { + if (value === void 0) return Object.freeze([]); + if (!Array.isArray(value) || value.length > 64) throw new TypeError("ACP extension methods must be a bounded allowlist"); + for (const method of value) { + if (typeof method !== "string" || !/^[A-Za-z_][A-Za-z0-9._/-]{1,255}$/.test(method) || !method.includes("/") || /^(?:session|fs|terminal|elicitation|notifications)\//.test(method)) throw new TypeError("ACP extension method cannot override a standard method"); + } + return Object.freeze([...new Set(value)]); +} +function mergeRuntimeClientCapabilities(base, extra) { + if (extra === void 0) return base; + const additional = boundedExtensionRecord(extra); + return { ...additional, ...base, _meta: { ...additional._meta, ...base._meta } }; +} var AcpClient = class { options; connection; @@ -4105,7 +4145,9 @@ var AcpClient = class { cwd: asAbsoluteCwd(options.cwd), authPolicy: options.authPolicy ?? "skip", permissionPolicy: snapshotPermissionPolicy(options.permissionPolicy), - elicitationModes: normalizeElicitationModes(options.elicitationModes) + elicitationModes: normalizeElicitationModes(options.elicitationModes), + extensionMethods: closedExtensionMethods(options.extensionMethods), + clientCapabilities: options.clientCapabilities === void 0 ? void 0 : boundedExtensionRecord(options.clientCapabilities) }; this.eventHandlers = { onAcpMessage: this.options.onAcpMessage, @@ -4216,9 +4258,21 @@ var AcpClient = class { this.lastAgentExit = void 0; this.lastKnownPid = child.pid ?? void 0; this.attachAgentLifecycleObservers(child); + if (this.options.onAgentSpawn) { + try { + if (!child.pid) throw new Error("ACPX agent spawn did not expose a valid process id."); + await this.options.onAgentSpawn({ pid: child.pid, startedAt: this.agentStartedAt }); + } catch (error) { + try { + child.kill("SIGKILL"); + } catch {} + throw error; + } + } const startupStderr = []; child.stderr.on("data", (chunk) => { this.captureStartupStderr(startupStderr, chunk); + this.options.onAgentStderr?.(String(chunk)); if (!this.options.verbose) return; process.stderr.write(chunk); }); @@ -4226,6 +4280,7 @@ var AcpClient = class { const output = Readable.toWeb(child.stdout); const stream = this.createTappedStream(createNdJsonMessageStream(this.options.agentCommand, input, output)); const connection = this.createConnection(stream, launch); + connection.signal.addEventListener("abort", () => stream.abortResponseDeliveries(), { once: true }); connection.signal.addEventListener("abort", () => { this.recordAgentExit("connection_close", child.exitCode ?? null, child.signalCode ?? null); }, { once: true }); @@ -4253,7 +4308,12 @@ var AcpClient = class { geminiAcp: isGeminiAcpCommand(spawnCommand, args), copilotAcp: isCopilotAcpCommand(spawnCommand, args), claudeAcp: isClaudeAcpCommand(spawnCommand, args), - spawnOptions: buildAgentSpawnOptions(this.options.cwd, this.options.authCredentials, this.options.sessionOptions?.env) + spawnOptions: buildAgentSpawnOptions( + this.options.spawnCwd ?? this.options.cwd, + this.options.authCredentials, + this.options.sessionOptions?.env, + this.options.spawnEnvironment + ) }; } logAgentLaunch(plan) { @@ -4279,10 +4339,17 @@ var AcpClient = class { } async spawnAgentProcess(plan) { const spawnCommand = buildAgentSpawnCommand(plan.spawnCommand, plan.args, process.platform, plan.spawnOptions.env); - const spawnedChild = spawn(spawnCommand.command, spawnCommand.args, { + const options = { ...plan.spawnOptions, windowsVerbatimArguments: spawnCommand.windowsVerbatimArguments - }); + }; + const spawnedChild = this.options.spawnAgent + ? this.options.spawnAgent({ + command: spawnCommand.command, + args: spawnCommand.args, + options + }) + : spawn(spawnCommand.command, spawnCommand.args, options); try { await waitForSpawn$1(spawnedChild); } catch (error) { @@ -4293,10 +4360,10 @@ var AcpClient = class { createConnection(stream, launch) { const app = client({ name: "acpx" }).onNotification(methods.client.session.update, async ({ params }) => { await this.handleSessionUpdate(params); - }).onNotification(methods.client.elicitation.complete, async () => {}).onRequest(methods.client.session.requestPermission, async ({ params }) => { - return await this.handlePermissionRequest(params); + }).onNotification(methods.client.elicitation.complete, async () => {}).onRequest(methods.client.session.requestPermission, async ({ params, requestId, signal }) => { + return await stream.withResponseDelivery(requestId, signal, (responseDelivery) => this.handlePermissionRequest(params, responseDelivery)); }).onRequest(methods.client.elicitation.create, async ({ params, requestId, signal }) => { - return await this.handleElicitationRequest(params, requestId, signal); + return await stream.withResponseDelivery(requestId, signal, (responseDelivery) => this.handleElicitationRequest(params, requestId, signal, responseDelivery)); }).onRequest(methods.client.fs.readTextFile, async ({ params }) => { return await this.handleReadTextFile(params); }).onRequest(methods.client.fs.writeTextFile, async ({ params }) => { @@ -4312,6 +4379,11 @@ var AcpClient = class { }).onRequest(methods.client.terminal.release, async ({ params }) => { return await this.handleReleaseTerminal(params); }); + for (const method of this.options.extensionMethods) { + if (this.options.onExtensionRequest) app.onRequest(method, boundedExtensionRecord, async ({ params, requestId, signal }) => { + return await stream.withResponseDelivery(requestId, signal, (responseDelivery) => this.handleExtensionRequest(method, params, requestId, signal, responseDelivery)); + }); + } if (launch.devinAcp) app.onRequest("_cognition.ai/request_diagnostics", (params) => { return params && typeof params === "object" && !Array.isArray(params) ? params : {}; }, async () => ({})); @@ -4333,12 +4405,12 @@ var AcpClient = class { async initializeProtocolConnection(connection, launch) { const initializePromise = connection.initialize({ protocolVersion: PROTOCOL_VERSION, - clientCapabilities: resolveClientCapabilities({ + clientCapabilities: mergeRuntimeClientCapabilities(resolveClientCapabilities({ devinAcp: launch.devinAcp, fs: this.options.fs !== false, terminal: this.options.terminal !== false, elicitationModes: this.options.elicitationModes ?? [] - }), + }), this.options.clientCapabilities), clientInfo: resolveClientInfo(launch.devinAcp) }); const initialized = launch.geminiAcp ? await withTimeout(initializePromise, resolveGeminiAcpStartupTimeoutMs()) : await initializePromise; @@ -4358,6 +4430,36 @@ var AcpClient = class { 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 @@ var AcpClient = class { 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 @@ var AcpClient = class { } }; return { + withResponseDelivery, abortResponseDeliveries, readable: new ReadableStream({ async start(controller) { const reader = base.readable.getReader(); try { @@ -4387,10 +4491,16 @@ var AcpClient = class { 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 @@ var AcpClient = class { onAcpOutputMessage()?.("outbound", message); onAcpMessage()?.("outbound", message); } + const receipt = id === void 0 ? void 0 : deliveries.get(id); const writer = base.writable.getWriter(); try { await writer.write(message); + if (receipt) { + if (!("result" in message) || receipt.expected === void 0 || JSON.stringify(message.result) !== receipt.expected) receipt.reject(new Error("ACP wrote an error or changed response instead of the selected resolution")); + else receipt.settle(); + } + } catch (error) { + receipt?.reject(new Error("ACP response write failed")); + abortResponseDeliveries(); + throw error; } finally { + if (receipt) { receipt.cleanup(); deliveries.delete(id); } writer.releaseLock(); } } }) @@ -4422,7 +4542,10 @@ var AcpClient = class { const createPromise = this.runConnectionRequest(() => connection.newSession({ cwd: sessionCwd, mcpServers: this.options.mcpServers ?? [], - _meta: buildClaudeCodeOptionsMeta(this.options.sessionOptions, claudeAcp) + _meta: this.isGrokBuildAcpCommand() + ? { rules: typeof this.options.sessionOptions?.systemPrompt === "string" + ? this.options.sessionOptions.systemPrompt : this.options.sessionOptions?.systemPrompt?.append ?? "" } + : buildClaudeCodeOptionsMeta(this.options.sessionOptions, claudeAcp) })); result = claudeAcp ? await withTimeout(createPromise, resolveClaudeAcpSessionCreateTimeoutMs()) : await createPromise; } catch (error) { @@ -4864,7 +4987,7 @@ var AcpClient = class { } selectAuthMethod(methods) { for (const method of methods) { - const envCredential = readEnvCredential(method.id); + const envCredential = readEnvCredential(method.id, resolveAgentEnvironment(this.options.spawnEnvironment)); if (envCredential) return { methodId: method.id, credential: envCredential, @@ -4890,7 +5013,7 @@ var AcpClient = class { } readAgentSpecificEnvCredential(methodId) { if (!this.isGrokBuildAcpCommand() || methodId !== "xai.api_key") return; - const value = process.env.XAI_API_KEY; + const value = (resolveAgentEnvironment(this.options.spawnEnvironment)).XAI_API_KEY; return typeof value === "string" && value.trim().length > 0 ? value : void 0; } selectAgentManagedAuthMethod(methodId) { @@ -4915,9 +5038,9 @@ var AcpClient = class { await connection.authenticate({ methodId: selected.methodId }); this.log(`authenticated with method ${selected.methodId} (${selected.source})`); } - async handlePermissionRequest(params) { + async handlePermissionRequest(params, responseDelivery) { if (this.cancellingSessionIds.has(params.sessionId)) return cancelledPermissionResponse(); - const hostResponse = await this.tryHandlePermissionRequestWithHost(params); + const hostResponse = await this.tryHandlePermissionRequestWithHost(params, responseDelivery); if (hostResponse) return hostResponse; const { response, recorded } = await this.resolvePermissionRequestFromMode(params); if (!recorded) { @@ -4926,14 +5049,49 @@ var AcpClient = class { } return response; } - async handleElicitationRequest(request, requestId, requestSignal) { + async handleExtensionRequest(method, params, requestId, requestSignal, responseDelivery) { + const active = this.activePrompt; + const connection = this.connection; + if (!this.options.extensionMethods.includes(method) || !this.options.onExtensionRequest) throw new Error("ACP extension method is not enabled"); + if (!active || !connection || this.closing || (params.sessionId !== void 0 && params.sessionId !== active.sessionId)) throw new Error("ACP extension request has no matching active prompt"); + if (!(typeof requestId === "string" || typeof requestId === "number" && Number.isFinite(requestId))) throw new Error("ACP extension request id is invalid"); + const signal = AbortSignal.any([requestSignal, connection.signal, active.elicitationController.signal]); + signal.throwIfAborted(); + let onAbort; + const aborted = new Promise((_, reject) => { + onAbort = () => reject(new Error("ACP extension request was cancelled")); + signal.addEventListener("abort", onAbort, { once: true }); + }); + try { + const result = await Promise.race([Promise.resolve().then(() => this.options.onExtensionRequest(method, boundedExtensionRecord(params), { requestId, signal, responseDelivery })), aborted]); + signal.throwIfAborted(); + if (this.activePrompt !== active || this.connection !== connection) throw new Error("ACP extension request outlived its prompt"); + return boundedExtensionRecord(result); + } finally { + signal.removeEventListener("abort", onAbort); + } + } + 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 @@ var AcpClient = class { isElicitationSessionCancelling(sessionId) { return this.closing || this.cancellingSessionIds.has(sessionId); } - async tryHandlePermissionRequestWithHost(params) { + async tryHandlePermissionRequestWithHost(params, responseDelivery) { if (!this.options.onPermissionRequest) return; const signal = this.cancellationSignalForSession(params.sessionId); try { @@ -4978,7 +5136,7 @@ var AcpClient = class { sessionId: params.sessionId, raw: params, inferredKind: inferToolKind(params) - }, { signal }); + }, { signal, responseDelivery }); return this.hostPermissionDecisionResponse(params, signal, decision); } catch (error) { return this.hostPermissionErrorResponse(params, signal, error); @@ -5028,6 +5186,12 @@ var AcpClient = class { attachAgentLifecycleObservers(child) { child.once("exit", (exitCode, signal) => { this.recordAgentExit("process_exit", exitCode, signal); + this.options.onAgentExit?.({ + pid: child.pid ?? this.lastKnownPid, + exitCode, + signal, + exitedAt: isoNow$1() + }); }); child.once("close", (exitCode, signal) => { this.recordAgentExit("process_close", exitCode, signal); @@ -5771,10 +5935,20 @@ function recordSessionUpdate(conversation, state, notification, timestamp = isoN trimConversationForRuntime(conversation); return acpx; } +function parsePaperclipPiReceipt(raw, costField = "cost_usd") { + const receipt = asRecord$4(raw); + if (receipt?.provenance !== "assistant_message_receipts" && receipt?.provenance !== "assistant_message_and_compaction_receipts") return; + const cost = receipt[costField]; + if (cost !== void 0 && !isNonNegativeFiniteNumber(cost)) return; + return { provenance: receipt.provenance, ...cost !== void 0 ? { cost_usd: cost } : {} }; +} function recordPromptResponseUsage(conversation, usage, promptMessageId, timestamp = isoNow()) { const tokenUsage = sourceToTokenUsage(usage); - if (!tokenUsage) return false; - applyTokenUsage(conversation, tokenUsage, promptMessageId); + const receipt = parsePaperclipPiReceipt(asRecord$4(asRecord$4(usage)?._meta)?.paperclipPi, "costUsd"); + if (!tokenUsage && !receipt) return false; + if (tokenUsage) applyTokenUsage(conversation, tokenUsage, promptMessageId); + const userId = promptMessageId ?? lastUserMessageId(conversation); + if (receipt && userId) conversation.request_token_usage[userId] = { ...conversation.request_token_usage[userId], paperclip_pi: receipt }; updateConversationTimestamp(conversation, timestamp); trimConversationForRuntime(conversation); return true; @@ -6435,6 +6609,27 @@ async function withConnectedSession(options) { //#region src/runtime/engine/prompt-turn.ts const SESSION_REPLY_IDLE_MS = 1e3; const SESSION_REPLY_DRAIN_TIMEOUT_MS = 5e3; +const TYPED_SESSION_FAILURE_CATEGORIES = /* @__PURE__ */ new Set([ + "connection", + "access", + "limit", + "service", + "request", + "unknown" +]); +function typedTerminalSessionFailureCategory(response) { + if (response === null || typeof response !== "object" || Array.isArray(response)) return null; + const meta = response._meta; + if (meta === null || typeof meta !== "object" || Array.isArray(meta)) return null; + const jetbrains = meta.jetbrains; + if (jetbrains === null || typeof jetbrains !== "object" || Array.isArray(jetbrains)) return null; + const air = jetbrains.air; + if (air === null || typeof air !== "object" || Array.isArray(air)) return null; + if (!Number.isInteger(air.version) || air.version < 1) return null; + const failure = air.sessionFailure; + if (failure === null || typeof failure !== "object" || Array.isArray(failure) || failure.severity !== "error") return null; + return typeof failure.category === "string" && TYPED_SESSION_FAILURE_CATEGORIES.has(failure.category) ? failure.category : "unknown"; +} async function runPromptTurn(params) { try { const promptPromise = params.client.prompt(params.sessionId, params.prompt, params.onPromptRequestStarted, params.onElicitation); @@ -6444,7 +6639,22 @@ async function runPromptTurn(params) { idleMs: SESSION_REPLY_IDLE_MS, timeoutMs: SESSION_REPLY_DRAIN_TIMEOUT_MS }).catch(() => {}); + // A failed terminal result may still include authoritative billable usage. recordPromptResponseUsage(params.conversation, response.usage, params.promptMessageId); + 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; }) => Promise; + /** Ephemeral initialize capabilities; mandatory base capabilities take precedence. */ + clientCapabilities?: Record; + /** 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, context: { + requestId: string | number; + signal: AbortSignal; + responseDelivery?: Promise; + }) => Promise>; + onExtensionNotification?: (method: string, params: Record) => void; + /** Ephemeral allowlisted environment evaluated immediately before child spawn. */ + spawnEnvironment?: () => Record; + /** 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; + 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; }; type AcpElicitationResponseMeta = { _meta?: Record | null; @@ -116,6 +119,35 @@ type AcpClientOptions = { }; env?: Record; }; + /** Ephemeral initialize capabilities; mandatory base capabilities take precedence. */ + clientCapabilities?: Record; + /** 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, context: { + requestId: string | number; + signal: AbortSignal; + responseDelivery?: Promise; + }) => Promise>; + onExtensionNotification?: (method: string, params: Record) => void; + /** Ephemeral child environment factory; its return value is never persisted. */ + spawnEnvironment?: () => Record; + /** 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; + 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; }) => Promise; }; declare const SESSION_RECORD_SCHEMA: "acpx.session.v1"; @@ -263,12 +297,15 @@ type SessionRecord = { lastAgentDisconnectReason?: string; protocolVersion?: number; agentCapabilities?: AgentCapabilities; + agentGoalCapability?: Record; title?: string | null; messages: SessionMessage[]; updated_at: string; cumulative_token_usage: SessionTokenUsage; cumulative_cost?: SessionUsageCost; - request_token_usage: Record; + request_token_usage: Record; acpx?: SessionAcpxState; importedFrom?: SessionImportedFrom; };