api-proxy-fork.spec.ts 6.7 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163
  1. /** Session-fork boundaries, lineage, and inherited model routing. */
  2. import { describe, expect, it } from 'vitest'
  3. import { Context } from 'cordis'
  4. import AgentRegistry, { agentEvents } from '@deepseek-ai/dsh-agent'
  5. import type { Agent, AgentHandle, CreateAgentOptions } from '@deepseek-ai/dsh-agent'
  6. import { createUserMessage, ReasoningEffortId } from '@deepseek-ai/dsh-llm'
  7. import type { LlmCallConfig } from '@deepseek-ai/dsh-llm'
  8. import SessionStore from '@deepseek-ai/dsh-session'
  9. import type { Session, SessionId } from '@deepseek-ai/dsh-session'
  10. import SystemPrompt from '@deepseek-ai/dsh-system-prompt'
  11. import UserInteractionService from '@deepseek-ai/dsh-user-interaction'
  12. import type { RpcRequest } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
  13. import { RpcId } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
  14. import { createApiProxy } from '@deepseek-ai/dsh-host-apiproxy'
  15. const sid = (id: string): SessionId => id as SessionId
  16. let nextRpc = 1
  17. function request<P>(payload: P): RpcRequest<P> {
  18. return { rpcId: RpcId(`fork-${String(nextRpc++)}`), payload }
  19. }
  20. async function composed(): Promise<Context> {
  21. const ctx = new Context()
  22. await ctx.plugin(SessionStore)
  23. await ctx.plugin(SystemPrompt, { persona: '' })
  24. await ctx.plugin(AgentRegistry)
  25. await ctx.plugin(UserInteractionService)
  26. ctx.provide('workspace', { list: () => [] } as never)
  27. ctx.agents.setFactory({
  28. createAgent: async (ownerCtx: Context, options: CreateAgentOptions): Promise<AgentHandle> => {
  29. const session = ctx.sessions.create(options.sessionId, {
  30. ...options.seed === undefined ? {} : { seed: [...options.seed] },
  31. ...options.meta === undefined ? {} : { meta: options.meta },
  32. })
  33. const agent = {} as Agent
  34. const agentCtx = ownerCtx.extend({ agent })
  35. Object.assign(agent, { id: session.id, session, status: 'idle', ctx: agentCtx })
  36. await options.setup?.(agentCtx)
  37. ctx.agents.register(agent)
  38. return { agent, dispose: () => Promise.resolve() }
  39. },
  40. resume: () => Promise.reject(new Error('fork test sources are live')),
  41. })
  42. return ctx
  43. }
  44. function liveAgent(ctx: Context, id: string, turns: number, openTail = false): Session {
  45. const session = ctx.sessions.create(sid(id), { meta: { cwd: '/proj' } })
  46. for (let turn = 1; turn <= turns; turn++) {
  47. session.append('turn/start', { turn, trigger: { kind: 'message', source: { kind: 'user' } } })
  48. session.append('user/message', createUserMessage({
  49. content: [{ type: 'text', text: `prompt ${String(turn)}` }],
  50. source: { kind: 'user' },
  51. }), { surfaceOp: 'append' })
  52. session.append('turn/end', { turn, reason: { kind: 'completed' } })
  53. }
  54. if (openTail) {
  55. session.append('turn/start', { turn: turns + 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  56. session.append('user/message', createUserMessage({
  57. content: [{ type: 'text', text: 'open prompt' }],
  58. source: { kind: 'user' },
  59. }), { surfaceOp: 'append' })
  60. }
  61. ctx.agents.register({ id: session.id, session, status: 'idle', ctx } as Agent)
  62. return session
  63. }
  64. const api = (ctx: Context) => createApiProxy(ctx, {
  65. provider: 'default-provider',
  66. model: 'default-model',
  67. cwd: '/tmp',
  68. workspaceRoot: '/tmp',
  69. })
  70. describe('sessions.fork', () => {
  71. it('cuts at the anchored completed turn and records lineage and cwd', async () => {
  72. const ctx = await composed()
  73. const source = liveAgent(ctx, 'session-source', 2)
  74. const response = await api(ctx).sessions.fork(request({ sessionId: source.id, atSeq: 1 }))
  75. expect(response.result.ok).toBe(true)
  76. if (!response.result.ok) return
  77. const child = ctx.sessions.get(response.result.value.sessionId)
  78. expect(child?.events.map(event => event.type)).toEqual([
  79. 'turn/start', 'user/message', 'turn/end', 'session/end-seed',
  80. ])
  81. expect(child?.header.parentSession).toBe(source.id)
  82. expect(child?.header.cwd).toBe('/proj')
  83. await ctx.fiber.dispose()
  84. })
  85. it('uses the last completed turn only for omitted and past-end anchors', async () => {
  86. const ctx = await composed()
  87. const source = liveAgent(ctx, 'session-tail', 2, true)
  88. const proxy = api(ctx)
  89. const expectedTypes = [
  90. 'turn/start', 'user/message', 'turn/end',
  91. 'turn/start', 'user/message', 'turn/end',
  92. 'session/end-seed',
  93. ]
  94. const omitted = await proxy.sessions.fork(request({ sessionId: source.id }))
  95. expect(omitted.result.ok).toBe(true)
  96. if (omitted.result.ok) {
  97. expect(ctx.sessions.get(omitted.result.value.sessionId)?.events.map(event => event.type))
  98. .toEqual(expectedTypes)
  99. }
  100. const pastEnd = await proxy.sessions.fork(request({ sessionId: source.id, atSeq: 999 }))
  101. expect(pastEnd.result.ok).toBe(true)
  102. if (pastEnd.result.ok) {
  103. expect(ctx.sessions.get(pastEnd.result.value.sessionId)?.events.map(event => event.type))
  104. .toEqual(expectedTypes)
  105. }
  106. await ctx.fiber.dispose()
  107. })
  108. it('rejects an in-log anchor whose turn is still open', async () => {
  109. const ctx = await composed()
  110. const source = liveAgent(ctx, 'session-open', 1, true)
  111. const anchor = source.events.at(-1)?.seq ?? 0
  112. const response = await api(ctx).sessions.fork(request({ sessionId: source.id, atSeq: anchor }))
  113. expect(response.result).toMatchObject({
  114. ok: false,
  115. error: { code: 'fork-unavailable', details: { sessionId: source.id } },
  116. })
  117. if (!response.result.ok) expect(response.result.error.message).toMatch(/has not completed/)
  118. await ctx.fiber.dispose()
  119. })
  120. it('installs the latest logged model target before the child can run', async () => {
  121. const ctx = await composed()
  122. const source = liveAgent(ctx, 'session-routed', 1)
  123. source.append('request/header', {
  124. header: {
  125. config: {
  126. provider: 'inherited-provider',
  127. model: 'inherited-model',
  128. reasoningEffort: ReasoningEffortId('high'),
  129. },
  130. },
  131. reason: 'initial',
  132. })
  133. const response = await api(ctx).sessions.fork(request({ sessionId: source.id }))
  134. expect(response.result.ok).toBe(true)
  135. if (!response.result.ok) return
  136. const child = ctx.agents.get(response.result.value.sessionId)
  137. if (child === undefined) throw new Error('fork did not publish the child agent')
  138. const assembly = await child.ctx.systemPrompt.assemble()
  139. expect(assembly.variables).toMatchObject({
  140. provider: 'inherited-provider',
  141. model: 'inherited-model',
  142. })
  143. const fallback: LlmCallConfig = { provider: 'default-provider', model: 'default-model' }
  144. await expect(agentEvents(child.ctx, child).waterfall(
  145. 'agent/request', 1, 0, new AbortController().signal, () => Promise.resolve(fallback),
  146. )).resolves.toMatchObject({
  147. provider: 'inherited-provider',
  148. model: 'inherited-model',
  149. reasoningEffort: 'high',
  150. })
  151. await ctx.fiber.dispose()
  152. })
  153. })