projection.spec.ts 16 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321
  1. /**
  2. * The `sessionStats` projection unit: mounting the plugin beside the
  3. * projection registry serves whole-log counts and wall times folded from step
  4. * boundaries, chunks, tool pairs, and assembled messages; compositions
  5. * without the registry are unaffected; unmounting the plugin removes the key
  6. * (HMR safety). The two counting regressions pinned here are the reasons the
  7. * fold counts step boundaries instead of assistant messages: a cancelled step
  8. * never assembles a message but still counts, and a max-tokens usage-host
  9. * message (empty content) adds no extra step. Wall-time math runs against the
  10. * exported definition directly, where event times are controlled.
  11. */
  12. import { describe, expect, it } from 'vitest'
  13. import { Context } from '@deepseek-ai/cordis'
  14. import { createMessage } from '@deepseek-ai/dsh-llm'
  15. import SessionStore, { SessionId } from '@deepseek-ai/dsh-session'
  16. import type { Session, SessionEvent } from '@deepseek-ai/dsh-session'
  17. import SessionProjectionRegistry from '@deepseek-ai/dsh-session-projection'
  18. import * as SessionStatsPlugin from '@deepseek-ai/dsh-session-stats'
  19. import { sessionStatsProjectionDefinition } from '@deepseek-ai/dsh-session-stats/src/projection.ts'
  20. import type { SessionStatsProjection } from '@deepseek-ai/dsh-session-stats/types'
  21. async function harness(withStatsPlugin: boolean): Promise<{ ctx: Context; session: Session }> {
  22. const ctx = new Context()
  23. await ctx.plugin(SessionStore)
  24. await ctx.plugin(SessionProjectionRegistry)
  25. if (withStatsPlugin) await ctx.plugin(SessionStatsPlugin)
  26. return { ctx, session: ctx.sessions.create(SessionId('counted')) }
  27. }
  28. /** Close one step; returns the counted `step/end` seq. */
  29. function closeStep(session: Session, turn: number, step: number): number {
  30. session.append('step/start', { turn, step })
  31. return session.append('step/end', { turn, step }).seq
  32. }
  33. /** Append the max-tokens usage-host shape: an assistant/message with empty content. */
  34. function appendEmptyAssistantMessage(session: Session, turn: number, step: number): void {
  35. session.append('assistant/message', {
  36. turn,
  37. step,
  38. message: createMessage({
  39. role: 'assistant',
  40. content: [],
  41. source: { kind: 'model', provider: 'mock', model: 'mock' },
  42. }),
  43. }, { surfaceOp: 'append', sourceEventSeqs: [] })
  44. }
  45. /** The all-zero projection value plus overrides, for exact fold expectations. */
  46. function totals(overrides: Partial<SessionStatsProjection> = {}): SessionStatsProjection {
  47. return {
  48. turns: 0, steps: 0, llmMs: 0, toolMs: 0, ttftMs: 0, ttftSteps: 0, decodeMs: 0, decodeTokens: 0,
  49. ...overrides,
  50. }
  51. }
  52. describe('sessionStats projection unit (registry drive)', () => {
  53. it('serves zero figures on the empty log', async () => {
  54. const { ctx, session } = await harness(true)
  55. expect(ctx.sessionProjections.snapshot(session).values.sessionStats).toEqual(totals())
  56. })
  57. it('counts distinct turns and closed steps and notifies the change feed with the causing seq', async () => {
  58. const { ctx, session } = await harness(true)
  59. const changes: { key: string; value: unknown; seq: number }[] = []
  60. ctx.sessionProjections.onChanged((_session, key, value, seq) => {
  61. changes.push({ key, value, seq })
  62. })
  63. session.append('turn/start', { turn: 1 })
  64. const firstSeq = closeStep(session, 1, 1)
  65. const secondSeq = closeStep(session, 1, 2)
  66. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  67. session.append('turn/start', { turn: 2 })
  68. const thirdSeq = closeStep(session, 2, 1)
  69. session.append('turn/end', { turn: 2, reason: { kind: 'completed' } })
  70. // Boundary events that carry no figure change (turn/start, empty-prune
  71. // turn/end, user input) fold to the same reference and stay silent;
  72. // step/start opens a boundary (internal state) and step/end commits the
  73. // counts, so each closed step notifies twice with the step/end value last.
  74. const counted = changes.filter(change => (change.value as SessionStatsProjection).steps > 0
  75. || change.seq === firstSeq)
  76. expect(changes.every(change => change.key === 'sessionStats')).toBe(true)
  77. expect(counted.map(change => ({ seq: change.seq, value: change.value }))).toContainEqual(
  78. { seq: firstSeq, value: totals({ turns: 1, steps: 1 }) },
  79. )
  80. expect(changes.at(-1)).toEqual({ key: 'sessionStats', value: totals({ turns: 2, steps: 3 }), seq: thirdSeq })
  81. const snapshot = ctx.sessionProjections.snapshot(session)
  82. expect(snapshot.values.sessionStats).toEqual(totals({ turns: 2, steps: 3 }))
  83. expect(snapshot.asOfSeq).toBe(session.seq - 1)
  84. expect(changes.map(change => change.seq)).toContain(secondSeq)
  85. })
  86. it('does not count a rejected or empty turn that closes with no step', async () => {
  87. const { ctx, session } = await harness(true)
  88. session.append('turn/start', { turn: 1 })
  89. session.append('turn/end', { turn: 1, reason: { kind: 'blocked' } })
  90. expect(ctx.sessionProjections.snapshot(session).values.sessionStats).toEqual(totals())
  91. })
  92. it('counts a cancelled step that closed without an assistant message', async () => {
  93. // Regression: an aborted stream never assembles assistant/message, but the
  94. // loop's finally still appends step/end — the step happened and counts.
  95. const { ctx, session } = await harness(true)
  96. session.append('turn/start', { turn: 1 })
  97. closeStep(session, 1, 1)
  98. session.append('turn/end', { turn: 1, reason: { kind: 'aborted', reason: { kind: 'legacy' } } })
  99. expect(ctx.sessionProjections.snapshot(session).values.sessionStats)
  100. .toMatchObject({ turns: 1, steps: 1 })
  101. })
  102. it('adds no extra step for a max-tokens usage-host assistant message', async () => {
  103. // Regression: the empty-content assistant/message exists only to host
  104. // usage and is excluded from the surface; the step counts once, from its
  105. // step/end, while the message contributes only its model wall time.
  106. const { ctx, session } = await harness(true)
  107. session.append('turn/start', { turn: 1 })
  108. session.append('step/start', { turn: 1, step: 1 })
  109. appendEmptyAssistantMessage(session, 1, 1)
  110. session.append('step/end', { turn: 1, step: 1 })
  111. session.append('turn/end', { turn: 1, reason: { kind: 'max-tokens' } })
  112. expect(ctx.sessionProjections.snapshot(session).values.sessionStats)
  113. .toMatchObject({ turns: 1, steps: 1, ttftSteps: 0, decodeTokens: 0 })
  114. })
  115. it('folds steps already in the log when the plugin mounts late (lazy cell build)', async () => {
  116. const { ctx, session } = await harness(false)
  117. session.append('turn/start', { turn: 1 })
  118. closeStep(session, 1, 1)
  119. closeStep(session, 1, 2)
  120. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  121. await ctx.plugin(SessionStatsPlugin)
  122. expect(ctx.sessionProjections.snapshot(session).values.sessionStats)
  123. .toMatchObject({ turns: 1, steps: 2 })
  124. })
  125. it('has no sessionStats key without the plugin, and drops it when the plugin unloads (HMR safety)', async () => {
  126. const { ctx, session } = await harness(false)
  127. expect('sessionStats' in ctx.sessionProjections.snapshot(session).values).toBe(false)
  128. const fiber = await ctx.plugin(SessionStatsPlugin)
  129. session.append('turn/start', { turn: 1 })
  130. closeStep(session, 1, 1)
  131. expect(ctx.sessionProjections.snapshot(session).values.sessionStats)
  132. .toMatchObject({ turns: 1, steps: 1 })
  133. await fiber.dispose()
  134. expect('sessionStats' in ctx.sessionProjections.snapshot(session).values).toBe(false)
  135. })
  136. })
  137. /** Build one synthetic committed event with a controlled timestamp. */
  138. function at(time: number, type: string, data: unknown): SessionEvent {
  139. return { type, seq: time, time, data } as unknown as SessionEvent
  140. }
  141. /** Fold a synthetic event list through the definition and view the result. */
  142. function fold(events: readonly SessionEvent[]): SessionStatsProjection {
  143. const state = events.reduce<Parameters<typeof sessionStatsProjectionDefinition.apply>[0]>(
  144. (folded, event) => sessionStatsProjectionDefinition.apply(folded, event),
  145. sessionStatsProjectionDefinition.init(),
  146. )
  147. return sessionStatsProjectionDefinition.wire.view(state)
  148. }
  149. describe('sessionStats wall-time fold (controlled timestamps)', () => {
  150. const message = createMessage({
  151. role: 'assistant',
  152. content: [{ type: 'text', text: 'answer' }],
  153. source: { kind: 'model', provider: 'mock', model: 'mock' },
  154. })
  155. it('accrues model, first-token, and decode time from one fully recorded step', () => {
  156. expect(fold([
  157. at(1_000, 'step/start', { turn: 1, step: 1 }),
  158. at(1_800, 'assistant/chunk', { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: 'a' } }),
  159. at(4_800, 'assistant/message', { turn: 1, step: 1, message, usage: { inputTokens: 10, outputTokens: 60 } }),
  160. at(4_900, 'step/end', { turn: 1, step: 1 }),
  161. ])).toEqual(totals({
  162. turns: 1, steps: 1, llmMs: 3_800, ttftMs: 800, ttftSteps: 1, decodeMs: 3_000, decodeTokens: 60,
  163. }))
  164. })
  165. it('keeps the first attempt token boundary across an in-step retry (window resetForRetry parity)', () => {
  166. expect(fold([
  167. at(1_000, 'step/start', { turn: 1, step: 1 }),
  168. at(1_200, 'assistant/chunk', { turn: 1, step: 1, chunk: { type: 'reasoning-delta', index: 0, text: 'x' } }),
  169. at(2_000, 'llm/retry', { turn: 1, step: 1 }),
  170. at(3_000, 'assistant/chunk', { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: 'y' } }),
  171. at(5_000, 'assistant/message', { turn: 1, step: 1, message }),
  172. at(5_100, 'step/end', { turn: 1, step: 1 }),
  173. ])).toEqual(totals({ turns: 1, steps: 1, llmMs: 4_000, ttftMs: 200, ttftSteps: 1 }))
  174. })
  175. it('ignores empty deltas, non-token chunks, and chunks outside the open step', () => {
  176. expect(fold([
  177. // Chunk before any step/start: no open boundary.
  178. at(500, 'assistant/chunk', { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: 'stray' } }),
  179. at(1_000, 'step/start', { turn: 1, step: 1 }),
  180. at(1_100, 'assistant/chunk', { turn: 1, step: 1, chunk: { type: 'block-start', index: 0, blockType: 'text' } }),
  181. at(1_200, 'assistant/chunk', { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: '' } }),
  182. at(1_300, 'assistant/chunk', { turn: 2, step: 9, chunk: { type: 'text-delta', index: 0, text: 'other' } }),
  183. at(1_400, 'assistant/chunk', { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: 'first' } }),
  184. at(2_000, 'assistant/message', { turn: 1, step: 1, message }),
  185. at(2_100, 'step/end', { turn: 1, step: 1 }),
  186. ])).toEqual(totals({ turns: 1, steps: 1, llmMs: 1_000, ttftMs: 400, ttftSteps: 1 }))
  187. })
  188. it('uses non-empty Tool-call names or arguments as the first token', () => {
  189. expect(fold([
  190. at(1_000, 'step/start', { turn: 1, step: 1 }),
  191. at(1_100, 'assistant/chunk', {
  192. turn: 1,
  193. step: 1,
  194. chunk: { type: 'tool-call-delta', index: 0, id: 'call-1', argumentsDelta: '' },
  195. }),
  196. at(1_200, 'assistant/chunk', {
  197. turn: 1,
  198. step: 1,
  199. chunk: { type: 'tool-call-delta', index: 0, id: 'call-1', name: 'read', argumentsDelta: '' },
  200. }),
  201. at(2_000, 'assistant/message', { turn: 1, step: 1, message }),
  202. at(2_100, 'step/end', { turn: 1, step: 1 }),
  203. ])).toEqual(totals({ turns: 1, steps: 1, llmMs: 1_000, ttftMs: 200, ttftSteps: 1 }))
  204. expect(fold([
  205. at(1_000, 'step/start', { turn: 1, step: 1 }),
  206. at(1_300, 'assistant/chunk', {
  207. turn: 1,
  208. step: 1,
  209. chunk: { type: 'tool-call-delta', index: 0, id: 'call-1', argumentsDelta: '{' },
  210. }),
  211. at(2_000, 'assistant/message', { turn: 1, step: 1, message }),
  212. at(2_100, 'step/end', { turn: 1, step: 1 }),
  213. ])).toEqual(totals({ turns: 1, steps: 1, llmMs: 1_000, ttftMs: 300, ttftSteps: 1 }))
  214. })
  215. it('leaves a cancelled step untimed: counted by step/end, no assembled message to accrue from', () => {
  216. expect(fold([
  217. at(1_000, 'step/start', { turn: 1, step: 1 }),
  218. at(1_500, 'assistant/chunk', { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: 'partial' } }),
  219. at(2_000, 'step/end', { turn: 1, step: 1 }),
  220. ])).toEqual(totals({ turns: 1, steps: 1 }))
  221. })
  222. it('pairs tool wall time by callId, ignores orphan results, and prunes leftovers at turn/end', () => {
  223. const result = (callId: string): unknown =>
  224. ({ turn: 1, step: 1, message: { source: { kind: 'tool', callId } } })
  225. const paired = fold([
  226. at(1_000, 'step/start', { turn: 1, step: 1 }),
  227. at(1_100, 'tool/call', { turn: 1, step: 1, callId: 'a', name: 'read', arguments: '{}' }),
  228. at(1_200, 'tool/call', { turn: 1, step: 1, callId: 'b', name: 'read', arguments: '{}' }),
  229. // Out-of-order settlement pairs by id, not adjacency.
  230. at(4_200, 'tool/result', result('b')),
  231. at(1_600, 'tool/result', result('a')),
  232. at(5_000, 'tool/result', result('ghost')),
  233. at(5_100, 'step/end', { turn: 1, step: 1 }),
  234. ])
  235. expect(paired).toEqual(totals({ turns: 1, steps: 1, toolMs: 3_500 }))
  236. // An unresolved call is dropped at turn/end; a later result cannot pair.
  237. const pruned = fold([
  238. at(1_000, 'step/start', { turn: 1, step: 1 }),
  239. at(1_100, 'tool/call', { turn: 1, step: 1, callId: 'orphan', name: 'read', arguments: '{}' }),
  240. at(2_000, 'step/end', { turn: 1, step: 1 }),
  241. at(2_100, 'turn/end', { turn: 1, reason: { kind: 'aborted', reason: { kind: 'legacy' } } }),
  242. at(9_000, 'tool/result', result('orphan')),
  243. ])
  244. expect(pruned).toEqual(totals({ turns: 1, steps: 1 }))
  245. })
  246. it('pairs only own pendingCalls keys: a prototype-name callId without a recorded call stays unmatched', () => {
  247. const result = (callId: string): unknown =>
  248. ({ turn: 1, step: 1, message: { source: { kind: 'tool', callId } } })
  249. // Crash recovery (TOOL_NOT_STARTED) emits results with no preceding
  250. // tool/call; a provider-minted callId colliding with an Object prototype
  251. // property must read as absent, not as an inherited function that would
  252. // fold toolMs to NaN and fail the value schema.
  253. expect(fold([
  254. at(1_000, 'step/start', { turn: 1, step: 1 }),
  255. at(1_500, 'tool/result', result('toString')),
  256. at(2_000, 'step/end', { turn: 1, step: 1 }),
  257. ])).toEqual(totals({ turns: 1, steps: 1 }))
  258. // The same name pairs normally once its call is recorded.
  259. expect(fold([
  260. at(1_000, 'step/start', { turn: 1, step: 1 }),
  261. at(1_100, 'tool/call', { turn: 1, step: 1, callId: 'constructor', name: 'read', arguments: '{}' }),
  262. at(1_600, 'tool/result', result('constructor')),
  263. at(2_000, 'step/end', { turn: 1, step: 1 }),
  264. ])).toEqual(totals({ turns: 1, steps: 1, toolMs: 500 }))
  265. })
  266. it('skips decode for an invalid usage report and ignores a duplicate assembled message', () => {
  267. const events = [
  268. at(1_000, 'step/start', { turn: 1, step: 1 }),
  269. at(1_400, 'assistant/chunk', { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: 'a' } }),
  270. // A malformed provider report: guarded like the window fold guards node usage.
  271. at(2_000, 'assistant/message', { turn: 1, step: 1, message, usage: { inputTokens: 1, outputTokens: -5 } }),
  272. ]
  273. expect(fold([...events, at(2_100, 'step/end', { turn: 1, step: 1 })]))
  274. .toEqual(totals({ turns: 1, steps: 1, llmMs: 1_000, ttftMs: 400, ttftSteps: 1 }))
  275. // The first message closed the step boundary; a defensive duplicate finds
  276. // no open step and folds to the same reference.
  277. const state = events.reduce<Parameters<typeof sessionStatsProjectionDefinition.apply>[0]>(
  278. (folded, event) => sessionStatsProjectionDefinition.apply(folded, event),
  279. sessionStatsProjectionDefinition.init(),
  280. )
  281. expect(sessionStatsProjectionDefinition.apply(
  282. state,
  283. at(2_050, 'assistant/message', { turn: 1, step: 1, message }),
  284. )).toBe(state)
  285. })
  286. it('accrues nothing for unrelated events and clamps negative clock skew to zero', () => {
  287. const state = sessionStatsProjectionDefinition.init()
  288. const untouched = sessionStatsProjectionDefinition.apply(state, at(1, 'user/message', { content: [] }))
  289. expect(untouched).toBe(state)
  290. expect(fold([
  291. at(2_000, 'step/start', { turn: 1, step: 1 }),
  292. at(1_000, 'assistant/message', { turn: 1, step: 1, message }),
  293. at(2_100, 'step/end', { turn: 1, step: 1 }),
  294. ])).toEqual(totals({ turns: 1, steps: 1 }))
  295. })
  296. })