index.ts 11 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314
  1. /** Session Remote owner: cold reads, explicit Agent commands, and live control state. */
  2. import { Context } from '@deepseek-ai/cordis'
  3. import z from '@deepseek-ai/schemastery'
  4. import { errorChain } from '@deepseek-ai/dsh-llm'
  5. import type { SessionEvent, SessionHeader, SessionId } from '@deepseek-ai/dsh-session'
  6. import type { SessionObservation } from '@deepseek-ai/dsh-session-query'
  7. import { Remote, TypertRemoteService } from '@deepseek-ai/dsh-typert-protocol'
  8. import {
  9. ApiSessionAgentController,
  10. inspectApiSession,
  11. type ApiSessionAgentResult,
  12. } from './agent.ts'
  13. import { SessionCommandController } from './commands.ts'
  14. import { SessionControlController } from './control.ts'
  15. import { SessionHistoryController } from './history.ts'
  16. import { ApiSessionList, DEFAULT_COLD_BLANK_PROBE_MAX_BYTES } from './list.ts'
  17. import { installModelSelectionProjection } from './model-selection-projection.ts'
  18. import type {
  19. SessionAttachmentRequest,
  20. SessionAttachmentValue,
  21. SessionCancelRequest,
  22. SessionCancelValue,
  23. SessionControlFrame,
  24. SessionCreateRequest,
  25. SessionCreateValue,
  26. SessionFollowFrame,
  27. SessionFollowRequest,
  28. SessionForkRequest,
  29. SessionForkValue,
  30. SessionListRequest,
  31. SessionListValue,
  32. SessionPage,
  33. SessionPageRequest,
  34. SessionPromptRequest,
  35. SessionPromptValue,
  36. SessionRenameRequest,
  37. SessionRenameValue,
  38. SessionSearchRequest,
  39. SessionSearchValue,
  40. SessionSelectModelRequest,
  41. SessionSelectModelValue,
  42. SessionUpdateQueueRequest,
  43. SessionUpdateQueueValue,
  44. } from './types.ts'
  45. export type * from './types.ts'
  46. export { ApiSessionNotFound } from './agent.ts'
  47. declare module '@deepseek-ai/cordis' {
  48. interface Context {
  49. /** Host Session business API and Remote namespace owner. */
  50. sessionController: SessionController
  51. }
  52. }
  53. /** Session Controller deployment policy. */
  54. export interface Config {
  55. /** Maximum cold Session artifact size eligible for one full projection observation. */
  56. readonly coldBlankProbeMaxBytes?: number
  57. }
  58. /** Host service backing the generated `ctx.remote.session` namespace. */
  59. export class SessionController extends TypertRemoteService {
  60. static inject = [
  61. 'agentDefaultModel',
  62. 'agents',
  63. 'attachments',
  64. 'llm',
  65. 'sessions',
  66. 'sessionProjections',
  67. 'sessionQuery',
  68. 'typert',
  69. 'workspaceRegistry',
  70. ]
  71. static Config: z<Config> = z.object({
  72. coldBlankProbeMaxBytes: z.natural().default(DEFAULT_COLD_BLANK_PROBE_MAX_BYTES),
  73. })
  74. private readonly agents: ApiSessionAgentController
  75. private readonly commands: SessionCommandController
  76. private readonly controlState: SessionControlController
  77. private readonly history: SessionHistoryController
  78. private readonly listState: ApiSessionList
  79. private readonly promotions = new Set<Promise<void>>()
  80. /**
  81. * @param ctx - Host context containing the Session capability assembly.
  82. * @param config - cold-list observation policy.
  83. */
  84. constructor(ctx: Context, config: Config) {
  85. super(ctx, 'sessionController', { namespace: 'session' })
  86. installModelSelectionProjection(ctx)
  87. this.agents = new ApiSessionAgentController(ctx)
  88. this.commands = new SessionCommandController(ctx, this.agents, process.cwd())
  89. this.controlState = new SessionControlController(ctx)
  90. // Registered before history so reverse-order teardown closes every
  91. // follower before waiting for already-admitted promotions.
  92. ctx.effect(() => async () => {
  93. await Promise.allSettled([...this.promotions])
  94. }, 'session-controller.promotions')
  95. this.history = new SessionHistoryController(ctx, (observation) => { this.promote(observation) })
  96. this.listState = new ApiSessionList(
  97. ctx,
  98. config.coldBlankProbeMaxBytes ?? DEFAULT_COLD_BLANK_PROBE_MAX_BYTES,
  99. )
  100. ctx.on('session/created', (session) => {
  101. ctx.emit('api-session/added', this.listState.summaryFor(session))
  102. })
  103. ctx.on('session/disposed', (session) => {
  104. ctx.emit('api-session/removed', session.id)
  105. })
  106. ctx.on('agent/status', ({ agent, status }) => {
  107. ctx.emit('api-session/status', agent.id, status === 'running')
  108. })
  109. ctx.on('agent/error', ({ agent, error }) => {
  110. ctx.emit('api-session/error', agent.id, errorChain(error))
  111. })
  112. ctx.on('session/event', (session, event) => {
  113. if (event.type === 'request/header') {
  114. const agent = ctx.agents.get(session.id)
  115. if (agent?.session === session) this.agents.consumeSelection(
  116. agent,
  117. event.data.header.config.provider,
  118. event.data.header.config.model,
  119. event.data.header.config.reasoningEffort,
  120. )
  121. }
  122. if (event.type !== 'user/message' || event.data.source.kind !== 'user') return
  123. ctx.emit('api-session/activity', session.id, event.time)
  124. })
  125. }
  126. private promote(observation: SessionObservation): void {
  127. const sessionId = observation.header.id
  128. const task = (async () => {
  129. using ownedObservation = observation
  130. const result = await this.agents.resolveObservedAgent(ownedObservation)
  131. if ('error' in result) this.ctx.emit('api-session/error', sessionId, result.error.message)
  132. })().catch((error: unknown) => {
  133. this.ctx.logger.error(`session-controller: background activation for "${sessionId}" failed: ${errorChain(error)}`)
  134. })
  135. this.promotions.add(task)
  136. void task.finally(() => { this.promotions.delete(task) })
  137. }
  138. /**
  139. * Resolve or resume one ordinary Session for another Host API domain.
  140. * @param sessionId - Session identity whose Agent owns the operation.
  141. * @returns the live Agent or the stable Session-domain failure.
  142. */
  143. resolveAgent(sessionId: SessionId): Promise<ApiSessionAgentResult> {
  144. return this.agents.resolveAgent(sessionId)
  145. }
  146. /**
  147. * Inspect one attached or persisted Session without activating its Agent.
  148. * @param sessionId - durable Session identity.
  149. * @param signal - optional caller cancellation for persistence reads.
  150. * @returns the current attached state or persisted header and event prefix.
  151. */
  152. inspect(
  153. sessionId: SessionId,
  154. signal?: AbortSignal,
  155. ): Promise<{ meta: SessionHeader; events: SessionEvent[] }> {
  156. const attached = this.ctx.sessions.get(sessionId)
  157. if (attached !== undefined) {
  158. return Promise.resolve({ meta: attached.header, events: [...attached.events] })
  159. }
  160. return inspectApiSession(this.ctx, sessionId, signal)
  161. }
  162. /**
  163. * Read all visible Session rows without resuming an Agent.
  164. * @param _request - reserved empty list request.
  165. * @param signal - cancellation for persistence reads.
  166. * @returns visible Session summaries ordered by activity.
  167. */
  168. @Remote('list')
  169. async list(_request: SessionListRequest, signal: AbortSignal): Promise<SessionListValue> {
  170. return { items: await this.listState.list(signal) }
  171. }
  172. /**
  173. * Search visible Session content without resuming an Agent.
  174. * @param request - literal message-content query.
  175. * @param signal - cancellation for list and search reads.
  176. * @returns authorized bounded Session search results.
  177. */
  178. @Remote('search')
  179. search(request: SessionSearchRequest, signal: AbortSignal): Promise<SessionSearchValue> {
  180. return this.listState.search(request.query, signal)
  181. }
  182. /**
  183. * Create or idempotently adopt one ordinary Session.
  184. * @param request - requested identity, location, and Agent preset.
  185. * @returns the Session identity and resolved preset when configured.
  186. */
  187. @Remote('create')
  188. create(request: SessionCreateRequest): Promise<SessionCreateValue> {
  189. return this.commands.create(request)
  190. }
  191. /**
  192. * Select one Session-local model after explicitly resuming the Session.
  193. * @param request - Session identity and requested model selection.
  194. * @returns the normalized selection installed for the Session.
  195. */
  196. @Remote('selectModel')
  197. selectModel(request: SessionSelectModelRequest): Promise<SessionSelectModelValue> {
  198. return this.commands.selectModel(request)
  199. }
  200. /**
  201. * Rename one Session after explicitly resuming it.
  202. * @param request - Session identity and proposed title.
  203. * @returns the accepted title and durable event sequence.
  204. */
  205. @Remote('rename')
  206. rename(request: SessionRenameRequest): Promise<SessionRenameValue> {
  207. return this.commands.rename(request)
  208. }
  209. /**
  210. * Fork one cold-readable completed-turn prefix into a new Session.
  211. * @param request - source Session and optional event anchor.
  212. * @returns the new Session identity.
  213. */
  214. @Remote('fork')
  215. fork(request: SessionForkRequest): Promise<SessionForkValue> {
  216. return this.commands.fork(request)
  217. }
  218. /**
  219. * Admit one prompt after explicitly resuming its Session.
  220. * @param request - Session identity, prompt content, source metadata, and delivery mode.
  221. * @param signal - caller cancellation before prompt admission begins.
  222. * @returns acknowledgement that the Agent accepted the prompt.
  223. */
  224. @Remote('prompt')
  225. prompt(request: SessionPromptRequest, signal: AbortSignal): Promise<SessionPromptValue> {
  226. signal.throwIfAborted()
  227. return this.commands.prompt(request)
  228. }
  229. /**
  230. * Read one image proven reachable from the addressed Session log.
  231. * @param request - Session and attachment identities used for authorization.
  232. * @returns the durable attachment reference and base64-encoded bytes.
  233. */
  234. @Remote('attachment')
  235. attachment(request: SessionAttachmentRequest): Promise<SessionAttachmentValue> {
  236. return this.commands.attachment(request)
  237. }
  238. /**
  239. * Mutate one still-pending queue occurrence on a live Agent.
  240. * @param request - Session, queue item, and requested mutation.
  241. * @returns acknowledgement that the queue mutation was applied.
  242. */
  243. @Remote('updateQueue')
  244. updateQueue(request: SessionUpdateQueueRequest): SessionUpdateQueueValue {
  245. return this.commands.updateQueue(request)
  246. }
  247. /**
  248. * Cancel one active Agent turn without dropping its pending inbox.
  249. * @param request - Session whose active Agent turn is cancelled.
  250. * @returns acknowledgement that cancellation was requested.
  251. */
  252. @Remote('cancel')
  253. cancel(request: SessionCancelRequest): SessionCancelValue {
  254. return this.commands.cancel(request)
  255. }
  256. /**
  257. * Read one cold-safe, message-aligned Session history page.
  258. * @param request - durable address, backward cursor, and page budget.
  259. * @param signal - cancellation for persistence reads.
  260. * @returns one chronological page.
  261. */
  262. @Remote('page')
  263. page(request: SessionPageRequest, signal: AbortSignal): Promise<SessionPage> {
  264. return this.history.page(request, signal)
  265. }
  266. /**
  267. * Follow one Session log from its opening or resume cursor.
  268. * @param request - durable address and last committed sequence already held by the caller.
  269. * @param signal - cancellation owned by the Remote stream carrier.
  270. * @returns a complete opening snapshot followed by gap-free event frames.
  271. */
  272. @Remote({ mode: 'stream' })
  273. follow(request: SessionFollowRequest, signal: AbortSignal): AsyncIterable<SessionFollowFrame> {
  274. return this.history.follow(request, signal)
  275. }
  276. /**
  277. * Stream a complete live-control baseline followed by replacement frames.
  278. * @param signal - cancellation owned by the Remote stream carrier.
  279. * @returns one complete baseline followed by live replacement frames.
  280. */
  281. @Remote({ mode: 'stream' })
  282. control(signal: AbortSignal): AsyncIterable<SessionControlFrame> {
  283. return this.controlState.control(signal)
  284. }
  285. }
  286. export { buildModelCatalog } from './catalog.ts'
  287. export default SessionController