返回源码地图

packages/api/session-controller/src/index.ts

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

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

1/** Session Remote owner: cold reads, explicit Agent commands, and live control state. */
2
3import { hostname } from 'node:os'
4import { resolve } from 'node:path'
5import type { Agent } from '@deepseek-ai/dsh-agent'
6import type {} from '@deepseek-ai/dsh-fs'
7import { Context } from '@deepseek-ai/cordis'
8import z from '@deepseek-ai/schemastery'
9import { errorChain, ReasoningEffortId } from '@deepseek-ai/dsh-llm'
10import type {} from '@deepseek-ai/dsh-client-file-upload'
11import { canOpenNativePath, nativeFileManager, nativeFileApplications, openNativeFileApplication, openNativeAssociatedPath, revealNativePath } from '@deepseek-ai/dsh-native-command'
12import type { SessionId } from '@deepseek-ai/dsh-session'
13import type { SessionInspection } from '@deepseek-ai/dsh-session-persistence'
14import { SessionQueryError, type SessionObservation } from '@deepseek-ai/dsh-session-query'
15import { Remote, RemoteError, TypertRemoteService } from '@deepseek-ai/dsh-typert-protocol'
16import {
17 ApiSessionAgentController,
18 inspectApiSession,
19 type ApiSessionAgentResult,
20} from './agent.ts'
21import { SessionCommandController } from './commands.ts'
22import { SessionControlController } from './control.ts'
23import { SessionHistoryController } from './history.ts'
24import { SessionFileReferences } from './file-references.ts'
25import { ApiSessionList } from './list.ts'
26import { buildModelCatalog, hasProviderApiKey } from './catalog.ts'
27import { installModelSelectionProjection } from './model-selection-projection.ts'
28import { SessionSkillCatalog } from './skill-catalog.ts'
29import { SessionMediaReferences } from './media-references.ts'
30import { ArchivedSessionGate } from './archived-session-gate.ts'
31import type {
32 ModelCatalog,
33 SessionWorkspacePathApplication,
34 SessionAttachmentRequest,
35 SessionAttachmentValue,
36 SessionCancelRequest,
37 SessionCancelValue,
38 SessionControlFrame,
39 SessionCreateRequest,
40 SessionCreateValue,
41 SessionFollowFrame,
42 SessionFollowRequest,
43 SessionForkRequest,
44 SessionForkValue,
45 SessionListRequest,
46 SessionListValue,
47 SessionOpenWorkspacePathRequest,
48 SessionOpenWorkspacePathValue,
49 SessionPage,
50 SessionPageRequest,
51 SessionPromptRequest,
52 SessionPromptValue,
53 SessionRenameRequest,
54 SessionRenameValue,
55 SessionSearchRequest,
56 SessionSearchValue,
57 SessionSelectModelRequest,
58 SessionSelectModelValue,
59 SessionProjectionsRequest,
60 SessionProjectionsValue,
61 SessionProjectionValues,
62 SessionUpdateQueueRequest,
63 SessionUpdateQueueValue,
64} from './types.ts'
65
66export type * from './types.ts'
67export { ApiSessionNotFound } from './agent.ts'
68export { SessionFileReferences } from './file-references.ts'
69export { SessionSkillCatalog } from './skill-catalog.ts'
70
71declare module '@deepseek-ai/cordis' {
72 interface Context {
73 /** Host Session business API and Remote namespace owner. */
74 sessionController: SessionController
75 }
76}
77
78/** Session Controller deployment policy. */
79export interface Config {
80 /** Override platform desktop-opener detection. */
81 readonly nativeOpen?: boolean
82 /** Positive integral milliseconds of list work before yielding between complete rows. */
83 readonly listWorkSliceMs?: number
84}
85
86/** Deployment policy after schema defaults have been applied. */
87type ResolvedConfig = Config & { readonly listWorkSliceMs: number }
88
89/** Host integrations replaceable by direct unit tests. */
90export interface SessionControllerInternals {
91 /** Native default-application handoff. */
92 readonly openPath?: (path: string, signal: AbortSignal) => Promise<void>
93 /** Native file-association query. */
94 readonly fileApplications?: typeof nativeFileApplications
95 /** Explicit registered-application handoff. */
96 readonly openFileApplication?: typeof openNativeFileApplication
97 /** Native file-manager handoff. */
98 readonly revealPath?: (path: string, signal: AbortSignal) => Promise<void>
99 /** Native handoff availability probe. */
100 readonly canOpenPath?: () => boolean
101}
102
103/** Host service backing the generated `ctx.remote.session` namespace. */
104export class SessionController extends TypertRemoteService {
105 static inject = [
106 'agentDefaultModel',
107 'agents',
108 'attachments',
109 'fileUploads',
110 'fs',
111 'llm',
112 'sessions',
113 'sessionProjections',
114 'sessionQuery',
115 'typert',
116 'workspaceRegistry',
117 ]
118
119 static Config: z<Config, ResolvedConfig> = z.object({
120 nativeOpen: z.boolean(),
121 listWorkSliceMs: z.natural().min(1).default(16),
122 })
123
124 private readonly agents: ApiSessionAgentController
125 private readonly commands: SessionCommandController
126 private readonly controlState: SessionControlController
127 private readonly history: SessionHistoryController
128 private readonly listState: ApiSessionList
129 private readonly openPath: (path: string, signal: AbortSignal) => Promise<void>
130 private readonly fileApplications: typeof nativeFileApplications
131 private readonly openFileApplication: typeof openNativeFileApplication
132 private readonly revealPath: (path: string, signal: AbortSignal) => Promise<void>
133 private readonly canOpenPath: () => boolean
134 private readonly promotions = new Set<Promise<void>>()
135
136 /**
137 * @param ctx - Host context containing the Session capability assembly.
138 * @param config - native-opener and list-scheduling deployment policy.
139 * @param internals - host integrations replaceable by direct unit tests.
140 */
141 constructor(ctx: Context, config: Config, internals: SessionControllerInternals = {}) {
142 super(ctx, 'sessionController', { namespace: 'session' })
143 const resolved = SessionController.Config(config)
144 installModelSelectionProjection(ctx)
145 this.agents = new ApiSessionAgentController(ctx)
146 this.commands = new SessionCommandController(ctx, this.agents, process.cwd())
147 ctx.effect(() => ctx.fileUploads.registerAgentResolver(async (sessionId) => {
148 const result = await this.agents.resolveAgent(sessionId)
149 if ('error' in result) throw result.error
150 return result.agent
151 }), 'session-controller: file-upload Agent resolver')
152 this.controlState = new SessionControlController(ctx)
153 // Registered before history so reverse-order teardown closes every
154 // follower before waiting for already-admitted promotions.
155 ctx.effect(() => async () => {
156 await Promise.allSettled([...this.promotions])
157 }, 'session-controller.promotions')
158 this.history = new SessionHistoryController(ctx, (observation) => { this.promote(observation) })
159 this.listState = new ApiSessionList(ctx, resolved.listWorkSliceMs)
160 this.fileApplications = internals.fileApplications ?? nativeFileApplications
161 this.openFileApplication = internals.openFileApplication ?? openNativeFileApplication
162 this.openPath = internals.openPath ?? openNativeAssociatedPath
163 this.revealPath = internals.revealPath ?? revealNativePath
164 this.canOpenPath = internals.canOpenPath
165 ?? (() => config.nativeOpen ?? (internals.openPath !== undefined || canOpenNativePath()))
166 ctx.plugin(SessionFileReferences)
167 ctx.plugin(SessionMediaReferences)
168 ctx.plugin(SessionSkillCatalog)
169 // An archived Session, or a subagent descendant of one, runs no model step
170 // until it is restored; what it still runs is stopped by the owners that
171 // answer the Workspace registry's archive-admission events.
172 ctx.plugin(ArchivedSessionGate)
173
174 ctx.on('session/created', (session) => {
175 ctx.emit('api-session/added', this.listState.summaryFor(session))
176 })
177 ctx.on('session/disposed', (session) => {
178 ctx.emit('api-session/removed', session.id)
179 })
180 const publishAgentAvailability = ({ agent }: { agent: Agent }): undefined => {
181 if (ctx.sessions.get(agent.id) === agent.session) {
182 ctx.emit('api-session/added', this.listState.summaryFor(agent.session))
183 }
184 }
185 ctx.on('agent/created', publishAgentAvailability)
186 ctx.on('agent/disposed', publishAgentAvailability)
187 ctx.on('agent/status', ({ agent, status }) => {
188 ctx.emit('api-session/status', agent.id, status === 'running')
189 })
190 ctx.on('agent/error', ({ agent, error }) => {
191 ctx.emit('api-session/error', agent.id, errorChain(error))
192 })
193 ctx.on('session/event', (session, event) => {
194 if (event.type === 'request/header') {
195 const agent = ctx.agents.get(session.id)
196 if (agent?.session === session) this.agents.consumeSelection(
197 agent,
198 event.data.header.config.provider,
199 event.data.header.config.model,
200 event.data.header.config.reasoningEffort,
201 )
202 }
203 if (event.type !== 'user/message' || event.data.source.kind !== 'user') return
204 ctx.emit('api-session/activity', session.id, event.time)
205 })
206 }
207
208 private promote(observation: SessionObservation): void {
209 const sessionId = observation.header.id
210 const task = (async () => {
211 using ownedObservation = observation
212 const result = await this.agents.resolveObservedAgent(ownedObservation)
213 if ('error' in result) this.ctx.emit('api-session/error', sessionId, result.error.message)
214 })().catch((error: unknown) => {
215 this.ctx.logger.error(`session-controller: background activation for "${sessionId}" failed: ${errorChain(error)}`)
216 })
217 this.promotions.add(task)
218 void task.finally(() => { this.promotions.delete(task) })
219 }
220
221 /**
222 * Resolve or resume one ordinary Session for another Host API domain.
223 * @param sessionId - Session identity whose Agent owns the operation.
224 * @returns the live Agent or the stable Session-domain failure.
225 */
226 resolveAgent(sessionId: SessionId): Promise<ApiSessionAgentResult> {
227 return this.agents.resolveAgent(sessionId)
228 }
229
230 /**
231 * Inspect one attached or persisted Session without activating its Agent.
232 * @param sessionId - durable Session identity.
233 * @param signal - optional caller cancellation for persistence reads.
234 * @returns the current attached state or persisted header and event prefix.
235 */
236 inspect(
237 sessionId: SessionId,
238 signal?: AbortSignal,
239 ): Promise<SessionInspection> {
240 const attached = this.ctx.sessions.get(sessionId)
241 if (attached !== undefined) {
242 return Promise.resolve({
243 meta: attached.header,
244 inheritedEventCount: attached.inheritedEventCount,
245 // oxlint-disable-next-line typescript/no-deprecated -- Existing Session history read; migration deferred.
246 events: attached.snapshotEvents(),
247 })
248 }
249 return inspectApiSession(this.ctx, sessionId, signal)
250 }
251
252 /**
253 * Read all visible Session rows without resuming an Agent.
254 * @param _request - reserved empty list request.
255 * @param signal - cancellation for persistence reads and summary generation.
256 * @returns visible Session summaries ordered by activity.
257 */
258 @Remote('list')
259 async list(_request: SessionListRequest, signal: AbortSignal): Promise<SessionListValue> {
260 return { items: await this.listState.list(signal) }
261 }
262
263 /**
264 * Search visible Session content without resuming an Agent.
265 * @param request - literal message-content query.
266 * @param signal - cancellation for list and search reads.
267 * @returns authorized bounded Session search results.
268 */
269 @Remote('search')
270 search(request: SessionSearchRequest, signal: AbortSignal): Promise<SessionSearchValue> {
271 return this.listState.search(request.query, signal)
272 }
273
274 /**
275 * Create or idempotently adopt one ordinary Session.
276 * @param request - requested identity, location, and Agent preset.
277 * @returns the Session identity and resolved preset when configured.
278 */
279 @Remote('create')
280 create(request: SessionCreateRequest): Promise<SessionCreateValue> {
281 return this.commands.create(request)
282 }
283
284 /**
285 * Select one Session-local model after explicitly resuming the Session; save the default in the background.
286 * @param request - Session identity and requested model selection.
287 * @returns the normalized selection installed for the Session, without waiting for default persistence.
288 */
289 @Remote('selectModel')
290 selectModel(request: SessionSelectModelRequest): Promise<SessionSelectModelValue> {
291 return this.commands.selectModel(request)
292 }
293
294 /**
295 * Select the first available account model after login when no provider API key is configured.
296 * @returns after saving the first available model or retaining the existing default.
297 */
298 @Remote
299 async initializeDefaultModel(): Promise<void> {
300 const provider = 'deepseek-account'
301 if (await hasProviderApiKey(this.ctx)) return
302 const catalog = await buildModelCatalog(this.ctx)
303 const model = catalog.groups.find(group => group.id === provider)?.models[0]
304 if (model === undefined) throw new RemoteError('session/provider-models-unavailable',
305 `provider "${provider}" has no available models`, { provider })
306 const selection = { provider, model: model.id,
307 ...model.reasoning?.defaultEffort === undefined ? {} : { reasoningEffort: ReasoningEffortId(model.reasoning.defaultEffort) },
308 }
309 await this.ctx.agentDefaultModel.saveSelection(selection)
310 }
311
312 /**
313 * Describe every currently routable model for Host-generation selectors.
314 * @returns provider-grouped models, the deployment default, and isolated provider failures.
315 */
316 @Remote('modelCatalog')
317 modelCatalog(): Promise<ModelCatalog> {
318 return buildModelCatalog(this.ctx)
319 }
320
321 /**
322 * Report whether this deployment can hand a Session workspace path to a native desktop.
323 * @returns true when the matching open operation is available.
324 */
325 @Remote
326 canOpenWorkspacePath(): boolean {
327 return this.canOpenPath()
328 }
329
330 /**
331 * Describe the serving desktop for authenticated file-action routes.
332 * @returns Host name, configured availability, and platform-specific file-manager behavior.
333 */
334 workspaceDesktop(): { name: string; available: boolean; fileManager: 'finder' | 'explorer' | 'directory' | null } {
335 const fileManager = nativeFileManager()
336 return { name: hostname(), available: fileManager !== null && this.canOpenPath(), fileManager }
337 }
338
339 /**
340 * Verify one path through the composed filesystem and open it on the Host desktop.
341 * @param request - path after best-effort Session workspace resolution.
342 * @param signal - caller lifetime; abort terminates the native command.
343 * @returns confirmation after the native opener accepts the path.
344 * @throws RemoteError when the request is invalid, has no verified Host mapping, is cancelled, or the opener fails.
345 */
346 @Remote('openWorkspacePath')
347 async openWorkspacePath(
348 request: SessionOpenWorkspacePathRequest,
349 signal: AbortSignal,
350 ): Promise<SessionOpenWorkspacePathValue> {
351 try {
352 const path = await this.verifyDesktopPath(request.path, signal)
353 if (request.action === 'reveal') await this.revealPath(path, signal)
354 else if (request.application !== undefined) await this.openFileApplication(path, request.application, signal)
355 else await this.openPath(path, signal)
356 return { opened: true }
357 } catch (error: unknown) {
358 if (signal.aborted) throw new RemoteError('gateway/cancelled', 'path open was aborted', {})
359 if (error instanceof RemoteError) throw error
360 throw new RemoteError(
361 'gateway/internal',
362 'path open failed',
363 {},
364 { cause: error },
365 )
366 }
367 }
368
369 /**
370 * Query current file handlers on the serving desktop without activating an Agent.
371 * @param request - file path in Host filesystem syntax.
372 * @param signal - caller lifetime, propagated to filesystem and desktop queries.
373 * @returns OS application names, icons, and default selection; empty when desktop opening is unavailable.
374 * @throws RemoteError when the path is invalid, the query is cancelled, or native discovery fails.
375 */
376 @Remote('workspacePathApplications')
377 async workspacePathApplications(
378 request: { readonly path: string }, signal: AbortSignal,
379 ): Promise<readonly SessionWorkspacePathApplication[]> {
380 if (!this.canOpenPath()) return []
381 try {
382 const path = await this.verifyDesktopPath(request.path, signal)
383 return await this.fileApplications(path, signal)
384 } catch (error: unknown) {
385 if (signal.aborted) throw new RemoteError('gateway/cancelled', 'application query was aborted', {})
386 if (error instanceof RemoteError) throw error
387 throw new RemoteError('gateway/internal', 'file application query failed', {}, { cause: error })
388 }
389 }
390
391 private async verifyDesktopPath(path: string, signal: AbortSignal): Promise<string> {
392 if (path.length === 0) throw new RemoteError('gateway/bad-request', 'A non-empty file path is required', {})
393 signal.throwIfAborted()
394 const hostPath = resolve(path)
395 const { fs } = this.ctx
396 const mapped = fs.processPathFromHostPath(hostPath)
397 if (mapped === undefined || fs.processPath(await fs.resolve(mapped, { signal })) !== hostPath) {
398 throw new RemoteError('gateway/bad-request', 'Path has no verified Host path', {})
399 }
400 signal.throwIfAborted()
401 return hostPath
402 }
403
404 /**
405 * Rename one Session after explicitly resuming it.
406 * @param request - Session identity and proposed title.
407 * @returns the accepted title and durable event sequence.
408 */
409 @Remote('rename')
410 rename(request: SessionRenameRequest): Promise<SessionRenameValue> {
411 return this.commands.rename(request)
412 }
413
414 /**
415 * Fork one cold-readable exact event prefix into a new Session. An omitted
416 * boundary selects the latest completed-turn prefix; an open cut receives
417 * synthetic fork closers.
418 * @param request - source Session and optional exact inclusive event boundary.
419 * @returns the new Session identity.
420 */
421 @Remote('fork')
422 fork(request: SessionForkRequest): Promise<SessionForkValue> {
423 return this.commands.fork(request)
424 }
425
426 /**
427 * Admit one prompt after explicitly resuming its Session.
428 * @param request - Session identity, prompt content, source metadata, and delivery mode.
429 * @param signal - caller cancellation before prompt admission begins.
430 * @returns acknowledgement that the Agent accepted the prompt.
431 */
432 @Remote('prompt')
433 prompt(request: SessionPromptRequest, signal: AbortSignal): Promise<SessionPromptValue> {
434 signal.throwIfAborted()
435 return this.commands.prompt(request)
436 }
437
438 /**
439 * Read one image proven reachable from the addressed Session log.
440 * @param request - Session and attachment identities used for authorization.
441 * @returns the durable attachment reference and base64-encoded bytes.
442 */
443 @Remote('attachment')
444 attachment(request: SessionAttachmentRequest): Promise<SessionAttachmentValue> {
445 return this.commands.attachment(request)
446 }
447
448 /**
449 * Mutate one still-pending queue occurrence, resuming a cold Agent first.
450 * @param request - Session, queue item, and requested mutation.
451 * @returns acknowledgement that the queue mutation was applied.
452 */
453 @Remote('updateQueue')
454 updateQueue(request: SessionUpdateQueueRequest): Promise<SessionUpdateQueueValue> {
455 return this.commands.updateQueue(request)
456 }
457
458 /**
459 * Cancel one active Agent turn without dropping its pending inbox.
460 * @param request - Session whose active Agent turn is cancelled.
461 * @returns acknowledgement that cancellation was requested.
462 */
463 @Remote('cancel')
464 cancel(request: SessionCancelRequest): SessionCancelValue {
465 return this.commands.cancel(request)
466 }
467
468 /**
469 * Read one cold-safe, message-aligned Session history page.
470 * @param request - durable address, backward cursor, and page budget.
471 * @param signal - cancellation for persistence reads.
472 * @returns one chronological page.
473 */
474 @Remote('page')
475 page(request: SessionPageRequest, signal: AbortSignal): Promise<SessionPage> {
476 return this.history.page(request, signal)
477 }
478
479 /**
480 * Follow one Session log from its opening or resume cursor.
481 * @param request - durable address and last committed sequence already held by the caller.
482 * @param signal - cancellation owned by the Remote stream carrier.
483 * @returns a complete opening snapshot followed by gap-free durable event
484 * frames and optional cursorless assistant-stream frames.
485 */
486 @Remote({ mode: 'stream' })
487 follow(request: SessionFollowRequest, signal: AbortSignal): AsyncIterable<SessionFollowFrame> {
488 return this.history.follow(request, signal)
489 }
490
491 /**
492 * Read all registered projections without activating an Agent.
493 * @param request - Session whose current values are required.
494 * @param signal - cancellation for the Session observation.
495 * @returns complete baseline, or null when the Session does not exist.
496 */
497 @Remote('projections')
498 async projections(request: SessionProjectionsRequest, signal: AbortSignal): Promise<SessionProjectionsValue> {
499 const { sessionId } = request
500 if (sessionId.length === 0) {
501 throw new RemoteError('gateway/bad-request', 'sessionId must not be empty', {})
502 }
503 try {
504 using observation = await this.ctx.sessionQuery.observeSession(sessionId, { signal })
505 const projections = observation.projections
506 if (projections === undefined) {
507 throw new RemoteError('session/projections-unavailable', 'Session projections are unavailable', {})
508 }
509 return { asOfSeq: projections.asOfSeq, values: projections.values as SessionProjectionValues }
510 } catch (error: unknown) {
511 if (error instanceof SessionQueryError && error.code === 'SESSION_QUERY_SESSION_NOT_FOUND') return null
512 if (signal.aborted
513 || (error instanceof SessionQueryError && error.code === 'SESSION_QUERY_ABORTED')) {
514 throw new RemoteError('gateway/cancelled', 'Session projection read was cancelled', {}, { cause: error })
515 }
516 if (error instanceof RemoteError) throw error
517 throw new RemoteError('gateway/internal', 'Session projection read failed', {}, { cause: error })
518 }
519 }
520
521 /**
522 * Stream a complete live-control baseline followed by replacement frames.
523 * @param signal - cancellation owned by the Remote stream carrier.
524 * @returns one complete baseline followed by live replacement frames.
525 */
526 @Remote({ mode: 'stream' })
527 control(signal: AbortSignal): AsyncIterable<SessionControlFrame> {
528 return this.controlState.control(signal)
529 }
530
531
532}
533
534export { buildModelCatalog }
535export default SessionController