index.ts 18 KB

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