json-stream.spec.ts 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421
  1. /** The `--json` run projection: commit-point emission, ordering, bounding, and disposal. */
  2. import { describe, expect, it } from 'vitest'
  3. import type { Context } from '@deepseek-ai/cordis'
  4. import type { Agent } from '@deepseek-ai/dsh-agent'
  5. import type { Session, SessionEvent } from '@deepseek-ai/dsh-session'
  6. import { boundJsonLine, MAX_STRING_BYTES, projectJsonRun, type JsonProjectionOptions } from '../src/json-stream.ts'
  7. interface ProjectionHarness {
  8. readonly lines: string[]
  9. readonly projection: ReturnType<typeof projectJsonRun>
  10. readonly agent: Agent
  11. readonly session: Session
  12. readonly parsed: () => Record<string, unknown>[]
  13. emitSession(event: SessionEvent): void
  14. emitRawSession(session: unknown, event: SessionEvent): void
  15. }
  16. /** One committed assistant message carrying the given content blocks. */
  17. function assistantMessage(content: unknown[], usage?: unknown): SessionEvent {
  18. return {
  19. type: 'assistant/message',
  20. data: {
  21. stream: [],
  22. turn: 1,
  23. step: 1,
  24. ...usage === undefined ? {} : { usage },
  25. message: {
  26. role: 'assistant',
  27. content,
  28. source: { kind: 'model', provider: 'p', model: 'm' },
  29. },
  30. },
  31. } as unknown as SessionEvent
  32. }
  33. /** One discarded attempt whose stream reports the given usage sample. */
  34. function attemptWithUsage(usage: unknown): SessionEvent {
  35. return {
  36. type: 'assistant/attempt',
  37. data: { turn: 1, step: 1, stream: [{ type: 'chunk', chunk: { type: 'usage', usage } }] },
  38. } as unknown as SessionEvent
  39. }
  40. /** One step boundary event closing the accumulated usage window. */
  41. function stepEnd(): SessionEvent {
  42. return { type: 'step/end', data: { turn: 1, step: 1 } } as unknown as SessionEvent
  43. }
  44. /** One tool result event with the given surface placement. */
  45. function toolResult(
  46. callId: string,
  47. content: unknown[],
  48. surfaceOp: 'append' | { op: 'replace'; startSeq: number; endSeq: number } = 'append',
  49. ): SessionEvent {
  50. return {
  51. type: 'tool/result',
  52. surfaceOp,
  53. data: {
  54. turn: 1,
  55. step: 1,
  56. message: { content: [{ type: 'tool-result', toolCallId: callId, content, isError: false }] },
  57. },
  58. } as unknown as SessionEvent
  59. }
  60. /** Drive the projector through a minimal Context and Agent double. */
  61. function harness(
  62. options: JsonProjectionOptions = {},
  63. agentId = 'session-1',
  64. cwd: string | null = '/',
  65. ): ProjectionHarness {
  66. const lines: string[] = []
  67. const sessionListeners = new Set<(session: unknown, event: SessionEvent) => void>()
  68. const ctx = {
  69. on(name: string, handler: unknown) {
  70. if (name === 'session/event') sessionListeners.add(handler as never)
  71. return () => { sessionListeners.delete(handler as never) }
  72. },
  73. } as unknown as Context
  74. const session = {} as Session
  75. const agent = { id: agentId, session } as unknown as Agent
  76. const projection = projectJsonRun(ctx, agent, {
  77. write: (chunk: string) => { lines.push(chunk); return true },
  78. }, { ...options, ...cwd === null ? {} : { cwd } })
  79. return {
  80. lines,
  81. projection,
  82. agent,
  83. session,
  84. parsed: () => lines.map(line => JSON.parse(line) as Record<string, unknown>),
  85. emitSession: (event) => { for (const listener of sessionListeners) listener(session, event) },
  86. emitRawSession: (rawSession, event) => { for (const listener of sessionListeners) listener(rawSession, event) },
  87. }
  88. }
  89. describe('--json projection', () => {
  90. it('opens with the session event before any observed event', () => {
  91. const test = harness()
  92. expect(test.parsed()).toEqual([{ type: 'session', sessionId: 'session-1', cwd: '/' }])
  93. })
  94. it('defaults the reported cwd and the per-string cap', () => {
  95. const test = harness({}, 'session-1', null)
  96. expect(test.parsed()[0]).toEqual({ type: 'session', sessionId: 'session-1', cwd: process.cwd() })
  97. test.emitSession(assistantMessage([{ type: 'text', text: 'x'.repeat(9000) }]))
  98. expect(test.parsed()[1]?.truncated).toBe(true)
  99. expect((test.parsed()[1]?.text as string).length).toBe(8 * 1024)
  100. })
  101. it('reports turn and step boundaries with the turn-end reason', () => {
  102. const test = harness()
  103. test.emitSession({ type: 'turn/start', data: { turn: 1 } } as unknown as SessionEvent)
  104. test.emitSession({ type: 'step/start', data: { turn: 1, step: 1 } } as unknown as SessionEvent)
  105. test.emitSession({
  106. type: 'turn/end',
  107. data: { turn: 1, reason: { kind: 'completed' } },
  108. } as unknown as SessionEvent)
  109. expect(test.parsed().slice(1)).toEqual([
  110. { type: 'status', phase: 'turn_start', turn: 1 },
  111. { type: 'status', phase: 'step_start', turn: 1, step: 1 },
  112. { type: 'status', phase: 'turn_end', turn: 1, reason: { kind: 'completed' } },
  113. ])
  114. })
  115. it('projects committed reasoning and text in content order', () => {
  116. const test = harness()
  117. test.emitSession(assistantMessage([
  118. { type: 'reasoning', text: 'think' },
  119. { type: 'tool-call', id: 'c1', name: 'bash', arguments: '{}' },
  120. { type: 'text', text: 'answer' },
  121. ]))
  122. expect(test.parsed().slice(1)).toEqual([
  123. { type: 'thinking', text: 'think' },
  124. { type: 'text', text: 'answer' },
  125. ])
  126. })
  127. it('ignores a discarded attempt so retried content never reaches the stream', () => {
  128. const test = harness()
  129. test.emitSession({
  130. type: 'assistant/attempt',
  131. data: { turn: 1, step: 1, stream: [] },
  132. } as unknown as SessionEvent)
  133. test.emitSession({ type: 'session/title', data: { title: 'ignored' } } as unknown as SessionEvent)
  134. expect(test.parsed().map(event => event.type)).toEqual(['session'])
  135. })
  136. it('attaches step usage to step_end and omits it when absent', () => {
  137. const test = harness()
  138. test.emitSession({ type: 'step/end', data: { turn: 1, step: 1 } } as unknown as SessionEvent)
  139. test.emitSession(assistantMessage([{ type: 'text', text: 'x' }], { inputTokens: 3, outputTokens: 4 }))
  140. test.emitSession({ type: 'step/end', data: { turn: 1, step: 2 } } as unknown as SessionEvent)
  141. const events = test.parsed()
  142. expect(events[1]).toEqual({ type: 'status', phase: 'step_end', turn: 1, step: 1 })
  143. expect(events[3]).toEqual({
  144. type: 'status', phase: 'step_end', turn: 1, step: 2,
  145. usage: { inputTokens: 3, outputTokens: 4 },
  146. })
  147. })
  148. it('reports tool results as completed or errored and keeps the final event last', () => {
  149. const test = harness()
  150. test.emitSession(toolResult('c1', [{ type: 'text', text: 'a.txt' }]))
  151. test.emitSession({
  152. type: 'tool/result',
  153. surfaceOp: 'append',
  154. data: {
  155. turn: 1,
  156. step: 1,
  157. message: {
  158. content: [{
  159. type: 'tool-result',
  160. toolCallId: 'c2',
  161. content: [{ type: 'image' }, { type: 'text' }, { type: 'text', text: 'boom' }],
  162. isError: true,
  163. }],
  164. },
  165. },
  166. } as unknown as SessionEvent)
  167. test.projection.finish('done')
  168. const events = test.parsed()
  169. expect(events[1]).toEqual({ type: 'tool_result', callId: 'c1', status: 'completed', result: 'a.txt' })
  170. expect(events[2]).toEqual({ type: 'tool_result', callId: 'c2', status: 'error', result: 'boom' })
  171. expect(events.at(-1)).toEqual({ type: 'final', text: 'done' })
  172. })
  173. it('skips a compaction replacement of an older tool result', () => {
  174. const test = harness()
  175. test.emitSession(toolResult('old', [{ type: 'text', text: 'history' }], { op: 'replace', startSeq: 1, endSeq: 2 }))
  176. expect(test.parsed().map(event => event.type)).toEqual(['session'])
  177. })
  178. it('keeps non-JSON tool arguments raw and bounds nested values', () => {
  179. const test = harness({ maxStringBytes: 11 }, 's1')
  180. test.emitSession({
  181. type: 'tool/call',
  182. data: { turn: 1, step: 1, callId: 'raw', name: 'bash', arguments: 'not json' },
  183. } as unknown as SessionEvent)
  184. test.emitSession({
  185. type: 'tool/call',
  186. data: {
  187. turn: 1,
  188. step: 1,
  189. callId: 'nested',
  190. name: 'bash',
  191. arguments: JSON.stringify({ items: ['éééééé', null, true], n: 1 }),
  192. },
  193. } as unknown as SessionEvent)
  194. const events = test.parsed()
  195. expect(events[1]).toEqual({ type: 'tool_call', callId: 'raw', tool: 'bash', input: 'not json' })
  196. expect(events[2]).toEqual({
  197. type: 'tool_call',
  198. callId: 'nested',
  199. tool: 'bash',
  200. input: { items: ['ééééé', null, true], n: 1 },
  201. truncated: true,
  202. })
  203. })
  204. it('keeps raw arguments that JSON cannot round-trip, such as an overflowing number', () => {
  205. const test = harness({}, 's1')
  206. test.emitSession({
  207. type: 'tool/call',
  208. data: { turn: 1, step: 1, callId: 'inf', name: 'bash', arguments: '{"n":1e400}' },
  209. } as unknown as SessionEvent)
  210. expect(test.parsed()[1]).toEqual({ type: 'tool_call', callId: 'inf', tool: 'bash', input: '{"n":1e400}' })
  211. })
  212. it('keeps a literal __proto__ key and bounds over-long object keys', () => {
  213. const proto = harness({ maxStringBytes: 32 }, 's1')
  214. proto.emitSession({
  215. type: 'tool/call',
  216. data: {
  217. turn: 1,
  218. step: 1,
  219. callId: 'p',
  220. name: 'bash',
  221. arguments: '{"__proto__":{"polluted":true},"a":1}',
  222. },
  223. } as unknown as SessionEvent)
  224. const protoInput = proto.parsed()[1]?.input as Record<string, unknown>
  225. expect(Object.keys(protoInput)).toEqual(['__proto__', 'a'])
  226. expect(protoInput['__proto__']).toEqual({ polluted: true })
  227. const longKey = harness({ maxStringBytes: 8 }, 's1')
  228. longKey.emitSession({
  229. type: 'tool/call',
  230. data: {
  231. turn: 1,
  232. step: 1,
  233. callId: 'k',
  234. name: 'bash',
  235. arguments: JSON.stringify({ ['k'.repeat(20)]: 1 }),
  236. },
  237. } as unknown as SessionEvent)
  238. const longEvent = longKey.parsed()[1] as { input: Record<string, unknown>; truncated?: boolean }
  239. expect(Object.keys(longEvent.input)).toEqual(['kkkkkkkk'])
  240. expect(longEvent.truncated).toBe(true)
  241. })
  242. it('drops a split trailing multibyte character when truncating', () => {
  243. const test = harness({ maxStringBytes: 5 }, 's1')
  244. test.emitSession(assistantMessage([{ type: 'text', text: 'ééé' }]))
  245. expect(test.parsed()[1]).toEqual({ type: 'text', text: 'éé', truncated: true })
  246. })
  247. it('normalizes empty tool arguments to an empty object like the executor', () => {
  248. const test = harness({}, 's1')
  249. test.emitSession({
  250. type: 'tool/call',
  251. data: { turn: 1, step: 1, callId: 'empty', name: 'bash', arguments: '' },
  252. } as unknown as SessionEvent)
  253. expect(test.parsed()[1]).toEqual({ type: 'tool_call', callId: 'empty', tool: 'bash', input: {} })
  254. })
  255. it('bounds one whole event line, dropping structured fields before scalars', () => {
  256. const input = Array.from({ length: 20_000 }, (_, index) => index)
  257. const line = boundJsonLine({ type: 'tool_call', callId: 'c', tool: 'bash', input }, MAX_STRING_BYTES, 1024)
  258. expect(Buffer.byteLength(line, 'utf8')).toBeLessThanOrEqual(1024)
  259. expect(JSON.parse(line)).toEqual({ type: 'tool_call', callId: 'c', tool: 'bash', truncated: true })
  260. })
  261. it('reduces a payload to its type when even its scalars exceed the line cap', () => {
  262. const line = boundJsonLine({ type: 'text', text: 'x'.repeat(2000) }, 4096, 64)
  263. expect(JSON.parse(line)).toEqual({ type: 'text', truncated: true })
  264. expect(JSON.parse(boundJsonLine({ type: 'text', text: 'ok' }))).toEqual({ type: 'text', text: 'ok' })
  265. })
  266. it('caps an error line even when control characters expand under JSON escaping', () => {
  267. const line = boundJsonLine({ type: 'error', message: '\u0000'.repeat(MAX_STRING_BYTES) })
  268. expect(Buffer.byteLength(line, 'utf8') + 1).toBeLessThanOrEqual(32 * 1024)
  269. expect(JSON.parse(line)).toEqual({ type: 'error', truncated: true })
  270. })
  271. it('reserves the trailing newline inside the whole-line cap', () => {
  272. // A 64-byte line exactly fills a 64-byte cap; the writer's newline must
  273. // force the scalar fallback rather than write 65 bytes.
  274. const line = boundJsonLine({ type: 'text', text: 'x'.repeat(39) }, 4096, 64)
  275. expect(Buffer.byteLength(line, 'utf8') + 1).toBeLessThanOrEqual(64)
  276. expect(JSON.parse(line)).toEqual({ type: 'text', truncated: true })
  277. })
  278. it('cuts a payload that nests past the depth budget instead of overflowing the stack', () => {
  279. let deepArray: unknown = 'leaf'
  280. for (let level = 0; level < 200; level += 1) deepArray = [deepArray]
  281. expect(JSON.parse(boundJsonLine({ type: 'tool_call', input: deepArray })))
  282. .toMatchObject({ type: 'tool_call', truncated: true })
  283. let deepObject: unknown = 'leaf'
  284. for (let level = 0; level < 200; level += 1) deepObject = { next: deepObject }
  285. expect(JSON.parse(boundJsonLine({ type: 'tool_call', input: deepObject })))
  286. .toMatchObject({ type: 'tool_call', truncated: true })
  287. })
  288. it('sums retried attempt usage into the step total and drops an unshared bucket', () => {
  289. const test = harness()
  290. test.emitSession(attemptWithUsage({ inputTokens: 10, outputTokens: 2, totalTokens: 12, cacheReadTokens: 4 }))
  291. test.emitSession(assistantMessage(
  292. [{ type: 'text', text: 'ok' }],
  293. { inputTokens: 3, outputTokens: 1, totalTokens: 4 },
  294. ))
  295. test.emitSession(stepEnd())
  296. expect(test.parsed().at(-1)).toEqual({
  297. type: 'status', phase: 'step_end', turn: 1, step: 1,
  298. usage: { inputTokens: 13, outputTokens: 3, totalTokens: 16 },
  299. })
  300. })
  301. it('sums every optional bucket both attempts report', () => {
  302. const test = harness()
  303. const sample = {
  304. inputTokens: 2, outputTokens: 1, totalTokens: 3,
  305. cacheReadTokens: 1, cacheWriteTokens: 1, reasoningTokens: 1,
  306. }
  307. test.emitSession(attemptWithUsage(sample))
  308. test.emitSession(attemptWithUsage(sample))
  309. test.emitSession(stepEnd())
  310. expect(test.parsed().at(-1)).toMatchObject({
  311. usage: {
  312. inputTokens: 4, outputTokens: 2, totalTokens: 6,
  313. cacheReadTokens: 2, cacheWriteTokens: 2, reasoningTokens: 2,
  314. },
  315. })
  316. })
  317. it('drops a bucket only the later attempt reports', () => {
  318. const test = harness()
  319. test.emitSession(attemptWithUsage({ inputTokens: 1, outputTokens: 1 }))
  320. test.emitSession(attemptWithUsage({ inputTokens: 1, outputTokens: 1, cacheWriteTokens: 2 }))
  321. test.emitSession(stepEnd())
  322. const usage = (test.parsed().at(-1) as { usage: Record<string, unknown> }).usage
  323. expect(usage).toEqual({ inputTokens: 2, outputTokens: 2 })
  324. expect(usage).not.toHaveProperty('cacheWriteTokens')
  325. })
  326. it('reads a message usage sample from its stream when the field is absent', () => {
  327. const test = harness()
  328. test.emitSession({
  329. type: 'assistant/message',
  330. data: {
  331. turn: 1,
  332. step: 1,
  333. stream: [{ type: 'chunk', chunk: { type: 'usage', usage: { inputTokens: 5, outputTokens: 1 } } }],
  334. message: {
  335. role: 'assistant',
  336. content: [{ type: 'text', text: 'hi' }],
  337. source: { kind: 'model', provider: 'p', model: 'm' },
  338. },
  339. },
  340. } as unknown as SessionEvent)
  341. test.emitSession(stepEnd())
  342. expect(test.parsed().at(-1)).toMatchObject({ usage: { inputTokens: 5, outputTokens: 1 } })
  343. })
  344. it('omits the step total when a later attempt reports no usage sample', () => {
  345. const test = harness()
  346. test.emitSession(attemptWithUsage({ inputTokens: 7, outputTokens: 3 }))
  347. test.emitSession(assistantMessage([{ type: 'text', text: 'ok' }]))
  348. test.emitSession(stepEnd())
  349. expect(test.parsed().at(-1)).toEqual({ type: 'status', phase: 'step_end', turn: 1, step: 1 })
  350. })
  351. it('omits the step total when an earlier attempt reports no usage sample', () => {
  352. const test = harness()
  353. test.emitSession({
  354. type: 'assistant/attempt',
  355. data: { turn: 1, step: 1, stream: [] },
  356. } as unknown as SessionEvent)
  357. test.emitSession(assistantMessage([{ type: 'text', text: 'ok' }], { inputTokens: 1, outputTokens: 1 }))
  358. test.emitSession(stepEnd())
  359. expect(test.parsed().at(-1)).toEqual({ type: 'status', phase: 'step_end', turn: 1, step: 1 })
  360. })
  361. it('omits the step total when no attempt reports a usage sample', () => {
  362. const test = harness()
  363. test.emitSession(stepEnd())
  364. expect(test.parsed().at(-1)).toEqual({ type: 'status', phase: 'step_end', turn: 1, step: 1 })
  365. })
  366. it('writes the terminal final event without bounding its answer', () => {
  367. const test = harness({ maxStringBytes: 4 })
  368. test.projection.finish('abcdefgh')
  369. expect(test.parsed().at(-1)).toEqual({ type: 'final', text: 'abcdefgh' })
  370. })
  371. it('ignores events from another Session and stops writing after dispose', () => {
  372. const test = harness()
  373. test.emitRawSession({}, assistantMessage([{ type: 'text', text: 'foreign' }]))
  374. test.projection.dispose()
  375. test.emitSession(assistantMessage([{ type: 'text', text: 'late' }]))
  376. test.projection.finish('ignored')
  377. expect(test.parsed().map(event => event.type)).toEqual(['session'])
  378. })
  379. it('bounds one projected payload through the line writer', () => {
  380. expect(JSON.parse(boundJsonLine({ type: 'error', message: 'x'.repeat(20) }, 8, 4096)))
  381. .toEqual({ type: 'error', message: 'xxxxxxxx', truncated: true })
  382. expect(JSON.parse(boundJsonLine({ type: 'error', message: 'ok' }))).toEqual({ type: 'error', message: 'ok' })
  383. })
  384. })