telemetry.spec.ts 19 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443
  1. import { createToolResultMessage, createUserMessage } from '@deepseek-ai/dsh-llm'
  2. /**
  3. * Coordinator semantics against a bare fake backend — the RFC's named unit
  4. * tier for the seam: adoption (fresh, seeded, re-adoption via the handoff
  5. * cursor), the fixed chunk projection, deep-copy isolation, turn-latency and
  6. * dispose-ordering pins, failure containment, and the `agent/error` relay.
  7. */
  8. import { describe, expect, it, vi } from 'vitest'
  9. import { Context } from 'cordis'
  10. import SessionStore, { SessionId, type Session, type SessionEvent } from '@deepseek-ai/dsh-session'
  11. import type { Agent } from '@deepseek-ai/dsh-agent'
  12. import { TelemetryCoordinator, type TelemetryBackend, type TelemetryRecord } from '../src/index.ts'
  13. declare module '@deepseek-ai/dsh-session' {
  14. interface SessionEventMap {
  15. /**
  16. * Test-only merged event proving unknown types flow through unchanged.
  17. * @mode emit
  18. * @param payload - opaque test payload
  19. */
  20. 'telemetry-test/opaque': { payload: { nested: string[] } }
  21. }
  22. }
  23. class FakeBackend implements TelemetryBackend {
  24. records: TelemetryRecord[] = []
  25. calls: string[] = []
  26. emitError: Error | undefined
  27. rejectSeq: number | undefined
  28. shutdownError: Error | undefined
  29. shutdownResolved = false
  30. emit(record: TelemetryRecord): void {
  31. if (this.emitError) throw this.emitError
  32. if (this.rejectSeq !== undefined && record.attributes['event.seq'] === this.rejectSeq) {
  33. throw new Error(`backend rejected seq ${this.rejectSeq}`)
  34. }
  35. this.records.push(record)
  36. this.calls.push(`emit:${String(record.attributes['event.seq'] ?? record.attributes['telemetry.op'])}`)
  37. }
  38. flush = vi.fn()
  39. async shutdown(): Promise<void> {
  40. this.calls.push('shutdown')
  41. await new Promise(resolve => setTimeout(resolve, 5))
  42. if (this.shutdownError) throw this.shutdownError
  43. this.shutdownResolved = true
  44. }
  45. ledger(): TelemetryRecord[] {
  46. return this.records.filter(r => r.channel === 'ledger')
  47. }
  48. }
  49. async function setup(backend: FakeBackend = new FakeBackend()) {
  50. const ctx = new Context()
  51. await ctx.plugin(SessionStore)
  52. const fiber = await ctx.plugin({
  53. name: 'fake-telemetry',
  54. inject: ['sessions'],
  55. apply: (inner: Context) => void new TelemetryCoordinator(inner, backend),
  56. })
  57. return { ctx, backend, fiber }
  58. }
  59. function liveSession(ctx: Context, id = `s-${Math.random().toString(36).slice(2)}`): Session {
  60. return ctx.sessions.create(SessionId(id), { meta: {} })
  61. }
  62. function appendTurn(session: Session): void {
  63. session.append('turn/start', { turn: 1 })
  64. session.append('user/message', createUserMessage({
  65. content: [{ type: 'text', text: 'hello' }], source: { kind: 'user' },
  66. }), { surfaceOp: 'append' })
  67. }
  68. describe('TelemetryCoordinator capture', () => {
  69. it('hands every appended event over with envelope identity and cloned body', async () => {
  70. const { ctx, backend } = await setup()
  71. const session = liveSession(ctx, 'cap')
  72. appendTurn(session)
  73. const start = backend.ledger()[0]!
  74. const message = backend.ledger()[1]!
  75. expect(start.attributes).toMatchObject({ 'session.id': 'cap', 'event.type': 'turn/start', 'event.seq': 0 })
  76. expect(start.time).toBe(session.events[0]!.time)
  77. expect(start.severity).toBe('info')
  78. expect(message.attributes['event.seq']).toBe(1)
  79. // Deep-copy isolation: mutating the handed-off body never reaches the log.
  80. ;(message.body as { content: { text: string }[] }).content[0]!.text = 'tampered'
  81. const logged = session.events[1] as SessionEvent<'user/message'>
  82. expect(logged.data.content[0]).toMatchObject({ text: 'hello' })
  83. })
  84. it('stamps header facts on every record when present', async () => {
  85. const { ctx, backend } = await setup()
  86. const parent = SessionId('parent')
  87. const session = ctx.sessions.create(SessionId('child'), { meta: { cwd: '/tmp/proj', parentSession: parent } })
  88. appendTurn(session)
  89. for (const record of backend.ledger()) {
  90. expect(record.attributes['session.cwd']).toBe('/tmp/proj')
  91. expect(record.attributes['session.parent_id']).toBe('parent')
  92. }
  93. })
  94. it('maps outcome flags to severity, unknown types falling through as info', async () => {
  95. const { ctx, backend } = await setup()
  96. const session = liveSession(ctx)
  97. session.append('turn/start', { turn: 1 })
  98. session.append('tool/result', {
  99. turn: 1, step: 1,
  100. message: createToolResultMessage({
  101. callId: 'c1' as never,
  102. content: [],
  103. isError: true,
  104. }),
  105. }, { surfaceOp: 'append' })
  106. session.append('tool/result', {
  107. turn: 1, step: 1,
  108. message: createToolResultMessage({
  109. callId: 'c2' as never,
  110. content: [],
  111. isError: false,
  112. }),
  113. }, { surfaceOp: 'append' })
  114. session.append('telemetry-test/opaque', { payload: { nested: [] } })
  115. session.append('turn/end', { turn: 1, reason: { kind: 'error', error: { message: 'boom', code: 'UNKNOWN' } } })
  116. const severities = backend.ledger().map(r => [r.attributes['event.type'], r.severity])
  117. expect(severities).toEqual([
  118. ['turn/start', 'info'],
  119. ['tool/result', 'error'],
  120. ['tool/result', 'info'],
  121. ['telemetry-test/opaque', 'info'],
  122. ['turn/end', 'error'],
  123. ])
  124. })
  125. it('passes unknown merged event types through unchanged', async () => {
  126. const { ctx, backend } = await setup()
  127. const session = liveSession(ctx)
  128. session.append('telemetry-test/opaque', { payload: { nested: ['a', 'b'] } })
  129. const record = backend.ledger()[0]!
  130. expect(record.attributes['event.type']).toBe('telemetry-test/opaque')
  131. expect(record.severity).toBe('info')
  132. expect(record.body).toEqual({ payload: { nested: ['a', 'b'] } })
  133. })
  134. it('ships only the first chunk of each (turn, step), per session', async () => {
  135. const { ctx, backend } = await setup()
  136. const a = liveSession(ctx, 'a')
  137. const b = liveSession(ctx, 'b')
  138. const chunk = (s: Session, turn: number, step: number, text: string) =>
  139. s.append('assistant/chunk', { turn, step, chunk: { type: 'text-delta', index: 0, text } })
  140. chunk(a, 1, 1, 'a11-first')
  141. chunk(a, 1, 1, 'a11-second')
  142. chunk(a, 1, 2, 'a12-first')
  143. chunk(b, 1, 1, 'b11-first')
  144. chunk(b, 1, 1, 'b11-second')
  145. const shipped = backend.ledger().map(r => [r.attributes['session.id'], (r.body as { chunk: { text: string } }).chunk.text])
  146. expect(shipped).toEqual([
  147. ['a', 'a11-first'],
  148. ['a', 'a12-first'],
  149. ['b', 'b11-first'],
  150. ])
  151. })
  152. })
  153. describe('TelemetryCoordinator adoption', () => {
  154. it('exports an unpublished suffix without re-exporting constructor history', async () => {
  155. const backend = new FakeBackend()
  156. const ctx = new Context()
  157. await ctx.plugin(SessionStore)
  158. const parent = liveSession(ctx, 'seed-parent')
  159. appendTurn(parent)
  160. await ctx.plugin({
  161. name: 'fake-telemetry',
  162. inject: ['sessions'],
  163. apply: (inner: Context) => void new TelemetryCoordinator(inner, backend),
  164. })
  165. const child = ctx.sessions.prepare(SessionId('seeded'), { seed: [...parent.events], meta: {} })
  166. child.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  167. ctx.sessions.enter(child)
  168. ctx.sessions.announce(child)
  169. const seqs = backend.ledger().map(r => [r.attributes['session.id'], r.attributes['event.seq']])
  170. expect(seqs).toEqual(expect.arrayContaining([['seed-parent', 0], ['seed-parent', 1]]))
  171. // 2 end-seed, 3 turn/end: both this lifecycle's own writes, while
  172. // inherited 0-1 stay with the parent stream.
  173. expect(seqs.filter(([id]) => id === 'seeded')).toEqual([['seeded', 2], ['seeded', 3]])
  174. })
  175. it('resume shape: a full-log seed exports only its own end-seed and rebuilds the chunk projection', async () => {
  176. const backend = new FakeBackend()
  177. const ctx = new Context()
  178. await ctx.plugin(SessionStore)
  179. const donor = ctx.sessions.create(SessionId('donor'), { meta: {} })
  180. donor.append('turn/start', { turn: 1 })
  181. donor.append('assistant/chunk', { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: 'first' } })
  182. const resumed = ctx.sessions.create(SessionId('resumed'), { seed: [...donor.events], meta: {} })
  183. await ctx.plugin({
  184. name: 'fake-telemetry',
  185. inject: ['sessions'],
  186. apply: (inner: Context) => void new TelemetryCoordinator(inner, backend),
  187. })
  188. const ofResumed = () => backend.ledger()
  189. .filter(r => r.attributes['session.id'] === 'resumed')
  190. .map(r => r.attributes['event.seq'])
  191. // Nothing inherited is re-exported; seq 2 is this session's own first
  192. // write — the end-seed event its constructor appended after the seed.
  193. expect(ofResumed()).toEqual([2])
  194. // The seed fed the projection: the (turn 1, step 1) first chunk already
  195. // shipped from the original process, so its continuation is re-dropped…
  196. resumed.append('assistant/chunk', { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: 'continuation' } })
  197. expect(ofResumed()).toEqual([2])
  198. // …while a new step's first chunk exports normally.
  199. resumed.append('assistant/chunk', { turn: 1, step: 2, chunk: { type: 'text-delta', index: 0, text: 'next step' } })
  200. expect(ofResumed()).toEqual([2, 4])
  201. })
  202. it('stamps session.seed_length from the header so receivers can stitch fork streams', async () => {
  203. const backend = new FakeBackend()
  204. const ctx = new Context()
  205. await ctx.plugin(SessionStore)
  206. const parent = liveSession(ctx, 'stitch-parent')
  207. appendTurn(parent)
  208. const child = ctx.sessions.create(SessionId('stitch-child'), {
  209. seed: [...parent.events],
  210. meta: { parentSession: SessionId('stitch-parent'), seedLength: 2 },
  211. })
  212. await ctx.plugin({
  213. name: 'fake-telemetry',
  214. inject: ['sessions'],
  215. apply: (inner: Context) => void new TelemetryCoordinator(inner, backend),
  216. })
  217. child.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  218. const record = backend.ledger().find(r => r.attributes['session.id'] === 'stitch-child')!
  219. expect(record.attributes['session.parent_id']).toBe('stitch-parent')
  220. expect(record.attributes['session.seed_length']).toBe(2)
  221. })
  222. it('adopts exactly once when created fires after the sweep', async () => {
  223. const backend = new FakeBackend()
  224. const ctx = new Context()
  225. await ctx.plugin(SessionStore)
  226. // The enter/announce window: prepare+enter puts the session in the store
  227. // (visible to the constructor sweep) before `session/created` fires, so a
  228. // coordinator loaded inside that window sees the session twice — sweep
  229. // first, created second. The second adoption must be a no-op.
  230. const session = ctx.sessions.prepare(SessionId('overlap'))
  231. appendTurn(session)
  232. ctx.sessions.enter(session)
  233. await ctx.plugin({
  234. name: 'fake-telemetry',
  235. inject: ['sessions'],
  236. apply: (inner: Context) => void new TelemetryCoordinator(inner, backend),
  237. })
  238. expect(backend.ledger()).toHaveLength(2)
  239. ctx.sessions.announce(session)
  240. expect(backend.ledger()).toHaveLength(2)
  241. })
  242. it('resumes from the handoff cursor across a reload, re-dropping mid-step chunks', async () => {
  243. const backend = new FakeBackend()
  244. const { ctx, fiber } = await setup(backend)
  245. const session = liveSession(ctx, 'hmr')
  246. session.append('turn/start', { turn: 1 })
  247. session.append('assistant/chunk', { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: 'first' } })
  248. expect(backend.ledger()).toHaveLength(2)
  249. await fiber.dispose()
  250. // The reload window: appends while no telemetry listener is registered.
  251. session.append('assistant/chunk', { turn: 1, step: 1, chunk: { type: 'text-delta', index: 0, text: 'mid-step continuation' } })
  252. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  253. const second = new FakeBackend()
  254. await ctx.plugin({
  255. name: 'fake-telemetry-2',
  256. inject: ['sessions'],
  257. apply: (inner: Context) => void new TelemetryCoordinator(inner, second),
  258. })
  259. // Only the window events past the cursor are re-handed, and the mid-step
  260. // continuation is re-dropped because ≤cursor events rebuilt the projection.
  261. expect(second.ledger().map(r => r.attributes['event.type'])).toEqual(['turn/end'])
  262. })
  263. it('replays past a record the backend rejects: one event withheld, the rest adopted', async () => {
  264. const backend = new FakeBackend()
  265. const ctx = new Context()
  266. await ctx.plugin(SessionStore)
  267. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => {})
  268. const session = liveSession(ctx, 'partial')
  269. appendTurn(session)
  270. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  271. // The backend rejects exactly the middle historical event: fail-closed
  272. // must withhold THAT record only — an adoption replay that dies on the
  273. // first contained failure would silently skip the rest of the log while
  274. // the session stays marked adopted.
  275. backend.rejectSeq = 1
  276. await ctx.plugin({
  277. name: 'fake-telemetry',
  278. inject: ['sessions'],
  279. apply: (inner: Context) => void new TelemetryCoordinator(inner, backend),
  280. })
  281. expect(backend.ledger().map(r => r.attributes['event.seq'])).toEqual([0, 2])
  282. expect(warn).toHaveBeenCalled()
  283. })
  284. it('re-hands the full log when no cursor survived (fresh session object)', async () => {
  285. const backend = new FakeBackend()
  286. const ctx = new Context()
  287. await ctx.plugin(SessionStore)
  288. const session = liveSession(ctx, 'fresh')
  289. appendTurn(session)
  290. await ctx.plugin({
  291. name: 'fake-telemetry',
  292. inject: ['sessions'],
  293. apply: (inner: Context) => void new TelemetryCoordinator(inner, backend),
  294. })
  295. expect(backend.ledger().map(r => r.attributes['event.seq'])).toEqual([0, 1])
  296. })
  297. })
  298. describe('TelemetryCoordinator lifecycle and containment', () => {
  299. it('forwards session/flush as a hint without awaiting backend work', async () => {
  300. const { ctx, backend } = await setup()
  301. const session = liveSession(ctx)
  302. let settled = false
  303. backend.flush.mockImplementation(() => {
  304. // The backend may kick off arbitrary async work; the loop's parallel must not wait for it.
  305. void new Promise(resolve => setTimeout(resolve, 50)).then(() => { settled = true })
  306. })
  307. await ctx.parallel('session/flush', session)
  308. expect(backend.flush).toHaveBeenCalledTimes(1)
  309. expect(settled).toBe(false)
  310. })
  311. it('ignores flush hints for sessions it never adopted', async () => {
  312. const { ctx, backend } = await setup()
  313. const stranger = ctx.sessions.prepare(SessionId('stranger'), { meta: {} })
  314. await ctx.parallel('session/flush', stranger)
  315. expect(backend.flush).not.toHaveBeenCalled()
  316. })
  317. it('emits no marker for a session whose announcement was vetoed before adoption', async () => {
  318. const backend = new FakeBackend()
  319. const ctx = new Context()
  320. await ctx.plugin(SessionStore)
  321. // A listener registered BEFORE the coordinator vetoes publication: the
  322. // store still emits the paired `session/disposed` for rollback, but the
  323. // coordinator never saw `session/created` — a marker for a session the
  324. // receiver saw no activity from would be noise, not signal.
  325. ctx.on('session/created', () => {
  326. throw new Error('vetoed by an earlier listener')
  327. })
  328. await ctx.plugin({
  329. name: 'fake-telemetry',
  330. inject: ['sessions'],
  331. apply: (inner: Context) => void new TelemetryCoordinator(inner, backend),
  332. })
  333. expect(() => ctx.sessions.create(SessionId('vetoed'), { meta: {} })).toThrow('vetoed')
  334. expect(backend.records.filter(r => r.channel === 'ops')).toHaveLength(0)
  335. })
  336. it('emits each adopted session’s shutdown record before awaiting backend shutdown', async () => {
  337. const { ctx, backend, fiber } = await setup()
  338. liveSession(ctx, 's1')
  339. liveSession(ctx, 's2')
  340. await fiber.dispose()
  341. expect(backend.calls).toEqual(['emit:shutdown', 'emit:shutdown', 'shutdown'])
  342. expect(backend.shutdownResolved).toBe(true)
  343. const ops = backend.records.filter(r => r.channel === 'ops')
  344. expect(ops.map(r => r.attributes['session.id']).sort()).toEqual(['s1', 's2'])
  345. expect(ops.every(r => r.attributes['telemetry.op'] === 'shutdown' && r.severity === 'info')).toBe(true)
  346. expect(ops.every(r => !('event.seq' in r.attributes) && !('event.type' in r.attributes))).toBe(true)
  347. })
  348. it('emits the shutdown marker at the session’s own disposal edge, then retires it', async () => {
  349. const { ctx, backend, fiber } = await setup()
  350. liveSession(ctx, 'survivor')
  351. // A session owned by its own fiber: disposing the fiber detaches it from
  352. // the store and emits `session/disposed` — the authoritative termination
  353. // edge. The marker must ride THAT edge (receivers classify a session with
  354. // activity and no marker as crashed, so a normally closed session in a
  355. // long-running host must not look like a crash), and the session retires
  356. // from the adopted set so unload neither retains it nor re-marks it.
  357. const owner = await ctx.plugin(Object.assign((inner: Context) => {
  358. inner.sessions.create(SessionId('ephemeral'), { meta: {} })
  359. }, { inject: ['sessions'] }))
  360. await owner.dispose()
  361. const atEdge = backend.records.filter(r => r.channel === 'ops')
  362. expect(atEdge.map(r => r.attributes['session.id'])).toEqual(['ephemeral'])
  363. expect(atEdge[0]!.attributes['telemetry.op']).toBe('shutdown')
  364. await fiber.dispose()
  365. const ops = backend.records.filter(r => r.channel === 'ops')
  366. expect(ops.map(r => r.attributes['session.id'])).toEqual(['ephemeral', 'survivor'])
  367. })
  368. it('warns instead of throwing when backend shutdown fails', async () => {
  369. const backend = new FakeBackend()
  370. backend.shutdownError = new Error('exporter unreachable')
  371. const { ctx, fiber } = await setup(backend)
  372. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => {})
  373. liveSession(ctx)
  374. await expect(fiber.dispose()).resolves.not.toThrow()
  375. expect(warn.mock.calls.some(args => String(args[0]).includes('shutdown failed'))).toBe(true)
  376. })
  377. it('contains emit failures: the append succeeds and capture heals', async () => {
  378. const { ctx, backend } = await setup()
  379. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => {})
  380. const session = liveSession(ctx)
  381. backend.emitError = new Error('backend broke')
  382. expect(() => session.append('turn/start', { turn: 1 })).not.toThrow()
  383. expect(warn).toHaveBeenCalled()
  384. backend.emitError = undefined
  385. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  386. expect(backend.ledger().map(r => r.attributes['event.type'])).toEqual(['turn/end'])
  387. })
  388. it.each([
  389. ['Error values', new TypeError('adapter exploded'), 'TypeError', 'adapter exploded'],
  390. ['non-Error values', 'plain failure', 'Error', 'plain failure'],
  391. ])('relays agent/error %s as an ops record with normalized identity', async (_label, error, name, message) => {
  392. const { ctx, backend } = await setup()
  393. const session = liveSession(ctx, 'erring')
  394. // Only the members the relay reads; the full Agent surface is irrelevant here.
  395. const agent = { id: 'agent-1', session } as Agent
  396. ctx.emit('agent/error', agent, 3, 2, error)
  397. const record = backend.records.find(r => r.channel === 'ops')!
  398. expect(record.severity).toBe('error')
  399. expect(record.attributes).toMatchObject({
  400. 'telemetry.op': 'agent-error',
  401. 'session.id': 'erring',
  402. 'agent.id': 'agent-1',
  403. 'error.name': name,
  404. turn: 3,
  405. step: 2,
  406. })
  407. expect(record.body).toEqual({ name, message })
  408. })
  409. })