index.ts 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532
  1. /**
  2. * Concrete agent-loop plugin: creates scoped ReactLoopAgents, publishes them
  3. * through the agent/session registries, and owns their ordered teardown.
  4. *
  5. * @module @deepseek-ai/dsh-agent-loop
  6. */
  7. import { Context, FiberState, Service } from 'cordis'
  8. import { randomUUID } from 'node:crypto'
  9. import z from 'schemastery'
  10. import { createScope } from '@deepseek-ai/dsh-scope'
  11. import type { Scope } from '@deepseek-ai/dsh-scope'
  12. import { agentEvents } from '@deepseek-ai/dsh-agent'
  13. import type {
  14. AgentFactory,
  15. AgentHandle,
  16. AgentOptions,
  17. CreateAgentOptions,
  18. ResumeAgentOptions,
  19. SessionStartSource,
  20. } from '@deepseek-ai/dsh-agent'
  21. import type {} from '@deepseek-ai/dsh-llm'
  22. import { SessionId } from '@deepseek-ai/dsh-session'
  23. import type { Session, SessionHeader } from '@deepseek-ai/dsh-session'
  24. import type {} from '@deepseek-ai/dsh-system-prompt'
  25. import type {} from '@deepseek-ai/dsh-tools'
  26. import type { SessionPersistence } from '@deepseek-ai/dsh-session-persistence'
  27. import {
  28. bindReactLoopAgentContext,
  29. prepareReactLoopAgent,
  30. ReactLoopAgent,
  31. } from './agent.ts'
  32. import type { PreparedReactLoopAgent } from './agent.ts'
  33. export { ReactLoopAgent } from './agent.ts'
  34. /** Fiber states that cannot own or serve a new lifecycle. */
  35. const INACTIVE_STATES: ReadonlySet<FiberState> = new Set([
  36. FiberState.UNLOADING,
  37. FiberState.DISPOSED,
  38. FiberState.FAILED,
  39. ])
  40. /** Factory-level ownership of every preparing or live transaction. */
  41. class FactoryOwnership {
  42. private accepting = true
  43. private transactions = new Set<AgentCreationTransaction>()
  44. private startupTasks = new Set<Promise<void>>()
  45. constructor(private readonly fiber: Context['fiber']) {}
  46. isActive(): boolean {
  47. return this.accepting && !INACTIVE_STATES.has(this.fiber.state)
  48. }
  49. track(transaction: AgentCreationTransaction): () => void {
  50. this.transactions.add(transaction)
  51. return () => { this.transactions.delete(transaction) }
  52. }
  53. /** Join config startup work that begins before an agent transaction exists. */
  54. trackStartup(task: Promise<void>): void {
  55. this.startupTasks.add(task)
  56. const forget = () => { this.startupTasks.delete(task) }
  57. void task.then(forget, forget)
  58. }
  59. async dispose(): Promise<void> {
  60. this.accepting = false
  61. const reason = new Error('agent loop is not active')
  62. await Promise.all([
  63. ...[...this.transactions].map(transaction => transaction.disposeForFactory(reason)),
  64. ...this.startupTasks,
  65. ])
  66. }
  67. }
  68. /** Build the public cancellation error while preserving a caller-supplied cause. */
  69. function signalAbortError(id: SessionId, signal: AbortSignal): Error {
  70. if (signal.reason instanceof Error) return signal.reason
  71. return new Error(`agent "${id}" creation aborted`, { cause: signal.reason })
  72. }
  73. /**
  74. * One create/resume transaction from caller ownership through unpublished
  75. * setup, rollback-covered publication, and final quiescent teardown.
  76. *
  77. * The class deliberately owns the state machine in one place. Registries only
  78. * arbitrate identity at their final `enter()` calls; before that point every
  79. * resource is private to this transaction.
  80. */
  81. class AgentCreationTransaction {
  82. private active = true
  83. private failure: Error | undefined
  84. private readonly deactivation = Promise.withResolvers<void>()
  85. private readonly publication = Promise.withResolvers<void>()
  86. private readonly torndown = Promise.withResolvers<void>()
  87. private readonly wrapperCompletion = Promise.withResolvers<void>()
  88. private preparing: Promise<void> | undefined
  89. private driver: PreparedReactLoopAgent | undefined
  90. private scope: Scope | undefined
  91. private session: Session | undefined
  92. private lifecycleDispose: (() => Promise<void> | void) | undefined
  93. private detachSession: (() => void) | undefined
  94. private detachAgent: (() => void) | undefined
  95. private publishing = false
  96. private cleanupTask: Promise<void> | undefined
  97. private ownerFollowing = true
  98. private readonly ownerDispose: () => Promise<void> | void
  99. private readonly untrackFactory: () => void
  100. private readonly abortListener: (() => void) | undefined
  101. readonly ownerAgent: Context['agent']
  102. readonly ownerFiber: Context['fiber']
  103. constructor(
  104. private readonly loopCtx: Context,
  105. private readonly ownerCtx: Context,
  106. private readonly ownership: FactoryOwnership,
  107. readonly id: SessionId,
  108. signal?: AbortSignal,
  109. ) {
  110. ownerCtx.fiber.assertActive()
  111. this.ownerAgent = ownerCtx.agent
  112. this.ownerFiber = ownerCtx.fiber
  113. if (!ownership.isActive()) throw new Error('agent loop is not active')
  114. this.ownerDispose = ownerCtx.effect(() => () => {
  115. if (!this.ownerFollowing) return
  116. return this.dispose(new Error(`agent "${id}" setup aborted: owner disposed during setup`))
  117. }, `agentLoop.owner(${id})`)
  118. this.untrackFactory = ownership.track(this)
  119. if (signal === undefined) {
  120. this.abortListener = undefined
  121. } else {
  122. this.abortListener = () => {
  123. /* v8 ignore next 3 -- transaction teardown contains callback/driver failures; rejection is a future-drift backstop. */
  124. void this.dispose(signalAbortError(id, signal)).catch((error: unknown) => {
  125. this.loopCtx.logger.error(error)
  126. })
  127. }
  128. signal.addEventListener('abort', this.abortListener, { once: true })
  129. if (signal.aborted) this.deactivate(signalAbortError(id, signal))
  130. }
  131. this.signal = signal
  132. }
  133. private readonly signal: AbortSignal | undefined
  134. /** Whether caller, provider, and optional parent-agent ownership remain live. */
  135. isActive(): boolean {
  136. return this.active
  137. && this.ownership.isActive()
  138. && this.ownerFiber.uid !== null
  139. && !INACTIVE_STATES.has(this.ownerFiber.state)
  140. && this.ownerAgent?.status !== 'disposed'
  141. }
  142. /** Fail synchronously at every real lifecycle boundary after deactivation. */
  143. assertActive(): void {
  144. if (this.isActive()) return
  145. if (!this.ownership.isActive()) throw new Error('agent loop is not active')
  146. throw this.failure ?? new Error(`agent "${this.id}" setup aborted: owner disposed during setup`)
  147. }
  148. /** Race an external async operation against structural/signal deactivation. */
  149. async waitFor<T>(operation: PromiseLike<T> | T): Promise<T> {
  150. this.assertActive()
  151. return await Promise.race([
  152. Promise.resolve(operation),
  153. this.deactivation.promise.then(() => {
  154. /* v8 ignore next -- deactivate() assigns failure before resolving deactivation. */
  155. throw this.failure ?? new Error(`agent "${this.id}" creation deactivated`)
  156. }),
  157. ])
  158. }
  159. /** Construct the driver and scope, then install their complete ordered lifecycle. */
  160. prepare(options: AgentOptions, session: Session): ReactLoopAgent {
  161. this.assertActive()
  162. const gate = Promise.withResolvers<void>()
  163. this.preparing = gate.promise
  164. try {
  165. this.session = session
  166. const driver = prepareReactLoopAgent(this.loopCtx, this.id, options, session)
  167. this.driver = driver
  168. const agent = driver.agent
  169. const scope = createScope(this.loopCtx, agent)
  170. this.scope = scope
  171. bindReactLoopAgentContext(agent, scope.ctx.extend({ agent }))
  172. this.installLifecycle(scope, driver)
  173. this.assertActive()
  174. return agent
  175. } catch (error: unknown) {
  176. if (!this.isActive() && error instanceof Error && /inactive context/.test(error.message)) {
  177. throw this.failure ?? this.disposalReason()
  178. }
  179. throw error
  180. } finally {
  181. gate.resolve()
  182. this.preparing = undefined
  183. }
  184. }
  185. /** Register the exact scope disposer inside the ordered transaction effect. */
  186. private installLifecycle(scope: Scope, driver: PreparedReactLoopAgent): void {
  187. this.lifecycleDispose = this.ownerCtx.effect(function* (this: AgentCreationTransaction) {
  188. // First yielded, disposed last.
  189. yield () => { this.finish() }
  190. yield scope.rawDispose
  191. yield () => {
  192. this.detachSession?.()
  193. this.detachSession = undefined
  194. }
  195. yield () => {
  196. this.detachAgent?.()
  197. this.detachAgent = undefined
  198. }
  199. // Last yielded, disposed first.
  200. yield () => {
  201. this.deactivate(this.disposalReason())
  202. if (this.publishing) {
  203. return this.publication.promise.then(() => driver.dispose())
  204. }
  205. return driver.dispose()
  206. }
  207. }.bind(this), `agentLoop.lifecycle(${this.id})`)
  208. }
  209. /** Publish the exact prepared objects and start the driver. */
  210. publish(source: SessionStartSource): AgentHandle {
  211. this.assertActive()
  212. const driver = this.driver
  213. /* v8 ignore next -- publish() is private and every caller invokes prepare() first. */
  214. if (driver === undefined) throw new Error(`agent "${this.id}" is not prepared`)
  215. const agent = driver.agent
  216. const session = this.session
  217. /* v8 ignore next -- prepare() assigns the session before it can produce the driver above. */
  218. if (session === undefined) throw new Error(`agent "${this.id}" has no prepared session`)
  219. this.publishing = true
  220. try {
  221. this.detachSession = agent.ctx.sessions.enter(session)
  222. this.detachAgent = this.loopCtx.agents.enter(agent, this.ownerAgent)
  223. agent.ctx.sessions.announce(session)
  224. this.assertActive()
  225. this.loopCtx.agents.announce(agent)
  226. this.assertActive()
  227. driver.markPublished()
  228. agentEvents(this.loopCtx, agent).emit('agent/session-start', source)
  229. this.assertActive()
  230. driver.startDriver()
  231. return { agent, dispose: () => this.dispose() }
  232. } finally {
  233. this.publishing = false
  234. this.publication.resolve()
  235. }
  236. }
  237. /** Mark the transaction inactive exactly once and wake load/setup races. */
  238. private deactivate(reason: Error): void {
  239. if (!this.active) return
  240. this.active = false
  241. this.failure = reason
  242. this.deactivation.resolve()
  243. }
  244. /** Choose the structural cause when an owner/factory effect starts teardown first. */
  245. private disposalReason(): Error {
  246. if (this.failure !== undefined) return this.failure
  247. if (!this.ownership.isActive()) return new Error('agent loop is not active')
  248. if (this.ownerFiber.uid === null || INACTIVE_STATES.has(this.ownerFiber.state) || this.ownerAgent?.status === 'disposed') {
  249. return new Error(`agent "${this.id}" setup aborted: owner disposed during setup`)
  250. }
  251. return new Error(`agent "${this.id}" lifecycle disposed`)
  252. }
  253. /** Complete ownership bookkeeping after every resource reached quiescence. */
  254. private finish(): void {
  255. this.untrackFactory()
  256. this.ownerFollowing = false
  257. void this.ownerDispose()
  258. this.torndown.resolve()
  259. }
  260. /**
  261. * Deactivate and quiesce this transaction. The promise is memoized because
  262. * Cordis effect disposers are single-shot while handles promise shared
  263. * quiescence to every racing owner.
  264. */
  265. dispose(reason = new Error(`agent "${this.id}" lifecycle disposed`)): Promise<void> {
  266. this.deactivate(reason)
  267. return (this.cleanupTask ??= (async () => {
  268. if (this.preparing !== undefined) await this.preparing
  269. if (this.lifecycleDispose !== undefined) {
  270. await this.lifecycleDispose()
  271. await this.torndown.promise
  272. return
  273. }
  274. try {
  275. await this.driver?.dispose()
  276. } finally {
  277. try {
  278. await this.scope?.dispose()
  279. } finally {
  280. this.finish()
  281. }
  282. }
  283. })())
  284. }
  285. /** Mark the public create/resume continuation settled and detach its creation-only signal. */
  286. finishWrapper(): void {
  287. if (this.signal !== undefined && this.abortListener !== undefined) {
  288. this.signal.removeEventListener('abort', this.abortListener)
  289. }
  290. this.wrapperCompletion.resolve()
  291. }
  292. /** Factory shutdown joins both resource teardown and the public wrapper's deactivation continuation. */
  293. async disposeForFactory(reason: Error): Promise<void> {
  294. await this.dispose(reason)
  295. await this.wrapperCompletion.promise
  296. }
  297. }
  298. declare module 'cordis' {
  299. interface Context {
  300. agentLoop: AgentLoop
  301. }
  302. }
  303. /** Plugin configuration for declarative startup agents. */
  304. export interface Config {
  305. /** Agents created or resumed at plugin startup. */
  306. agents: (AgentOptions & {
  307. /** Stable config label used in logs and as the fresh combined-id prefix. */
  308. id: string
  309. /** Optional stable identity; remounts resume its materialized history, while first use creates it fresh. */
  310. sessionId?: SessionId
  311. /** Optional workspace for a fresh session. */
  312. cwd?: string
  313. /** Persisted session to resume instead of creating a fresh session. */
  314. resumeSessionId?: SessionId
  315. })[]
  316. }
  317. /** Concrete ReactLoopAgent factory and driver service. */
  318. export class AgentLoop extends Service implements AgentFactory {
  319. static inject = ['agents', 'sessions', 'llm', 'tools', 'systemPrompt']
  320. /** Runtime schema for declarative agents. */
  321. static Config = z.object({
  322. agents: z.array(z.object({
  323. id: z.string().required(),
  324. sessionId: z.string(),
  325. model: z.string(),
  326. cwd: z.string(),
  327. resumeSessionId: z.string(),
  328. })).default([]),
  329. }) as unknown as z<Config>
  330. private readonly ownership: FactoryOwnership
  331. /** Plain holder prevents Cordis from re-tracing the factory's dependency context through a caller shadow. */
  332. private readonly runtime: { ctx: Context }
  333. constructor(ctx: Context, public config: Config) {
  334. super(ctx, 'agentLoop')
  335. this.ownership = new FactoryOwnership(ctx.fiber)
  336. this.runtime = { ctx }
  337. ctx.effect(() => () => this.ownership.dispose(), 'agentLoop.transactions()')
  338. ctx.effect(() => ctx.agents.setFactory(this), 'agentLoop.setFactory()')
  339. ctx.systemPrompt.variable('model', context => context.agent?.options.model)
  340. ctx.systemPrompt.variable('cwd', context => context.agent?.session.header.cwd)
  341. for (const { id, sessionId, cwd, resumeSessionId, ...options } of config.agents) {
  342. const meta = cwd === undefined ? {} : { cwd }
  343. if (resumeSessionId === undefined || resumeSessionId === '') {
  344. const configuredId = sessionId ?? SessionId(`${id}-session-${randomUUID()}`)
  345. const persistence = sessionId === undefined ? undefined : ctx.get('sessionPersistence')
  346. if (persistence === undefined) {
  347. this.create(configuredId, options, meta)
  348. } else {
  349. const startup = this.restoreOrCreateConfigured(ctx, persistence, configuredId, options, meta).catch((error: unknown) => {
  350. ctx.logger.warn(`agent "${id}": config-driven restore of "${configuredId}" failed: ${String(error)}`)
  351. })
  352. this.ownership.trackStartup(startup)
  353. }
  354. continue
  355. }
  356. if (sessionId !== undefined) {
  357. throw new Error(`agent "${id}": sessionId and resumeSessionId are mutually exclusive`)
  358. }
  359. ctx.effect(() => {
  360. const fiber = ctx.inject(['sessionPersistence'], (childCtx: Context) => {
  361. void this.resumeWith(ctx, childCtx.sessionPersistence, {
  362. resumeSessionId,
  363. agentOptions: options,
  364. }).catch((error: unknown) => {
  365. ctx.logger.warn(`agent "${id}": config-driven resume of "${resumeSessionId}" failed: ${String(error)}`)
  366. })
  367. })
  368. return fiber.dispose
  369. }, `agentLoop.resume(${id})`)
  370. }
  371. }
  372. /** Restore a materialized exact config identity on remount, or create it on first use. */
  373. private async restoreOrCreateConfigured(
  374. ownerCtx: Context,
  375. persistence: SessionPersistence,
  376. sessionId: SessionId,
  377. agentOptions: AgentOptions,
  378. meta: Pick<SessionHeader, 'cwd'>,
  379. ): Promise<void> {
  380. const exists = (await persistence.list()).some(header => header.id === sessionId)
  381. if (!this.ownership.isActive()) return
  382. if (exists) {
  383. await this.resumeWith(ownerCtx, persistence, { resumeSessionId: sessionId, agentOptions })
  384. return
  385. }
  386. this.create(sessionId, agentOptions, meta)
  387. }
  388. /**
  389. * Create an agent and session under one caller-supplied identity, owned by
  390. * the accessing fiber. Constructor-driven config calls mint a fresh combined
  391. * id before entering this boundary.
  392. * @param id - shared agent/session identity.
  393. * @param options - concrete loop options.
  394. * @param meta - optional fresh-session workspace metadata.
  395. * @returns the published running agent.
  396. */
  397. create(id: SessionId, options: AgentOptions = {}, meta: Pick<SessionHeader, 'cwd'> = {}): ReactLoopAgent {
  398. const loopCtx = this.runtime.ctx
  399. const transaction = new AgentCreationTransaction(loopCtx, this.ctx, this.ownership, id)
  400. try {
  401. const session = loopCtx.sessions.prepare(id, { meta })
  402. const agent = transaction.prepare(options, session)
  403. transaction.publish('startup')
  404. return agent
  405. } catch (error: unknown) {
  406. void transaction.dispose(error instanceof Error ? error : new Error(String(error)))
  407. throw error
  408. } finally {
  409. transaction.finishWrapper()
  410. }
  411. }
  412. /**
  413. * Create an owned agent on a caller-supplied session id.
  414. * @param ownerCtx - caller context that structurally owns the transaction.
  415. * @param options - identities, session seed/metadata, loop options, setup, and cancellation.
  416. * @returns the published handle.
  417. */
  418. async createAgent(ownerCtx: Context, options: CreateAgentOptions): Promise<AgentHandle> {
  419. const transaction = new AgentCreationTransaction(
  420. this.runtime.ctx,
  421. ownerCtx,
  422. this.ownership,
  423. options.sessionId,
  424. options.signal,
  425. )
  426. try {
  427. const session = this.runtime.ctx.sessions.prepare(options.sessionId, {
  428. ...options.seed === undefined ? {} : { seed: options.seed },
  429. ...options.meta === undefined ? {} : { meta: options.meta },
  430. })
  431. const agent = transaction.prepare(options.agentOptions ?? {}, session)
  432. await transaction.waitFor(options.setup?.(agent.ctx))
  433. transaction.assertActive()
  434. return transaction.publish('startup')
  435. } catch (error: unknown) {
  436. await transaction.dispose(error instanceof Error ? error : new Error(String(error)))
  437. throw error
  438. } finally {
  439. transaction.finishWrapper()
  440. }
  441. }
  442. /**
  443. * Resume an owned agent from the configured persistence service.
  444. * @param ownerCtx - caller context that owns load, setup, and the live lifecycle.
  445. * @param options - persisted identity, loop options, setup, and cancellation.
  446. * @returns the published handle.
  447. */
  448. async resume(ownerCtx: Context, options: ResumeAgentOptions): Promise<AgentHandle> {
  449. const persistence = this.runtime.ctx.get('sessionPersistence')
  450. if (persistence === undefined) {
  451. throw new Error('cannot resume: session persistence is not configured (load a dsh-session-persistence backend)')
  452. }
  453. return this.resumeWith(ownerCtx, persistence, options)
  454. }
  455. /** Resume through an explicit persistence handle used by the deferred config path. */
  456. private async resumeWith(
  457. ownerCtx: Context,
  458. persistence: SessionPersistence,
  459. options: ResumeAgentOptions,
  460. ): Promise<AgentHandle> {
  461. const transaction = new AgentCreationTransaction(
  462. this.runtime.ctx,
  463. ownerCtx,
  464. this.ownership,
  465. options.resumeSessionId,
  466. options.signal,
  467. )
  468. try {
  469. const loaded = await transaction.waitFor(persistence.load(options.resumeSessionId))
  470. transaction.assertActive()
  471. const session = this.runtime.ctx.sessions.prepare(options.resumeSessionId, {
  472. seed: loaded.events,
  473. meta: {
  474. createdAt: loaded.meta.createdAt,
  475. ...loaded.meta.cwd === undefined ? {} : { cwd: loaded.meta.cwd },
  476. ...loaded.meta.parentSession === undefined ? {} : { parentSession: loaded.meta.parentSession },
  477. ...loaded.meta.seedLength === undefined ? {} : { seedLength: loaded.meta.seedLength },
  478. },
  479. })
  480. const agent = transaction.prepare(options.agentOptions ?? {}, session)
  481. await transaction.waitFor(options.setup?.(agent.ctx))
  482. transaction.assertActive()
  483. return transaction.publish('resume')
  484. } catch (error: unknown) {
  485. await transaction.dispose(error instanceof Error ? error : new Error(String(error)))
  486. throw error
  487. } finally {
  488. transaction.finishWrapper()
  489. }
  490. }
  491. }
  492. export default AgentLoop