api-proxy-fork.spec.ts 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288
  1. /** Session-fork boundaries, lineage, and inherited model routing. */
  2. import { describe, expect, it, vi } 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, SessionEvent, SessionHeader, 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 { Workspace } from '@deepseek-ai/dsh-workspace'
  13. import type { RpcRequest } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
  14. import { RpcId } from '@deepseek-ai/dsh-host-apiproxy/api/rpc'
  15. import { createApiProxy } from '@deepseek-ai/dsh-host-apiproxy'
  16. const sid = (id: string): SessionId => id as SessionId
  17. let nextRpc = 1
  18. function request<P>(payload: P): RpcRequest<P> {
  19. return { rpcId: RpcId(`fork-${String(nextRpc++)}`), payload }
  20. }
  21. async function composed(workspaces: readonly Workspace[] = []): Promise<Context> {
  22. const ctx = new Context()
  23. await ctx.plugin(SessionStore)
  24. await ctx.plugin(SystemPrompt, { persona: '' })
  25. await ctx.plugin(AgentRegistry)
  26. await ctx.plugin(UserInteractionService)
  27. ctx.provide('workspace', { list: () => workspaces } as never)
  28. ctx.agents.setFactory({
  29. createAgent: async (ownerCtx: Context, options: CreateAgentOptions): Promise<AgentHandle> => {
  30. const session = ctx.sessions.create(options.sessionId, {
  31. ...options.seed === undefined ? {} : { seed: [...options.seed] },
  32. ...options.meta === undefined ? {} : { meta: options.meta },
  33. })
  34. const agent = {} as Agent
  35. const agentCtx = ownerCtx.extend({ agent })
  36. Object.assign(agent, { id: session.id, session, status: 'idle', ctx: agentCtx })
  37. await options.setup?.(agentCtx)
  38. ctx.agents.register(agent)
  39. return { agent, dispose: () => Promise.resolve() }
  40. },
  41. resume: () => Promise.reject(new Error('fork test sources are live')),
  42. })
  43. return ctx
  44. }
  45. /** Tail turn appended after the completed ones: left open, or closed as aborted (a stopped turn). */
  46. type Tail = 'none' | 'open' | 'aborted'
  47. function liveAgent(
  48. ctx: Context,
  49. id: string,
  50. turns: number,
  51. tail: Tail = 'none',
  52. lineage: { parentSession?: SessionId; origin?: 'subagent' } = {},
  53. ): Session {
  54. const session = ctx.sessions.create(sid(id), { meta: { cwd: '/proj', ...lineage } })
  55. for (let turn = 1; turn <= turns; turn++) {
  56. session.append('turn/start', { turn, trigger: { kind: 'message', source: { kind: 'user' } } })
  57. session.append('user/message', createUserMessage({
  58. content: [{ type: 'text', text: `prompt ${String(turn)}` }],
  59. source: { kind: 'user' },
  60. }), { surfaceOp: 'append' })
  61. session.append('turn/end', { turn, reason: { kind: 'completed' } })
  62. }
  63. if (tail !== 'none') {
  64. session.append('turn/start', { turn: turns + 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  65. session.append('user/message', createUserMessage({
  66. content: [{ type: 'text', text: 'open prompt' }],
  67. source: { kind: 'user' },
  68. }), { surfaceOp: 'append' })
  69. if (tail === 'aborted') session.append('turn/end', { turn: turns + 1, reason: { kind: 'aborted' } })
  70. }
  71. ctx.agents.register({ id: session.id, session, status: 'idle', ctx } as Agent)
  72. return session
  73. }
  74. const api = (ctx: Context) => createApiProxy(ctx, {
  75. provider: 'default-provider',
  76. model: 'default-model',
  77. cwd: '/tmp',
  78. workspaceRoot: '/tmp',
  79. })
  80. describe('sessions.fork', () => {
  81. it('cuts at the anchored completed turn and records lineage and cwd', async () => {
  82. const ctx = await composed()
  83. const source = liveAgent(ctx, 'session-source', 2)
  84. const response = await api(ctx).sessions.fork(request({ sessionId: source.id, atSeq: 1 }))
  85. expect(response.result.ok).toBe(true)
  86. if (!response.result.ok) return
  87. const child = ctx.sessions.get(response.result.value.sessionId)
  88. expect(child?.events.map(event => event.type)).toEqual([
  89. 'turn/start', 'user/message', 'turn/end', 'session/end-seed',
  90. ])
  91. expect(child?.header.parentSession).toBe(source.id)
  92. expect(child?.header.cwd).toBe('/proj')
  93. await ctx.fiber.dispose()
  94. })
  95. it('attaches a subagent fork to its nearest workspace-owning ancestor', async () => {
  96. const accounted: SessionId[] = []
  97. const attachSession = vi.fn<(sessionId: SessionId) => Promise<void>>()
  98. .mockResolvedValue(undefined)
  99. const workspace = {
  100. sessionIds: accounted,
  101. attachSession,
  102. } as unknown as Workspace
  103. const ctx = await composed([workspace])
  104. const owner = liveAgent(ctx, 'session-owner', 1)
  105. accounted.push(owner.id)
  106. const child = liveAgent(ctx, 'session-child', 1, 'none', {
  107. parentSession: owner.id,
  108. origin: 'subagent',
  109. })
  110. const grandchild = liveAgent(ctx, 'session-grandchild', 1, 'none', {
  111. parentSession: child.id,
  112. origin: 'subagent',
  113. })
  114. ctx.provide('sessionQuery', {
  115. traceSession: vi.fn(() => Promise.resolve({
  116. target: { header: grandchild.header, live: true, persisted: false },
  117. ancestors: [
  118. { header: child.header, live: true, persisted: false },
  119. { header: owner.header, live: true, persisted: false },
  120. ],
  121. descendants: [],
  122. complete: true,
  123. root: { header: owner.header, live: true, persisted: false },
  124. })),
  125. } as never)
  126. const response = await api(ctx).sessions.fork(request({ sessionId: grandchild.id }))
  127. expect(response.result.ok).toBe(true)
  128. if (!response.result.ok) return
  129. expect(attachSession).toHaveBeenCalledWith(response.result.value.sessionId)
  130. expect(ctx.sessions.get(response.result.value.sessionId)?.header).toMatchObject({
  131. parentSession: grandchild.id,
  132. cwd: '/proj',
  133. })
  134. expect(ctx.sessions.get(response.result.value.sessionId)?.header.origin).toBeUndefined()
  135. await ctx.fiber.dispose()
  136. })
  137. it('forks a persisted subagent without resuming its Agent', async () => {
  138. const ctx = await composed()
  139. const sourceId = sid('session-cold-subagent')
  140. const parentId = sid('session-cold-parent')
  141. const header: SessionHeader = {
  142. version: 0,
  143. id: sourceId,
  144. createdAt: 1,
  145. cwd: '/proj',
  146. parentSession: parentId,
  147. origin: 'subagent',
  148. }
  149. const events = [
  150. { type: 'turn/start', seq: 0, time: 1, data: { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } } },
  151. {
  152. type: 'user/message',
  153. seq: 1,
  154. time: 2,
  155. data: createUserMessage({ content: [{ type: 'text', text: 'work' }], source: { kind: 'user' } }),
  156. surfaceOp: 'append',
  157. },
  158. { type: 'turn/end', seq: 2, time: 3, data: { turn: 1, reason: { kind: 'completed' } } },
  159. ] as SessionEvent[]
  160. ctx.provide('sessionPersistence', {
  161. list: () => Promise.resolve([header]),
  162. inspect: () => Promise.resolve({ meta: header, events }),
  163. } as never)
  164. ctx.provide('sessionQuery', {
  165. traceSession: () => Promise.resolve({
  166. target: { header, live: false, persisted: true },
  167. ancestors: [],
  168. descendants: [],
  169. complete: true,
  170. root: { header, live: false, persisted: true },
  171. }),
  172. } as never)
  173. const resume = vi.spyOn(ctx.agents, 'resume')
  174. const response = await api(ctx).sessions.fork(request({ sessionId: sourceId }))
  175. expect(response.result.ok).toBe(true)
  176. if (!response.result.ok) return
  177. expect(resume).not.toHaveBeenCalled()
  178. expect(ctx.agents.get(sourceId)).toBeUndefined()
  179. expect(ctx.sessions.get(response.result.value.sessionId)?.header).toMatchObject({
  180. parentSession: sourceId,
  181. cwd: '/proj',
  182. })
  183. expect(ctx.sessions.get(response.result.value.sessionId)?.header.origin).toBeUndefined()
  184. await ctx.fiber.dispose()
  185. })
  186. it('uses the last completed turn only for omitted and past-end anchors', async () => {
  187. const ctx = await composed()
  188. const source = liveAgent(ctx, 'session-tail', 2, 'open')
  189. const proxy = api(ctx)
  190. const expectedTypes = [
  191. 'turn/start', 'user/message', 'turn/end',
  192. 'turn/start', 'user/message', 'turn/end',
  193. 'session/end-seed',
  194. ]
  195. const omitted = await proxy.sessions.fork(request({ sessionId: source.id }))
  196. expect(omitted.result.ok).toBe(true)
  197. if (omitted.result.ok) {
  198. expect(ctx.sessions.get(omitted.result.value.sessionId)?.events.map(event => event.type))
  199. .toEqual(expectedTypes)
  200. }
  201. const pastEnd = await proxy.sessions.fork(request({ sessionId: source.id, atSeq: 999 }))
  202. expect(pastEnd.result.ok).toBe(true)
  203. if (pastEnd.result.ok) {
  204. expect(ctx.sessions.get(pastEnd.result.value.sessionId)?.events.map(event => event.type))
  205. .toEqual(expectedTypes)
  206. }
  207. await ctx.fiber.dispose()
  208. })
  209. it('cuts through an aborted turn: stopped is closed, not open', async () => {
  210. const ctx = await composed()
  211. const source = liveAgent(ctx, 'session-aborted', 1, 'aborted')
  212. // What a stopped message's fork button anchors on: the frozen node sits
  213. // one event before its turn/end, floored client-side to that event's seq.
  214. const anchor = (source.events.at(-1)?.seq ?? 0) - 1
  215. const response = await api(ctx).sessions.fork(request({ sessionId: source.id, atSeq: anchor }))
  216. expect(response.result.ok).toBe(true)
  217. if (!response.result.ok) return
  218. expect(ctx.sessions.get(response.result.value.sessionId)?.events.map(event => event.type)).toEqual([
  219. 'turn/start', 'user/message', 'turn/end',
  220. 'turn/start', 'user/message', 'turn/end',
  221. 'session/end-seed',
  222. ])
  223. await ctx.fiber.dispose()
  224. })
  225. it('rejects an in-log anchor whose turn is still open', async () => {
  226. const ctx = await composed()
  227. const source = liveAgent(ctx, 'session-open', 1, 'open')
  228. const anchor = source.events.at(-1)?.seq ?? 0
  229. const response = await api(ctx).sessions.fork(request({ sessionId: source.id, atSeq: anchor }))
  230. expect(response.result).toMatchObject({
  231. ok: false,
  232. error: { code: 'fork-unavailable', details: { sessionId: source.id } },
  233. })
  234. if (!response.result.ok) expect(response.result.error.message).toMatch(/has not completed/)
  235. await ctx.fiber.dispose()
  236. })
  237. it('installs the latest logged model target before the child can run', async () => {
  238. const ctx = await composed()
  239. const source = liveAgent(ctx, 'session-routed', 1)
  240. source.append('request/header', {
  241. header: {
  242. config: {
  243. provider: 'inherited-provider',
  244. model: 'inherited-model',
  245. reasoningEffort: ReasoningEffortId('high'),
  246. },
  247. },
  248. reason: 'initial',
  249. })
  250. const response = await api(ctx).sessions.fork(request({ sessionId: source.id }))
  251. expect(response.result.ok).toBe(true)
  252. if (!response.result.ok) return
  253. const child = ctx.agents.get(response.result.value.sessionId)
  254. if (child === undefined) throw new Error('fork did not publish the child agent')
  255. const assembly = await child.ctx.systemPrompt.assemble()
  256. expect(assembly.variables).toMatchObject({
  257. provider: 'inherited-provider',
  258. model: 'inherited-model',
  259. })
  260. const fallback: LlmCallConfig = { provider: 'default-provider', model: 'default-model' }
  261. await expect(agentEvents(child.ctx, child).waterfall(
  262. 'agent/request', 1, 0, new AbortController().signal, () => Promise.resolve(fallback),
  263. )).resolves.toMatchObject({
  264. provider: 'inherited-provider',
  265. model: 'inherited-model',
  266. reasoningEffort: 'high',
  267. })
  268. await ctx.fiber.dispose()
  269. })
  270. })