telemetry.spec.ts 24 KB

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