resume.spec.ts 37 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857
  1. import { createUserMessage } from '@deepseek-ai/dsh-llm'
  2. import { afterEach, describe, expect, it } from 'vitest'
  3. import { Context } from 'cordis'
  4. import { mkdtemp, rm } from 'node:fs/promises'
  5. import { tmpdir } from 'node:os'
  6. import { join } from 'node:path'
  7. import LlmService from '@deepseek-ai/dsh-llm'
  8. import SessionStore, { SESSION_FORMAT_VERSION, Session, SessionId } from '@deepseek-ai/dsh-session'
  9. import type { SessionEvent } from '@deepseek-ai/dsh-session'
  10. import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
  11. import ToolRegistry from '@deepseek-ai/dsh-tools'
  12. import AgentRegistry, { type Agent } from '@deepseek-ai/dsh-agent'
  13. import SessionPersistenceJsonl from '@deepseek-ai/dsh-session-persistence-jsonl'
  14. import AgentLoop from '@deepseek-ai/dsh-agent-loop'
  15. import { MockAdapter, textResponse } from './mock-adapter.ts'
  16. const dirs: string[] = []
  17. afterEach(async () => { for (const d of dirs.splice(0)) await rm(d, { recursive: true, force: true }) })
  18. async function persistentHarness(adapter: MockAdapter): Promise<{ ctx: Context; root: string }> {
  19. const root = await mkdtemp(join(tmpdir(), 'dsh-resume-'))
  20. dirs.push(root)
  21. return { ctx: await mountPersistentHarness(root, adapter), root }
  22. }
  23. async function mountPersistentHarness(root: string, adapter: MockAdapter): Promise<Context> {
  24. const ctx = new Context()
  25. await ctx.plugin(LlmService)
  26. await ctx.plugin(SessionStore)
  27. await ctx.plugin(SystemPrompt)
  28. await ctx.plugin(ToolRegistry)
  29. await ctx.plugin(AgentRegistry)
  30. await ctx.plugin(AgentLoop, { agents: [] })
  31. await ctx.plugin(SessionPersistenceJsonl, { root })
  32. ctx.llm.registerAdapter(['mock'], adapter)
  33. return ctx
  34. }
  35. async function persistSession(sessionId: SessionId): Promise<string> {
  36. const { ctx, root } = await persistentHarness(new MockAdapter([textResponse('seed')]))
  37. // Persistence deliberately has no artifact for a truly empty session. A
  38. // balanced completed turn is the smallest resumable log and avoids running
  39. // the model merely to construct this lifecycle fixture.
  40. const seed: SessionEvent[] = [
  41. { type: 'turn/start', seq: 0, time: 1, data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } } },
  42. { type: 'turn/end', seq: 1, time: 2, data: { turn: 1, reason: { kind: 'completed' } } },
  43. ]
  44. const session = ctx.sessions.create(sessionId, { seed })
  45. await ctx.sessions.flush(session)
  46. await ctx.fiber.dispose()
  47. return root
  48. }
  49. function waitForIdle(ctx: Context, agent: Agent): Promise<void> {
  50. return new Promise((resolve) => {
  51. const dispose = ctx.on('agent/status', (subject, status) => {
  52. if (subject === agent && status === 'idle') { dispose(); resolve() }
  53. })
  54. })
  55. }
  56. /** Fail a lifecycle regression promptly instead of waiting for Vitest's suite timeout. */
  57. async function promptly<T>(task: Promise<T>): Promise<T> {
  58. const timeout = Promise.withResolvers<never>()
  59. const timer = setTimeout(() => { timeout.reject(new Error('lifecycle task did not settle promptly')) }, 1000)
  60. try {
  61. return await Promise.race([task, timeout.promise])
  62. } finally {
  63. clearTimeout(timer)
  64. }
  65. }
  66. /** Throw an arbitrary callback value to exercise the public unknown-error boundary. */
  67. function throwUnknown(value: unknown): never {
  68. throw value
  69. }
  70. describe('the session-persistence Agent Note: AgentLoop factory create/resume', () => {
  71. it('resumes a session persisted before messages gained identities', async () => {
  72. const sessionId = SessionId('pre-identity-resume')
  73. const first = await persistentHarness(new MockAdapter([]))
  74. await first.ctx.sessionPersistence.create({
  75. version: SESSION_FORMAT_VERSION,
  76. id: sessionId,
  77. createdAt: 1,
  78. })
  79. await first.ctx.sessionPersistence.append(sessionId, [
  80. { type: 'turn/start', seq: 0, time: 1, data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } } },
  81. {
  82. type: 'user/message',
  83. seq: 1,
  84. time: 2,
  85. data: { content: [{ type: 'text', text: 'old question' }], source: { kind: 'user' } },
  86. surfaceOp: 'append',
  87. },
  88. { type: 'step/start', seq: 2, time: 3, data: { turn: 1, step: 1 } },
  89. {
  90. type: 'assistant/message',
  91. seq: 3,
  92. time: 4,
  93. data: {
  94. turn: 1,
  95. step: 1,
  96. content: [{ type: 'text', text: 'old answer' }],
  97. provenance: { provider: 'mock', model: 'mock' },
  98. },
  99. surfaceOp: 'append',
  100. },
  101. { type: 'step/end', seq: 4, time: 5, data: { turn: 1, step: 1 } },
  102. { type: 'turn/end', seq: 5, time: 6, data: { turn: 1, reason: { kind: 'completed' } } },
  103. ] as unknown as SessionEvent[])
  104. await first.ctx.fiber.dispose()
  105. const ctx = await mountPersistentHarness(first.root, new MockAdapter([textResponse('new answer')]))
  106. const handle = await ctx.agents.resume({
  107. resumeSessionId: sessionId,
  108. agentOptions: { provider: 'mock', model: 'mock' },
  109. })
  110. expect(handle.agent.session.deriveMessages()).toMatchObject([
  111. { id: `legacy-message:${sessionId}:1`, role: 'user' },
  112. { id: `legacy-message:${sessionId}:3`, role: 'assistant' },
  113. ])
  114. handle.agent.followup(createUserMessage({
  115. content: [{ type: 'text', text: 'new question' }],
  116. source: { kind: 'user' },
  117. }))
  118. await waitForIdle(ctx, handle.agent)
  119. expect(handle.agent.session.deriveMessages()).toHaveLength(4)
  120. expect(handle.agent.session.events.at(-1)).toMatchObject({
  121. type: 'turn/end',
  122. data: { reason: { kind: 'completed' } },
  123. })
  124. await handle.dispose()
  125. await ctx.fiber.dispose()
  126. })
  127. it('normalizes a non-Error resume publication failure for rollback and rethrows it', async () => {
  128. const sessionId = SessionId('unknown-resume-failure-s')
  129. const root = await persistSession(sessionId)
  130. const ctx = await mountPersistentHarness(root, new MockAdapter([textResponse('next')]))
  131. const failure = { source: 'resume' }
  132. ctx.on('session/created', () => throwUnknown(failure))
  133. await expect(ctx.agents.resume({
  134. resumeSessionId: sessionId,
  135. })).rejects.toBe(failure)
  136. expect(ctx.agents.get(SessionId('unknown-resume-failure'))).toBeUndefined()
  137. expect(ctx.sessions.get(sessionId)).toBeUndefined()
  138. await ctx.fiber.dispose()
  139. })
  140. it('createAgent uses the caller-supplied sessionId (not ${id}-session)', async () => {
  141. const adapter = new MockAdapter([textResponse('hi')])
  142. const { ctx } = await persistentHarness(adapter)
  143. const { agent } = await ctx.agents.create({ sessionId: SessionId('custom-session'), meta: { cwd: '/w' } })
  144. expect(agent.session.id).toBe('custom-session')
  145. expect(agent.session.header.cwd).toBe('/w')
  146. await ctx.fiber.dispose()
  147. })
  148. it('createAgent rejects a duplicate identity without orphaning a session', async () => {
  149. const adapter = new MockAdapter([textResponse('hi')])
  150. const { ctx } = await persistentHarness(adapter)
  151. const sessionId = SessionId('sess-a')
  152. await ctx.agents.create({ sessionId })
  153. await expect(ctx.agents.create({ sessionId })).rejects.toThrow(/already exists/)
  154. expect(ctx.sessions.list()).toHaveLength(1)
  155. await ctx.fiber.dispose()
  156. })
  157. it('resume cannot crash-repair a turn owned by a live agent', async () => {
  158. const { ctx } = await persistentHarness(new MockAdapter([textResponse('unused')]))
  159. const sessionId = SessionId('live-resume-race')
  160. const first = (await ctx.agents.create({ sessionId })).agent
  161. first.session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  162. await ctx.sessions.flush(first.session)
  163. await expect(ctx.agents.resume({ resumeSessionId: sessionId }))
  164. .rejects.toThrow(/live turn is open/)
  165. first.session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  166. await ctx.sessions.flush(first.session)
  167. const loaded = await ctx.sessionPersistence.load(sessionId)
  168. expect(loaded.events.map(event => event.type)).toEqual(['turn/start', 'turn/end'])
  169. expect(loaded.events.at(-1)).toMatchObject({
  170. type: 'turn/end',
  171. data: { reason: { kind: 'completed' } },
  172. })
  173. await ctx.fiber.dispose()
  174. })
  175. it('createAgent works without meta (no cwd)', async () => {
  176. const adapter = new MockAdapter([textResponse('hi')])
  177. const { ctx } = await persistentHarness(adapter)
  178. const { agent } = await ctx.agents.create({ sessionId: SessionId('nometa-session') })
  179. expect(agent.session.id).toBe('nometa-session')
  180. expect(agent.session.header.cwd).toBeUndefined()
  181. await ctx.fiber.dispose()
  182. })
  183. it('resume of a session with no cwd carries an undefined cwd header', async () => {
  184. // Lifecycle 1: create a no-cwd session and run a turn.
  185. const adapter1 = new MockAdapter([textResponse('a')])
  186. const { ctx: ctx1, root } = await persistentHarness(adapter1)
  187. const a1 = (await ctx1.agents.create({ sessionId: SessionId('nocwd-sess') })).agent
  188. a1.followup(createUserMessage({ content: [{ type: 'text', text: 'q' }], source: { kind: 'user' } }))
  189. await waitForIdle(ctx1, a1)
  190. await ctx1.fiber.dispose()
  191. // Lifecycle 2: resume it; the header cwd stays undefined (no-cwd branch).
  192. const adapter2 = new MockAdapter([textResponse('b')])
  193. const ctx2 = new Context()
  194. await ctx2.plugin(LlmService)
  195. await ctx2.plugin(SessionStore)
  196. await ctx2.plugin(SystemPrompt)
  197. await ctx2.plugin(ToolRegistry)
  198. await ctx2.plugin(AgentRegistry)
  199. await ctx2.plugin(AgentLoop, { agents: [] })
  200. await ctx2.plugin(SessionPersistenceJsonl, { root })
  201. ctx2.llm.registerAdapter(['mock'], adapter2)
  202. const a2 = (await ctx2.agents.resume({ resumeSessionId: SessionId('nocwd-sess') })).agent
  203. expect(a2.session.header.cwd).toBeUndefined()
  204. await ctx2.fiber.dispose()
  205. })
  206. it('agent/session-start fires "startup" for createAgent and "resume" for resume()', async () => {
  207. // Lifecycle 1: a fresh createAgent emits session-start with source 'startup'.
  208. const adapter1 = new MockAdapter([textResponse('a')])
  209. const { ctx: ctx1, root } = await persistentHarness(adapter1)
  210. const sources1: string[] = []
  211. ctx1.on('agent/session-start', (_agent, source) => void sources1.push(source))
  212. const a1 = (await ctx1.agents.create({ sessionId: SessionId('start-sess') })).agent
  213. expect(sources1).toEqual(['startup'])
  214. a1.followup(createUserMessage({ content: [{ type: 'text', text: 'q' }], source: { kind: 'user' } }))
  215. await waitForIdle(ctx1, a1)
  216. await ctx1.fiber.dispose()
  217. // Lifecycle 2: resuming the persisted session emits session-start 'resume'.
  218. const adapter2 = new MockAdapter([textResponse('b')])
  219. const ctx2 = new Context()
  220. await ctx2.plugin(LlmService)
  221. await ctx2.plugin(SessionStore)
  222. await ctx2.plugin(SystemPrompt)
  223. await ctx2.plugin(ToolRegistry)
  224. await ctx2.plugin(AgentRegistry)
  225. await ctx2.plugin(AgentLoop, { agents: [] })
  226. await ctx2.plugin(SessionPersistenceJsonl, { root })
  227. ctx2.llm.registerAdapter(['mock'], adapter2)
  228. const sources2: string[] = []
  229. ctx2.on('agent/session-start', (_agent, source) => void sources2.push(source))
  230. await ctx2.agents.resume({ resumeSessionId: SessionId('start-sess') })
  231. expect(sources2).toEqual(['resume'])
  232. await ctx2.fiber.dispose()
  233. })
  234. it('resume awaits setup while unpublished, then publishes a fully composed world in order', async () => {
  235. const sessionId = SessionId('resume-setup-success')
  236. const root = await persistSession(sessionId)
  237. const ctx = await mountPersistentHarness(root, new MockAdapter([textResponse('next')]))
  238. const gate = Promise.withResolvers<undefined>()
  239. const setupStarted = Promise.withResolvers<undefined>()
  240. const order: string[] = []
  241. ctx.on('session/created', (session) => {
  242. expect(ctx.sessions.get(session.id)).toBe(session)
  243. expect(ctx.agents.get(sessionId)?.session).toBe(session)
  244. order.push('session/created')
  245. })
  246. ctx.on('agent/created', (agent) => {
  247. expect(agent.status).toBe('idle')
  248. order.push('agent/created')
  249. })
  250. ctx.on('agent/session-start', (agent) => {
  251. expect(() => { agent.cancel({ kind: 'user' }) }).not.toThrow()
  252. order.push('agent/session-start')
  253. })
  254. const resuming = ctx.agents.resume({
  255. resumeSessionId: sessionId,
  256. agentOptions: { provider: 'mock', model: 'mock' },
  257. setup: async (agentCtx) => {
  258. expect(agentCtx.agent?.id).toBe(sessionId)
  259. // The two persisted events plus the end-seed marker.
  260. expect(agentCtx.agent?.session.events).toHaveLength(3)
  261. agentCtx.on('session/created', () => void order.push('setup-listener:session/created'))
  262. agentCtx.on('agent/created', () => void order.push('setup-listener:agent/created'))
  263. order.push('setup:start')
  264. setupStarted.resolve(undefined)
  265. await gate.promise
  266. order.push('setup:end')
  267. return {
  268. commit: () => {
  269. expect(ctx.agents.get(sessionId)).toBeUndefined()
  270. expect(ctx.sessions.get(sessionId)).toBeUndefined()
  271. order.push('setup:commit')
  272. },
  273. }
  274. },
  275. })
  276. await setupStarted.promise
  277. expect(ctx.agents.get(sessionId)).toBeUndefined()
  278. expect(ctx.sessions.get(sessionId)).toBeUndefined()
  279. expect(order).toEqual(['setup:start'])
  280. gate.resolve(undefined)
  281. const handle = await resuming
  282. expect(order).toEqual([
  283. 'setup:start',
  284. 'setup:end',
  285. 'setup:commit',
  286. 'session/created',
  287. 'setup-listener:session/created',
  288. 'agent/created',
  289. 'setup-listener:agent/created',
  290. 'agent/session-start',
  291. ])
  292. await handle.dispose()
  293. await ctx.fiber.dispose()
  294. })
  295. it('successful resume disposal retires its caller-owned transaction effects', async () => {
  296. const sessionId = SessionId('resume-retired-effects-s')
  297. const root = await persistSession(sessionId)
  298. const ctx = await mountPersistentHarness(root, new MockAdapter([textResponse('next')]))
  299. const handle = await ctx.agents.resume({
  300. resumeSessionId: sessionId,
  301. agentOptions: { provider: 'mock', model: 'mock' },
  302. })
  303. const transactionLabels = [`agentLoop.lifecycle(${sessionId})`]
  304. expect(ctx.fiber.getEffects().map(effect => effect.label)).toEqual(expect.arrayContaining(transactionLabels))
  305. await handle.dispose()
  306. expect(ctx.fiber.getEffects().filter(effect => transactionLabels.includes(effect.label))).toEqual([])
  307. await ctx.fiber.dispose()
  308. })
  309. it('resume setup rejection publishes nothing, unwinds, and releases the identity', async () => {
  310. const sessionId = SessionId('resume-setup-reject')
  311. const root = await persistSession(sessionId)
  312. const ctx = await mountPersistentHarness(root, new MockAdapter([textResponse('next')]))
  313. const published: string[] = []
  314. ctx.on('session/created', () => void published.push('session/created'))
  315. ctx.on('agent/created', () => void published.push('agent/created'))
  316. ctx.on('agent/session-start', () => void published.push('agent/session-start'))
  317. await expect(ctx.agents.resume({
  318. resumeSessionId: sessionId,
  319. agentOptions: { provider: 'mock', model: 'mock' },
  320. setup: async () => {
  321. await Promise.resolve()
  322. throw new Error('resume setup failed')
  323. },
  324. })).rejects.toThrow('resume setup failed')
  325. expect(published).toEqual([])
  326. expect(ctx.agents.get(sessionId)).toBeUndefined()
  327. expect(ctx.sessions.get(sessionId)).toBeUndefined()
  328. const retry = await ctx.agents.resume({
  329. resumeSessionId: sessionId,
  330. agentOptions: { provider: 'mock', model: 'mock' },
  331. })
  332. await retry.dispose()
  333. await ctx.fiber.dispose()
  334. })
  335. it('resume setup commit rejection publishes nothing and releases the identity', async () => {
  336. const sessionId = SessionId('resume-setup-commit-reject')
  337. const root = await persistSession(sessionId)
  338. const ctx = await mountPersistentHarness(root, new MockAdapter([textResponse('next')]))
  339. const published: string[] = []
  340. ctx.on('session/created', () => void published.push('session/created'))
  341. ctx.on('agent/created', () => void published.push('agent/created'))
  342. await expect(ctx.agents.resume({
  343. resumeSessionId: sessionId,
  344. agentOptions: { provider: 'mock', model: 'mock' },
  345. setup: () => ({
  346. commit: () => { throw new Error('resume setup commit failed') },
  347. }),
  348. })).rejects.toThrow('resume setup commit failed')
  349. expect(published).toEqual([])
  350. expect(ctx.agents.get(sessionId)).toBeUndefined()
  351. expect(ctx.sessions.get(sessionId)).toBeUndefined()
  352. const retry = await ctx.agents.resume({
  353. resumeSessionId: sessionId,
  354. agentOptions: { provider: 'mock', model: 'mock' },
  355. })
  356. await retry.dispose()
  357. await ctx.fiber.dispose()
  358. })
  359. it('owner unload aborts resume setup and cannot publish after the callback settles', async () => {
  360. const sessionId = SessionId('resume-setup-owner-unload')
  361. const root = await persistSession(sessionId)
  362. const ctx = await mountPersistentHarness(root, new MockAdapter([textResponse('next')]))
  363. const gate = Promise.withResolvers<undefined>()
  364. const setupStarted = Promise.withResolvers<undefined>()
  365. const published: string[] = []
  366. ctx.on('session/created', () => void published.push('session/created'))
  367. ctx.on('agent/created', () => void published.push('agent/created'))
  368. let resuming!: ReturnType<typeof ctx.agents.resume>
  369. const owner = await ctx.plugin(Object.assign((inner: Context) => {
  370. resuming = inner.agents.resume({
  371. resumeSessionId: sessionId,
  372. agentOptions: { provider: 'mock', model: 'mock' },
  373. setup: async () => {
  374. setupStarted.resolve(undefined)
  375. await gate.promise
  376. },
  377. })
  378. }, { inject: ['agents'] }))
  379. await setupStarted.promise
  380. await owner.dispose()
  381. await expect(resuming).rejects.toThrow(/owner disposed during setup/)
  382. expect(published).toEqual([])
  383. expect(ctx.agents.get(sessionId)).toBeUndefined()
  384. expect(ctx.sessions.get(sessionId)).toBeUndefined()
  385. gate.resolve(undefined)
  386. await Promise.resolve()
  387. expect(published).toEqual([])
  388. await ctx.fiber.dispose()
  389. })
  390. it('owner unload aborts a never-settling persistence load, releases the identity, and blocks late publication', async () => {
  391. const sessionId = SessionId('resume-load-owner-unload')
  392. const root = await persistSession(sessionId)
  393. const ctx = await mountPersistentHarness(root, new MockAdapter([textResponse('next')]))
  394. const snapshot = await ctx.sessionPersistence.load(sessionId)
  395. const lateLoad = Promise.withResolvers<typeof snapshot>()
  396. const loadStarted = Promise.withResolvers<undefined>()
  397. let loads = 0
  398. ctx.sessionPersistence.load = (id) => {
  399. expect(id).toBe(sessionId)
  400. loads += 1
  401. if (loads === 1) {
  402. loadStarted.resolve(undefined)
  403. return lateLoad.promise
  404. }
  405. return Promise.resolve(structuredClone(snapshot))
  406. }
  407. const published: string[] = []
  408. ctx.on('session/created', () => void published.push('session/created'))
  409. ctx.on('agent/created', () => void published.push('agent/created'))
  410. ctx.on('agent/session-start', () => void published.push('agent/session-start'))
  411. let resuming!: ReturnType<typeof ctx.agents.resume>
  412. const owner = await ctx.plugin(Object.assign((inner: Context) => {
  413. resuming = inner.agents.resume({ resumeSessionId: sessionId, agentOptions: { provider: 'mock', model: 'mock' } })
  414. }, { inject: ['agents'] }))
  415. await loadStarted.promise
  416. const rejection = expect(promptly(resuming)).rejects.toThrow(/owner disposed during setup/)
  417. await promptly(owner.dispose())
  418. expect(published).toEqual([])
  419. expect(ctx.agents.get(sessionId)).toBeUndefined()
  420. expect(ctx.sessions.get(sessionId)).toBeUndefined()
  421. // owner.dispose() awaited transaction settlement, so the same identities
  422. // can be reused before awaiting the public rejection.
  423. const retry = await promptly(ctx.agents.resume({ resumeSessionId: sessionId, agentOptions: { provider: 'mock', model: 'mock' } }))
  424. await rejection
  425. expect(loads).toBe(2)
  426. expect(published).toEqual(['session/created', 'agent/created', 'agent/session-start'])
  427. // Settlement of the abandoned backend promise cannot resume the old
  428. // transaction or emit a second publication after the retry owns the ids.
  429. lateLoad.resolve(structuredClone(snapshot))
  430. await Promise.resolve()
  431. await Promise.resolve()
  432. expect(ctx.agents.get(sessionId)).toBe(retry.agent)
  433. expect(ctx.sessions.get(sessionId)).toBe(retry.agent.session)
  434. expect(published).toEqual(['session/created', 'agent/created', 'agent/session-start'])
  435. await retry.dispose()
  436. await ctx.fiber.dispose()
  437. })
  438. it('AgentLoop unload aborts persistence load and awaits wrapper settlement', async () => {
  439. const sessionId = SessionId('resume-load-factory-unload')
  440. const root = await persistSession(sessionId)
  441. const ctx = new Context()
  442. await ctx.plugin(LlmService)
  443. await ctx.plugin(SessionStore)
  444. await ctx.plugin(SystemPrompt)
  445. await ctx.plugin(ToolRegistry)
  446. await ctx.plugin(AgentRegistry)
  447. const loopFiber = await ctx.plugin(AgentLoop, { agents: [] })
  448. await ctx.plugin(SessionPersistenceJsonl, { root })
  449. ctx.llm.registerAdapter(['mock'], new MockAdapter([textResponse('next')]))
  450. const snapshot = await ctx.sessionPersistence.load(sessionId)
  451. const lateLoad = Promise.withResolvers<typeof snapshot>()
  452. const loadStarted = Promise.withResolvers<undefined>()
  453. ctx.sessionPersistence.load = (id) => {
  454. expect(id).toBe(sessionId)
  455. loadStarted.resolve(undefined)
  456. return lateLoad.promise
  457. }
  458. const published: string[] = []
  459. ctx.on('session/created', () => void published.push('session/created'))
  460. ctx.on('agent/created', () => void published.push('agent/created'))
  461. const resuming = ctx.agents.resume({ resumeSessionId: sessionId, agentOptions: { provider: 'mock', model: 'mock' } })
  462. await loadStarted.promise
  463. const rejection = expect(promptly(resuming)).rejects.toThrow(/agent loop is not active/)
  464. await promptly(loopFiber.dispose())
  465. await rejection
  466. expect(published).toEqual([])
  467. expect(ctx.agents.get(sessionId)).toBeUndefined()
  468. expect(ctx.sessions.get(sessionId)).toBeUndefined()
  469. lateLoad.resolve(structuredClone(snapshot))
  470. await Promise.resolve()
  471. await Promise.resolve()
  472. expect(published).toEqual([])
  473. await ctx.fiber.dispose()
  474. })
  475. it('resume of a forked session preserves the lineage, seed boundary, and delegation depth in the header', async () => {
  476. // Lifecycle 1: persist a FORKED session (carries parentSession + seedLength
  477. // in its header) by creating it with a complete-turn seed — the write path
  478. // materializes the fork (header + seed) on disk.
  479. const seed: SessionEvent[] = [
  480. { type: 'turn/start', seq: 0, time: 1, data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } } },
  481. { type: 'turn/end', seq: 1, time: 2, data: { turn: 1, reason: { kind: 'completed' } } },
  482. ]
  483. const adapter1 = new MockAdapter([textResponse('a')])
  484. const { ctx: ctx1, root } = await persistentHarness(adapter1)
  485. const forked = ctx1.sessions.create(SessionId('forked-sess'), {
  486. seed,
  487. meta: { cwd: '/w', parentSession: SessionId('parent-sess'), seedLength: seed.length, delegationDepth: 1 },
  488. })
  489. await ctx1.sessions.flush(forked)
  490. await ctx1.fiber.dispose()
  491. // Lifecycle 2: resume it; the parentSession + seedLength header survives the
  492. // round-trip (exercises resume's parentSession- and seedLength-present
  493. // branches). seedLength must come from the PERSISTED header, not from the
  494. // resume seed length (which is the whole stored log, not the original
  495. // boundary).
  496. const adapter2 = new MockAdapter([textResponse('b')])
  497. const ctx2 = new Context()
  498. await ctx2.plugin(LlmService)
  499. await ctx2.plugin(SessionStore)
  500. await ctx2.plugin(SystemPrompt)
  501. await ctx2.plugin(ToolRegistry)
  502. await ctx2.plugin(AgentRegistry)
  503. await ctx2.plugin(AgentLoop, { agents: [] })
  504. await ctx2.plugin(SessionPersistenceJsonl, { root })
  505. ctx2.llm.registerAdapter(['mock'], adapter2)
  506. const a2 = (await ctx2.agents.resume({ resumeSessionId: SessionId('forked-sess') })).agent
  507. expect(a2.session.header.parentSession).toBe('parent-sess')
  508. expect(a2.session.header.cwd).toBe('/w')
  509. expect(a2.session.header.seedLength).toBe(seed.length)
  510. // The recursion budget survives resume — a dropped depth would let a
  511. // resumed child delegate as if it were top-level.
  512. expect(a2.session.header.delegationDepth).toBe(1)
  513. await ctx2.fiber.dispose()
  514. })
  515. it('an idle inject() survives persist + resume without a synthetic turn', async () => {
  516. const adapter1 = new MockAdapter([textResponse('answer')])
  517. const { ctx: ctx1, root } = await persistentHarness(adapter1)
  518. const a1 = (await ctx1.agents.create({ sessionId: SessionId('inject-sess'), meta: { cwd: '/w' } })).agent
  519. a1.followup(createUserMessage({ content: [{ type: 'text', text: 'q' }], source: { kind: 'user' } }))
  520. await waitForIdle(ctx1, a1)
  521. a1.inject(createUserMessage({ content: [{ type: 'text', text: 'background task 42 finished' }], source: { kind: 'plugin', plugin: 'tool-bash' } }))
  522. await a1.whenIdle()
  523. await ctx1.fiber.dispose()
  524. // Lifecycle 2: resume; the injected context is still in the derived history.
  525. const adapter2 = new MockAdapter([textResponse('next')])
  526. const ctx2 = new Context()
  527. await ctx2.plugin(LlmService)
  528. await ctx2.plugin(SessionStore)
  529. await ctx2.plugin(SystemPrompt)
  530. await ctx2.plugin(ToolRegistry)
  531. await ctx2.plugin(AgentRegistry)
  532. await ctx2.plugin(AgentLoop, { agents: [] })
  533. await ctx2.plugin(SessionPersistenceJsonl, { root })
  534. ctx2.llm.registerAdapter(['mock'], adapter2)
  535. const a2 = (await ctx2.agents.resume({ resumeSessionId: SessionId('inject-sess') })).agent
  536. const flat = JSON.stringify(a2.session.deriveMessages())
  537. expect(flat).toContain('background task 42 finished')
  538. await ctx2.fiber.dispose()
  539. })
  540. it('resume reloads a persisted session: history + turn numbering continue, no duplicate seqs', async () => {
  541. // Lifecycle 1: run one full turn, persisting it.
  542. const adapter1 = new MockAdapter([textResponse('first answer')])
  543. const { ctx: ctx1, root } = await persistentHarness(adapter1)
  544. const a1 = (await ctx1.agents.create({ sessionId: SessionId('sess-resume'), meta: { cwd: '/w' } })).agent
  545. a1.followup(createUserMessage({ content: [{ type: 'text', text: 'first question' }], source: { kind: 'user' } }))
  546. await waitForIdle(ctx1, a1)
  547. const events1 = [...a1.session.events]
  548. const seqs1 = events1.map(e => e.seq)
  549. expect(seqs1).toEqual([...seqs1].sort((x, y) => x - y)) // contiguous
  550. await ctx1.fiber.dispose()
  551. // Lifecycle 2: a brand-new context over the SAME root; resume the session.
  552. const adapter2 = new MockAdapter([textResponse('second answer')])
  553. const ctx2 = new Context()
  554. await ctx2.plugin(LlmService)
  555. await ctx2.plugin(SessionStore)
  556. await ctx2.plugin(SystemPrompt)
  557. await ctx2.plugin(ToolRegistry)
  558. await ctx2.plugin(AgentRegistry)
  559. await ctx2.plugin(AgentLoop, { agents: [] })
  560. await ctx2.plugin(SessionPersistenceJsonl, { root })
  561. ctx2.llm.registerAdapter(['mock'], adapter2)
  562. const a2 = (await ctx2.agents.resume({ resumeSessionId: SessionId('sess-resume') })).agent
  563. // The resumed session carries the prior history…
  564. expect(a2.session.id).toBe('sess-resume')
  565. // …followed by one end-seed event marking the constructor seed.
  566. expect(a2.session.events.length).toBe(events1.length + 1)
  567. expect(a2.session.firstLiveSeq).toBe(events1.length)
  568. expect(a2.session.events.at(-1)?.type).toBe('session/end-seed')
  569. const replay = Session.create(SessionId('replay'), events1)
  570. expect(a2.session.deriveMessages()).toEqual(replay.deriveMessages())
  571. // …and a new turn continues numbering (turn 2) with contiguous seqs.
  572. a2.followup(createUserMessage({ content: [{ type: 'text', text: 'second question' }], source: { kind: 'user' } }))
  573. await waitForIdle(ctx2, a2)
  574. const allSeqs = a2.session.events.map(e => e.seq)
  575. expect(allSeqs).toEqual(allSeqs.map((_, i) => i)) // 0..N contiguous, no duplicates
  576. const turnStarts = a2.session.events.filter(e => e.type === 'turn/start')
  577. expect(turnStarts.map(e => e.type === 'turn/start' && e.data.turn)).toEqual([1, 2])
  578. await ctx2.fiber.dispose()
  579. })
  580. it('resume rejects when session persistence is not configured', async () => {
  581. // A harness WITHOUT the persistence plugin.
  582. const adapter = new MockAdapter([textResponse('x')])
  583. const ctx = new Context()
  584. await ctx.plugin(LlmService)
  585. await ctx.plugin(SessionStore)
  586. await ctx.plugin(SystemPrompt)
  587. await ctx.plugin(ToolRegistry)
  588. await ctx.plugin(AgentRegistry)
  589. await ctx.plugin(AgentLoop, { agents: [] })
  590. ctx.llm.registerAdapter(['mock'], adapter)
  591. await expect(ctx.agents.resume({ resumeSessionId: SessionId('nope') }))
  592. .rejects.toThrow(/session persistence is not configured/)
  593. await ctx.fiber.dispose()
  594. })
  595. })
  596. describe('creation and resume cancellation edges', () => {
  597. it('rejects create() with a pre-aborted signal, including a non-Error reason', async () => {
  598. const { ctx } = await persistentHarness(new MockAdapter([]))
  599. const errorReason = new AbortController()
  600. errorReason.abort(new Error('caller gave up'))
  601. await expect(promptly(ctx.agents.create({
  602. sessionId: SessionId('pre-aborted-error'),
  603. agentOptions: { provider: 'mock', model: 'mock' },
  604. signal: errorReason.signal,
  605. }))).rejects.toThrow('caller gave up')
  606. // A non-Error reason is wrapped into the creation-aborted error.
  607. const stringReason = new AbortController()
  608. stringReason.abort('operator string reason')
  609. await expect(promptly(ctx.agents.create({
  610. sessionId: SessionId('pre-aborted-string'),
  611. agentOptions: { provider: 'mock', model: 'mock' },
  612. signal: stringReason.signal,
  613. }))).rejects.toThrow(/creation aborted/)
  614. expect(ctx.agents.get(SessionId('pre-aborted-error'))).toBeUndefined()
  615. expect(ctx.agents.get(SessionId('pre-aborted-string'))).toBeUndefined()
  616. await ctx.fiber.dispose()
  617. })
  618. it('a non-Error abort reason arriving during setup is wrapped for the caller', async () => {
  619. const { ctx } = await persistentHarness(new MockAdapter([]))
  620. const controller = new AbortController()
  621. const setupEntered = Promise.withResolvers<undefined>()
  622. const setupGate = Promise.withResolvers<undefined>()
  623. const creating = ctx.agents.create({
  624. sessionId: SessionId('setup-string-abort'),
  625. agentOptions: { provider: 'mock', model: 'mock' },
  626. signal: controller.signal,
  627. async setup() {
  628. setupEntered.resolve(undefined)
  629. await setupGate.promise
  630. },
  631. })
  632. await setupEntered.promise
  633. controller.abort('mid-setup string reason')
  634. setupGate.resolve(undefined)
  635. await expect(promptly(creating)).rejects.toThrow(/creation aborted/)
  636. expect(ctx.agents.get(SessionId('setup-string-abort'))).toBeUndefined()
  637. await ctx.fiber.dispose()
  638. })
  639. it('resume with a pre-aborted caller signal rejects out of the load race', async () => {
  640. const sessionId = SessionId('resume-pre-aborted')
  641. const root = await persistSession(sessionId)
  642. const ctx = await mountPersistentHarness(root, new MockAdapter([]))
  643. const controller = new AbortController()
  644. controller.abort(new Error('resume abandoned'))
  645. await expect(promptly(ctx.agents.resume({
  646. resumeSessionId: sessionId,
  647. agentOptions: { provider: 'mock', model: 'mock' },
  648. signal: controller.signal,
  649. }))).rejects.toThrow('resume abandoned')
  650. expect(ctx.agents.get(sessionId)).toBeUndefined()
  651. await ctx.fiber.dispose()
  652. })
  653. it('factory teardown during a hung resume load rejects with loop-inactive', async () => {
  654. const sessionId = SessionId('resume-loop-teardown')
  655. const root = await persistSession(sessionId)
  656. const ctx = await mountPersistentHarness(root, new MockAdapter([]))
  657. const snapshot = await ctx.sessionPersistence.load(sessionId)
  658. const gate = Promise.withResolvers<typeof snapshot>()
  659. const loadStarted = Promise.withResolvers<undefined>()
  660. ctx.sessionPersistence.load = () => {
  661. loadStarted.resolve(undefined)
  662. return gate.promise
  663. }
  664. const resuming = ctx.agents.resume({
  665. resumeSessionId: sessionId,
  666. agentOptions: { provider: 'mock', model: 'mock' },
  667. })
  668. await loadStarted.promise
  669. // Resolve the load only after teardown began: the post-load ownership
  670. // check, not the abort race, must reject the wrapper.
  671. const rejection = expect(promptly(resuming)).rejects.toThrow()
  672. const disposal = ctx.fiber.dispose()
  673. gate.resolve(structuredClone(snapshot))
  674. await rejection
  675. await disposal
  676. })
  677. })
  678. describe('configured-start failure edges', () => {
  679. it('a non-Error mid-load abort reason is wrapped for the resume caller', async () => {
  680. const sessionId = SessionId('resume-string-mid-abort')
  681. const root = await persistSession(sessionId)
  682. const ctx = await mountPersistentHarness(root, new MockAdapter([]))
  683. const gate = Promise.withResolvers<never>()
  684. gate.promise.catch(() => undefined)
  685. const loadStarted = Promise.withResolvers<undefined>()
  686. ctx.sessionPersistence.load = () => {
  687. loadStarted.resolve(undefined)
  688. return gate.promise
  689. }
  690. const controller = new AbortController()
  691. const resuming = ctx.agents.resume({
  692. resumeSessionId: sessionId,
  693. agentOptions: { provider: 'mock', model: 'mock' },
  694. signal: controller.signal,
  695. })
  696. await loadStarted.promise
  697. controller.abort('operator string reason')
  698. await expect(promptly(resuming)).rejects.toThrow(/creation aborted/)
  699. expect(ctx.agents.get(sessionId)).toBeUndefined()
  700. await ctx.fiber.dispose()
  701. })
  702. it('a failing exact-id restore over an existing artifact stays loud', async () => {
  703. const sessionId = SessionId('config-existing-corrupt')
  704. const root = await persistSession(sessionId)
  705. const ctx = await mountPersistentHarness(root, new MockAdapter([]))
  706. // The artifact exists (list reports it) but its load fails: this is
  707. // corruption, not first creation — the failure must be reported, and no
  708. // fresh same-id session may shadow the broken one.
  709. ctx.sessionPersistence.load = () => Promise.reject(new Error('artifact corrupt'))
  710. const configured = new Context()
  711. await configured.plugin(LlmService)
  712. await configured.plugin(SessionStore)
  713. await configured.plugin(SystemPrompt)
  714. await configured.plugin(ToolRegistry)
  715. await configured.plugin(AgentRegistry)
  716. await configured.plugin(SessionPersistenceJsonl, { root })
  717. configured.llm.registerAdapter(['mock'], new MockAdapter([]))
  718. configured.sessionPersistence.load = id => ctx.sessionPersistence.load(id)
  719. const configFailures: unknown[] = []
  720. configured.on('agent-loop/config-start-failed', (_id, error) => { configFailures.push(error) })
  721. const configWarnings: string[] = []
  722. const configWarn = configured.logger.warn.bind(configured.logger)
  723. configured.logger.warn = ((...args: unknown[]) => {
  724. if (typeof args[0] === 'string') configWarnings.push(args[0])
  725. return (configWarn as (...a: unknown[]) => unknown)(...args)
  726. }) as typeof configured.logger.warn
  727. const loop = await configured.plugin(AgentLoop, {
  728. agents: [{ id: 'main', sessionId, provider: 'mock', model: 'mock' }],
  729. })
  730. await expect.poll(() => configFailures.length).toBe(1)
  731. expect(configFailures[0]).toBeInstanceOf(Error)
  732. expect((configFailures[0] as Error).message).toBe('artifact corrupt')
  733. expect(configWarnings.some(w => w.includes('config-driven restore'))).toBe(true)
  734. expect(configured.agents.get(sessionId)).toBeUndefined()
  735. await loop.dispose()
  736. await configured.fiber.dispose()
  737. await ctx.fiber.dispose()
  738. })
  739. it('suppresses a configured-resume failure that lands after teardown', async () => {
  740. const sessionId = SessionId('config-late-resume-failure')
  741. const root = await persistSession(sessionId)
  742. const ctx = await mountPersistentHarness(root, new MockAdapter([]))
  743. const gate = Promise.withResolvers<never>()
  744. gate.promise.catch(() => undefined)
  745. const loadStarted = Promise.withResolvers<undefined>()
  746. ctx.sessionPersistence.load = () => {
  747. loadStarted.resolve(undefined)
  748. return gate.promise
  749. }
  750. const failures: unknown[] = []
  751. ctx.on('agent-loop/config-start-failed', (_id, error) => { failures.push(error) })
  752. const configured = new Context()
  753. await configured.plugin(LlmService)
  754. await configured.plugin(SessionStore)
  755. await configured.plugin(SystemPrompt)
  756. await configured.plugin(ToolRegistry)
  757. await configured.plugin(AgentRegistry)
  758. await configured.plugin(SessionPersistenceJsonl, { root })
  759. configured.llm.registerAdapter(['mock'], new MockAdapter([]))
  760. configured.sessionPersistence.load = id => ctx.sessionPersistence.load(id)
  761. configured.on('agent-loop/config-start-failed', (_id, error) => { failures.push(error) })
  762. const loop = await configured.plugin(AgentLoop, {
  763. agents: [{ id: 'main', resumeSessionId: sessionId, provider: 'mock', model: 'mock' }],
  764. })
  765. await loadStarted.promise
  766. const disposal = loop.dispose()
  767. gate.reject(new Error('late backend failure'))
  768. await disposal
  769. await new Promise(r => setTimeout(r, 20))
  770. // Ownership deactivated before the failure landed: the report is dropped.
  771. expect(failures).toEqual([])
  772. await configured.fiber.dispose()
  773. await ctx.fiber.dispose()
  774. })
  775. })