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 {} from '@deepseek-ai/dsh-agent-execution'
  14. import type {
  15. AgentFactory,
  16. AgentHandle,
  17. AgentId,
  18. AgentOptions,
  19. CreateAgentOptions,
  20. ResumeAgentOptions,
  21. SessionStartSource,
  22. } from '@deepseek-ai/dsh-agent'
  23. import type {} from '@deepseek-ai/dsh-llm'
  24. import { SessionId } from '@deepseek-ai/dsh-session'
  25. import type { Session, SessionHeader } from '@deepseek-ai/dsh-session'
  26. import type {} from '@deepseek-ai/dsh-system-prompt'
  27. import type {} from '@deepseek-ai/dsh-tools'
  28. import type { SessionPersistence } from '@deepseek-ai/dsh-session-persistence'
  29. import {
  30. bindReactLoopAgentContext,
  31. prepareReactLoopAgent,
  32. ReactLoopAgent,
  33. } from './agent.ts'
  34. import type { PreparedReactLoopAgent } from './agent.ts'
  35. export { ReactLoopAgent } from './agent.ts'
  36. /** Fiber states that cannot own or serve a new lifecycle. */
  37. const INACTIVE_STATES: ReadonlySet<FiberState> = new Set([
  38. FiberState.UNLOADING,
  39. FiberState.DISPOSED,
  40. FiberState.FAILED,
  41. ])
  42. /** Factory-level ownership of every preparing or live transaction. */
  43. class FactoryOwnership {
  44. private accepting = true
  45. private transactions = new Set<AgentCreationTransaction>()
  46. constructor(private readonly fiber: Context['fiber']) {}
  47. isActive(): boolean {
  48. return this.accepting && !INACTIVE_STATES.has(this.fiber.state)
  49. }
  50. track(transaction: AgentCreationTransaction): () => void {
  51. this.transactions.add(transaction)
  52. return () => { this.transactions.delete(transaction) }
  53. }
  54. async dispose(): Promise<void> {
  55. this.accepting = false
  56. const reason = new Error('agent loop is not active')
  57. await Promise.all(
  58. [...this.transactions].map(transaction => transaction.disposeForFactory(reason)),
  59. )
  60. }
  61. }
  62. /** Build the public cancellation error while preserving a caller-supplied cause. */
  63. function signalAbortError(id: AgentId, signal: AbortSignal): Error {
  64. if (signal.reason instanceof Error) return signal.reason
  65. return new Error(`agent "${id}" creation aborted`, { cause: signal.reason })
  66. }
  67. /**
  68. * Caller-owned create/resume transaction through rollback-covered publication
  69. * and quiescent teardown. Resources remain private until the final registry
  70. * entry arbitrates identity.
  71. */
  72. class AgentCreationTransaction {
  73. private active = true
  74. private failure: Error | undefined
  75. private readonly deactivation = Promise.withResolvers<void>()
  76. private readonly publication = Promise.withResolvers<void>()
  77. private readonly torndown = Promise.withResolvers<void>()
  78. private readonly wrapperCompletion = Promise.withResolvers<void>()
  79. private preparing: Promise<void> | undefined
  80. private driver: PreparedReactLoopAgent | undefined
  81. private scope: Scope | undefined
  82. private session: Session | undefined
  83. private lifecycleDispose: (() => Promise<void> | void) | undefined
  84. private detachSession: (() => void) | undefined
  85. private detachAgent: (() => void) | undefined
  86. private publishing = false
  87. private cleanupTask: Promise<void> | undefined
  88. private ownerFollowing = true
  89. private readonly ownerDispose: () => Promise<void> | void
  90. private readonly untrackFactory: () => void
  91. private readonly abortListener: (() => void) | undefined
  92. readonly ownerAgent: Context['agent']
  93. readonly ownerFiber: Context['fiber']
  94. constructor(
  95. private readonly loopCtx: Context,
  96. private readonly ownerCtx: Context,
  97. private readonly ownership: FactoryOwnership,
  98. readonly id: AgentId,
  99. signal?: AbortSignal,
  100. ) {
  101. ownerCtx.fiber.assertActive()
  102. this.ownerAgent = ownerCtx.agent
  103. this.ownerFiber = ownerCtx.fiber
  104. if (!ownership.isActive()) throw new Error('agent loop is not active')
  105. this.ownerDispose = ownerCtx.effect(() => () => {
  106. if (!this.ownerFollowing) return
  107. return this.dispose(new Error(`agent "${id}" setup aborted: owner disposed during setup`))
  108. }, `agentLoop.owner(${id})`)
  109. this.untrackFactory = ownership.track(this)
  110. if (signal === undefined) {
  111. this.abortListener = undefined
  112. } else {
  113. this.abortListener = () => {
  114. /* v8 ignore next 3 -- transaction teardown contains callback/driver failures; rejection is a future-drift backstop. */
  115. void this.dispose(signalAbortError(id, signal)).catch((error: unknown) => {
  116. this.loopCtx.logger.error(error)
  117. })
  118. }
  119. signal.addEventListener('abort', this.abortListener, { once: true })
  120. if (signal.aborted) this.deactivate(signalAbortError(id, signal))
  121. }
  122. this.signal = signal
  123. }
  124. private readonly signal: AbortSignal | undefined
  125. /** Whether caller, provider, and optional parent-agent ownership remain live. */
  126. isActive(): boolean {
  127. return this.active
  128. && this.ownership.isActive()
  129. && this.ownerFiber.uid !== null
  130. && !INACTIVE_STATES.has(this.ownerFiber.state)
  131. && this.ownerAgent?.status !== 'disposed'
  132. }
  133. /** Fail synchronously at every real lifecycle boundary after deactivation. */
  134. assertActive(): void {
  135. if (this.isActive()) return
  136. if (!this.ownership.isActive()) throw new Error('agent loop is not active')
  137. throw this.failure ?? new Error(`agent "${this.id}" setup aborted: owner disposed during setup`)
  138. }
  139. /** Race an external async operation against structural/signal deactivation. */
  140. async waitFor<T>(operation: PromiseLike<T> | T): Promise<T> {
  141. this.assertActive()
  142. return await Promise.race([
  143. Promise.resolve(operation),
  144. this.deactivation.promise.then(() => {
  145. /* v8 ignore next -- deactivate() assigns failure before resolving deactivation. */
  146. throw this.failure ?? new Error(`agent "${this.id}" creation deactivated`)
  147. }),
  148. ])
  149. }
  150. /** Construct the driver and scope, then install their complete ordered lifecycle. */
  151. prepare(options: AgentOptions, session: Session): ReactLoopAgent {
  152. this.assertActive()
  153. const gate = Promise.withResolvers<void>()
  154. this.preparing = gate.promise
  155. try {
  156. this.session = session
  157. const driver = prepareReactLoopAgent(this.loopCtx, this.id, options, session)
  158. this.driver = driver
  159. const agent = driver.agent
  160. const scope = createScope(this.loopCtx, agent)
  161. this.scope = scope
  162. bindReactLoopAgentContext(agent, scope.ctx.extend({ agent }))
  163. this.installLifecycle(scope, driver)
  164. this.assertActive()
  165. return agent
  166. } catch (error: unknown) {
  167. if (!this.isActive() && error instanceof Error && /inactive context/.test(error.message)) {
  168. throw this.failure ?? this.disposalReason()
  169. }
  170. throw error
  171. } finally {
  172. gate.resolve()
  173. this.preparing = undefined
  174. }
  175. }
  176. /** Register the exact scope disposer inside the ordered transaction effect. */
  177. private installLifecycle(scope: Scope, driver: PreparedReactLoopAgent): void {
  178. this.lifecycleDispose = this.ownerCtx.effect(function* (this: AgentCreationTransaction) {
  179. // First yielded, disposed last.
  180. yield () => { this.finish() }
  181. yield scope.rawDispose
  182. yield () => {
  183. this.detachSession?.()
  184. this.detachSession = undefined
  185. }
  186. yield () => {
  187. this.detachAgent?.()
  188. this.detachAgent = undefined
  189. }
  190. // Last yielded, disposed first.
  191. yield () => {
  192. this.deactivate(this.disposalReason())
  193. if (this.publishing) {
  194. return this.publication.promise.then(() => driver.dispose())
  195. }
  196. return driver.dispose()
  197. }
  198. }.bind(this), `agentLoop.lifecycle(${this.id})`)
  199. }
  200. /** Publish the exact prepared objects and start the driver. */
  201. publish(source: SessionStartSource): AgentHandle {
  202. this.assertActive()
  203. const driver = this.driver
  204. /* v8 ignore next -- publish() is private and every caller invokes prepare() first. */
  205. if (driver === undefined) throw new Error(`agent "${this.id}" is not prepared`)
  206. const agent = driver.agent
  207. const session = this.session
  208. /* v8 ignore next -- prepare() assigns the session before it can produce the driver above. */
  209. if (session === undefined) throw new Error(`agent "${this.id}" has no prepared session`)
  210. this.publishing = true
  211. try {
  212. this.detachSession = agent.ctx.sessions.enter(session)
  213. this.detachAgent = this.loopCtx.agents.enter(agent)
  214. agent.ctx.sessions.announce(session)
  215. this.assertActive()
  216. this.loopCtx.agents.announce(agent)
  217. this.assertActive()
  218. driver.markPublished()
  219. agentEvents(this.loopCtx, agent).emit('agent/session-start', source)
  220. this.assertActive()
  221. driver.startDriver()
  222. return { agent, dispose: () => this.dispose() }
  223. } finally {
  224. this.publishing = false
  225. this.publication.resolve()
  226. }
  227. }
  228. /** Mark the transaction inactive exactly once and wake load/setup races. */
  229. private deactivate(reason: Error): void {
  230. if (!this.active) return
  231. this.active = false
  232. this.failure = reason
  233. this.deactivation.resolve()
  234. }
  235. /** Choose the structural cause when an owner/factory effect starts teardown first. */
  236. private disposalReason(): Error {
  237. if (this.failure !== undefined) return this.failure
  238. if (!this.ownership.isActive()) return new Error('agent loop is not active')
  239. if (this.ownerFiber.uid === null || INACTIVE_STATES.has(this.ownerFiber.state) || this.ownerAgent?.status === 'disposed') {
  240. return new Error(`agent "${this.id}" setup aborted: owner disposed during setup`)
  241. }
  242. return new Error(`agent "${this.id}" lifecycle disposed`)
  243. }
  244. /** Complete ownership bookkeeping after every resource reached quiescence. */
  245. private finish(): void {
  246. this.untrackFactory()
  247. this.ownerFollowing = false
  248. void this.ownerDispose()
  249. this.torndown.resolve()
  250. }
  251. /**
  252. * Deactivate and quiesce this transaction. The promise is memoized because
  253. * Cordis effect disposers are single-shot while handles promise shared
  254. * quiescence to every racing owner.
  255. */
  256. dispose(reason = new Error(`agent "${this.id}" lifecycle disposed`)): Promise<void> {
  257. this.deactivate(reason)
  258. return (this.cleanupTask ??= (async () => {
  259. if (this.preparing !== undefined) await this.preparing
  260. if (this.lifecycleDispose !== undefined) {
  261. await this.lifecycleDispose()
  262. await this.torndown.promise
  263. return
  264. }
  265. try {
  266. await this.driver?.dispose()
  267. } finally {
  268. try {
  269. await this.scope?.dispose()
  270. } finally {
  271. this.finish()
  272. }
  273. }
  274. })())
  275. }
  276. /** Mark the public create/resume continuation settled and detach its creation-only signal. */
  277. finishWrapper(): void {
  278. if (this.signal !== undefined && this.abortListener !== undefined) {
  279. this.signal.removeEventListener('abort', this.abortListener)
  280. }
  281. this.wrapperCompletion.resolve()
  282. }
  283. /** Factory shutdown joins both resource teardown and the public wrapper's deactivation continuation. */
  284. async disposeForFactory(reason: Error): Promise<void> {
  285. await this.dispose(reason)
  286. await this.wrapperCompletion.promise
  287. }
  288. }
  289. declare module 'cordis' {
  290. interface Context {
  291. agentLoop: AgentLoop
  292. }
  293. }
  294. /** Plugin configuration for declarative startup agents. */
  295. export interface Config {
  296. /** Agents created or resumed at plugin startup. */
  297. agents: (AgentOptions & {
  298. /** Registry identity for the live agent. */
  299. id: AgentId
  300. /** Optional workspace for a fresh session. */
  301. cwd?: string
  302. /** Persisted session to resume instead of creating a fresh session. */
  303. resumeSessionId?: SessionId
  304. })[]
  305. }
  306. /** Concrete ReactLoopAgent factory and driver service. */
  307. export class AgentLoop extends Service implements AgentFactory {
  308. static inject = ['agents', 'agentExecution', 'sessions', 'llm', 'tools', 'systemPrompt']
  309. /** Runtime schema for declarative agents. */
  310. static Config = z.object({
  311. agents: z.array(z.object({
  312. id: z.string().required(),
  313. provider: z.string(),
  314. model: z.string(),
  315. cwd: z.string(),
  316. resumeSessionId: z.string(),
  317. })).default([]),
  318. }) as unknown as z<Config>
  319. private readonly ownership: FactoryOwnership
  320. /** Plain holder prevents Cordis from re-tracing the factory's dependency context through a caller shadow. */
  321. private readonly runtime: { ctx: Context }
  322. constructor(ctx: Context, public config: Config) {
  323. super(ctx, 'agentLoop')
  324. this.ownership = new FactoryOwnership(ctx.fiber)
  325. this.runtime = { ctx }
  326. ctx.effect(() => () => this.ownership.dispose(), 'agentLoop.transactions()')
  327. ctx.effect(() => ctx.agents.setFactory(this), 'agentLoop.setFactory()')
  328. ctx.systemPrompt.variable('provider', context => context.agent?.options.provider)
  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