persistence.spec.ts 19 KB

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