session-history-journal.host.spec.ts 13 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304
  1. /** Raw Session journal transport and message-aligned pagination coverage. */
  2. import { describe, expect, it, vi } from 'vitest'
  3. import { Context } from '@deepseek-ai/cordis'
  4. import AgentRegistry from '@deepseek-ai/dsh-agent'
  5. import SessionStore from '@deepseek-ai/dsh-session'
  6. import { CallId, createMessage, createToolResultMessage, createUserMessage } from '@deepseek-ai/dsh-llm'
  7. import type { Session, SessionEvent, SessionId } from '@deepseek-ai/dsh-session'
  8. import { SessionHistoryController } from '@deepseek-ai/dsh-api-session-controller/src/history.ts'
  9. import type { SessionFollowFrame } from '@deepseek-ai/dsh-api-session-controller/types'
  10. import { createSessionTestRemote, installSessionReadTestServices } from './test-remote.ts'
  11. /** Append a production-shaped human prompt to the session surface. */
  12. function appendUserText(session: Session, text: string): SessionEvent {
  13. return session.append('user/message', createUserMessage({
  14. content: [{ type: 'text', text }], source: { kind: 'user' },
  15. }), { surfaceOp: 'append' })
  16. }
  17. /** Append a production-shaped assistant message to the session surface. */
  18. function appendAssistantText(session: Session, text: string, step: number): SessionEvent {
  19. return session.append('assistant/message', {
  20. turn: 1,
  21. step,
  22. message: createMessage({
  23. role: 'assistant',
  24. content: [{ type: 'text', text }],
  25. source: { kind: 'model', provider: 'p', model: 'm' },
  26. }),
  27. }, { surfaceOp: 'append' })
  28. }
  29. /**
  30. * Append a plugin-owned log-only event. The host proxy is projection-only, so it
  31. * declares no compaction vocabulary; the cast writes the real event shape without
  32. * depending on the owning package.
  33. */
  34. function appendExtension(session: Session, type: string, data: unknown): SessionEvent {
  35. return (session.append as unknown as (type: string, data: unknown) => SessionEvent)(type, data)
  36. }
  37. async function harness(): Promise<{ ctx: Context }> {
  38. const ctx = new Context()
  39. await ctx.plugin(SessionStore)
  40. await ctx.plugin(AgentRegistry)
  41. installSessionReadTestServices(ctx)
  42. return { ctx }
  43. }
  44. /** Drain one Session follow until `count` event frames arrive. */
  45. async function collect(
  46. iterable: AsyncIterable<SessionFollowFrame>,
  47. count: number,
  48. abort: AbortController,
  49. ): Promise<SessionFollowFrame[]> {
  50. const frames: SessionFollowFrame[] = []
  51. for await (const frame of iterable) {
  52. frames.push(frame)
  53. if (frames.filter(candidate => candidate.type === 'event').length >= count) abort.abort()
  54. }
  55. return frames
  56. }
  57. /** Open follow and wait until its cursor is fixed before appending fixtures. */
  58. async function openFollow(
  59. history: SessionHistoryController,
  60. sessionId: SessionId,
  61. signal: AbortSignal,
  62. ): Promise<AsyncIterable<SessionFollowFrame>> {
  63. const iterator = history.follow({
  64. address: { kind: 'session', sessionId },
  65. }, signal)[Symbol.asyncIterator]()
  66. await expect(iterator.next()).resolves.toMatchObject({
  67. done: false,
  68. value: { type: 'snapshot' },
  69. })
  70. return { [Symbol.asyncIterator]: () => iterator }
  71. }
  72. describe('Session history raw journal', () => {
  73. it('follows raw tool events and preserves result metadata without a Tools service', async () => {
  74. const { ctx } = await harness()
  75. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  76. const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
  77. const abort = new AbortController()
  78. const stream = await openFollow(history, session.id, abort.signal)
  79. const collected = collect(stream, 2, abort)
  80. const call = session.append('tool/call', {
  81. turn: 1, step: 1, callId: CallId('raw-call'), name: 'custom', arguments: '{malformed',
  82. })
  83. const result = session.append('tool/result', {
  84. turn: 1, step: 1,
  85. message: createToolResultMessage({
  86. callId: CallId('raw-call'),
  87. content: [{ type: 'text', text: 'raw output' }],
  88. isError: false,
  89. }),
  90. meta: { nested: { count: 2 }, paths: ['a.ts', 'b.ts'] },
  91. }, { surfaceOp: 'append' })
  92. const frames = await collected
  93. expect(frames).toEqual([
  94. { type: 'event', event: call },
  95. { type: 'event', event: result },
  96. ])
  97. expect((frames[1] as Extract<SessionFollowFrame, { type: 'event' }>).event.data)
  98. .toMatchObject({ meta: { nested: { count: 2 }, paths: ['a.ts', 'b.ts'] } })
  99. })
  100. it('follows live results without rescanning Session history', async () => {
  101. const { ctx } = await harness()
  102. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  103. const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
  104. const abort = new AbortController()
  105. const stream = await openFollow(history, session.id, abort.signal)
  106. const iterator = stream[Symbol.asyncIterator]()
  107. session.append('tool/call', {
  108. turn: 1, step: 1, callId: CallId('live-fast'), name: 'term', arguments: '{"cmd":"pwd"}',
  109. })
  110. await expect(iterator.next()).resolves.toMatchObject({
  111. value: { type: 'event', event: { type: 'tool/call', data: { callId: 'live-fast' } } },
  112. })
  113. const events = vi.spyOn(session, 'events', 'get').mockImplementation(() => {
  114. throw new Error('live result rescanned Session history')
  115. })
  116. try {
  117. session.append('tool/result', {
  118. turn: 1, step: 1,
  119. message: createToolResultMessage({
  120. callId: CallId('live-fast'),
  121. content: [{ type: 'text', text: 'ok' }],
  122. isError: false,
  123. }),
  124. }, { surfaceOp: 'append' })
  125. await expect(iterator.next()).resolves.toMatchObject({
  126. value: { type: 'event', event: { type: 'tool/result', data: { message: { source: { callId: 'live-fast' } } } } },
  127. })
  128. } finally {
  129. events.mockRestore()
  130. abort.abort()
  131. await iterator.next()
  132. await ctx.fiber.dispose()
  133. }
  134. })
  135. it('serves raw call and result entries without parsing tool arguments', async () => {
  136. const { ctx } = await harness()
  137. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  138. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  139. const start = session.append('turn/start', { turn: 1 })
  140. const call = session.append('tool/call', {
  141. turn: 1, step: 1, callId: CallId('history-call'), name: 'custom', arguments: '{broken',
  142. })
  143. const result = session.append('tool/result', {
  144. turn: 1, step: 1,
  145. message: createToolResultMessage({
  146. callId: CallId('history-call'),
  147. content: [{ type: 'text', text: 'failed raw output' }],
  148. isError: true,
  149. }),
  150. meta: { persisted: true, count: 3 },
  151. }, { surfaceOp: 'append' })
  152. const response = await remote.page({
  153. address: { kind: 'session', sessionId: session.id },
  154. throughSeq: session.seq - 1,
  155. })
  156. expect(response.ok).toBe(true)
  157. if (!response.ok) throw new Error('unreachable')
  158. expect(response.value.events).toEqual([
  159. { event: start },
  160. { event: call },
  161. { event: result },
  162. ])
  163. })
  164. it('counts only append-origin messages toward maxMessages and keeps each compaction summary with its replacement', async () => {
  165. const { ctx } = await harness()
  166. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  167. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  168. session.append('turn/start', { turn: 1 })
  169. const first = appendUserText(session, 'first prompt')
  170. appendAssistantText(session, 'first reply', 1)
  171. const third = appendUserText(session, 'second prompt')
  172. appendAssistantText(session, 'second reply', 2)
  173. const shadowed = [...session.surface.nodes]
  174. // A compaction transaction: a log-only summary record immediately followed by the
  175. // replacement that shadows the range.
  176. const summary = appendExtension(session, 'compaction/summary', {
  177. summary: [{ type: 'text', text: 'summary' }],
  178. shadowedRange: { start: shadowed[0], end: shadowed.at(-1) },
  179. shadowedSeqs: shadowed,
  180. shadowedTokenCount: 0,
  181. provider: 'p',
  182. model: 'm',
  183. })
  184. session.append('user/message', createUserMessage({
  185. content: [{ type: 'text', text: '<context_checkpoint>summary</context_checkpoint>' }],
  186. source: { kind: 'plugin', plugin: 'compact' },
  187. }), {
  188. surfaceOp: { op: 'replace', start: shadowed[0] as number, end: shadowed.at(-1) as number },
  189. sourceEventSeqs: [...shadowed, summary.seq],
  190. })
  191. const response = await remote.page({
  192. address: { kind: 'session', sessionId: session.id },
  193. throughSeq: session.seq - 1,
  194. maxMessages: 2,
  195. })
  196. if (!response.ok) throw new Error('unreachable')
  197. const page = response.value.events.map(entry => entry.event)
  198. // Two append-origin messages fill the page even though a replacement copy of
  199. // the same event type sits in the window: the copy is model-only.
  200. const messages = page.filter(event => event.type === 'user/message' || event.type === 'assistant/message')
  201. expect(messages.map(event => event.seq)).toEqual([third.seq, third.seq + 1, third.seq + 3])
  202. expect(page.some(event => event.seq === first.seq)).toBe(false)
  203. expect(response.value.hasMore).toBe(true)
  204. // The range stays contiguous, so the checkpoint's summary record is readable on
  205. // the same page as the checkpoint itself.
  206. const summaryIndex = page.findIndex(event => event.seq === summary.seq)
  207. expect(summaryIndex).toBeGreaterThan(-1)
  208. expect(page[summaryIndex + 1]?.seq).toBe(summary.seq + 1)
  209. expect(page.map(event => event.seq)).toEqual(page.map((_event, index) => third.seq + index))
  210. })
  211. it('paginates a message with many provenance sources without variadic argument expansion', async () => {
  212. const { ctx } = await harness()
  213. const remote = createSessionTestRemote(ctx, { defaultModelSelection: () => ({ provider: 'p', model: 'm' }), cwd: '/tmp' })
  214. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  215. session.append('turn/start', { turn: 1 })
  216. const sources = Array.from({ length: 128 }, (_unused, index) => session.append('assistant/chunk', {
  217. turn: 1,
  218. step: 1,
  219. chunk: { type: 'text-delta', index, text: 'x' },
  220. }).seq)
  221. const message = session.append('assistant/message', {
  222. turn: 1,
  223. step: 1,
  224. message: createMessage({
  225. role: 'assistant',
  226. content: [{ type: 'text', text: 'x'.repeat(sources.length) }],
  227. source: { kind: 'model', provider: 'p', model: 'm' },
  228. }),
  229. }, { surfaceOp: 'append', sourceEventSeqs: sources })
  230. const scalarMin = Math.min
  231. const min = vi.spyOn(Math, 'min').mockImplementation((...values) => {
  232. if (values.length > 2) throw new RangeError('variadic minimum rejected by regression harness')
  233. return scalarMin(...values)
  234. })
  235. try {
  236. const response = await remote.page({
  237. address: { kind: 'session', sessionId: session.id },
  238. throughSeq: message.seq,
  239. maxMessages: 1,
  240. })
  241. if (!response.ok) throw new Error('unreachable')
  242. expect(response.value.events.map(entry => entry.event.seq)).toEqual([...sources, message.seq])
  243. expect(response.value.hasMore).toBe(true)
  244. } finally {
  245. min.mockRestore()
  246. }
  247. })
  248. it('follows a result after turn/end without reading the addressed Session log', async () => {
  249. const { ctx } = await harness()
  250. const session = ctx.sessions.create(undefined, { meta: { cwd: '/workspace' } })
  251. const history = new SessionHistoryController(ctx, (observation) => { observation[Symbol.dispose]() })
  252. const abort = new AbortController()
  253. const stream = await openFollow(history, session.id, abort.signal)
  254. const iterator = stream[Symbol.asyncIterator]()
  255. session.append('turn/start', { turn: 1 })
  256. await expect(iterator.next()).resolves.toMatchObject({ value: { event: { type: 'turn/start' } } })
  257. session.append('tool/call', { turn: 1, step: 1, callId: CallId('c-late'), name: 'term', arguments: '{"cmd":"tail"}' })
  258. await expect(iterator.next()).resolves.toMatchObject({ value: { event: { type: 'tool/call' } } })
  259. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  260. await expect(iterator.next()).resolves.toMatchObject({ value: { event: { type: 'turn/end' } } })
  261. const events = vi.spyOn(session, 'events', 'get').mockImplementation(() => {
  262. throw new Error('live result rescanned Session history')
  263. })
  264. try {
  265. const result = session.append('tool/result', {
  266. turn: 1, step: 1,
  267. message: createToolResultMessage({
  268. callId: CallId('c-late'),
  269. content: [{ type: 'text', text: 'ok' }],
  270. isError: false,
  271. }),
  272. }, { surfaceOp: 'append' })
  273. await expect(iterator.next()).resolves.toEqual({
  274. done: false,
  275. value: { type: 'event', event: result },
  276. })
  277. } finally {
  278. events.mockRestore()
  279. abort.abort()
  280. await iterator.next()
  281. await ctx.fiber.dispose()
  282. }
  283. })
  284. })