1
0

agent.ts 8.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181
  1. /**
  2. * The concrete Agent implementation: LoopAgent plus its inbox. Everything
  3. * observable happens through session events and the agent/* event taxonomy —
  4. * plugins never need this class.
  5. *
  6. * @module dsh-agent-loop/agent
  7. */
  8. import type { Context } from 'cordis'
  9. import type { AgentId, AgentOptions, AgentStatus, SendOptions } from '@deepseek-ai/dsh-agent'
  10. import type { Agent } from '@deepseek-ai/dsh-agent'
  11. import type { ContentBlock, MessageSource } from '@deepseek-ai/dsh-llm'
  12. import type { Session } from '@deepseek-ai/dsh-session'
  13. import { Inbox } from './inbox.ts'
  14. import { isTurnOpen, lastTurnNumber, runLoop } from './loop.ts'
  15. /**
  16. * The concrete {@link Agent} implementation owned by the agent-loop plugin.
  17. *
  18. * Owns the inbox (queued + steering FIFOs), the per-step AbortController, and
  19. * the loop driver. Everything observable happens through session events and
  20. * the agent/* event taxonomy — plugins never need this class.
  21. */
  22. export class LoopAgent implements Agent {
  23. readonly inbox = new Inbox()
  24. private _status: AgentStatus = 'idle'
  25. private currentAbort: AbortController | undefined
  26. private disposed: Promise<void>
  27. private resolveDisposed!: () => void
  28. /** Resolves when the driver loop has fully exited (tests/disposal). */
  29. done: Promise<void> = Promise.resolve()
  30. constructor(
  31. private ctx: Context,
  32. public readonly id: AgentId,
  33. public readonly options: AgentOptions,
  34. public readonly session: Session,
  35. ) {
  36. const { promise, resolve } = Promise.withResolvers<void>()
  37. this.disposed = promise
  38. this.resolveDisposed = resolve
  39. }
  40. get status(): AgentStatus {
  41. return this._status
  42. }
  43. private setStatus(status: AgentStatus): void {
  44. if (this._status === status || this._status === 'disposed') return
  45. this._status = status
  46. this.ctx.emit('agent/status', this, status)
  47. }
  48. private resolveSource(options?: SendOptions): MessageSource {
  49. return options?.source ?? { kind: 'user' }
  50. }
  51. send(content: ContentBlock[], options?: SendOptions): void {
  52. if (this._status === 'disposed') throw new Error(`agent "${this.id}" is disposed`)
  53. const source = this.resolveSource(options)
  54. this.inbox.enqueue({ content, source })
  55. this.ctx.emit('agent/queued', this, content, { source, steering: false })
  56. }
  57. steer(content: ContentBlock[], options?: SendOptions): void {
  58. if (this._status === 'disposed') throw new Error(`agent "${this.id}" is disposed`)
  59. if (this._status !== 'running') { this.send(content, options); return }
  60. const source = this.resolveSource(options)
  61. this.inbox.steer({ content, source })
  62. this.ctx.emit('agent/queued', this, content, { source, steering: true })
  63. }
  64. inject(content: ContentBlock[], options?: SendOptions): void {
  65. if (this._status === 'disposed') throw new Error(`agent "${this.id}" is disposed`)
  66. const source = this.resolveSource(options)
  67. if (isTurnOpen(this.session)) {
  68. // A turn is open in the LOG (decided from the log, not agent status —
  69. // status can be `running` with no turn open): the context/message is
  70. // turn-enclosed by that turn, so append it directly.
  71. this.session.append('context/message', { content, source })
  72. return
  73. }
  74. // No turn open: wrap the injection in a one-shot turn so every event stays
  75. // turn-enclosed (the durability/replay boundary is the turn).
  76. const turn = lastTurnNumber(this.session) + 1
  77. // Once turn/start enters the log, a turn/end is OWED no matter what — even
  78. // if a throwing `session/event` listener escapes from the turn/start append
  79. // (Session.append pushes the event BEFORE notifying listeners) or the
  80. // context/message append throws (non-serializable content, throwing
  81. // listener). The finally re-checks the log via isTurnOpen() and closes the
  82. // turn if one was actually opened, so the log never carries a permanently
  83. // open injection turn that would corrupt later turns/replay. (If the
  84. // turn/start append throws BEFORE pushing — non-serializable trigger, which
  85. // can't happen for our fixed trigger — no turn was opened and none is owed.)
  86. try {
  87. this.session.append('turn/start', { turn, trigger: { kind: 'injection', source } })
  88. this.session.append('context/message', { content, source })
  89. } finally {
  90. // Close the turn if turn/start made it into the log. Contain a throwing
  91. // turn/end listener: Session.append pushes before notifying, so a throw
  92. // here still leaves turn/end in the log (the turn is balanced) — swallow
  93. // it so it neither replaces the original exception nor skips the flush
  94. // decision below. (It surfaces through the flush path is not needed; the
  95. // turn-balance contract is what matters and it holds.)
  96. if (isTurnOpen(this.session)) {
  97. try {
  98. this.session.append('turn/end', { turn, reason: { kind: 'completed' } })
  99. } catch {
  100. // turn/end is already in the log (pushed before the listener threw),
  101. // so the turn is balanced; the throw is the listener's bug.
  102. }
  103. }
  104. // Decide the durability checkpoint from the LOG, not a flag: a turn was
  105. // recorded iff this turn's turn/start is logged (it may have been closed
  106. // by a throwing-listener turn/end above, which still counts). A
  107. // `turnRecorded` boolean set after append('turn/end') would be skipped by
  108. // a throwing turn/end listener, losing the flush for a balanced in-memory
  109. // turn (crash before the next turn/dispose would drop the idle injection).
  110. const turnRecorded = this.session.events.some(e => e.type === 'turn/start' && e.data.turn === turn)
  111. // Checkpoint the one-shot turn for durability, exactly as the loop does at
  112. // every turn/end. The loop is NOT running (we are idle), so nothing else
  113. // will flush this turn. Fire-and-forget with error containment: inject()
  114. // is synchronous, and a persistence backend failing must not throw into
  115. // the caller (e.g. a tool-bash task-done callback). Disposal still drains
  116. // independently, so a slow flush is safe. A flush failure is reported via
  117. // agent/error (step 0 — the idle-injection convention, there is no real
  118. // step) AND the logger, mirroring the loop's post-turn/end flush path so
  119. // plugins monitoring agent/error see idle-injection persistence failures
  120. // too. A throwing agent/error listener is contained.
  121. if (turnRecorded) {
  122. void Promise.resolve(this.ctx.parallel('session/flush', this.session)).catch((error: unknown) => {
  123. const err = error instanceof Error ? error : new Error(String(error))
  124. this.ctx.logger.warn(`agent "${this.id}": flush after idle injection failed: ${err.message}`)
  125. try {
  126. this.ctx.emit('agent/error', this, turn, 0, err)
  127. } catch {
  128. // contained: the failure is already logged; a throwing agent/error
  129. // listener must not escape this fire-and-forget catch.
  130. }
  131. })
  132. }
  133. }
  134. }
  135. abort(reason?: string): void {
  136. this.currentAbort?.abort(reason ?? 'aborted')
  137. }
  138. /**
  139. * Start the driver loop. Returns a disposer: calling it sets status to
  140. * `disposed`, emits `agent/status('disposed')`, resolves the disposed
  141. * promise (unblocking the idle wait), and aborts the current request if
  142. * any. The returned `agent.done` promise resolves once the loop exits.
  143. */
  144. start(): () => void {
  145. this.done = runLoop(this.ctx, this, {
  146. setStatus: (status) => { this.setStatus(status) },
  147. setAbort: controller => void (this.currentAbort = controller),
  148. disposed: this.disposed,
  149. isDisposed: () => this._status === 'disposed',
  150. })
  151. // The disposer must be infallible: it runs inside the fiber's LIFO
  152. // disposal chain, where a throw would skip later disposers (e.g. the
  153. // registry unregistration) and leave `done` pending forever.
  154. return () => {
  155. if (this._status === 'disposed') return
  156. this._status = 'disposed'
  157. this.resolveDisposed()
  158. this.currentAbort?.abort('disposed')
  159. // setStatus refuses transitions out of 'disposed', so emit directly —
  160. // 'disposed' is part of the agent/status contract. Guarded: a throwing
  161. // listener must not break the disposal chain.
  162. try {
  163. this.ctx.emit('agent/status', this, 'disposed')
  164. } catch {
  165. // listener error during disposal — nothing safe left to do with it
  166. }
  167. }
  168. }
  169. }