index.ts 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347
  1. /**
  2. * THE concrete agent plugin: creates ReactLoopAgents, runs their loops, and
  3. * registers them in ctx.agents. Deliberately thin — every behavior beyond
  4. * "call the model, run the tools, repeat" belongs to plugins on the event
  5. * taxonomy.
  6. *
  7. * @module @deepseek-ai/dsh-agent-loop
  8. */
  9. import { Context, Service } from 'cordis'
  10. import { randomUUID } from 'node:crypto'
  11. import z from 'schemastery'
  12. import type { AgentFactory, AgentHandle, AgentId, AgentOptions, CreateAgentOptions, ResumeAgentOptions, SessionStartSource } from '@deepseek-ai/dsh-agent'
  13. import type {} from '@deepseek-ai/dsh-llm'
  14. import { SessionId } from '@deepseek-ai/dsh-session'
  15. import type { Session } from '@deepseek-ai/dsh-session'
  16. import type {} from '@deepseek-ai/dsh-system-prompt'
  17. import type {} from '@deepseek-ai/dsh-tools'
  18. import type { SessionPersistence } from '@deepseek-ai/dsh-session-persistence'
  19. import { ReactLoopAgent } from './agent.ts'
  20. export { ReactLoopAgent } from './agent.ts'
  21. export { Inbox, type InboxMessage } from './inbox.ts'
  22. export { runLoop } from './loop.ts'
  23. declare module 'cordis' {
  24. interface Context {
  25. agentLoop: AgentLoop
  26. }
  27. }
  28. /**
  29. * Plugin config: the agents to create — or resume, via `resumeSessionId` —
  30. * declaratively at startup, so a cordis.yml deployment needs no code.
  31. */
  32. export interface Config {
  33. /** Agents created from configuration at startup. */
  34. agents: (AgentOptions & {
  35. id: AgentId
  36. /**
  37. * If set, the config agent RESUMES this persisted session id instead of
  38. * starting a fresh `${id}-session-<uuid>`. Sourced from an env var in
  39. * cordis.yml (`resumeSessionId: !!js process.env.RESUME_SESSION_ID`), so a
  40. * demo can continue a prior conversation without code changes. Requires a
  41. * `dsh-session-persistence` backend; the resume is deferred until that
  42. * service is available (via `ctx.inject`) and the loaded session's events
  43. * seed the live session so history continues.
  44. *
  45. * The schema accepts a plain string at runtime (cordis.yml values are
  46. * untyped); the brand is compile-time only — the config format is the
  47. * boundary where an id enters, so the TYPE declares the brand here.
  48. */
  49. resumeSessionId?: SessionId
  50. })[]
  51. }
  52. /**
  53. * The agent-loop plugin (`ctx.agentLoop`): creates {@link ReactLoopAgent}s, runs
  54. * their loops, and registers them in `ctx.agents`. Also implements the
  55. * {@link AgentFactory} seam, so plugins create/resume agents through
  56. * `ctx.agents` (the interface) without depending on this concrete package.
  57. *
  58. * The loop itself is deliberately thin — every behavior beyond "call the
  59. * model, run the tools, repeat" belongs to plugins listening on the event
  60. * taxonomy declared in @deepseek-ai/dsh-agent.
  61. */
  62. export class AgentLoop extends Service implements AgentFactory {
  63. static inject = ['agents', 'sessions', 'llm', 'tools', 'systemPrompt']
  64. // The schema validates plain strings (cordis.yml config values are untyped at
  65. // runtime); the {@link Config} TYPE declares the branded `id`/`resumeSessionId`
  66. // because the config format is the boundary where an id enters. The brand is a
  67. // zero-cost compile-time cast, so the runtime schema stays string-based and we
  68. // assert the branded view once here — the single schema boundary.
  69. static Config = z.object({
  70. agents: z.array(z.object({
  71. id: z.string().required(),
  72. model: z.string(),
  73. resumeSessionId: z.string(),
  74. })).default([]),
  75. }) as unknown as z<Config>
  76. constructor(ctx: Context, public config: Config) {
  77. super(ctx, 'agentLoop')
  78. // Provide the agent-creation factory to the registry (effect-scoped: the
  79. // slot is cleared on dispose).
  80. ctx.effect(() => this.ctx.agents.setFactory(this), 'agentLoop.setFactory()')
  81. // The prompt variables the shipped loop provides, registered once. The
  82. // sections themselves (`harness:identity`, `deployment:persona`) belong to
  83. // dsh-system-prompt — they must survive a swapped loop plugin — but
  84. // `{{model}}`/`{{cwd}}` are runtime facts of the agents THIS loop drives:
  85. // it assembles with `{ agent }` each step (loop.ts), and the variables
  86. // project the agent's configured model and its session workspace from that
  87. // context. A provider returns undefined when the fact is absent
  88. // (renderPrompt then rejects a persona that claims it — fail loud).
  89. ctx.systemPrompt.variable('model', context => context.agent?.options.model)
  90. ctx.systemPrompt.variable('cwd', context => context.agent?.session.header.cwd)
  91. for (const { id, resumeSessionId, ...options } of config.agents) {
  92. if (resumeSessionId !== undefined && resumeSessionId !== '') {
  93. // Resume a prior session instead of starting fresh. resume() needs
  94. // `ctx.sessionPersistence`, which may load AFTER this plugin (cordis.yml
  95. // lists the backend later). `ctx.inject(['sessionPersistence'], cb)`
  96. // runs `cb` with a child ctx once the service exists; the child reads
  97. // the persistence and hands it to resumeWith (which uses this.ctx — the
  98. // parent — for sessions/registry, all in AgentLoop's static inject). A
  99. // failed resume is contained + logged: startup must not crash.
  100. ctx.effect(() => {
  101. const fiber = this.ctx.inject(['sessionPersistence'], (childCtx: Context) => {
  102. void this.resumeWith(childCtx.sessionPersistence, { agentId: id, resumeSessionId, agentOptions: options })
  103. .catch((error: unknown) => {
  104. this.ctx.logger.warn(`agent "${id}": config-driven resume of "${resumeSessionId}" failed: ${String(error)}`)
  105. })
  106. })
  107. return () => void fiber.dispose()
  108. }, `agentLoop.resume(${id})`)
  109. } else {
  110. this.create(id, options)
  111. }
  112. }
  113. }
  114. /**
  115. * Config-driven create: an agent on a FRESH, non-colliding session id per run
  116. * (`${id}-session-<uuid>`, no cwd). Used for `cordis.yml`-configured agents
  117. * and as the shared core for the programmatic factory {@link createAgent}.
  118. *
  119. * Why a per-run id, not a fixed `${id}-session`: once a durable persistence
  120. * backend is loaded, a fixed id collides on the second run — the backend
  121. * refuses to re-create an id whose log already exists on disk (the SessionId
  122. * is the identity). A fresh id means each run is a new session.
  123. *
  124. * TODO(demo): each run starting a brand-new session is fine for demos but is
  125. * NOT real conversation continuity. A production config-driven agent needs a
  126. * deliberate resume-or-create policy (resume the prior session if one exists,
  127. * else start fresh) or an explicit caller-chosen session id — revisit when the
  128. * UI/ACP path owns session selection.
  129. * @param id - the agent id; also seeds the generated session id.
  130. * @param options - loop options (model, limits, …); defaults applied per option.
  131. * @returns the running agent, owned by the calling fiber (no handle).
  132. */
  133. create(id: AgentId, options: AgentOptions = {}): ReactLoopAgent {
  134. this.assertAgentIdFree(id)
  135. // Config/programmatic path: prepare the session and let start() fold its
  136. // lifecycle into the agent's composite effect (so a fiber unload tears the
  137. // session + agent down as one ordered chain, capturing the loop's closing
  138. // flush). The whole effect is owned by THIS fiber; no AgentHandle is needed.
  139. const session = this.ctx.sessions.prepare(SessionId(`${id}-session-${randomUUID()}`), { meta: {} })
  140. const { agent } = this.start(id, options, session, 'startup')
  141. return agent
  142. }
  143. /**
  144. * Programmatic factory create ({@link AgentFactory}): an agent on a
  145. * caller-supplied `sessionId` (NOT `${id}-session`), with optional session
  146. * metadata (validated `cwd`, lineage) and an optional `seed` event prefix. The
  147. * ACP bridge uses this so the client-generated session id becomes the
  148. * live/persisted session id; the in-process FORK subagent backend passes a
  149. * `seed` (a balanced completed-turn prefix of the parent's log) so the child
  150. * starts with the parent's context. Returns an {@link AgentHandle} the owner
  151. * disposes to tear down exactly this agent.
  152. * @param options - agent id, caller-supplied session id, optional seed/meta,
  153. * and agent options.
  154. * @returns the handle whose dispose tears down exactly this agent.
  155. */
  156. createAgent(options: CreateAgentOptions): AgentHandle {
  157. // Check the agent id BEFORE preparing the session: register() would reject a
  158. // duplicate id only AFTER the session enters the store, leaving an orphaned
  159. // live session (and lazy persistence state) that blocks reuse of that id.
  160. this.assertAgentIdFree(options.agentId)
  161. const session = this.ctx.sessions.prepare(options.sessionId, {
  162. ...options.seed !== undefined ? { seed: options.seed } : {},
  163. meta: options.meta ?? {},
  164. })
  165. // A seeded (forked) create is still a fresh start, NOT a resume — `resume`
  166. // is reserved for reloading a PERSISTED session via resume()/resumeWith().
  167. return this.startOwned(options.agentId, options.agentOptions ?? {}, session, 'startup')
  168. }
  169. /**
  170. * Resume an agent on a persisted session ({@link AgentFactory}). Loads the
  171. * session log + metadata via `ctx.sessionPersistence`, reconstructs the live
  172. * session with the loaded events (so `lastTurnNumber`/`deriveMessages`
  173. * continue), and starts a fresh agent on it. The live session id is the
  174. * resumed id, NOT `${agentId}-session`.
  175. *
  176. * Requires `ctx.sessionPersistence`; rejects with a clear error if it is not
  177. * configured. NOT hard-injected (that would make non-persistent demos pend
  178. * forever) — callers that need resume (ACP) inject `sessionPersistence`, so
  179. * by the time this runs the service exists.
  180. * @param options - the persisted session id to reload, plus agent id/options.
  181. * @returns the handle for the agent resumed on the reconstructed session.
  182. */
  183. async resume(options: ResumeAgentOptions): Promise<AgentHandle> {
  184. // Read the service through `ctx.get('sessionPersistence')` — a direct
  185. // global-store lookup keyed by the isolate symbol — NOT
  186. // `this.ctx.sessionPersistence`. AgentLoop deliberately does NOT inject
  187. // `sessionPersistence` (injecting it would pend non-persistent demos
  188. // forever). The `ctx.<name>` property proxy resolves a service by an
  189. // ancestor-only walk of the current fiber's parent chain; from AgentLoop's
  190. // own fiber (which lacks the inject) that walk never reaches the sibling
  191. // backend fiber and throws "cannot get property … without inject". Worse,
  192. // when the call arrives via a traceable shadow (e.g. the ACP bridge child
  193. // fiber → `ctx.agents.resume()` → `this.factory.resume()`), the walk starts
  194. // at the shadow's origin fiber and fails the same way. `ctx.get(name)`
  195. // sidesteps the fiber walk entirely (a store lookup by the global isolate
  196. // key), so resume works from any caller fiber. It is strict by default: a
  197. // backend that is not ACTIVE (absent, or mid-teardown) reads as undefined
  198. // and we reject below, rather than handing back an unusable handle.
  199. const persistence = this.ctx.get('sessionPersistence')
  200. if (persistence === undefined) {
  201. throw new Error('cannot resume: session persistence is not configured (load a dsh-session-persistence backend)')
  202. }
  203. return this.resumeWith(persistence, options)
  204. }
  205. /**
  206. * Resume against an EXPLICIT persistence handle. Factored out of {@link resume}
  207. * so the config-driven path can pass the handle it obtained from a
  208. * `ctx.inject(['sessionPersistence'], …)` child context: `this.ctx` (the
  209. * service's own fiber) did not inject `sessionPersistence`, so reading it
  210. * there from inside the inject child trips the cordis inject guard. The
  211. * sessions store + registry are still read through `this.ctx` (both are in
  212. * AgentLoop's static inject, so they resolve fine).
  213. */
  214. private async resumeWith(persistence: SessionPersistence, options: ResumeAgentOptions): Promise<AgentHandle> {
  215. this.assertAgentIdFree(options.agentId)
  216. const { meta, events } = await persistence.load(options.resumeSessionId)
  217. // Re-check the agent id AFTER the await: the pre-load check above can go
  218. // stale while load() is pending (a concurrent resume/create may register the
  219. // same id). Re-checking immediately before prepare()/start keeps the
  220. // "no orphaned session on a duplicate id" guarantee under concurrency.
  221. this.assertAgentIdFree(options.agentId)
  222. // Reconstruct the live session with the FULL persisted header (createdAt,
  223. // cwd, lineage) so resume preserves identity, not just the cwd. The seed
  224. // events make lastTurnNumber/deriveMessages continue; the backend already
  225. // has state (cursor) from the load above, so onCreated is a no-op and the
  226. // seed is not re-persisted. prepare() (not create()) so the session
  227. // lifecycle folds into the agent's composite effect (ordered teardown).
  228. const session = this.ctx.sessions.prepare(options.resumeSessionId, {
  229. seed: events,
  230. meta: {
  231. createdAt: meta.createdAt,
  232. ...meta.cwd !== undefined ? { cwd: meta.cwd } : {},
  233. ...meta.parentSession !== undefined ? { parentSession: meta.parentSession } : {},
  234. // Reconstruct the seed boundary from the persisted header, NOT from
  235. // `events.length` (the resume seeds the WHOLE stored log).
  236. ...meta.seedLength !== undefined ? { seedLength: meta.seedLength } : {},
  237. },
  238. })
  239. return this.startOwned(options.agentId, options.agentOptions ?? {}, session, 'resume')
  240. }
  241. /**
  242. * Reject a duplicate agent id BEFORE the session is entered into the store, so
  243. * a failed factory call never leaves an orphaned live session (and lazy
  244. * persistence state) behind. `register()` enforces the same uniqueness, but
  245. * only after the session has already entered the store.
  246. */
  247. private assertAgentIdFree(id: AgentId): void {
  248. if (this.ctx.agents.get(id) !== undefined) {
  249. throw new Error(`agent "${id}" is already registered`)
  250. }
  251. }
  252. /**
  253. * Shared: construct a ReactLoopAgent over a PREPARED (not-yet-entered)
  254. * session, then build the ONE composite effect that owns the whole agent
  255. * lifecycle — session entry, registry registration, and the loop. Keeping all
  256. * three in a SINGLE effect (not sibling effects) is load-bearing: a fiber
  257. * unload disposes sibling effects CONCURRENTLY (`Promise.all`), which would
  258. * race the session detach against the loop's closing flush and drop the
  259. * closing `turn/end`. Inside one effect the disposers run as an ORDERED LIFO
  260. * chain — the runtime awaits each disposer's returned promise before the next:
  261. *
  262. * yield session-detach (disposed LAST — detach onAppend + remove entry)
  263. * yield register (disposed 2nd — unregister)
  264. * yield stop-and-drain (disposed FIRST — request loop stop, await agent.done)
  265. *
  266. * So on teardown: the loop is stopped and AWAITED to exit (its final
  267. * `session/flush` + `turn/end` fire through the still-attached `onAppend`),
  268. * THEN the agent is unregistered, THEN the session is detached — capturing the
  269. * closing events before detach, whether the trigger is the handle's `dispose()`
  270. * OR a fiber unload. Rollback safety: each yield runs before the next mutation,
  271. * so a throwing `session/created`/`agent/created` listener unwinds the
  272. * already-yielded disposers instead of leaking.
  273. *
  274. * `source` says why the session began ({@link SessionStartSource}); it is
  275. * emitted as `agent/session-start` once, AFTER the agent is registered (so a
  276. * listener can resolve the agent via `ctx.agents.get(id)` and `inject()` into
  277. * it) and BEFORE the loop starts its first turn. The emit is contained: a
  278. * throwing session-start listener must not abort agent construction — it is
  279. * logged, and the agent still starts. (Unlike a turn-boundary throw, there is
  280. * no open turn here to balance; the durable evidence of a session-start hook
  281. * is whatever it `inject()`ed.)
  282. *
  283. * Returns the agent plus the composite effect's disposer (`disposeAgent`).
  284. */
  285. private start(
  286. id: AgentId, options: AgentOptions, session: Session, source: SessionStartSource,
  287. ): { agent: ReactLoopAgent; disposeAgent: () => Promise<void> } {
  288. const agent = new ReactLoopAgent(this.ctx, id, options, session)
  289. const dispose = this.ctx.effect(function* (this: AgentLoop) {
  290. yield this.ctx.sessions.enter(session)
  291. this.ctx.sessions.announce(session)
  292. yield this.ctx.agents.register(agent)
  293. // Fire AFTER register (a listener can ctx.agents.get(id) + inject()) and
  294. // BEFORE the loop's first turn. Contained: a throwing listener is logged,
  295. // never aborts construction (no open turn to balance here).
  296. try {
  297. this.ctx.emit('agent/session-start', agent, source)
  298. } catch (error: unknown) {
  299. this.ctx.logger.warn(`agent "${id}": agent/session-start listener threw: ${String(error)}`)
  300. }
  301. const stop = agent.start()
  302. // Disposed FIRST (LIFO): request loop stop (sync), then AWAIT the loop's
  303. // actual exit so its closing flush lands while onAppend (yielded above,
  304. // disposed later) is still attached.
  305. yield async () => { stop(); await agent.done }
  306. }.bind(this), 'agentLoop.start()')
  307. return { agent, disposeAgent: async () => { await dispose() } }
  308. }
  309. /**
  310. * Build an {@link AgentHandle} for a PREPARED session + a fresh agent. The
  311. * handle's `dispose()` runs the composite effect's disposer (see
  312. * {@link start}) — which stops the loop, awaits its exit (final flush
  313. * captured), unregisters the agent, and detaches the session, in that order.
  314. * The same composite effect is what a fiber unload disposes, so both teardown
  315. * triggers honor the ordering identically.
  316. *
  317. * `dispose()` is MEMOIZED: the underlying cordis effect disposer is
  318. * single-shot (a second call returns immediately because the effect's epoch is
  319. * already cleared, NOT awaiting the in-flight teardown), so concurrent/repeated
  320. * `dispose()` calls would otherwise resolve before the first call's
  321. * `await agent.done` + final flush completed. Memoizing the promise makes every
  322. * caller observe the SAME quiescence boundary, honoring the
  323. * `AgentHandle.dispose(): Promise<void>` contract (mirrors the ACP `quiesce()`
  324. * helper).
  325. */
  326. private startOwned(id: AgentId, options: AgentOptions, session: Session, source: SessionStartSource): AgentHandle {
  327. const { agent, disposeAgent } = this.start(id, options, session, source)
  328. let disposing: Promise<void> | undefined
  329. return { agent, dispose: () => (disposing ??= disposeAgent()) }
  330. }
  331. }
  332. export default AgentLoop