fork.spec.ts 10 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240
  1. import { describe, expect, it } from 'vitest'
  2. import { Context } from 'cordis'
  3. import { CallId } from '@deepseek-ai/dsh-llm'
  4. import SessionStore, { Session, SessionForkError, SessionId } from '@deepseek-ai/dsh-session'
  5. import type { SessionEvent, TurnEndReason } from '@deepseek-ai/dsh-session'
  6. async function setup(): Promise<{ ctx: Context; sessions: SessionStore }> {
  7. const ctx = new Context()
  8. await ctx.plugin(SessionStore)
  9. return { ctx, sessions: ctx.sessions }
  10. }
  11. function appendClosedTurn(
  12. session: Session,
  13. turn: number,
  14. text = `hello ${turn}`,
  15. reason: TurnEndReason = { kind: 'completed' },
  16. ): void {
  17. session.append('turn/start', { turn, trigger: { kind: 'message', source: { kind: 'user' } } })
  18. session.append('user/message', {
  19. content: [{ type: 'text', text }],
  20. source: { kind: 'user' },
  21. }, { surfaceOp: 'append' })
  22. session.append('turn/end', { turn, reason })
  23. }
  24. function appendOpenTurn(session: Session, turn: number): void {
  25. session.append('turn/start', { turn, trigger: { kind: 'message', source: { kind: 'user' } } })
  26. session.append('user/message', {
  27. content: [{ type: 'text', text: `open ${turn}` }],
  28. source: { kind: 'user' },
  29. }, { surfaceOp: 'append' })
  30. }
  31. function firstUserMessage(events: readonly SessionEvent[]): SessionEvent<'user/message'> {
  32. const event = events.find((e): e is SessionEvent<'user/message'> => e.type === 'user/message')
  33. if (event === undefined) throw new Error('missing user/message')
  34. return event
  35. }
  36. function lastSeq(session: Session): number {
  37. const event = session.events.at(-1)
  38. if (event === undefined) throw new Error('missing last event')
  39. return event.seq
  40. }
  41. describe('SessionStore.fork', () => {
  42. it('forks an empty live session as an empty child with lineage metadata', async () => {
  43. const { ctx, sessions } = await setup()
  44. const source = ctx.sessions.create(SessionId('empty-parent'), { meta: { cwd: '/workspace' } })
  45. const child = sessions.fork(source, undefined, SessionId('empty-child'))
  46. expect(child.events).toEqual([])
  47. expect(child.header).toMatchObject({
  48. id: SessionId('empty-child'),
  49. cwd: '/workspace',
  50. parentSession: SessionId('empty-parent'),
  51. seedLength: 0,
  52. })
  53. })
  54. it('forks the latest completed boundary by default and deep-clones seed events', async () => {
  55. const { ctx, sessions } = await setup()
  56. const source = ctx.sessions.create(SessionId('parent'), { meta: { cwd: '/workspace' } })
  57. appendClosedTurn(source, 1, 'hello')
  58. const child = sessions.fork(SessionId('parent'), undefined, SessionId('child'))
  59. expect(child.events).toEqual(source.events)
  60. expect(child.events).not.toBe(source.events)
  61. expect(child.events[1]).not.toBe(source.events[1])
  62. firstUserMessage(child.events).data.content[0] = { type: 'text', text: 'child mutation' }
  63. expect(firstUserMessage(source.events).data.content).toEqual([{ type: 'text', text: 'hello' }])
  64. expect(child.header).toMatchObject({
  65. id: SessionId('child'),
  66. cwd: '/workspace',
  67. parentSession: SessionId('parent'),
  68. seedLength: source.events.length,
  69. })
  70. })
  71. it('forks from an earlier turn boundary even when the source currently has an open tail', async () => {
  72. const { ctx, sessions } = await setup()
  73. const source = ctx.sessions.create(SessionId('parent'), { meta: { cwd: '/workspace' } })
  74. appendClosedTurn(source, 1, 'first')
  75. const firstBoundary = lastSeq(source)
  76. appendClosedTurn(source, 2, 'second')
  77. appendOpenTurn(source, 3)
  78. const child = sessions.fork(source, firstBoundary, SessionId('child-from-first'))
  79. expect(child.events).toEqual(source.events.slice(0, firstBoundary + 1))
  80. expect(child.header.seedLength).toBe(firstBoundary + 1)
  81. expect(child.deriveMessages()).toEqual([{ role: 'user', content: [{ type: 'text', text: 'first' }] }])
  82. })
  83. it('accepts every turn/end reason as an explicit fork boundary', async () => {
  84. const { ctx, sessions } = await setup()
  85. const reasons: TurnEndReason[] = [
  86. { kind: 'completed' },
  87. { kind: 'aborted', reason: 'cancelled by user' },
  88. { kind: 'error', step: 1, message: 'model failed', code: 'MODEL' },
  89. { kind: 'disposed' },
  90. { kind: 'max-tokens' },
  91. { kind: 'interrupted' },
  92. ]
  93. for (const reason of reasons) {
  94. const source = ctx.sessions.create(SessionId(`parent-${reason.kind}`))
  95. appendClosedTurn(source, 1, reason.kind, reason)
  96. const child = sessions.fork(source, lastSeq(source), SessionId(`child-${reason.kind}`))
  97. expect(child.events.at(-1)?.type).toBe('turn/end')
  98. expect(child.header.seedLength).toBe(source.events.length)
  99. }
  100. })
  101. it('rejects invalid boundaries before creating a child', async () => {
  102. const { ctx, sessions } = await setup()
  103. const empty = ctx.sessions.create(SessionId('empty'))
  104. expect(() => sessions.fork(empty, 0, SessionId('empty-child')))
  105. .toThrow(new SessionForkError('fork boundary 0 does not exist in session "empty" (last seq: none)', 'INVALID_BOUNDARY'))
  106. expect(ctx.sessions.get(SessionId('empty-child'))).toBeUndefined()
  107. const source = ctx.sessions.create(SessionId('parent'))
  108. appendClosedTurn(source, 1)
  109. expect(() => sessions.fork(source, -1, SessionId('negative')))
  110. .toThrow(/non-negative safe integer/)
  111. expect(() => sessions.fork(source, 0.5, SessionId('fraction')))
  112. .toThrow(/non-negative safe integer/)
  113. expect(() => sessions.fork(source, Number.MAX_SAFE_INTEGER + 1, SessionId('unsafe')))
  114. .toThrow(/non-negative safe integer/)
  115. expect(() => sessions.fork(source, source.seq, SessionId('past-end')))
  116. .toThrow(new SessionForkError(`fork boundary ${source.seq} does not exist in session "parent" (last seq: ${source.seq - 1})`, 'INVALID_BOUNDARY'))
  117. })
  118. it('rejects a corrupted live source whose array index no longer matches event seq', async () => {
  119. const { ctx, sessions } = await setup()
  120. const source = ctx.sessions.create(SessionId('corrupt-parent'))
  121. appendClosedTurn(source, 1)
  122. const mutableLog = (source as unknown as { log: SessionEvent[] }).log
  123. mutableLog[2] = { ...mutableLog[2]!, seq: 99 }
  124. expect(() => sessions.fork(source, 2, SessionId('corrupt-child')))
  125. .toThrow(new SessionForkError('fork boundary 2 does not match a contiguous event seq in session "corrupt-parent"', 'INVALID_BOUNDARY'))
  126. expect(ctx.sessions.get(SessionId('corrupt-child'))).toBeUndefined()
  127. })
  128. it('rejects an unknown live session id', async () => {
  129. const { sessions } = await setup()
  130. expect(() => sessions.fork(SessionId('missing')))
  131. .toThrow(new SessionForkError('session "missing" not found', 'SESSION_NOT_FOUND'))
  132. })
  133. it('rejects a detached Session object that is not live in ctx.sessions', async () => {
  134. const { sessions } = await setup()
  135. const detached = new Session(SessionId('detached'))
  136. expect(() => sessions.fork(detached))
  137. .toThrow(new SessionForkError('session "detached" not found', 'SESSION_NOT_FOUND'))
  138. })
  139. it('rejects a stale Session object whose id is live on a different instance', async () => {
  140. const { ctx, sessions } = await setup()
  141. ctx.sessions.create(SessionId('same-id'))
  142. const stale = new Session(SessionId('same-id'))
  143. expect(() => sessions.fork(stale))
  144. .toThrow(new SessionForkError('session "same-id" is not the live store instance', 'SESSION_NOT_LIVE'))
  145. })
  146. it('rejects selected slices whose boundary is inside an open turn', async () => {
  147. const { ctx, sessions } = await setup()
  148. const cases: [string, (session: Session) => number][] = [
  149. ['turn/start', (session) => {
  150. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  151. return lastSeq(session)
  152. }],
  153. ['step/start', (session) => {
  154. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  155. session.append('step/start', { turn: 1, step: 1 })
  156. return lastSeq(session)
  157. }],
  158. ['user/message', (session) => {
  159. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  160. session.append('user/message', { content: [{ type: 'text', text: 'open' }], source: { kind: 'user' } }, { surfaceOp: 'append' })
  161. return lastSeq(session)
  162. }],
  163. ['assistant/message', (session) => {
  164. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  165. session.append('step/start', { turn: 1, step: 1 })
  166. session.append('assistant/message', { turn: 1, step: 1, content: [{ type: 'text', text: 'partial' }] }, { surfaceOp: 'append' })
  167. return lastSeq(session)
  168. }],
  169. ['tool/call', (session) => {
  170. const callId = CallId('call-open')
  171. session.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  172. session.append('step/start', { turn: 1, step: 1 })
  173. session.append('assistant/message', {
  174. turn: 1,
  175. step: 1,
  176. content: [{ type: 'tool-call', id: callId, name: 'bash', arguments: '{}' }],
  177. }, { surfaceOp: 'append' })
  178. session.append('tool/call', { turn: 1, step: 1, callId, name: 'bash', arguments: '{}' })
  179. return lastSeq(session)
  180. }],
  181. ]
  182. for (const [lastType, build] of cases) {
  183. const source = ctx.sessions.create(SessionId(`open-${lastType}`))
  184. const boundary = build(source)
  185. expect(() => sessions.fork(source, boundary))
  186. .toThrow(new SessionForkError(`fork boundary ${boundary} in session "open-${lastType}" must be turn/end, got ${lastType}`, 'OPEN_TURN'))
  187. }
  188. })
  189. it('rejects a child session id that is already live with a typed fork error', async () => {
  190. const { ctx, sessions } = await setup()
  191. const source = ctx.sessions.create(SessionId('parent'))
  192. appendClosedTurn(source, 1)
  193. ctx.sessions.create(SessionId('child'))
  194. expect(() => sessions.fork(source, undefined, SessionId('child')))
  195. .toThrow(new SessionForkError('session "child" already exists', 'SESSION_ALREADY_EXISTS'))
  196. })
  197. it('rejects a duplicate child session id before validating the boundary', async () => {
  198. const { ctx, sessions } = await setup()
  199. const source = ctx.sessions.create(SessionId('open-parent'))
  200. source.append('turn/start', { turn: 1, trigger: { kind: 'message', source: { kind: 'user' } } })
  201. ctx.sessions.create(SessionId('child'))
  202. expect(() => sessions.fork(source, undefined, SessionId('child')))
  203. .toThrow(new SessionForkError('session "child" already exists', 'SESSION_ALREADY_EXISTS'))
  204. })
  205. })