telemetry.spec.ts 28 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680
  1. import { createAssistantMessage, 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), lifecycle-suffix replay, 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 '@deepseek-ai/cordis'
  10. import SessionStore, {
  11. SESSION_FORMAT_VERSION,
  12. Session,
  13. SessionId,
  14. SessionLogOffset,
  15. type SessionEvent,
  16. } from '@deepseek-ai/dsh-session'
  17. import type { Agent } from '@deepseek-ai/dsh-agent'
  18. import {
  19. SessionTelemetryCoordinator,
  20. type SessionTelemetrySink,
  21. type SessionTelemetryCapture,
  22. type SessionTelemetryRecord,
  23. } from '../src/index.ts'
  24. declare module '@deepseek-ai/dsh-session/types' {
  25. interface SessionEventMap {
  26. /**
  27. * Test-only merged event proving unknown types flow through unchanged.
  28. * @mode emit
  29. * @param payload - opaque test payload
  30. */
  31. 'telemetry-test/opaque': { payload: { nested: string[] } }
  32. }
  33. }
  34. class FakeBackend implements SessionTelemetrySink {
  35. records: SessionTelemetryRecord[] = []
  36. calls: string[] = []
  37. emitError: Error | undefined
  38. rejectSeq: number | undefined
  39. shutdownError: Error | undefined
  40. shutdownResolved = false
  41. emit(record: SessionTelemetryRecord): void {
  42. if (this.emitError) throw this.emitError
  43. if (this.rejectSeq !== undefined && record.attributes['event.seq'] === this.rejectSeq) {
  44. throw new Error(`backend rejected seq ${this.rejectSeq}`)
  45. }
  46. this.records.push(record)
  47. this.calls.push(`emit:${String(record.attributes['event.seq'] ?? record.attributes['telemetry.op'])}`)
  48. }
  49. flush = vi.fn()
  50. async shutdown(): Promise<void> {
  51. this.calls.push('shutdown')
  52. await new Promise(resolve => setTimeout(resolve, 5))
  53. if (this.shutdownError) throw this.shutdownError
  54. this.shutdownResolved = true
  55. }
  56. ledger(): SessionTelemetryRecord[] {
  57. return this.records.filter(r => r.channel === 'ledger')
  58. }
  59. }
  60. async function setup(
  61. backend: FakeBackend = new FakeBackend(),
  62. capture: SessionTelemetryCapture = 'live',
  63. ) {
  64. const ctx = new Context()
  65. await ctx.plugin(SessionStore)
  66. let coordinator!: SessionTelemetryCoordinator
  67. const fiber = await ctx.plugin({
  68. name: 'fake-telemetry',
  69. inject: ['sessions'],
  70. apply: (inner: Context) => {
  71. coordinator = new SessionTelemetryCoordinator(inner, backend, { capture })
  72. },
  73. })
  74. return { ctx, backend, coordinator, fiber }
  75. }
  76. function liveSession(ctx: Context, id = `s-${Math.random().toString(36).slice(2)}`): Session {
  77. return ctx.sessions.create(SessionId(id), { meta: {} })
  78. }
  79. function appendTurn(session: Session): void {
  80. session.append('turn/start', { turn: 1 })
  81. session.append('user/message', createUserMessage({
  82. content: [{ type: 'text', text: 'hello' }], source: { kind: 'user' },
  83. }), { surfaceOp: 'append' })
  84. }
  85. function appendAssistantMessage(
  86. session: Session,
  87. turn: number,
  88. step: number,
  89. texts: readonly string[],
  90. time0 = 100,
  91. ): void {
  92. session.append('assistant/message', {
  93. turn,
  94. step,
  95. message: createAssistantMessage({
  96. content: [{ type: 'text', text: texts.join('') }],
  97. source: { provider: 'mock', model: 'mock' },
  98. }),
  99. stream: [
  100. { type: 'text-chunks', time0, index: 0, dt: texts.slice(1).map(() => 5), texts: [...texts] },
  101. { type: 'chunk', time: time0 + Math.max(0, texts.length - 1) * 5, chunk: { type: 'finish', reason: { kind: 'stop' } } },
  102. ],
  103. }, { surfaceOp: 'append' })
  104. }
  105. describe('SessionTelemetryCoordinator capture', () => {
  106. it('hands every appended event over with envelope identity and cloned body', async () => {
  107. const { ctx, backend } = await setup()
  108. const session = liveSession(ctx, 'cap')
  109. appendTurn(session)
  110. const start = backend.ledger()[0]!
  111. const message = backend.ledger()[1]!
  112. expect(start.attributes).toMatchObject({
  113. 'session.id': 'cap',
  114. 'session.format_version': SESSION_FORMAT_VERSION,
  115. 'event.type': 'turn/start',
  116. 'event.seq': 0,
  117. })
  118. expect(start.time).toBe(session.snapshotEvents()[0]!.time)
  119. expect(start.severity).toBe('info')
  120. expect(message.attributes['event.seq']).toBe(1)
  121. // Deep-copy isolation: mutating the handed-off body never reaches the log.
  122. ;(message.body as { content: { text: string }[] }).content[0]!.text = 'tampered'
  123. const logged = session.snapshotEvents()[1] as SessionEvent<'user/message'>
  124. expect(logged.data.content[0]).toMatchObject({ text: 'hello' })
  125. })
  126. it('captures a live request header without replaying previously withheld events', async () => {
  127. const { ctx, backend } = await setup()
  128. try {
  129. const session = liveSession(ctx, 'live-header')
  130. const disposeRule = ctx.on('session-telemetry/record', () => {
  131. throw new Error('withheld')
  132. })
  133. session.append('turn/start', { turn: 1 })
  134. disposeRule()
  135. session.append('request/header', {
  136. header: { config: { provider: 'mock', model: 'mock' } }, reason: 'initial',
  137. })
  138. expect(backend.ledger().map(record => record.attributes['event.type'])).toEqual(['request/header'])
  139. } finally {
  140. await ctx.fiber.dispose()
  141. }
  142. })
  143. it('stamps header facts on every record when present', async () => {
  144. const { ctx, backend } = await setup()
  145. const parent = SessionId('parent')
  146. const session = ctx.sessions.create(SessionId('child'), { meta: { cwd: '/tmp/proj', parentSession: parent } })
  147. appendTurn(session)
  148. for (const record of backend.ledger()) {
  149. expect(record.attributes['session.format_version']).toBe(SESSION_FORMAT_VERSION)
  150. expect(record.attributes['session.cwd']).toBe('/tmp/proj')
  151. expect(record.attributes['session.parent_id']).toBe('parent')
  152. }
  153. })
  154. it('maps outcome flags to severity, unknown types falling through as info', async () => {
  155. const { ctx, backend } = await setup()
  156. const session = liveSession(ctx)
  157. session.append('turn/start', { turn: 1 })
  158. session.append('tool/result', {
  159. turn: 1, step: 1,
  160. message: createToolResultMessage({
  161. callId: 'c1' as never,
  162. content: [],
  163. isError: true,
  164. }),
  165. }, { surfaceOp: 'append' })
  166. session.append('tool/result', {
  167. turn: 1, step: 1,
  168. message: createToolResultMessage({
  169. callId: 'c2' as never,
  170. content: [],
  171. isError: false,
  172. }),
  173. }, { surfaceOp: 'append' })
  174. session.append('telemetry-test/opaque', { payload: { nested: [] } })
  175. session.append('turn/end', { turn: 1, reason: { kind: 'error', error: { message: 'boom', code: 'UNKNOWN' } } })
  176. const severities = backend.ledger().map(r => [r.attributes['event.type'], r.severity])
  177. expect(severities).toEqual([
  178. ['turn/start', 'info'],
  179. ['tool/result', 'error'],
  180. ['tool/result', 'info'],
  181. ['telemetry-test/opaque', 'info'],
  182. ['turn/end', 'error'],
  183. ])
  184. })
  185. it('passes unknown merged event types through unchanged', async () => {
  186. const { ctx, backend } = await setup()
  187. const session = liveSession(ctx)
  188. session.append('telemetry-test/opaque', { payload: { nested: ['a', 'b'] } })
  189. const record = backend.ledger()[0]!
  190. expect(record.attributes['event.type']).toBe('telemetry-test/opaque')
  191. expect(record.severity).toBe('info')
  192. expect(record.body).toEqual({ payload: { nested: ['a', 'b'] } })
  193. })
  194. it('ships every assistant stream in canonical order with its complete body', async () => {
  195. const { ctx, backend } = await setup()
  196. const a = liveSession(ctx, 'a')
  197. const b = liveSession(ctx, 'b')
  198. appendAssistantMessage(a, 1, 1, ['a11-first', 'a11-second'], 110)
  199. appendAssistantMessage(a, 1, 2, ['a12-first'], 120)
  200. appendAssistantMessage(b, 1, 1, ['b11-first', 'b11-second'], 210)
  201. const shipped = backend.ledger().map(r => [
  202. r.attributes['session.id'],
  203. r.attributes['event.seq'],
  204. (r.body as SessionEvent<'assistant/message'>['data']).stream,
  205. ])
  206. expect(shipped).toEqual([
  207. ['a', 0, [
  208. { type: 'text-chunks', time0: 110, index: 0, dt: [5], texts: ['a11-first', 'a11-second'] },
  209. { type: 'chunk', time: 115, chunk: { type: 'finish', reason: { kind: 'stop' } } },
  210. ]],
  211. ['a', 1, [
  212. { type: 'text-chunks', time0: 120, index: 0, dt: [], texts: ['a12-first'] },
  213. { type: 'chunk', time: 120, chunk: { type: 'finish', reason: { kind: 'stop' } } },
  214. ]],
  215. ['b', 0, [
  216. { type: 'text-chunks', time0: 210, index: 0, dt: [5], texts: ['b11-first', 'b11-second'] },
  217. { type: 'chunk', time: 215, chunk: { type: 'finish', reason: { kind: 'stop' } } },
  218. ]],
  219. ])
  220. })
  221. })
  222. describe('SessionTelemetryCoordinator on-demand capture', () => {
  223. it('captures one canonical-log prefix at a time without following later events', async () => {
  224. const { ctx, backend, coordinator } = await setup(new FakeBackend(), 'on-demand')
  225. const session = liveSession(ctx, 'on-demand-prefix')
  226. appendTurn(session)
  227. appendAssistantMessage(session, 1, 1, ['first'], 100)
  228. const firstBoundary = session.snapshotEvents()[2]!.seq
  229. appendAssistantMessage(session, 1, 2, ['second'], 200)
  230. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  231. expect(backend.records).toEqual([])
  232. coordinator.captureSession(session, firstBoundary)
  233. expect(backend.ledger().map(record => record.attributes['event.type'])).toEqual([
  234. 'turn/start',
  235. 'user/message',
  236. 'assistant/message',
  237. ])
  238. expect(backend.ledger().map(record => record.attributes['event.seq'])).toEqual([0, 1, 2])
  239. expect(backend.ledger()[2]!.body).toMatchObject({
  240. stream: [
  241. { type: 'text-chunks', time0: 100, index: 0, dt: [], texts: ['first'] },
  242. { type: 'chunk', time: 100, chunk: { type: 'finish', reason: { kind: 'stop' } } },
  243. ],
  244. })
  245. expect(backend.ledger()).toHaveLength(3)
  246. coordinator.captureSession(session)
  247. coordinator.captureSession(session)
  248. expect(backend.ledger().map(record => record.attributes['event.type'])).toEqual([
  249. 'turn/start',
  250. 'user/message',
  251. 'assistant/message',
  252. 'assistant/message',
  253. 'turn/end',
  254. ])
  255. expect(backend.ledger().map(record => record.attributes['event.seq'])).toEqual([0, 1, 2, 3, 4])
  256. expect(backend.ledger()[3]!.body).toMatchObject({
  257. stream: [
  258. { type: 'text-chunks', time0: 200, index: 0, dt: [], texts: ['second'] },
  259. { type: 'chunk', time: 200, chunk: { type: 'finish', reason: { kind: 'stop' } } },
  260. ],
  261. })
  262. })
  263. it('includes inherited history only when explicitly requested, through the exact sequence', async () => {
  264. const ctx = new Context()
  265. const backend = new FakeBackend()
  266. try {
  267. await ctx.plugin(SessionStore)
  268. let coordinator!: SessionTelemetryCoordinator
  269. await ctx.plugin({
  270. name: 'history-telemetry',
  271. inject: ['sessions'],
  272. apply: (inner: Context) => {
  273. coordinator = new SessionTelemetryCoordinator(inner, backend, {
  274. capture: 'on-demand', includeHistory: true,
  275. })
  276. },
  277. })
  278. const parent = liveSession(ctx, 'history-parent')
  279. appendTurn(parent)
  280. const child = ctx.sessions.create(SessionId('history-child'), { seed: [...parent.snapshotEvents()] })
  281. const boundary = child.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  282. child.append('turn/start', { turn: 2 })
  283. expect(backend.records).toEqual([])
  284. coordinator.captureSession(child, boundary.seq)
  285. coordinator.captureSession(child, boundary.seq)
  286. expect(backend.ledger().map(record => record.attributes['event.seq'])).toEqual([0, 1, 2, 3])
  287. expect(backend.ledger().at(-1)?.attributes['event.seq']).toBe(boundary.seq)
  288. expect(backend.ledger().every(record => record.attributes['session.id'] === child.id)).toBe(true)
  289. } finally {
  290. await ctx.fiber.dispose()
  291. }
  292. })
  293. it('runs the currently mounted redaction policy during canonical-log capture', async () => {
  294. const { ctx, backend, coordinator } = await setup(new FakeBackend(), 'on-demand')
  295. const session = liveSession(ctx, 'on-demand-redacted')
  296. session.append('turn/start', { turn: 1 })
  297. const disposeRule = ctx.on('session-telemetry/record', (_record, next) => ({
  298. ...next(),
  299. body: { scrubbed: true },
  300. }))
  301. coordinator.captureSession(session)
  302. expect(backend.ledger()[0]!.body).toEqual({ scrubbed: true })
  303. disposeRule()
  304. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  305. coordinator.captureSession(session)
  306. expect(backend.ledger()[1]!.body).toEqual({ turn: 1, reason: { kind: 'completed' } })
  307. })
  308. it('contains each backend failure independently while replaying a prefix', async () => {
  309. const backend = new FakeBackend()
  310. backend.rejectSeq = 1
  311. const { ctx, coordinator } = await setup(backend, 'on-demand')
  312. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => {})
  313. const session = liveSession(ctx, 'on-demand-failure')
  314. appendTurn(session)
  315. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  316. coordinator.captureSession(session)
  317. expect(backend.ledger().map(record => record.attributes['event.seq'])).toEqual([0, 2])
  318. expect(warn).toHaveBeenCalled()
  319. })
  320. it('captures a pending prefix after coordinator reload without retained records', async () => {
  321. const first = new FakeBackend()
  322. const { ctx, fiber } = await setup(first, 'on-demand')
  323. const session = liveSession(ctx, 'on-demand-reload')
  324. session.append('turn/start', { turn: 1 })
  325. await fiber.dispose()
  326. expect(first.records).toEqual([])
  327. const second = new FakeBackend()
  328. let coordinator!: SessionTelemetryCoordinator
  329. await ctx.plugin({
  330. name: 'fake-telemetry-after-on-demand-reload',
  331. inject: ['sessions'],
  332. apply: (inner: Context) => {
  333. coordinator = new SessionTelemetryCoordinator(inner, second, { capture: 'on-demand' })
  334. },
  335. })
  336. coordinator.captureSession(session)
  337. expect(second.ledger().map(record => record.attributes['event.seq'])).toEqual([0])
  338. })
  339. it('registers no continuous capture, flush, or ops listeners', async () => {
  340. const { ctx, backend, coordinator, fiber } = await setup(new FakeBackend(), 'on-demand')
  341. const redact = vi.fn((_record: SessionTelemetryRecord, next: () => SessionTelemetryRecord) => next())
  342. ctx.on('session-telemetry/record', redact)
  343. const session = liveSession(ctx, 'on-demand-ledger-only')
  344. session.append('turn/start', { turn: 1 })
  345. await ctx.parallel('session/flush', session)
  346. const agent = { id: 'agent-1', session } as Agent
  347. ctx.emit('agent/error', { agent, turn: 1, step: 1, error: new Error('local only') })
  348. expect(backend.flush).not.toHaveBeenCalled()
  349. expect(backend.records).toEqual([])
  350. expect(redact).not.toHaveBeenCalled()
  351. coordinator.captureSession(session)
  352. expect(redact).toHaveBeenCalledTimes(1)
  353. await fiber.dispose()
  354. expect(backend.records.map(record => record.channel)).toEqual(['ledger'])
  355. })
  356. })
  357. describe('SessionTelemetryCoordinator adoption', () => {
  358. it('replays a new fork object from its constructor boundary without its inherited prefix', async () => {
  359. const backend = new FakeBackend()
  360. const ctx = new Context()
  361. await ctx.plugin(SessionStore)
  362. const parent = liveSession(ctx, 'seed-parent')
  363. appendTurn(parent)
  364. await ctx.plugin({
  365. name: 'fake-telemetry',
  366. inject: ['sessions'],
  367. apply: (inner: Context) => void new SessionTelemetryCoordinator(inner, backend),
  368. })
  369. const child = ctx.sessions.prepare(SessionId('seeded'), { seed: [...parent.snapshotEvents()], meta: {} })
  370. child.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  371. ctx.sessions.enter(child)
  372. ctx.sessions.announce(child)
  373. const seqs = backend.ledger().map(r => [r.attributes['session.id'], r.attributes['event.seq']])
  374. expect(seqs).toEqual(expect.arrayContaining([['seed-parent', 0], ['seed-parent', 1]]))
  375. expect(seqs.filter(([id]) => id === 'seeded')).toEqual([
  376. ['seeded', 2],
  377. ['seeded', 3],
  378. ])
  379. })
  380. it('replays a restored post-migration Session from its constructor boundary', async () => {
  381. const backend = new FakeBackend()
  382. const ctx = new Context()
  383. await ctx.plugin(SessionStore)
  384. const donor = Session.create(SessionId('donor'))
  385. donor.append('turn/start', { turn: 1 })
  386. appendAssistantMessage(donor, 1, 1, ['first'], 100)
  387. await ctx.plugin({
  388. name: 'fake-telemetry',
  389. inject: ['sessions'],
  390. apply: (inner: Context) => void new SessionTelemetryCoordinator(inner, backend),
  391. })
  392. // Session persistence migrates before it constructs the restored Session;
  393. // telemetry therefore receives a current-format object with the complete
  394. // migrated canonical seed.
  395. const resumed = ctx.sessions.prepare(SessionId('resumed'), {
  396. seed: structuredClone(donor.snapshotEvents()) as SessionEvent[],
  397. meta: {
  398. version: SESSION_FORMAT_VERSION,
  399. id: SessionId('resumed'),
  400. createdAt: 1,
  401. isSeeded: false,
  402. },
  403. inheritedEventCount: SessionLogOffset(0),
  404. eventState: 'detached',
  405. })
  406. ctx.sessions.enter(resumed)
  407. ctx.sessions.announce(resumed)
  408. const ofResumed = () => backend.ledger()
  409. .filter(r => r.attributes['session.id'] === 'resumed')
  410. expect(ofResumed().map(r => r.attributes['event.seq'])).toEqual([2])
  411. expect(ofResumed().every(r => r.attributes['session.format_version'] === SESSION_FORMAT_VERSION)).toBe(true)
  412. appendAssistantMessage(resumed, 1, 1, ['continuation'], 200)
  413. appendAssistantMessage(resumed, 1, 2, ['next step'], 300)
  414. expect(ofResumed().map(r => r.attributes['event.seq'])).toEqual([2, 3, 4])
  415. expect(ofResumed().map(r => (r.body as { stream?: { type: string; texts?: string[] }[] }).stream?.[0]?.texts)).toEqual([
  416. undefined,
  417. ['continuation'],
  418. ['next step'],
  419. ])
  420. })
  421. it('stamps session.seed_length from the exact Session cut so receivers can stitch fork streams', async () => {
  422. const backend = new FakeBackend()
  423. const ctx = new Context()
  424. await ctx.plugin(SessionStore)
  425. const parent = liveSession(ctx, 'stitch-parent')
  426. appendTurn(parent)
  427. const child = ctx.sessions.create(SessionId('stitch-child'), {
  428. seed: [...parent.snapshotEvents()],
  429. inheritedEventCount: SessionLogOffset(parent.snapshotEvents().length),
  430. meta: { parentSession: SessionId('stitch-parent'), isSeeded: true },
  431. })
  432. await ctx.plugin({
  433. name: 'fake-telemetry',
  434. inject: ['sessions'],
  435. apply: (inner: Context) => void new SessionTelemetryCoordinator(inner, backend),
  436. })
  437. child.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  438. const record = backend.ledger().find(r => r.attributes['session.id'] === 'stitch-child')!
  439. expect(record.attributes['session.parent_id']).toBe('stitch-parent')
  440. expect(record.attributes['session.seed_length']).toBe(2)
  441. })
  442. it('adopts exactly once when created fires after the sweep', async () => {
  443. const backend = new FakeBackend()
  444. const ctx = new Context()
  445. await ctx.plugin(SessionStore)
  446. // The enter/announce window: prepare+enter puts the session in the store
  447. // (visible to the constructor sweep) before `session/created` fires, so a
  448. // coordinator loaded inside that window sees the session twice — sweep
  449. // first, created second. The second adoption must be a no-op.
  450. const session = ctx.sessions.prepare(SessionId('overlap'))
  451. appendTurn(session)
  452. ctx.sessions.enter(session)
  453. await ctx.plugin({
  454. name: 'fake-telemetry',
  455. inject: ['sessions'],
  456. apply: (inner: Context) => void new SessionTelemetryCoordinator(inner, backend),
  457. })
  458. expect(backend.ledger()).toHaveLength(2)
  459. ctx.sessions.announce(session)
  460. expect(backend.ledger()).toHaveLength(2)
  461. })
  462. it('resumes from the handoff cursor across same-object re-adoption without duplicates', async () => {
  463. const backend = new FakeBackend()
  464. const { ctx, fiber } = await setup(backend)
  465. const session = liveSession(ctx, 'hmr')
  466. session.append('turn/start', { turn: 1 })
  467. appendAssistantMessage(session, 1, 1, ['first'], 100)
  468. expect(backend.ledger()).toHaveLength(2)
  469. await fiber.dispose()
  470. // The reload window: appends while no telemetry listener is registered.
  471. appendAssistantMessage(session, 1, 2, ['mid-step continuation'], 200)
  472. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  473. const second = new FakeBackend()
  474. await ctx.plugin({
  475. name: 'fake-telemetry-2',
  476. inject: ['sessions'],
  477. apply: (inner: Context) => void new SessionTelemetryCoordinator(inner, second),
  478. })
  479. // Only window events past the same object's cursor are re-handed.
  480. expect(second.ledger().map(r => [r.attributes['event.seq'], r.attributes['event.type']])).toEqual([
  481. [2, 'assistant/message'],
  482. [3, 'turn/end'],
  483. ])
  484. })
  485. it('replays past a record the backend rejects: one event withheld, the rest adopted', async () => {
  486. const backend = new FakeBackend()
  487. const ctx = new Context()
  488. await ctx.plugin(SessionStore)
  489. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => {})
  490. const session = liveSession(ctx, 'partial')
  491. appendTurn(session)
  492. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  493. // The backend rejects exactly the middle historical event: fail-closed
  494. // must withhold THAT record only — an adoption replay that dies on the
  495. // first contained failure would silently skip the rest of the log while
  496. // the session stays marked adopted.
  497. backend.rejectSeq = 1
  498. await ctx.plugin({
  499. name: 'fake-telemetry',
  500. inject: ['sessions'],
  501. apply: (inner: Context) => void new SessionTelemetryCoordinator(inner, backend),
  502. })
  503. expect(backend.ledger().map(r => r.attributes['event.seq'])).toEqual([0, 2])
  504. expect(warn).toHaveBeenCalled()
  505. })
  506. it('re-hands the full log when no cursor survived (fresh session object)', async () => {
  507. const backend = new FakeBackend()
  508. const ctx = new Context()
  509. await ctx.plugin(SessionStore)
  510. const session = liveSession(ctx, 'fresh')
  511. appendTurn(session)
  512. await ctx.plugin({
  513. name: 'fake-telemetry',
  514. inject: ['sessions'],
  515. apply: (inner: Context) => void new SessionTelemetryCoordinator(inner, backend),
  516. })
  517. expect(backend.ledger().map(r => r.attributes['event.seq'])).toEqual([0, 1])
  518. })
  519. })
  520. describe('SessionTelemetryCoordinator lifecycle and containment', () => {
  521. it('forwards session/flush as a hint without awaiting backend work', async () => {
  522. const { ctx, backend } = await setup()
  523. const session = liveSession(ctx)
  524. let settled = false
  525. backend.flush.mockImplementation(() => {
  526. // The backend may kick off arbitrary async work; the loop's parallel must not wait for it.
  527. void new Promise(resolve => setTimeout(resolve, 50)).then(() => { settled = true })
  528. })
  529. await ctx.parallel('session/flush', session)
  530. expect(backend.flush).toHaveBeenCalledTimes(1)
  531. expect(settled).toBe(false)
  532. })
  533. it('ignores flush hints for sessions it never adopted', async () => {
  534. const { ctx, backend } = await setup()
  535. const stranger = ctx.sessions.prepare(SessionId('stranger'), { meta: {} })
  536. await ctx.parallel('session/flush', stranger)
  537. expect(backend.flush).not.toHaveBeenCalled()
  538. })
  539. it('emits no marker for a session whose announcement was vetoed before adoption', async () => {
  540. const backend = new FakeBackend()
  541. const ctx = new Context()
  542. await ctx.plugin(SessionStore)
  543. // A listener registered BEFORE the coordinator vetoes publication: the
  544. // store still emits the paired `session/disposed` for rollback, but the
  545. // coordinator never saw `session/created` — a marker for a session the
  546. // receiver saw no activity from would be noise, not signal.
  547. ctx.on('session/created', () => {
  548. throw new Error('vetoed by an earlier listener')
  549. })
  550. await ctx.plugin({
  551. name: 'fake-telemetry',
  552. inject: ['sessions'],
  553. apply: (inner: Context) => void new SessionTelemetryCoordinator(inner, backend),
  554. })
  555. expect(() => ctx.sessions.create(SessionId('vetoed'), { meta: {} })).toThrow('vetoed')
  556. expect(backend.records.filter(r => r.channel === 'ops')).toHaveLength(0)
  557. })
  558. it('emits each adopted session’s shutdown record before awaiting backend shutdown', async () => {
  559. const { ctx, backend, fiber } = await setup()
  560. liveSession(ctx, 's1')
  561. liveSession(ctx, 's2')
  562. await fiber.dispose()
  563. expect(backend.calls).toEqual(['emit:shutdown', 'emit:shutdown', 'shutdown'])
  564. expect(backend.shutdownResolved).toBe(true)
  565. const ops = backend.records.filter(r => r.channel === 'ops')
  566. expect(ops.map(r => r.attributes['session.id']).sort()).toEqual(['s1', 's2'])
  567. expect(ops.every(r => r.attributes['telemetry.op'] === 'shutdown' && r.severity === 'info')).toBe(true)
  568. expect(ops.every(r => !('event.seq' in r.attributes) && !('event.type' in r.attributes))).toBe(true)
  569. })
  570. it('emits the shutdown marker at the session’s own disposal edge, then retires it', async () => {
  571. const { ctx, backend, fiber } = await setup()
  572. liveSession(ctx, 'survivor')
  573. // A session owned by its own fiber: disposing the fiber detaches it from
  574. // the store and emits `session/disposed` — the authoritative termination
  575. // edge. The marker must ride THAT edge (receivers classify a session with
  576. // activity and no marker as crashed, so a normally closed session in a
  577. // long-running host must not look like a crash), and the session retires
  578. // from the adopted set so unload neither retains it nor re-marks it.
  579. const owner = await ctx.plugin(Object.assign((inner: Context) => {
  580. inner.sessions.create(SessionId('ephemeral'), { meta: {} })
  581. }, { inject: ['sessions'] }))
  582. await owner.dispose()
  583. const atEdge = backend.records.filter(r => r.channel === 'ops')
  584. expect(atEdge.map(r => r.attributes['session.id'])).toEqual(['ephemeral'])
  585. expect(atEdge[0]!.attributes['telemetry.op']).toBe('shutdown')
  586. await fiber.dispose()
  587. const ops = backend.records.filter(r => r.channel === 'ops')
  588. expect(ops.map(r => r.attributes['session.id'])).toEqual(['ephemeral', 'survivor'])
  589. })
  590. it('warns instead of throwing when backend shutdown fails', async () => {
  591. const backend = new FakeBackend()
  592. backend.shutdownError = new Error('exporter unreachable')
  593. const { ctx, fiber } = await setup(backend)
  594. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => {})
  595. liveSession(ctx)
  596. await expect(fiber.dispose()).resolves.not.toThrow()
  597. expect(warn.mock.calls.some(args => String(args[0]).includes('shutdown failed'))).toBe(true)
  598. })
  599. it('contains emit failures: the append succeeds and capture heals', async () => {
  600. const { ctx, backend } = await setup()
  601. const warn = vi.spyOn(ctx.logger, 'warn').mockImplementation(() => {})
  602. const session = liveSession(ctx)
  603. backend.emitError = new Error('backend broke')
  604. expect(() => session.append('turn/start', { turn: 1 })).not.toThrow()
  605. expect(warn).toHaveBeenCalled()
  606. backend.emitError = undefined
  607. session.append('turn/end', { turn: 1, reason: { kind: 'completed' } })
  608. expect(backend.ledger().map(r => r.attributes['event.type'])).toEqual(['turn/end'])
  609. })
  610. it.each([
  611. ['Error values', new TypeError('adapter exploded'), 'TypeError', 'adapter exploded'],
  612. ['non-Error values', 'plain failure', 'Error', 'plain failure'],
  613. ])('relays agent/error %s as an ops record with normalized identity', async (_label, error, name, message) => {
  614. const { ctx, backend } = await setup()
  615. const session = liveSession(ctx, 'erring')
  616. // Only the members the relay reads; the full Agent surface is irrelevant here.
  617. const agent = { id: 'agent-1', session } as Agent
  618. ctx.emit('agent/error', { agent, turn: 3, step: 2, error })
  619. const record = backend.records.find(r => r.channel === 'ops')!
  620. expect(record.severity).toBe('error')
  621. expect(record.attributes).toMatchObject({
  622. 'telemetry.op': 'agent-error',
  623. 'session.id': 'erring',
  624. 'agent.id': 'agent-1',
  625. 'error.name': name,
  626. turn: 3,
  627. step: 2,
  628. })
  629. expect(record.body).toEqual({ name, message })
  630. })
  631. })