返回源码地图

packages/acp/acp/src/index.ts

main snapshot · da00f7f5358f · 正文引用章节 15 / 15 / 15;完整原文可核对,不声称全文件人工逐行审计

完整原文供逐行核对;页面收录不代表每行都经过人工语义审核。MIT 许可见 许可证。

1/**
2 * Automation-only Agent Client Protocol server over JSON-RPC stdio.
3 *
4 * The bridge exposes persistent harness sessions to trusted programmatic
5 * clients. It carries standard configuration, MCP mounts, prompt content,
6 * committed semantic updates, cancellation, and one-shot permission decisions;
7 * presentation and human-interaction features stay with the harness's UI modules.
8 *
9 * @module @deepseek-ai/dsh-acp
10 */
11
12import type { Context } from '@deepseek-ai/cordis'
13import { Buffer } from 'node:buffer'
14import { randomUUID } from 'node:crypto'
15import { realpath } from 'node:fs/promises'
16import { isAbsolute, resolve } from 'node:path'
17import { Readable, Writable } from 'node:stream'
18import Schema from '@deepseek-ai/schemastery'
19import { brandString } from '@deepseek-ai/dsh-brand'
20import { errorChain } from '@deepseek-ai/dsh-llm'
21import {
22 agent as createAcpAgentApp,
23 methods,
24 ndJsonStream,
25 PROTOCOL_VERSION,
26 RequestError,
27 type AgentContext,
28 type AuthenticateRequest,
29 type CancelNotification,
30 type CloseSessionRequest,
31 type CloseSessionResponse,
32 type InitializeRequest,
33 type InitializeResponse,
34 type ListSessionsRequest,
35 type ListSessionsResponse,
36 type NewSessionRequest,
37 type NewSessionResponse,
38 type PromptRequest,
39 type PromptResponse,
40 type RequestPermissionRequest,
41 type ResumeSessionRequest,
42 type ResumeSessionResponse,
43 type SetSessionConfigOptionRequest,
44 type SetSessionConfigOptionResponse,
45 type SessionNotification,
46 type Stream,
47} from '@agentclientprotocol/sdk'
48import type { ModelSelection } from '@deepseek-ai/dsh-agent'
49import type { SessionId } from '@deepseek-ai/dsh-session'
50import type {} from '@deepseek-ai/dsh-session-persistence'
51// Side-effect type import: declaration-merges the approval waterfall answered below.
52import type {} from '@deepseek-ai/dsh-user-approval'
53import { supportsAcpImagePrompts } from './content.ts'
54import { AcpMcpConfigError } from './mcp.ts'
55import { AcpModelConfigError } from './model-control.ts'
56import { AcpSession } from './session.ts'
57
58const DEFAULT_SESSION_LIST_PAGE_SIZE = 100
59
60export const name = 'acp'
61/** Core services required by the standard automation controls. */
62export const inject = ['agents', 'llm', 'sessionPersistence', 'sessions']
63
64/** Preserve invalid-parameter detail in the SDK wire error message. */
65function invalidParams(detail: string): RequestError {
66 return RequestError.invalidParams(undefined, detail)
67}
68
69/** Preserve failed-turn detail; plain handler errors become a generic wire internal error. */
70function internalError(detail: string): RequestError {
71 return RequestError.internalError(undefined, detail)
72}
73
74/** Plugin config: the provider/model selection used for each ACP-created agent. */
75export interface AcpConfig {
76 /** Provider route for created agents. */
77 provider?: string
78 /** Model name for created agents. */
79 model?: string
80 /** Maximum summaries returned by one session/list page. */
81 sessionListPageSize?: number
82 /** Runtime-only transport override; production uses stdio. */
83 stream?: Stream
84}
85
86export const Config: Schema<AcpConfig> = Schema.object({
87 provider: Schema.string(),
88 model: Schema.string(),
89 sessionListPageSize: Schema.natural().min(1).default(DEFAULT_SESSION_LIST_PAGE_SIZE),
90})
91
92/**
93 * Mount the automation-only ACP server.
94 * @param ctx - Cordis context carrying the agent factory and session events.
95 * @param config - Initial provider/model selection and optional test transport.
96 */
97export function apply(ctx: Context, config: AcpConfig): void {
98 // ACP handlers execute outside this plugin's injection scope, so capture the
99 // injected service during apply rather than reading it lazily in a callback.
100 const persistence = ctx.sessionPersistence
101 const logger = ctx.logger
102 const sessionListPageSize = resolveSessionListPageSize(config.sessionListPageSize)
103 const sessions = new Map<SessionId, AcpSession>()
104 const activating = new Set<SessionId>()
105 let closed = false
106 let imagePromptEnabled = false
107
108 /** Return the bridge-owned record for an agent, rejecting same-id impostors. */
109 const ownedRecord = (agent: Parameters<AcpSession['owns']>[0]): AcpSession | undefined => {
110 const record = sessions.get(agent.session.id)
111 return record?.owns(agent) === true ? record : undefined
112 }
113
114 const assertOpen = (): void => {
115 if (closed) throw internalError('the ACP bridge has been disposed')
116 }
117
118 const requireSession = (sessionId: SessionId): AcpSession => {
119 const record = sessions.get(sessionId)
120 if (record === undefined) throw invalidParams(`unknown session: ${sessionId}`)
121 return record
122 }
123
124 /** Send one ordered protocol update while containing transport-only failure. */
125 const notify = async (notification: SessionNotification): Promise<void> => {
126 try {
127 await conn.notify(methods.client.session.update, notification)
128 /* v8 ignore start -- the ACP SDK contains notification-handler failures; only a transport write failure reaches this guard. */
129 } catch (error: unknown) {
130 logger.warn(`acp: session/update failed: ${String(error)}`)
131 }
132 /* v8 ignore stop */
133 }
134
135 ctx.on('session/event', (session, event) => {
136 const record = sessions.get(session.header.id)
137 if (record?.ownsSession(session) === true) record.onSessionEvent(session, event)
138 })
139
140 ctx.on('agent/inbox/claimed', ({ agent, message, turn }) => {
141 ownedRecord(agent)?.onInboxClaimed(message, turn)
142 })
143
144 ctx.on('agent/error', ({ agent, turn, error }) => {
145 ownedRecord(agent)?.onAgentError(turn, error)
146 })
147
148 ctx.on('llm/adapters-updated', () => {
149 for (const record of sessions.values()) record.topologyChanged()
150 })
151
152 // Permission requests are a machine policy channel for ACP clients such as
153 // dsh-subagent-acp. The bridge offers one-shot choices only and never infers a
154 // durable grant from an unknown client response.
155 ctx.on('approval/request', (request, next) => {
156 const record = ownedRecord(request.agent)
157 if (record === undefined || request.callId === undefined) return next()
158 const callId = request.callId
159 return record.drainUpdates().then(() => {
160 const params: RequestPermissionRequest = {
161 sessionId: record.agent.session.id,
162 toolCall: { toolCallId: callId },
163 options: [
164 { optionId: 'allow-once', name: 'Allow once', kind: 'allow_once' },
165 { optionId: 'reject-once', name: 'Reject', kind: 'reject_once' },
166 ],
167 }
168 return conn.request(methods.client.session.requestPermission, params)
169 }).then(({ outcome }) => {
170 if (outcome.outcome === 'cancelled') return 'cancelled'
171 return outcome.optionId === 'allow-once' ? 'allowed-once' : 'rejected'
172 })
173 })
174
175 const implementation = {
176 async initialize(_params: InitializeRequest): Promise<InitializeResponse> {
177 // Single-version agent: the spec's "same version if supported, else
178 // the latest supported" both resolve to this server's one version.
179 imagePromptEnabled = await supportsAcpImagePrompts(ctx, config.provider, config.model)
180 return {
181 protocolVersion: PROTOCOL_VERSION,
182 agentInfo: { name: 'deepseek-harness-acp', version: '0.0.1' },
183 agentCapabilities: {
184 mcpCapabilities: { http: true },
185 promptCapabilities: { image: imagePromptEnabled, audio: false, embeddedContext: false },
186 sessionCapabilities: { close: {}, list: {}, resume: {} },
187 },
188 authMethods: [],
189 }
190 },
191
192 authenticate(_params: AuthenticateRequest): Promise<void> {
193 return Promise.resolve()
194 },
195
196 async newSession(params: NewSessionRequest, signal: AbortSignal): Promise<NewSessionResponse> {
197 assertOpen()
198 validateWorkspaceParams(params)
199 const sessionId = brandString<SessionId>(randomUUID())
200 // No preset composition: the ACP bundle keeps the model-facing rows in
201 // the host plane, so this agent reads them from the global layer. A
202 // deployment that configures a roster has to join one here first
203 // (@deepseek-ai/dsh-agent-preset-registry README, "Composing a child agent").
204 let record: AcpSession
205 try {
206 record = await AcpSession.create(ctx, {
207 sessionId,
208 cwd: params.cwd,
209 mcpServers: params.mcpServers,
210 agentOptions: agentOptions(config),
211 fallbackSelection: initialSelection(config),
212 signal,
213 notify,
214 })
215 } catch (error: unknown) {
216 if (error instanceof AcpMcpConfigError) throw invalidParams(error.message)
217 throw error
218 }
219 /* v8 ignore next 4 -- a real stdio close can race an in-flight create. */
220 if (closed) {
221 await record.close('connection closed during session/new')
222 throw internalError('connection closed during session/new')
223 }
224 sessions.set(sessionId, record)
225 try {
226 const configOptions = await record.configOptions(signal)
227 assertOpen()
228 // The attached log writer's flush materializes an empty session durably.
229 await ctx.sessions.flush(record.agent.session)
230 assertOpen()
231 return { sessionId, configOptions }
232 } catch (error: unknown) {
233 sessions.delete(sessionId)
234 await record.close('session/new activation failed')
235 throw error
236 }
237 },
238
239 async resumeSession(params: ResumeSessionRequest, signal: AbortSignal): Promise<ResumeSessionResponse> {
240 assertOpen()
241 validateWorkspaceParams(params)
242 const sessionId = brandString<SessionId>(params.sessionId)
243 if (sessions.has(sessionId) || activating.has(sessionId) || ctx.sessions.get(sessionId) !== undefined) {
244 throw invalidParams(`session is already active: ${sessionId}`)
245 }
246 activating.add(sessionId)
247 return (async (): Promise<ResumeSessionResponse> => {
248 const persisted = (await persistence.stat(sessionId, { signal }))?.header
249 if (persisted === undefined || persisted.origin === 'subagent' || persisted.parentSession !== undefined) {
250 throw invalidParams(`session is not resumable: ${sessionId}`)
251 }
252 if (!await sameDirectory(persisted.cwd, params.cwd)) {
253 throw invalidParams(`session cwd does not match: ${params.cwd}`)
254 }
255 let record: AcpSession
256 try {
257 record = await AcpSession.resume(ctx, {
258 sessionId,
259 cwd: params.cwd,
260 mcpServers: params.mcpServers ?? [],
261 agentOptions: agentOptions(config),
262 fallbackSelection: initialSelection(config),
263 signal,
264 notify,
265 })
266 } catch (error: unknown) {
267 if (error instanceof AcpMcpConfigError) throw invalidParams(error.message)
268 throw error
269 }
270 /* v8 ignore start -- the persisted header was checked before resume; the factory restores that exact header. */
271 if (!await sameDirectory(record.agent.session.header.cwd, params.cwd)) {
272 await record.close('session/resume cwd mismatch')
273 throw invalidParams(`session cwd does not match: ${params.cwd}`)
274 }
275 /* v8 ignore stop */
276 /* v8 ignore next 4 -- a real stdio close can race an in-flight resume. */
277 if (closed) {
278 await record.close('connection closed during session/resume')
279 throw internalError('connection closed during session/resume')
280 }
281 sessions.set(sessionId, record)
282 try {
283 return { configOptions: await record.configOptions(signal) }
284 } catch (error: unknown) {
285 sessions.delete(sessionId)
286 await record.close('session/resume option discovery failed')
287 throw error
288 }
289 })().finally(() => { activating.delete(sessionId) })
290 },
291
292 async listSessions(params: ListSessionsRequest, signal: AbortSignal): Promise<ListSessionsResponse> {
293 assertOpen()
294 if (params.cwd !== undefined && params.cwd !== null && !isAbsolute(params.cwd)) {
295 throw invalidParams(`cwd must be an absolute path: ${params.cwd}`)
296 }
297 let cursor: SessionListCursor | undefined
298 try {
299 cursor = decodeSessionListCursor(params.cursor)
300 } catch (error: unknown) {
301 throw invalidParams((error as Error).message)
302 }
303 const listed = await persistence.list({ signal })
304 const filtered = await Promise.all(listed.map(async ({ header }) => {
305 if (
306 sessions.has(header.id)
307 || activating.has(header.id)
308 || ctx.sessions.get(header.id) !== undefined
309 || header.origin === 'subagent'
310 || header.parentSession !== undefined
311 || header.cwd === undefined
312 || !isAbsolute(header.cwd)
313 ) return undefined
314 if (params.cwd !== undefined && params.cwd !== null && !await sameDirectory(header.cwd, params.cwd)) {
315 return undefined
316 }
317 return { sessionId: header.id, cwd: header.cwd, createdAt: header.createdAt }
318 }))
319 const entries = filtered
320 .filter((entry): entry is NonNullable<typeof entry> => entry !== undefined)
321 .sort((left, right) => right.createdAt - left.createdAt || compareSessionIds(left.sessionId, right.sessionId))
322 const remaining = cursor === undefined
323 ? entries
324 : entries.filter(entry => isAfterSessionListCursor(entry, cursor))
325 const page = remaining.slice(0, sessionListPageSize)
326 const next = remaining.length > page.length ? page.at(-1) : undefined
327 return {
328 sessions: page.map(({ sessionId, cwd }) => ({ sessionId, cwd })),
329 ...next === undefined ? {} : { nextCursor: encodeSessionListCursor(next) },
330 }
331 },
332
333 async setSessionConfigOption(
334 params: SetSessionConfigOptionRequest,
335 signal: AbortSignal,
336 ): Promise<SetSessionConfigOptionResponse> {
337 assertOpen()
338 const record = requireSession(brandString<SessionId>(params.sessionId))
339 try {
340 return { configOptions: await record.setConfig(params.configId, params.value, signal) }
341 } catch (error: unknown) {
342 if (error instanceof AcpModelConfigError) throw invalidParams(error.message)
343 throw error
344 }
345 },
346
347 async closeSession(params: CloseSessionRequest): Promise<CloseSessionResponse> {
348 assertOpen()
349 const sessionId = brandString<SessionId>(params.sessionId)
350 const record = requireSession(sessionId)
351 try {
352 await record.close('ACP session closed')
353 } catch (error: unknown) {
354 throw internalError(`session close failed: ${errorChain(error)}`)
355 } finally {
356 if (sessions.get(sessionId) === record) sessions.delete(sessionId)
357 }
358 return {}
359 },
360
361 async prompt(params: PromptRequest, requestSignal: AbortSignal): Promise<PromptResponse> {
362 assertOpen()
363 const record = requireSession(brandString<SessionId>(params.sessionId))
364 return record.prompt(params, imagePromptEnabled, requestSignal)
365 },
366
367 cancel(params: CancelNotification): Promise<void> {
368 sessions.get(brandString<SessionId>(params.sessionId))?.cancel()
369 return Promise.resolve()
370 },
371 }
372
373 /* v8 ignore next 4 -- production stdio wiring; tests inject config.stream. */
374 const stream: Stream = config.stream ?? ndJsonStream(
375 Writable.toWeb(process.stdout) as WritableStream<Uint8Array>,
376 Readable.toWeb(process.stdin) as ReadableStream<Uint8Array>,
377 )
378 const app = createAcpAgentApp({ name: 'deepseek-harness-acp' })
379 .onRequest(methods.agent.initialize, ({ params }) => implementation.initialize(params))
380 .onRequest(methods.agent.authenticate, async ({ params }) => {
381 await implementation.authenticate(params)
382 return {}
383 })
384 .onRequest(methods.agent.session.new, ({ params, signal }) => implementation.newSession(params, signal))
385 .onRequest(methods.agent.session.list, ({ params, signal }) => implementation.listSessions(params, signal))
386 .onRequest(methods.agent.session.resume, ({ params, signal }) => implementation.resumeSession(params, signal))
387 .onRequest(methods.agent.session.close, ({ params }) => implementation.closeSession(params))
388 .onRequest(methods.agent.session.setConfigOption, ({ params, signal }) => implementation.setSessionConfigOption(params, signal))
389 .onRequest(methods.agent.session.prompt, ({ params, signal }) => implementation.prompt(params, signal))
390 .onNotification(methods.agent.session.cancel, ({ params }) => implementation.cancel(params))
391 const connection = app.connect(stream)
392 const conn: AgentContext = connection.client
393
394 let quiescing: Promise<void> | undefined
395 const quiesce = (): Promise<void> => {
396 if (quiescing !== undefined) return quiescing
397 closed = true
398 const records = [...sessions.values()]
399 // AcpSession.close cancels synchronously before its first await, so every owned
400 // prompt stops before any descendant or persistence drain can block.
401 quiescing = (async () => {
402 const disposals = await Promise.allSettled(records.map(record => record.close('ACP bridge disposed')))
403 for (const record of records) {
404 /* v8 ignore next -- closed blocks concurrent handlers; each captured record remains mapped until this loop. */
405 if (sessions.get(record.agent.session.id) === record) sessions.delete(record.agent.session.id)
406 }
407 const failures: unknown[] = []
408 for (const result of disposals) {
409 if (result.status === 'rejected') failures.push(result.reason as unknown)
410 }
411 if (failures.length > 0) {
412 // The production consumer logs this AggregateError through `String`,
413 // which renders only its message. Embed every per-session diagnostic,
414 // including nested causes and aggregate members, in that message.
415 const detail = failures.map(failure => errorChain(failure)).join('; ')
416 throw new AggregateError(
417 failures,
418 `ACP agent teardown failed for ${failures.length} session(s): ${detail}`,
419 )
420 }
421 })()
422 return quiescing
423 }
424
425 /* v8 ignore start -- production transport rejection and teardown failure. */
426 void connection.closed
427 .catch((error: unknown) => {
428 logger.warn(`acp: connection closed with an error: ${String(error)}`)
429 })
430 .then(quiesce)
431 .catch((error: unknown) => {
432 logger.warn(`acp: connection-close teardown failed: ${String(error)}`)
433 })
434 /* v8 ignore stop */
435
436 ctx.effect(() => quiesce, 'acp.connection')
437}
438
439/**
440 * Build per-agent options from plugin config without assigning absent optional fields.
441 * @param config - ACP provider/model configuration.
442 * @returns the configured fields only.
443 */
444function agentOptions(config: AcpConfig): { provider?: string; model?: string } {
445 return {
446 ...config.provider !== undefined ? { provider: config.provider } : {},
447 ...config.model !== undefined ? { model: config.model } : {},
448 }
449}
450
451/** Initial session selection when both deployment fields are present. */
452function initialSelection(config: AcpConfig): ModelSelection | undefined {
453 return config.provider === undefined || config.model === undefined
454 ? undefined
455 : { provider: config.provider, model: config.model }
456}
457
458interface SessionListCursor {
459 createdAt: number
460 sessionId: string
461}
462
463/** Resolve and validate the deployment-owned session page limit. */
464function resolveSessionListPageSize(value: number | undefined): number {
465 const resolved = value ?? DEFAULT_SESSION_LIST_PAGE_SIZE
466 /* v8 ignore start -- Cordis applies the positive-integer Config schema; this protects direct apply callers. */
467 if (!Number.isSafeInteger(resolved) || resolved < 1) {
468 throw new Error('acp: sessionListPageSize must be a positive safe integer')
469 }
470 /* v8 ignore stop */
471 return resolved
472}
473
474/** Decode an opaque keyset cursor without assigning meaning to client metadata. */
475function decodeSessionListCursor(value: string | null | undefined): SessionListCursor | undefined {
476 if (value === undefined || value === null) return undefined
477 if (!/^[A-Za-z0-9_-]+$/.test(value)) throw new Error('session/list cursor is invalid')
478 try {
479 const decoded: unknown = JSON.parse(Buffer.from(value, 'base64url').toString('utf8'))
480 const createdAt: unknown = Array.isArray(decoded) ? decoded[0] : undefined
481 const sessionId: unknown = Array.isArray(decoded) ? decoded[1] : undefined
482 if (
483 !Array.isArray(decoded)
484 || decoded.length !== 2
485 || typeof createdAt !== 'number'
486 || !Number.isSafeInteger(createdAt)
487 || createdAt < 0
488 || typeof sessionId !== 'string'
489 || sessionId.length === 0
490 ) throw new Error('invalid cursor fields')
491 const canonical = Buffer.from(JSON.stringify(decoded), 'utf8').toString('base64url')
492 if (canonical !== value) throw new Error('non-canonical cursor')
493 return { createdAt, sessionId }
494 } catch (_invalidCursor) {
495 throw new Error('session/list cursor is invalid')
496 }
497}
498
499/** Encode the last returned ordering key as an opaque continuation token. */
500function encodeSessionListCursor(entry: SessionListCursor): string {
501 return Buffer.from(JSON.stringify([entry.createdAt, entry.sessionId]), 'utf8').toString('base64url')
502}
503
504/** Test whether an entry follows the cursor in newest-first list order. */
505function isAfterSessionListCursor(entry: SessionListCursor, cursor: SessionListCursor): boolean {
506 return entry.createdAt < cursor.createdAt
507 || (entry.createdAt === cursor.createdAt && compareSessionIds(entry.sessionId, cursor.sessionId) > 0)
508}
509
510/** Compare opaque session ids by stable UTF-8 bytes, independent of process locale. */
511function compareSessionIds(left: string, right: string): number {
512 return Buffer.compare(Buffer.from(left), Buffer.from(right))
513}
514
515/** Reject workspace features outside the automation contract. */
516function validateWorkspaceParams(params: { cwd: string; additionalDirectories?: string[] | null }): void {
517 if (!isAbsolute(params.cwd)) throw invalidParams(`cwd must be an absolute path: ${params.cwd}`)
518 if (
519 params.additionalDirectories !== undefined
520 && params.additionalDirectories !== null
521 && params.additionalDirectories.length > 0
522 ) {
523 throw invalidParams('additionalDirectories is not supported')
524 }
525}
526
527/** Compare existing directories by physical identity and missing paths lexically. */
528async function sameDirectory(left: string | undefined, right: string): Promise<boolean> {
529 if (left === undefined) return false
530 try {
531 const [realLeft, realRight] = await Promise.all([realpath(left), realpath(right)])
532 return realLeft === realRight
533 } catch (_unresolvablePath) {
534 return resolve(left) === resolve(right)
535 }
536}