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

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