agent.ts 21 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449
  1. /**
  2. * The concrete Agent implementation: ReactLoopAgent 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 { agentEvents } from '@deepseek-ai/dsh-agent'
  10. import type { AgentOptions, AgentStatus, SendOptions } from '@deepseek-ai/dsh-agent'
  11. import type { Agent } from '@deepseek-ai/dsh-agent'
  12. import { deepFreeze } from '@deepseek-ai/dsh-llm'
  13. import type { ContentBlock, MessageSource } from '@deepseek-ai/dsh-llm'
  14. import { snapshotJsonValue, type Session, type SessionId } from '@deepseek-ai/dsh-session'
  15. import { Inbox, type InboxMessage } from './inbox.ts'
  16. import { isTurnOpen, lastTurnNumber, runLoop } from './loop.ts'
  17. /** Sessions already claimed by a concrete driver construction. */
  18. const claimedDriverSessions = new WeakSet<Session>()
  19. /** Module-private driver entry: its symbol is absent from the package surface. */
  20. const startDriver = Symbol('dsh.agent-loop.start-driver')
  21. /** Module-private quiescent stop, valid both before and after driver start. */
  22. const stopDriver = Symbol('dsh.agent-loop.stop-driver')
  23. /** Module-private context binding for the mutually referential agent scope. */
  24. const bindContext = Symbol('dsh.agent-loop.bind-context')
  25. /** Module-private publication marker. */
  26. const publishAgent = Symbol('dsh.agent-loop.publish-agent')
  27. /** Factory-owned controls that can operate only on the agent created with them. */
  28. export interface PreparedReactLoopAgent {
  29. /** The unpublished concrete agent. */
  30. agent: ReactLoopAgent
  31. /** Mark the agent public so teardown emits its status lifecycle. */
  32. markPublished(): void
  33. /** Stop the prepared instance even when publication has not started its loop. */
  34. dispose(): Promise<void> | void
  35. /**
  36. * Start its driver after publication and session-start notification.
  37. * The returned disposer reaches quiescence for both the loop and every
  38. * fire-and-forget idle-injection flush the agent started.
  39. */
  40. startDriver(): () => Promise<void> | void
  41. }
  42. /**
  43. * Construct one concrete agent together with unforgeable, instance-bound
  44. * lifecycle controls. The package surface deliberately exposes neither source
  45. * subpaths nor this helper: setup code may identify the concrete class, but it
  46. * cannot publish or start the factory's unpublished instance.
  47. * @param ctx - the agent-loop service context used for driving and events.
  48. * @param id - the concrete agent identity.
  49. * @param options - loop options for the agent.
  50. * @param session - the prepared session the agent will own.
  51. * @returns the agent and closures bound only to that exact instance.
  52. */
  53. export function prepareReactLoopAgent(
  54. ctx: Context, id: SessionId, options: AgentOptions, session: Session,
  55. ): PreparedReactLoopAgent {
  56. if (claimedDriverSessions.has(session)) {
  57. throw new Error(`session "${session.id}" already has a concrete agent driver`)
  58. }
  59. const agent = new ReactLoopAgent(ctx, id, options, session)
  60. claimedDriverSessions.add(session)
  61. const dispose = () => agent[stopDriver]()
  62. return {
  63. agent,
  64. markPublished: () => { agent[publishAgent]() },
  65. dispose,
  66. startDriver: () => {
  67. agent[startDriver]()
  68. return dispose
  69. },
  70. }
  71. }
  72. /**
  73. * Install the concrete agent's scope context exactly once. Construction and
  74. * scope minting are mutually referential (the scope key is the agent), so the
  75. * factory performs this one post-construction binding before setup receives
  76. * the unpublished agent. The module-private binding rejects a second bind.
  77. * @param agent - the unpublished concrete agent to bind.
  78. * @param ctx - its fully extended agent scope context.
  79. */
  80. export function bindReactLoopAgentContext(agent: ReactLoopAgent, ctx: Context): void {
  81. agent[bindContext](ctx)
  82. }
  83. /**
  84. * The concrete {@link Agent} implementation owned by the agent-loop plugin.
  85. *
  86. * Owns the inbox (queued + steering FIFOs), the per-step AbortController, and
  87. * the loop driver. Everything observable happens through session events and
  88. * the agent/* event taxonomy — plugins never need this class.
  89. */
  90. export class ReactLoopAgent implements Agent {
  91. /** Queued + steering FIFOs; native-private so callers cannot bypass the public driving verbs. */
  92. readonly #inbox = new Inbox()
  93. /**
  94. * The agent's scope context ({@link Agent.ctx}), wired by the factory right
  95. * after the scope is minted — before the agent is registered, announced, or
  96. * driven, so no consumer can observe it unset. Definite-assignment (`!`)
  97. * expresses that two-phase construction: the agent object and its scope
  98. * context are mutually referential (the scope is keyed BY this agent), so
  99. * neither can exist strictly before the other.
  100. */
  101. private boundContext: Context | undefined
  102. /** The agent's scoped composition context, bound once by its factory. */
  103. get ctx(): Context {
  104. if (this.boundContext === undefined) throw new Error(`agent "${this.id}" context is not bound`)
  105. return this.boundContext
  106. }
  107. private _status: AgentStatus = 'idle'
  108. private currentAbort: AbortController | undefined
  109. /** Whether runLoop has been installed into {@link done}. */
  110. private driverStarted = false
  111. /** Whether registry publication began and status disposal is externally visible. */
  112. private published = false
  113. /**
  114. * Turn-scoped cancel marker, set by {@link cancel} and read/cleared by the
  115. * driver loop (via the LoopHandle) at every point a turn could start or
  116. * continue. Armed ONLY when there is something to cancel (a running turn, an
  117. * in-flight step, or queued/steering work), so an idle no-op cancel cannot
  118. * leave it set to wrongly drop a later prompt.
  119. */
  120. private cancelRequested = false
  121. /**
  122. * The resolved reason for the pending {@link cancel} (`reason ?? 'cancelled'`),
  123. * read by the driver loop's marker branches so a turn dropped in a
  124. * marker-only window (pre-step / continuation, where no `AbortController`
  125. * carries the reason) ends with the SAME `{kind:'aborted', reason}` the
  126. * mid-step abort path produces from `abort.signal.reason`. Without this the
  127. * caller's `cancel(reason)` would be silently replaced by the literal
  128. * 'cancelled' whenever the cancel landed outside a running step — making the
  129. * logged reason race-dependent and the public `reason?` param half-effective.
  130. */
  131. private cancelReason = 'cancelled'
  132. private disposed: Promise<void>
  133. private resolveDisposed!: () => void
  134. /** Resolves when the driver loop has fully exited (tests/disposal). */
  135. done: Promise<void> = Promise.resolve()
  136. /**
  137. * Pending {@link whenIdle} waiters, resolved by {@link settleIdleWaiters} when
  138. * the agent next settles out of `running`. Kept as internal agent state (NOT
  139. * an effect-scoped `ctx.on` listener) so a concurrent fiber disposal — which
  140. * runs the agent's own listeners' disposers — cannot drop the waiter before
  141. * the `disposed` transition fires and leave the promise hanging.
  142. */
  143. private idleWaiters: (() => void)[] = []
  144. /**
  145. * Durability checkpoints started by idle {@link inject} calls. `inject()` is
  146. * synchronous, so it cannot await them itself; the driver disposer drains
  147. * this set before the lifecycle unregisters the agent or detaches its session.
  148. */
  149. private pendingIdleFlushes = new Set<Promise<void>>()
  150. constructor(
  151. private loopCtx: Context,
  152. public readonly id: SessionId,
  153. public readonly options: AgentOptions,
  154. public readonly session: Session,
  155. ) {
  156. const { promise, resolve } = Promise.withResolvers<void>()
  157. this.disposed = promise
  158. this.resolveDisposed = resolve
  159. }
  160. get status(): AgentStatus {
  161. return this._status
  162. }
  163. private setStatus(status: AgentStatus): void {
  164. if (this._status === status || this._status === 'disposed') return
  165. this._status = status
  166. // Release quiescence waiters on a transition OUT of running BEFORE emitting
  167. // (the disposer handles the disposed transition separately). Settling first
  168. // means a throwing `agent/status` subscriber cannot starve a `whenIdle()`
  169. // waiter (docs/defensive-patterns.md "contain callback exceptions" — a lifecycle await must
  170. // not hang on one bad listener).
  171. if (status !== 'running') this.settleIdleWaiters()
  172. agentEvents(this.loopCtx, this).emit('agent/status', status)
  173. }
  174. /**
  175. * Resolve and clear all pending {@link whenIdle} waiters. Called on a
  176. * running→idle transition (from {@link setStatus}) and on disposal (from the
  177. * internal driver disposer, which chains `done` for true loop-exit quiescence).
  178. */
  179. private settleIdleWaiters(): void {
  180. const waiters = this.idleWaiters
  181. this.idleWaiters = []
  182. for (const resolve of waiters) resolve()
  183. }
  184. private resolveSource(options?: SendOptions): MessageSource {
  185. return options?.source ?? { kind: 'user' }
  186. }
  187. /**
  188. * Accept one public send/steer payload as the exact detached record shared by
  189. * the live notification and inbox. Lossless-JSON materialization reads every
  190. * nested field once; deep freeze prevents an observer from rewriting queued
  191. * work before the loop drains it.
  192. */
  193. private acceptInboxMessage(content: ContentBlock[], options?: SendOptions): InboxMessage {
  194. const source = this.resolveSource(options)
  195. const accepted = snapshotJsonValue({ content, source })
  196. if (accepted === undefined) {
  197. throw new TypeError('agent message content and source must be losslessly JSON-serializable')
  198. }
  199. return deepFreeze(accepted)
  200. }
  201. /** Reject a driving operation once teardown has synchronously closed the agent. */
  202. private assertNotDisposed(): void {
  203. if (this._status === 'disposed') throw new Error(`agent "${this.id}" is disposed`)
  204. }
  205. send(content: ContentBlock[], options?: SendOptions): void {
  206. this.assertNotDisposed()
  207. const accepted = this.acceptInboxMessage(content, options)
  208. this.#inbox.enqueue(accepted)
  209. const info = { source: accepted.source, steering: false } as const
  210. agentEvents(this.loopCtx, this).emit('agent/queued', accepted.content, info)
  211. }
  212. steer(content: ContentBlock[], options?: SendOptions): void {
  213. this.assertNotDisposed()
  214. if (this._status !== 'running') { this.send(content, options); return }
  215. const accepted = this.acceptInboxMessage(content, options)
  216. this.#inbox.steer(accepted)
  217. const info = { source: accepted.source, steering: true } as const
  218. agentEvents(this.loopCtx, this).emit('agent/queued', accepted.content, info)
  219. }
  220. inject(content: ContentBlock[], options?: SendOptions): void {
  221. this.assertNotDisposed()
  222. const source = this.resolveSource(options)
  223. if (isTurnOpen(this.session)) {
  224. // A turn is open in the LOG (decided from the log, not agent status —
  225. // status can be `running` with no turn open): the context/message is
  226. // turn-enclosed by that turn, so append it directly.
  227. this.session.append('context/message', { content, source }, { surfaceOp: 'append' })
  228. return
  229. }
  230. // No turn open: wrap the injection in a one-shot turn so every event stays
  231. // turn-enclosed (the durability/replay boundary is the turn).
  232. const turn = lastTurnNumber(this.session) + 1
  233. // Once turn/start enters the log, a turn/end is owed even if the message
  234. // append fails acceptance or pre-commit validation. The finally re-checks
  235. // the log and closes only a turn that actually opened; post-commit observers
  236. // are contained by Session and cannot create a false append failure.
  237. try {
  238. this.session.append('turn/start', { turn, trigger: { kind: 'injection', source } })
  239. this.session.append('context/message', { content, source }, { surfaceOp: 'append' })
  240. } finally {
  241. // Close the turn if turn/start made it into the log. A pre-commit veto
  242. // must escape rather than being mistaken for a committed turn/end.
  243. if (isTurnOpen(this.session)) {
  244. this.session.append('turn/end', { turn, reason: { kind: 'completed' } })
  245. }
  246. // Decide the durability checkpoint from the log: an accepted one-shot
  247. // turn must be flushed even when its message append was the failing step.
  248. const turnRecorded = this.session.events.some(e => e.type === 'turn/start' && e.data.turn === turn)
  249. // Checkpoint the one-shot turn for durability, exactly as the loop does at
  250. // every turn/end. The loop is NOT running (we are idle), so nothing else
  251. // will flush this turn. Fire-and-forget with error containment: inject()
  252. // is synchronous, and a persistence backend failing must not throw into
  253. // the caller (e.g. a tool-bash task-done callback). Disposal still drains
  254. // independently, so a slow flush is safe. The task is tracked until it
  255. // settles: driver disposal awaits every pending idle-injection checkpoint
  256. // before unregistering the agent or detaching the session. A flush failure
  257. // is reported via agent/error (step 0 — the idle-injection convention,
  258. // there is no real step) AND the logger, mirroring the loop's post-turn/end
  259. // flush path so plugins monitoring agent/error see idle-injection
  260. // persistence failures too. A throwing agent/error listener is contained.
  261. if (turnRecorded) {
  262. // Through the store's flush (the carrier owner), never a raw parallel.
  263. const flush = this.loopCtx.sessions.flush(this.session).catch((error: unknown) => {
  264. const rendered = renderThrown(error)
  265. const err = error instanceof Error ? error : new Error(rendered)
  266. this.loopCtx.logger.warn(`agent "${this.id}": flush after idle injection failed: ${rendered}`)
  267. agentEvents(this.loopCtx, this).emit('agent/error', turn, 0, err)
  268. })
  269. this.pendingIdleFlushes.add(flush)
  270. // Attach the same retirement callback to both settlement arms so even a
  271. // logger failure in the catch above cannot become an unhandled rejection.
  272. // Teardown uses allSettled for the same reason: a reporting failure must
  273. // not strand ownership.
  274. const retire = (): void => { this.pendingIdleFlushes.delete(flush) }
  275. void flush.then(retire, retire)
  276. }
  277. }
  278. }
  279. cancel(reason?: string): void {
  280. // Arm-gate: only mark a cancellation when there is actually work to cancel —
  281. // a running turn, an in-flight step, or queued/steering work. An idle cancel
  282. // with nothing pending is a true no-op; arming the marker then would wrongly
  283. // drop the NEXT legitimate prompt (the marker is consumed only at the loop's
  284. // turn-decision points, which an idle parked loop does not reach until woken
  285. // by a real send()). Note the gate canNOT be `status === 'running'` alone:
  286. // the pre-step window (a send() queued but the loop not yet flipped to
  287. // running) has status `idle` with `hasQueued` true, and the marker exists
  288. // precisely to cover it.
  289. if (this._status === 'running' || this.currentAbort !== undefined || this.#inbox.hasQueued || this.#inbox.hasSteering) {
  290. this.cancelRequested = true
  291. // Capture the resolved reason for the marker-only windows (pre-step /
  292. // continuation). The mid-step path reads it from abort.signal.reason
  293. // below; the marker path reads it via the LoopHandle's cancelReason().
  294. this.cancelReason = reason ?? 'cancelled'
  295. }
  296. // Drop all pending queued + steering work (un-started prompts never run; the
  297. // cancelled turn's steering is not re-enqueued). Cleared directly even when
  298. // the loop is parked in waitForQueued — there is no turn to stop and nothing
  299. // left for the parked loop to run, so no wake is needed.
  300. this.#inbox.clear()
  301. // Interrupt an in-flight step immediately (the running turn observes the
  302. // abort and ends `aborted`). The marker covers the windows where no step is
  303. // running (pre-step, continuation).
  304. this.currentAbort?.abort(reason ?? 'cancelled')
  305. }
  306. /**
  307. * Resolve once the agent has reached quiescence after settling out of
  308. * `running`. If it is already disposed, awaits {@link done} (the loop-exit
  309. * promise) — `agent/status('disposed')` fires in the disposer BEFORE the
  310. * driver loop has unwound, so it is NOT itself a quiescence signal. If it is
  311. * idle AND has no queued work, resolves immediately. Otherwise queues an
  312. * internal waiter (see {@link idleWaiters}) released on the next
  313. * running→idle/disposed transition, resolving on `idle` directly (the turn
  314. * fully ended) or chaining {@link done} on `disposed` (wait for the loop to
  315. * actually exit). Implements the {@link Agent.whenIdle} contract: a non-owner
  316. * quiescence-observation hook, distinct from teardown (a lifecycle owner stops
  317. * and unregisters via `AgentHandle.dispose()`, whose driver boundary awaits
  318. * both {@link done} and outstanding idle-injection flushes, not through this).
  319. */
  320. whenIdle(): Promise<void> {
  321. if (this._status === 'disposed') return this.done
  322. if (this._status !== 'running' && !this.#inbox.hasQueued) return Promise.resolve()
  323. // Register an internal waiter (resolved by settleIdleWaiters on the next
  324. // running→idle/disposed transition), NOT an effect-scoped `ctx.on` listener:
  325. // a concurrent fiber disposal runs this agent's listener disposers, which
  326. // could remove a `ctx.on` waiter before the `disposed` transition fires and
  327. // hang the promise. On disposal the disposer settles the waiter AND we chain
  328. // `done` here for true loop-exit quiescence (status flips to disposed before
  329. // the loop unwinds); a plain idle transition resolves directly.
  330. return new Promise<void>((resolve) => {
  331. this.idleWaiters.push(() => {
  332. resolve(this._status === 'disposed' ? this.done : undefined)
  333. })
  334. })
  335. }
  336. /** Bind the mutually referential scope context once. */
  337. private [bindContext](ctx: Context): void {
  338. if (this.boundContext !== undefined) throw new Error(`agent "${this.id}" context is already bound`)
  339. this.boundContext = ctx
  340. }
  341. /** Mark that public lifecycle publication began. */
  342. private [publishAgent](): void {
  343. this.published = true
  344. }
  345. /**
  346. * Start the driver loop. The prepared controller already owns its stable
  347. * disposer, so teardown can mark the agent disposed even in the narrow
  348. * publication window before this method runs.
  349. */
  350. [startDriver](): void {
  351. if (this._status === 'disposed') return
  352. this.driverStarted = true
  353. this.done = runLoop(this.loopCtx, this, {
  354. inbox: this.#inbox,
  355. setStatus: (status) => { this.setStatus(status) },
  356. setAbort: controller => void (this.currentAbort = controller),
  357. disposed: this.disposed,
  358. isDisposed: () => this._status === 'disposed',
  359. isCancelled: () => this.cancelRequested,
  360. cancelReason: () => this.cancelReason,
  361. clearCancel: () => { this.cancelRequested = false },
  362. // Settle whenIdle() waiters WITHOUT a status transition — the pre-step
  363. // cancel-skip path drops the about-to-run turn and re-parks without ever
  364. // flipping running→idle, so a waiter registered in the pre-step window
  365. // (status idle, hasQueued was true) would otherwise hang. This emits no
  366. // agent/status, so an ACP agent/status listener never sees a spurious idle
  367. // that would resolve a freshly-queued prompt as cancelled.
  368. settleIdle: () => { this.settleIdleWaiters() },
  369. })
  370. }
  371. /**
  372. * Quiescent stop shared by pre-start rollback and live teardown. It marks the
  373. * agent disposed synchronously, contains an unexpected loop rejection, and
  374. * drains every idle-injection flush before resolving.
  375. */
  376. private [stopDriver](): Promise<void> | void {
  377. if (this._status !== 'disposed') {
  378. this._status = 'disposed'
  379. this.resolveDisposed()
  380. // Release whenIdle waiters BEFORE the (guarded) event emit — they are
  381. // internal state that must settle even if a listener throws below. Each
  382. // waiter chains `done`, so it resolves only once the loop actually exits.
  383. this.settleIdleWaiters()
  384. this.currentAbort?.abort('disposed')
  385. // An unpublished rollback has no public status lifecycle to announce.
  386. // Once publication begins, disposed is part of the agent/status contract.
  387. if (this.published) {
  388. agentEvents(this.loopCtx, this).emit('agent/status', 'disposed')
  389. }
  390. }
  391. // Before runLoop starts there is normally nothing asynchronous to drain;
  392. // keep publication rollback synchronous so create() cannot throw while its
  393. // session/agent entries are still briefly live. A session-start listener
  394. // may have called inject(), however, so preserve
  395. // its durability checkpoint as a real quiescence boundary.
  396. if (!this.driverStarted && this.pendingIdleFlushes.size === 0) return
  397. return this.drainDriver()
  398. }
  399. /** Await the loop (when started) and every outstanding idle flush. */
  400. private async drainDriver(): Promise<void> {
  401. // An unexpected driver rejection must not skip registry/session/scope
  402. // cleanup. The normal loop contains turn failures itself; allSettled is the
  403. // final lifecycle backstop for anything outside those boundaries.
  404. await Promise.allSettled([this.done])
  405. // No new inject() can start after the synchronous disposed transition.
  406. // Loop because settled tasks retire themselves in promise reactions that
  407. // may run beside this continuation; either the set is empty or this waits
  408. // the exact remaining quiescence boundary. allSettled keeps a failure in
  409. // error reporting from skipping registry/session/scope disposers.
  410. while (this.pendingIdleFlushes.size > 0) {
  411. await Promise.allSettled([...this.pendingIdleFlushes])
  412. }
  413. }
  414. }
  415. /** Render an ordinary thrown value for the error event and log. */
  416. function renderThrown(value: unknown): string {
  417. return value instanceof Error ? value.message : String(value)
  418. }