persistence.spec.ts 20 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514
  1. import { afterEach, describe, expect, it, vi } from 'vitest'
  2. import { mkdtempSync } from 'node:fs'
  3. import { rm } from 'node:fs/promises'
  4. import { tmpdir } from 'node:os'
  5. import { join } from 'node:path'
  6. import { Context } from '@deepseek-ai/cordis'
  7. import type { Agent, AgentHandle } from '@deepseek-ai/dsh-agent'
  8. import AgentLoop from '@deepseek-ai/dsh-agent-loop'
  9. import { mountAgentLoopTestDependencies } from '@deepseek-ai/dsh-agent-loop-testkit'
  10. import { createUserMessage } from '@deepseek-ai/dsh-llm'
  11. import { SessionId, type SessionEvent } from '@deepseek-ai/dsh-session'
  12. import JsonlSessionPersistence from '@deepseek-ai/dsh-session-persistence-jsonl'
  13. import SubagentService, { snapshotSubagentDescriptor } from '@deepseek-ai/dsh-subagent'
  14. import * as SubagentSpawn from '@deepseek-ai/dsh-subagent-spawn-in-process'
  15. import { MockAdapter, textResponse } from '../../../core/agent-loop/tests/mock-adapter.ts'
  16. import TeamService, { TeamId, TeamMessageId } from '../src/index.ts'
  17. import type { TeamMailbox } from '../src/mailbox.ts'
  18. import { teamProjectionDefinition } from '../src/projection.ts'
  19. import type { TeamMemberSnapshot, TeamMessageSnapshot, TeamTaskSnapshot } from '../src/index.ts'
  20. import { TestSessionQuery } from './test-session-query.ts'
  21. const SIGNAL = new AbortController().signal
  22. const PERSISTENCE_TEST_TIMEOUT_MS = 15_000
  23. const roots: string[] = []
  24. const contexts = new Set<Context>()
  25. /** Detached durable Team read through the same projection definition as the service. */
  26. function durable(agent: Agent): {
  27. members: TeamMemberSnapshot[]
  28. tasks: TeamTaskSnapshot[]
  29. pendingMessages: TeamMessageSnapshot[]
  30. } {
  31. let projected = teamProjectionDefinition.init(agent.session.header)
  32. for (const event of agent.session.snapshotEvents()) projected = teamProjectionDefinition.apply(projected, event)
  33. if (projected.failure !== undefined) throw new Error(projected.failure)
  34. const state = projected
  35. return {
  36. members: state.members,
  37. tasks: state.tasks,
  38. pendingMessages: state.messages.filter(message => !state.delivered.includes(message.id)),
  39. }
  40. }
  41. /** Read one stored session's full event log through a short-lived read handle. */
  42. async function storedEvents(ctx: Context, id: SessionId): Promise<readonly SessionEvent[]> {
  43. const handle = await ctx.sessionPersistence.open(id, 'read')
  44. try {
  45. return (await handle.read()).events
  46. } finally {
  47. await handle.close()
  48. }
  49. }
  50. /** Await mailbox acknowledgements through their flush and dispatch completion. */
  51. async function settleMailbox(ctx: Context): Promise<void> {
  52. const { mailbox } = ctx.agentTeams as unknown as { readonly mailbox: TeamMailbox }
  53. await Promise.all(mailbox.pendingDispatches())
  54. }
  55. async function disposeContext(ctx: Context): Promise<void> {
  56. try {
  57. await ctx.fiber.dispose()
  58. } finally {
  59. contexts.delete(ctx)
  60. }
  61. }
  62. afterEach(async () => {
  63. const failures: unknown[] = []
  64. for (const ctx of [...contexts].reverse()) {
  65. try {
  66. await disposeContext(ctx)
  67. } catch (error: unknown) {
  68. failures.push(error)
  69. }
  70. }
  71. for (const root of roots.splice(0)) {
  72. try {
  73. await rm(root, { recursive: true, force: true, maxRetries: 5, retryDelay: 100 })
  74. } catch (error: unknown) {
  75. failures.push(error)
  76. }
  77. }
  78. if (failures.length > 0) throw new AggregateError(failures, 'Agent Teams persistence test cleanup failed')
  79. })
  80. interface PersistenceMount {
  81. readonly name: string
  82. mount(ctx: Context, root: string): Promise<{ dispose(): Promise<void> }>
  83. }
  84. const backends: PersistenceMount[] = [
  85. {
  86. name: 'JSONL',
  87. mount: async (ctx, root) => await ctx.plugin(JsonlSessionPersistence, {
  88. root: join(root, 'jsonl'),
  89. compression: 'none',
  90. }),
  91. },
  92. ]
  93. async function stack(
  94. backend: PersistenceMount,
  95. root: string,
  96. script: ConstructorParameters<typeof MockAdapter>[0],
  97. ) {
  98. const ctx = new Context()
  99. contexts.add(ctx)
  100. await mountAgentLoopTestDependencies(ctx)
  101. await backend.mount(ctx, root)
  102. await ctx.plugin(TestSessionQuery)
  103. await ctx.plugin(AgentLoop, { agents: [] })
  104. await ctx.plugin(SubagentService)
  105. await ctx.plugin(SubagentSpawn, { providerName: 'spawn' })
  106. await ctx.plugin(TeamService)
  107. const adapter = new MockAdapter(script)
  108. ctx.llm.registerAdapter(['mock'], adapter)
  109. return {
  110. ctx,
  111. adapter,
  112. dispose: async () => { await disposeContext(ctx) },
  113. }
  114. }
  115. function provisioning(childId: SessionId, name: string): TeamMemberSnapshot {
  116. return {
  117. id: childId,
  118. name,
  119. description: `${name} recovery`,
  120. provider: 'spawn',
  121. context: 'fresh',
  122. phase: 'provisioning',
  123. }
  124. }
  125. async function persistedChild(
  126. ctx: Context,
  127. rootId: SessionId,
  128. childId: SessionId,
  129. message: ReturnType<typeof createUserMessage>,
  130. ) {
  131. const descriptor = snapshotSubagentDescriptor({
  132. mode: 'continuable',
  133. provider: 'spawn',
  134. label: 'persisted child fixture',
  135. agentProvider: 'mock',
  136. agentModel: 'mock',
  137. })
  138. const child = ctx.sessions.create(childId, {
  139. meta: { parentSession: rootId, origin: 'subagent' },
  140. })
  141. child.append('subagent/descriptor', descriptor)
  142. child.append('agent/inbox/spliced', {
  143. target: 'next-turn',
  144. start: 0,
  145. inserted: [message],
  146. })
  147. // Live sessions persist only through an attached agent-loop writer; this
  148. // bare fixture session seeds its durable log directly for the cold restart.
  149. const handle = await ctx.sessionPersistence.create(child.header)
  150. await handle.append(child.snapshotEvents())
  151. await handle.close()
  152. return child
  153. }
  154. for (const backend of backends) {
  155. describe(`${backend.name} Agent Teams recovery`, () => {
  156. it('reconciles a persisted child to active and a missing child to durable failed', {
  157. timeout: PERSISTENCE_TEST_TIMEOUT_MS,
  158. }, async () => {
  159. const storageRoot = mkdtempSync(join(tmpdir(), `dsh-team-${backend.name.toLowerCase()}-`))
  160. roots.push(storageRoot)
  161. const first = await stack(backend, storageRoot, [textResponse('initial child answer')])
  162. const activeRootId = SessionId(`${backend.name.toLowerCase()}-active-root`)
  163. const failedRootId = SessionId(`${backend.name.toLowerCase()}-failed-root`)
  164. const childId = SessionId(`${backend.name.toLowerCase()}-child`)
  165. const activeRoot = await first.ctx.agentLoop.create(activeRootId, { provider: 'mock', model: 'mock' })
  166. const failedRoot = await first.ctx.agentLoop.create(failedRootId, { provider: 'mock', model: 'mock' })
  167. // Let each root's startup recovery observe the empty initial log before
  168. // simulating the crash-only provisioning prefix.
  169. await Promise.resolve()
  170. await Promise.resolve()
  171. activeRoot.session.append('team/member', {
  172. version: 2,
  173. teamId: TeamId(activeRoot.id),
  174. member: provisioning(childId, 'recoverable'),
  175. })
  176. failedRoot.session.append('team/member', {
  177. version: 2,
  178. teamId: TeamId(failedRoot.id),
  179. member: provisioning(SessionId(`${backend.name}-missing`), 'missing'),
  180. })
  181. await Promise.all([
  182. first.ctx.sessions.flush(activeRoot.session),
  183. first.ctx.sessions.flush(failedRoot.session),
  184. ])
  185. await first.ctx.subagents.startContinuable({
  186. childId,
  187. provider: 'spawn',
  188. label: 'recoverable recovery',
  189. request: {
  190. prompt: [{ type: 'text', text: 'persist before active edge' }],
  191. parent: activeRoot,
  192. },
  193. signal: SIGNAL,
  194. })
  195. await vi.waitFor(() => { expect(first.ctx.agents.get(childId)).toBeUndefined() }, { timeout: 5_000 })
  196. expect((await storedEvents(first.ctx, childId))
  197. .some(event => event.type === 'user/message')).toBe(true)
  198. await first.dispose()
  199. const second = await stack(backend, storageRoot, [textResponse('cold resumed answer')])
  200. const activeHandle = await second.ctx.agents.resume({
  201. resumeSessionId: activeRootId,
  202. agentOptions: { provider: 'mock', model: 'mock' },
  203. })
  204. const failedHandle = await second.ctx.agents.resume({
  205. resumeSessionId: failedRootId,
  206. agentOptions: { provider: 'mock', model: 'mock' },
  207. })
  208. await vi.waitFor(() => {
  209. expect(durable(activeHandle.agent).members[0]?.phase).toBe('active')
  210. const failedMember = durable(failedHandle.agent).members[0]
  211. expect(failedMember?.phase).toBe('failed')
  212. expect(failedMember?.error).toContain('child Session recovery failed')
  213. }, { timeout: 5_000 })
  214. const receipt = await second.ctx.agentTeams.sendMessage(activeHandle.agent, {
  215. target: 'recoverable',
  216. content: [{ type: 'text', text: 'resume after reconciliation' }],
  217. signal: SIGNAL,
  218. })
  219. expect(receipt.status).toBe('accepted')
  220. await vi.waitFor(() => { expect(second.ctx.agents.get(childId)).toBeUndefined() }, { timeout: 5_000 })
  221. await vi.waitFor(() => { expect(durable(activeHandle.agent).pendingMessages).toEqual([]) })
  222. await activeHandle.dispose()
  223. await failedHandle.dispose()
  224. await second.dispose()
  225. })
  226. it('reconciles a provisioning child whose initial prompt is durably pending', {
  227. timeout: PERSISTENCE_TEST_TIMEOUT_MS,
  228. }, async () => {
  229. const storageRoot = mkdtempSync(join(tmpdir(), `dsh-team-pending-${backend.name.toLowerCase()}-`))
  230. roots.push(storageRoot)
  231. const rootId = SessionId(`${backend.name.toLowerCase()}-pending-root`)
  232. const childId = SessionId(`${backend.name.toLowerCase()}-pending-child`)
  233. const first = await stack(backend, storageRoot, [])
  234. const root = await first.ctx.agentLoop.create(rootId, { provider: 'mock', model: 'mock' })
  235. await Promise.resolve()
  236. await Promise.resolve()
  237. root.session.append('team/member', {
  238. version: 2,
  239. teamId: TeamId(root.id),
  240. member: provisioning(childId, 'pending-worker'),
  241. })
  242. const initial = createUserMessage({
  243. content: [{ type: 'text', text: 'durably pending initial task' }],
  244. source: { kind: 'user' },
  245. })
  246. await persistedChild(first.ctx, rootId, childId, initial)
  247. await first.ctx.sessions.flush(root.session)
  248. await first.dispose()
  249. const second = await stack(backend, storageRoot, [])
  250. const rootHandle = await second.ctx.agents.resume({
  251. resumeSessionId: rootId,
  252. agentOptions: { provider: 'mock', model: 'mock' },
  253. })
  254. await vi.waitFor(() => {
  255. expect(durable(rootHandle.agent).members[0]?.phase).toBe('active')
  256. })
  257. expect(second.adapter.requests).toEqual([])
  258. const stored = await storedEvents(second.ctx, childId)
  259. expect(stored.some(event => event.type === 'agent/inbox/spliced'
  260. && event.data.inserted.some(message => message.id === initial.id))).toBe(true)
  261. await rootHandle.dispose()
  262. await second.dispose()
  263. })
  264. it('retries queued mail through cold-resume Steer after restart', {
  265. timeout: PERSISTENCE_TEST_TIMEOUT_MS,
  266. }, async () => {
  267. const storageRoot = mkdtempSync(join(tmpdir(), `dsh-team-mail-${backend.name.toLowerCase()}-`))
  268. roots.push(storageRoot)
  269. const rootId = SessionId(`${backend.name.toLowerCase()}-mail-root`)
  270. const first = await stack(backend, storageRoot, [textResponse('initial teammate answer')])
  271. const firstLead = await first.ctx.agentLoop.create(rootId, { provider: 'mock', model: 'mock' })
  272. const started = await first.ctx.agentTeams.spawnTeammate(firstLead, {
  273. name: 'mail-worker',
  274. description: 'mail recovery worker',
  275. prompt: [{ type: 'text', text: 'finish before restart' }],
  276. context: 'fresh',
  277. provider: 'spawn',
  278. signal: SIGNAL,
  279. })
  280. await vi.waitFor(() => { expect(first.ctx.agents.get(started.member.id)).toBeUndefined() }, { timeout: 5_000 })
  281. vi.spyOn(first.ctx.sessionPersistence, 'open')
  282. .mockRejectedValueOnce(new Error('temporary target read failure'))
  283. const queued = await first.ctx.agentTeams.sendMessage(firstLead, {
  284. target: 'mail-worker',
  285. content: [{ type: 'text', text: 'durable retry context' }],
  286. signal: SIGNAL,
  287. })
  288. expect(queued.status).toBe('queued')
  289. expect(durable(firstLead).pendingMessages.map(message => message.id)).toEqual([queued.messageId])
  290. await first.dispose()
  291. const second = await stack(backend, storageRoot, [textResponse('resumed teammate answer')])
  292. const rootHandle = await second.ctx.agents.resume({
  293. resumeSessionId: rootId,
  294. agentOptions: { provider: 'mock', model: 'mock' },
  295. })
  296. await vi.waitFor(() => { expect(second.ctx.agents.get(started.member.id)).toBeUndefined() }, { timeout: 5_000 })
  297. await vi.waitFor(() => { expect(durable(rootHandle.agent).pendingMessages).toEqual([]) })
  298. const child = await storedEvents(second.ctx, started.member.id)
  299. const peerIds = child.flatMap(event => event.type === 'user/message'
  300. && event.data.source.kind === 'team-message'
  301. ? [event.data.source.messageId]
  302. : [])
  303. expect(peerIds).toEqual([queued.messageId])
  304. await rootHandle.dispose()
  305. await second.dispose()
  306. })
  307. it('acknowledges target-recorded mail after restart without delivering it twice', {
  308. timeout: PERSISTENCE_TEST_TIMEOUT_MS,
  309. }, async () => {
  310. const storageRoot = mkdtempSync(join(tmpdir(), `dsh-team-dedup-${backend.name.toLowerCase()}-`))
  311. roots.push(storageRoot)
  312. const rootId = SessionId(`${backend.name.toLowerCase()}-dedup-root`)
  313. const messageId = TeamMessageId(`${backend.name.toLowerCase()}-recorded-message`)
  314. const first = await stack(backend, storageRoot, [textResponse('initial teammate answer')])
  315. const firstLead = await first.ctx.agentLoop.create(rootId, { provider: 'mock', model: 'mock' })
  316. const started = await first.ctx.agentTeams.spawnTeammate(firstLead, {
  317. name: 'dedup-worker',
  318. description: 'mail deduplication worker',
  319. prompt: [{ type: 'text', text: 'finish before the crash window' }],
  320. context: 'fresh',
  321. provider: 'spawn',
  322. signal: SIGNAL,
  323. })
  324. await vi.waitFor(() => { expect(first.ctx.agents.get(started.member.id)).toBeUndefined() }, { timeout: 5_000 })
  325. const targetHandle = await first.ctx.agents.resume({
  326. resumeSessionId: started.member.id,
  327. agentOptions: { provider: 'mock', model: 'mock' },
  328. })
  329. targetHandle.agent.session.append('user/message', createUserMessage({
  330. content: [
  331. { type: 'text', text: `Team message ${messageId} from lead:` },
  332. { type: 'text', text: 'already recorded before acknowledgement' },
  333. ],
  334. source: {
  335. kind: 'team-message',
  336. teamId: TeamId(rootId),
  337. messageId,
  338. senderId: rootId,
  339. senderName: 'lead',
  340. },
  341. }), { surfaceOp: 'append' })
  342. await first.ctx.sessions.flush(targetHandle.agent.session)
  343. // Finish the target observer before writing the crash-only queued prefix.
  344. await settleMailbox(first.ctx)
  345. await targetHandle.dispose()
  346. const queued: TeamMessageSnapshot = {
  347. id: messageId,
  348. senderId: rootId,
  349. senderName: 'lead',
  350. targetId: started.member.id,
  351. content: [{ type: 'text', text: 'already recorded before acknowledgement' }],
  352. }
  353. firstLead.session.append('team/message/queued', {
  354. version: 2,
  355. teamId: TeamId(rootId),
  356. message: queued,
  357. })
  358. await first.ctx.sessions.flush(firstLead.session)
  359. expect(durable(firstLead).pendingMessages.map(message => message.id)).toEqual([messageId])
  360. await first.dispose()
  361. const second = await stack(backend, storageRoot, [])
  362. const { mailbox } = second.ctx.agentTeams as unknown as { readonly mailbox: TeamMailbox }
  363. const flush = second.ctx.sessions.flush.bind(second.ctx.sessions)
  364. const checkpointEntered = Promise.withResolvers<undefined>()
  365. const releaseCheckpoint = Promise.withResolvers<undefined>()
  366. const delayedCheckpoint = vi.spyOn(second.ctx.sessions, 'flush').mockImplementation(async (session) => {
  367. if (session.id === rootId && session.snapshotEvents().some(event =>
  368. event.type === 'team/message/delivered' && event.data.messageId === messageId)) {
  369. checkpointEntered.resolve(undefined)
  370. await releaseCheckpoint.promise
  371. }
  372. return await flush(session)
  373. })
  374. let rootHandle: AgentHandle
  375. try {
  376. rootHandle = await second.ctx.agents.resume({
  377. resumeSessionId: rootId,
  378. agentOptions: { provider: 'mock', model: 'mock' },
  379. })
  380. await checkpointEntered.promise
  381. expect(durable(rootHandle.agent).pendingMessages).toEqual([])
  382. expect(mailbox.pendingDispatches().length).toBeGreaterThan(0)
  383. let settled = false
  384. const settlement = settleMailbox(second.ctx).then(() => { settled = true })
  385. await Promise.resolve()
  386. expect(settled).toBe(false)
  387. releaseCheckpoint.resolve(undefined)
  388. await settlement
  389. } finally {
  390. releaseCheckpoint.resolve(undefined)
  391. try {
  392. await settleMailbox(second.ctx)
  393. } finally {
  394. delayedCheckpoint.mockRestore()
  395. }
  396. }
  397. expect(second.ctx.agents.get(started.member.id)).toBeUndefined()
  398. expect(second.adapter.requests).toEqual([])
  399. const child = await storedEvents(second.ctx, started.member.id)
  400. const occurrences = child.filter(event => event.type === 'user/message'
  401. && event.data.source.kind === 'team-message'
  402. && event.data.source.messageId === messageId)
  403. expect(occurrences).toHaveLength(1)
  404. await rootHandle.dispose()
  405. await second.dispose()
  406. })
  407. it('acknowledges durably pending target mail without cold-resume duplication', {
  408. timeout: PERSISTENCE_TEST_TIMEOUT_MS,
  409. }, async () => {
  410. const storageRoot = mkdtempSync(join(tmpdir(), `dsh-team-inbox-${backend.name.toLowerCase()}-`))
  411. roots.push(storageRoot)
  412. const rootId = SessionId(`${backend.name.toLowerCase()}-inbox-root`)
  413. const childId = SessionId(`${backend.name.toLowerCase()}-inbox-child`)
  414. const messageId = TeamMessageId(`${backend.name.toLowerCase()}-pending-team-message`)
  415. const first = await stack(backend, storageRoot, [])
  416. const root = await first.ctx.agentLoop.create(rootId, { provider: 'mock', model: 'mock' })
  417. await Promise.resolve()
  418. await Promise.resolve()
  419. const provisioned = provisioning(childId, 'pending-mail-worker')
  420. const active: TeamMemberSnapshot = {
  421. ...provisioned,
  422. phase: 'active',
  423. }
  424. const queued: TeamMessageSnapshot = {
  425. id: messageId,
  426. senderId: rootId,
  427. senderName: 'lead',
  428. targetId: childId,
  429. content: [{ type: 'text', text: 'already durable in target inbox' }],
  430. }
  431. root.session.append('team/member', {
  432. version: 2,
  433. teamId: TeamId(root.id),
  434. member: provisioned,
  435. })
  436. root.session.append('team/member', {
  437. version: 2,
  438. teamId: TeamId(root.id),
  439. member: active,
  440. })
  441. root.session.append('team/message/queued', {
  442. version: 2,
  443. teamId: TeamId(root.id),
  444. message: queued,
  445. })
  446. const pending = createUserMessage({
  447. content: [{ type: 'text', text: 'already durable in target inbox' }],
  448. source: {
  449. kind: 'team-message',
  450. teamId: TeamId(rootId),
  451. messageId,
  452. senderId: rootId,
  453. senderName: 'lead',
  454. },
  455. })
  456. await persistedChild(first.ctx, rootId, childId, pending)
  457. await first.ctx.sessions.flush(root.session)
  458. await first.dispose()
  459. const second = await stack(backend, storageRoot, [])
  460. const rootHandle = await second.ctx.agents.resume({
  461. resumeSessionId: rootId,
  462. agentOptions: { provider: 'mock', model: 'mock' },
  463. })
  464. await vi.waitFor(() => {
  465. expect(durable(rootHandle.agent).pendingMessages).toEqual([])
  466. })
  467. await settleMailbox(second.ctx)
  468. expect(second.adapter.requests).toEqual([])
  469. expect(second.ctx.agents.get(childId)).toBeUndefined()
  470. const stored = await storedEvents(second.ctx, childId)
  471. const pendingCopies = stored.flatMap(event => event.type === 'agent/inbox/spliced'
  472. ? event.data.inserted.filter(message => message.source.kind === 'team-message'
  473. && message.source.messageId === messageId)
  474. : [])
  475. expect(pendingCopies).toHaveLength(1)
  476. await rootHandle.dispose()
  477. await second.dispose()
  478. })
  479. })
  480. }