control.ts 6.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210
  1. /** Live Session queue, jobs, and projection state with reconnect baselines. */
  2. import type { Context } from '@deepseek-ai/cordis'
  3. import type { Agent, InboxState } from '@deepseek-ai/dsh-agent'
  4. import { Deque } from '@deepseek-ai/dsh-deque'
  5. import type { JobSnapshot } from '@deepseek-ai/dsh-jobs'
  6. import type {
  7. Session, SessionId, UserMessage,
  8. } from '@deepseek-ai/dsh-session'
  9. import type { JsonValue } from '@deepseek-ai/dsh-util-values'
  10. import type {
  11. SessionControlBaseline,
  12. SessionControlFrame,
  13. SessionJob,
  14. SessionProjectionBaseline,
  15. SessionProjectionValues,
  16. SessionQueuedItem,
  17. } from './types.ts'
  18. /** Owns the Host-wide Session control stream. */
  19. export class SessionControlController {
  20. private readonly streams = new Set<ControlQueue>()
  21. /** @param ctx - Host context carrying live Agent, projection, and jobs services. */
  22. constructor(private readonly ctx: Context) {
  23. ctx.sessionProjections.onChanged((session, key, value, seq) => {
  24. this.broadcast({
  25. type: 'projection',
  26. sessionId: session.id,
  27. key,
  28. value: value as JsonValue,
  29. seq,
  30. })
  31. if (key !== 'inbox') return
  32. const agent = this.ctx.agents.get(session.id)
  33. if (agent?.session !== session) return
  34. this.broadcast({
  35. type: 'queue',
  36. sessionId: session.id,
  37. items: queueItemsFromInbox(value as InboxState),
  38. })
  39. })
  40. ctx.inject(['jobs'], (jobsCtx) => {
  41. jobsCtx.jobs.onJobsChanged((owner) => { this.onJobsChanged(owner) })
  42. })
  43. ctx.on('session/created', (session) => {
  44. const jobs = this.jobsFor(this.ctx.agents.get(session.id))
  45. if (jobs.length > 0) this.broadcast({ type: 'jobs', sessionId: session.id, jobs })
  46. })
  47. ctx.effect(() => () => {
  48. for (const stream of this.streams) stream.end()
  49. this.streams.clear()
  50. }, 'session-controller.control')
  51. }
  52. /**
  53. * Open one generation of Host-wide live control state.
  54. * @param signal - Remote stream cancellation.
  55. * @returns one complete baseline followed by live replacement frames.
  56. */
  57. async *control(signal: AbortSignal): AsyncIterable<SessionControlFrame> {
  58. signal.throwIfAborted()
  59. const queue = new ControlQueue()
  60. this.streams.add(queue)
  61. try {
  62. yield { type: 'baseline', value: this.baseline() }
  63. yield* queue.iterate(signal)
  64. } finally {
  65. this.streams.delete(queue)
  66. queue.end()
  67. }
  68. }
  69. private baseline(): SessionControlBaseline {
  70. const sessions = this.ctx.sessions.list()
  71. const queues = Object.create(null) as Record<SessionId, readonly SessionQueuedItem[]>
  72. const jobs = Object.create(null) as Record<SessionId, readonly SessionJob[]>
  73. for (const session of sessions) {
  74. const agent = this.ctx.agents.get(session.id)
  75. queues[session.id] = agent?.session === session ? queueItems(agent) : []
  76. jobs[session.id] = this.jobsFor(agent)
  77. }
  78. return {
  79. queues,
  80. jobs,
  81. projections: this.projectionBaseline(sessions),
  82. }
  83. }
  84. private projectionBaseline(
  85. sessions: readonly Session[],
  86. ): Readonly<Record<SessionId, SessionProjectionBaseline>> {
  87. const blocks = Object.create(null) as Record<SessionId, SessionProjectionBaseline>
  88. for (const session of sessions) {
  89. const snapshot = this.ctx.sessionProjections.snapshot(session)
  90. blocks[session.id] = {
  91. asOfSeq: snapshot.asOfSeq,
  92. // Every projection definition validates its value before snapshot publication.
  93. values: snapshot.values as SessionProjectionValues,
  94. }
  95. }
  96. return blocks
  97. }
  98. private onJobsChanged(owner: Agent | undefined): void {
  99. if (owner !== undefined) {
  100. this.broadcast({ type: 'jobs', sessionId: owner.id, jobs: this.jobsFor(owner) })
  101. return
  102. }
  103. for (const session of this.ctx.sessions.list()) {
  104. this.broadcast({
  105. type: 'jobs',
  106. sessionId: session.id,
  107. jobs: this.jobsFor(this.ctx.agents.get(session.id)),
  108. })
  109. }
  110. }
  111. private jobsFor(agent: Agent | undefined): SessionJob[] {
  112. const jobs = this.ctx.get('jobs')
  113. return jobs === undefined ? [] : jobs.list(agent).map(jobView)
  114. }
  115. private broadcast(frame: SessionControlFrame): void {
  116. for (const stream of this.streams) stream.push(frame)
  117. }
  118. }
  119. class ControlQueue {
  120. private readonly buffer = new Deque<SessionControlFrame>()
  121. private wake: (() => void) | undefined
  122. private done = false
  123. push(frame: SessionControlFrame): void {
  124. if (this.done) return
  125. this.buffer.pushBack(frame)
  126. const wake = this.wake
  127. this.wake = undefined
  128. wake?.()
  129. }
  130. end(): void {
  131. if (this.done) return
  132. this.done = true
  133. const wake = this.wake
  134. this.wake = undefined
  135. wake?.()
  136. }
  137. async *iterate(signal: AbortSignal): AsyncIterable<SessionControlFrame> {
  138. const onAbort = (): void => { this.end() }
  139. signal.addEventListener('abort', onAbort, { once: true })
  140. try {
  141. while (!this.done && !signal.aborted) {
  142. const frame = this.buffer.popFront()
  143. if (frame !== undefined) {
  144. yield frame
  145. continue
  146. }
  147. await new Promise<void>((resolve) => { this.wake = resolve })
  148. }
  149. while (this.buffer.size > 0 && !signal.aborted) yield this.buffer.popFront() as SessionControlFrame
  150. } finally {
  151. signal.removeEventListener('abort', onAbort)
  152. this.end()
  153. }
  154. }
  155. }
  156. function queueItems(agent: Agent): SessionQueuedItem[] {
  157. return queueItemsFromInbox({
  158. 'next-turn': agent.inbox.nextTurn,
  159. 'next-step': agent.inbox.nextStep,
  160. })
  161. }
  162. function queueItemsFromInbox(inbox: InboxState): SessionQueuedItem[] {
  163. return [
  164. ...inbox['next-turn'].map(message => ({
  165. id: message.id,
  166. placement: 'queued' as const,
  167. ...promptRpcId(message),
  168. message: { id: message.id, content: message.content as unknown as JsonValue[] },
  169. })),
  170. ...inbox['next-step'].map(message => ({
  171. id: message.id,
  172. placement: message.source.kind === 'user' ? 'steering' as const : 'context' as const,
  173. ...promptRpcId(message),
  174. message: { id: message.id, content: message.content as unknown as JsonValue[] },
  175. })),
  176. ]
  177. }
  178. /** Prompt-RPC identity carried by a browser-submitted message's user source. */
  179. function promptRpcId(message: UserMessage): Pick<SessionQueuedItem, 'rpcId'> {
  180. const source = message.source
  181. return source.kind === 'user' && 'rpcId' in source ? { rpcId: source.rpcId } : {}
  182. }
  183. function jobView(job: JobSnapshot): SessionJob {
  184. return {
  185. id: job.id,
  186. kind: job.kind,
  187. label: job.label,
  188. status: job.status,
  189. ...(job.detail === undefined ? {} : { detail: job.detail }),
  190. startedAt: job.startedAt,
  191. ...(job.finishedAt === undefined ? {} : { finishedAt: job.finishedAt }),
  192. }
  193. }