control.ts 6.7 KB

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