agent-continuation.worker.ts 6.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142
  1. /** Plain-Node measurements of active request history and cold tool-heavy continuation. */
  2. import { performance } from 'node:perf_hooks'
  3. import { scheduler } from 'node:timers/promises'
  4. import { Context } from '@deepseek-ai/cordis'
  5. import AgentLoop from '@deepseek-ai/dsh-agent-loop'
  6. import type { Agent, AgentHandle } from '@deepseek-ai/dsh-agent'
  7. import { mountAgentLoopTestDependencies } from '@deepseek-ai/dsh-agent-loop-testkit'
  8. import { createUserMessage, LlmAdapter } from '@deepseek-ai/dsh-llm'
  9. import type { GenerateOptions, LlmResolvedModelInfo, StreamChunk } from '@deepseek-ai/dsh-llm'
  10. import { SESSION_FORMAT_VERSION } from '@deepseek-ai/dsh-session'
  11. import JsonlSessionPersistence from '@deepseek-ai/dsh-session-persistence-jsonl'
  12. import { defineContentToolFixture } from '@deepseek-ai/dsh-tools'
  13. import { assertBuiltBenchmarkRuntime } from '../support/built-worker.ts'
  14. import { PARENT_ID, response, resultText, syntheticHistory, TIME_ZERO, WORKLOAD } from './workload.ts'
  15. /** Raw timing and retained-memory report from one isolated backend process. */
  16. export interface ContinuationReport {
  17. readonly totalMs: number
  18. readonly resumeMs: number
  19. readonly turnsMs: number
  20. readonly flushMs: number
  21. readonly cpuUserMs: number
  22. readonly cpuSystemMs: number
  23. readonly retainedHeapMb: number
  24. readonly peakRssMb: number
  25. readonly requests: number
  26. readonly toolCalls: number
  27. readonly events: number
  28. }
  29. class SyntheticAdapter extends LlmAdapter {
  30. requests = 0
  31. constructor(private readonly toolsPerTurn: number) { super() }
  32. override resolveModel(provider: string, model: string): Promise<LlmResolvedModelInfo> {
  33. return Promise.resolve({ provider, id: model, name: model })
  34. }
  35. async * stream(_options: GenerateOptions): AsyncIterable<StreamChunk> {
  36. const tools = this.toolsPerTurn > 0 && this.requests % 2 === 0 ? this.toolsPerTurn : 0
  37. const reply = response(100_000 + this.requests++, tools)
  38. yield* reply.chunks
  39. }
  40. }
  41. async function collectHeap(): Promise<number> {
  42. if (globalThis.gc === undefined) throw new Error('backend benchmark requires --expose-gc')
  43. globalThis.gc()
  44. await scheduler.yield()
  45. globalThis.gc()
  46. return process.memoryUsage().heapUsed / 1_048_576
  47. }
  48. async function seed(root: string): Promise<void> {
  49. const ctx = new Context()
  50. try {
  51. await mountAgentLoopTestDependencies(ctx)
  52. await ctx.plugin(JsonlSessionPersistence, { root, compression: 'zstd' })
  53. const handle = await ctx.sessionPersistence.create({
  54. version: SESSION_FORMAT_VERSION, id: PARENT_ID, createdAt: TIME_ZERO, cwd: '/bench', isSeeded: false,
  55. }, {})
  56. try {
  57. await handle.append(syntheticHistory(WORKLOAD.historyTurns))
  58. await handle.flush()
  59. } finally { await handle.close() }
  60. } finally { await ctx.fiber.dispose() }
  61. }
  62. async function runTurns(agent: Agent, turns: number): Promise<void> {
  63. for (let turn = 0; turn < turns; turn++) {
  64. agent.followup(createUserMessage({ content: [{ type: 'text', text: 'Continue synthetic task ' + String(turn) }], source: { kind: 'user' } }))
  65. await agent.whenIdle()
  66. }
  67. }
  68. async function measure(root: string, scenario: string): Promise<ContinuationReport> {
  69. const ctx = new Context()
  70. let handle: AgentHandle | undefined
  71. const toolHeavy = scenario === 'tool-continuation'
  72. const adapter = new SyntheticAdapter(toolHeavy ? WORKLOAD.toolsPerLiveTurn : 0)
  73. let toolCalls = 0
  74. try {
  75. await mountAgentLoopTestDependencies(ctx)
  76. await ctx.plugin(JsonlSessionPersistence, { root, compression: 'zstd' })
  77. await ctx.plugin(AgentLoop, { agents: [] })
  78. ctx.effect(() => ctx.llm.registerAdapter(['bench'], adapter))
  79. ctx.effect(() => ctx.tools.register(defineContentToolFixture({
  80. name: 'bench_tool', description: 'Read a bounded synthetic module.',
  81. parameters: { ordinal: { type: 'number', required: true } },
  82. isConcurrencySafe: () => true,
  83. execute(args) {
  84. toolCalls++
  85. return Promise.resolve([{ type: 'text', text: resultText(args.ordinal) }])
  86. },
  87. })))
  88. if (!toolHeavy) {
  89. handle = await ctx.agents.resume({ resumeSessionId: PARENT_ID, agentOptions: { provider: 'bench', model: 'bench' } })
  90. }
  91. const beforeHeap = await collectHeap()
  92. const cpuStart = process.cpuUsage()
  93. const start = performance.now()
  94. if (handle === undefined) {
  95. handle = await ctx.agents.resume({ resumeSessionId: PARENT_ID, agentOptions: { provider: 'bench', model: 'bench' } })
  96. }
  97. const resumed = performance.now()
  98. await runTurns(handle.agent, toolHeavy ? WORKLOAD.continuationTurns : WORKLOAD.requestTurns)
  99. const turnsDone = performance.now()
  100. await ctx.sessions.flush(handle.agent.session)
  101. const end = performance.now()
  102. const cpu = process.cpuUsage(cpuStart)
  103. const retainedHeapMb = (await collectHeap()) - beforeHeap
  104. if (adapter.requests !== (toolHeavy ? WORKLOAD.continuationTurns * 2 : WORKLOAD.requestTurns)
  105. || toolCalls !== (toolHeavy ? WORKLOAD.continuationTurns * WORKLOAD.toolsPerLiveTurn : 0)) {
  106. throw new Error('backend benchmark did not complete every requested model/tool step')
  107. }
  108. return {
  109. totalMs: end - start, resumeMs: resumed - start, turnsMs: turnsDone - resumed, flushMs: end - turnsDone,
  110. cpuUserMs: cpu.user / 1_000, cpuSystemMs: cpu.system / 1_000,
  111. retainedHeapMb, peakRssMb: process.resourceUsage().maxRSS / 1_024,
  112. requests: adapter.requests, toolCalls, events: handle.agent.session.seq,
  113. }
  114. } finally {
  115. await handle?.dispose()
  116. await ctx.fiber.dispose()
  117. }
  118. }
  119. assertBuiltBenchmarkRuntime(import.meta.url, Object.fromEntries([
  120. '@deepseek-ai/dsh-agent-loop', '@deepseek-ai/dsh-session', '@deepseek-ai/dsh-llm',
  121. '@deepseek-ai/dsh-tools', '@deepseek-ai/dsh-session-persistence-jsonl',
  122. ].map(name => [name, import.meta.resolve(name)])))
  123. const [root, scenario] = process.argv.slice(2)
  124. if (root === undefined || scenario === undefined || !['seed', 'request-history', 'tool-continuation'].includes(scenario)) {
  125. throw new Error('usage: agent-continuation.worker.js <root> <seed|request-history|tool-continuation>')
  126. }
  127. if (scenario === 'seed') {
  128. await seed(root)
  129. process.stdout.write(JSON.stringify({ seeded: true }) + '\n')
  130. } else {
  131. process.stdout.write(JSON.stringify(await measure(root, scenario)) + '\n')
  132. }