coordinator.ts 12 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301
  1. /**
  2. * Capture coordinator for the telemetry capability. Live capture subscribes to
  3. * the session firehose plus the one live-bus relay (`agent/error`). Both
  4. * capture paths build one logical record per canonical Session event and run
  5. * each through the
  6. * `session-telemetry/record` waterfall (deployment-mounted redaction rules;
  7. * pass-through when none), then hands the result to the backend. Live capture
  8. * follows the session firehose; on-demand capture replays the canonical log
  9. * only when requested. Every synchronous handler is self-contained so a
  10. * failing backend can never starve other subscribers (cordis `emit` is
  11. * stop-on-throw) or touch the agent loop. Composed by a backend in its
  12. * constructor.
  13. *
  14. * @module @deepseek-ai/dsh-session-telemetry/coordinator
  15. */
  16. import type { Context } from '@deepseek-ai/cordis'
  17. import {
  18. SessionSeq,
  19. SessionLogOffset,
  20. type Session,
  21. type SessionEvent,
  22. type SessionSeq as SessionSeqType,
  23. type SessionSeqCursor,
  24. } from '@deepseek-ai/dsh-session'
  25. import type { Agent } from '@deepseek-ai/dsh-agent'
  26. import type { SessionTelemetrySink, SessionTelemetryRecord, SessionTelemetrySeverity } from './index.ts'
  27. /** Whether capture follows live events or reads the canonical log only when requested. */
  28. export type SessionTelemetryCapture = 'live' | 'on-demand'
  29. /** Backend-selected capture mode and history policy. */
  30. export interface SessionTelemetryCaptureOptions {
  31. /** Follow live events, or wait for explicit capture; defaults to live. */
  32. capture?: SessionTelemetryCapture
  33. /** Include stored history before this lifecycle; defaults to false. */
  34. includeHistory?: boolean
  35. }
  36. /** One record ready for backend handoff. */
  37. interface PendingRecord {
  38. readonly record: SessionTelemetryRecord
  39. /** Ledger cursor advanced only after the backend accepts this record. */
  40. readonly seq?: SessionSeqType
  41. }
  42. /**
  43. * The handoff cursor: per session, the highest `seq` handed to a backend.
  44. * Deliberately MODULE-scope ambient state — a narrow, documented exception
  45. * to the registrations-are-effects discipline: cordis has no HMR
  46. * state-handover API, and keying by the `Session` object (which belongs to
  47. * the session store and outlives any telemetry fiber) is the only in-process
  48. * lifetime that lets a re-adopting fiber resume instead of re-handing
  49. * history. Entries die with their sessions; a missing entry safely means
  50. * "re-hand everything". Advanced only at emit time — the cursor marks
  51. * handed-off, not delivered.
  52. */
  53. const handoffCursor = new WeakMap<Session, SessionSeqCursor>()
  54. /**
  55. * Install the telemetry capture side onto a context for one backend.
  56. *
  57. * Live capture registers its own `session/created` / `session/event` /
  58. * `session/disposed` listener set plus the `agent/error` relay, all through
  59. * `ctx.effect()`/`ctx.on()` on the composing fiber, and sweeps already-live
  60. * sessions (a hot reload does not replay `session/created`). A `session/disposed` captures the session's `shutdown`
  61. * operational record at its own termination edge and retires it from the
  62. * adopted set. On-demand capture registers none of those continuous listeners;
  63. * {@link captureSession} reads the canonical log explicitly and never creates
  64. * operational records. Disposal captures shutdown markers for live-adopted
  65. * sessions, then awaits the backend's `shutdown()`; a failure there warns
  66. * instead of throwing — best-effort reporting must not fail application
  67. * teardown.
  68. */
  69. export class SessionTelemetryCoordinator {
  70. /**
  71. * Sessions adopted by THIS fiber and still live, for double-adoption
  72. * protection and the teardown sweep of unmarked sessions;
  73. * `session/disposed` marks and retires entries.
  74. */
  75. private readonly adopted = new Set<Session>()
  76. /**
  77. * @param ctx - the composing backend's context; listeners bind to its fiber.
  78. * @param backend - the backend receiving records; owned elsewhere, never disposed here beyond `shutdown()` forwarding.
  79. * @param options - capture mode and history policy.
  80. */
  81. constructor(
  82. private readonly ctx: Context,
  83. private readonly backend: SessionTelemetrySink,
  84. private readonly options: SessionTelemetryCaptureOptions = {},
  85. ) {
  86. if ((options.capture ?? 'live') === 'live') {
  87. ctx.on('session/created', (session) => {
  88. this.adopt(session)
  89. })
  90. // Capture the shutdown marker at the session's own termination edge,
  91. // then retire the only strong reference owned by this coordinator.
  92. ctx.on('session/disposed', (session) => {
  93. this.contain(() => {
  94. if (!this.adopted.delete(session)) return
  95. this.deliver(session, { record: this.redact(shutdownRecord(session)) })
  96. })
  97. })
  98. ctx.on('session/event', (session, event) => {
  99. this.contain(() => {
  100. this.captureEvent(session, event)
  101. })
  102. })
  103. // Parallel listeners are awaited by the loop at turn end; returning void
  104. // (not the SDK's flush promise) is the turn-latency contract.
  105. ctx.on('session/flush', (session) => {
  106. this.contain(() => {
  107. this.hintFlush(session)
  108. })
  109. })
  110. ctx.on('agent/error', ({ agent, turn, step, error }) => {
  111. this.contain(() => {
  112. this.relayAgentError(agent, turn, step, error)
  113. })
  114. })
  115. for (const session of ctx.sessions.list()) {
  116. this.adopt(session)
  117. }
  118. }
  119. ctx.effect(() => async () => {
  120. // Sessions still adopted here are alive through whole-application
  121. // teardown, so capture the marker before the backend quiesces.
  122. for (const session of this.adopted) {
  123. this.contain(() => {
  124. this.deliver(session, { record: this.redact(shutdownRecord(session)) })
  125. })
  126. }
  127. try {
  128. await this.backend.shutdown()
  129. } catch (error) {
  130. this.ctx.logger.warn(`telemetry: backend shutdown failed: ${String(error)}`)
  131. }
  132. }, 'telemetry capture')
  133. }
  134. /**
  135. * Copy, redact, and hand over the canonical session-log suffix after the handoff
  136. * cursor, optionally stopping at an inclusive sequence boundary. Redaction
  137. * runs during this call, so an on-demand caller retains no copied records
  138. * before requesting capture and uses the policy mounted at that time.
  139. * Backend and policy failures remain contained per event and do not starve
  140. * later events in the same replay.
  141. * @param session - session whose current canonical-log prefix may be handed over.
  142. * @param throughSeq - optional last sequence included in this capture.
  143. */
  144. captureSession(session: Session, throughSeq?: SessionSeqType): void {
  145. const cursor = handoffCursor.get(session)
  146. ?? (this.options.includeHistory === true || session.firstLiveSeq === 0 ? -1 : SessionSeq(session.firstLiveSeq - 1))
  147. // Containment is PER EVENT: one rejected record is withheld fail-closed
  148. // while the rest of the historical replay proceeds.
  149. for (const event of session.snapshotEvents(SessionLogOffset(cursor + 1))) {
  150. if (throughSeq !== undefined && event.seq > throughSeq) break
  151. this.contain(() => {
  152. this.captureEvent(session, event)
  153. })
  154. }
  155. }
  156. /**
  157. * Adopt a session and replay after its handoff cursor, then follow live events.
  158. * New objects include inherited and restored history only with includeHistory;
  159. * otherwise replay starts at the constructor boundary. Re-adopting the same
  160. * object resumes after its cursor.
  161. * @param session - the live session to adopt; a second adoption is a no-op.
  162. */
  163. private adopt(session: Session): void {
  164. if (this.adopted.has(session)) return
  165. this.adopted.add(session)
  166. this.captureSession(session)
  167. }
  168. /** Copy, redact, and hand one canonical event to the backend. */
  169. private captureEvent(session: Session, event: SessionEvent): void {
  170. this.deliver(session, {
  171. record: this.redact({
  172. channel: 'ledger',
  173. time: event.time,
  174. severity: severityOf(event),
  175. attributes: identityOf(session, event),
  176. // The canonical event object is mutable and the backend serializes
  177. // later; append-time validation guarantees this clone cannot throw.
  178. body: structuredClone(event.data),
  179. }),
  180. seq: event.seq,
  181. })
  182. }
  183. /**
  184. * Run the `session-telemetry/record` waterfall at capture time. The innermost `next`
  185. * passes the record through unchanged — this package ships no rules; exported
  186. * data is as clean as the listeners a deployment mounts. Callers run inside
  187. * {@link contain}, so a throwing rule withholds the record instead of
  188. * reaching the loop (fail-closed). On-demand capture invokes this waterfall
  189. * while reading the canonical session log, not when the event was appended.
  190. */
  191. private redact(record: SessionTelemetryRecord): SessionTelemetryRecord {
  192. return this.ctx.waterfall('session-telemetry/record', record, () => record)
  193. }
  194. /** Hand one redacted record to the backend, then advance its ledger cursor. */
  195. private deliver(session: Session, pending: PendingRecord): void {
  196. this.backend.emit(pending.record)
  197. if (pending.seq !== undefined) handoffCursor.set(session, pending.seq)
  198. }
  199. /** Forward the turn-end boundary to the backend's optional flush hint. */
  200. private hintFlush(session: Session): void {
  201. if (this.adopted.has(session)) this.backend.flush?.()
  202. }
  203. /** Relay one `agent/error` bus emission as an `agent-error` operational record. */
  204. private relayAgentError(agent: Agent, turn: number, step: number, error: unknown): void {
  205. const detail = errorDetail(error)
  206. this.deliver(agent.session, {
  207. record: this.redact({
  208. channel: 'ops',
  209. time: Date.now(),
  210. severity: 'error',
  211. attributes: {
  212. 'telemetry.op': 'agent-error',
  213. 'session.id': String(agent.session.id),
  214. 'agent.id': agent.id,
  215. 'error.name': detail.name,
  216. turn,
  217. step,
  218. },
  219. body: detail,
  220. }),
  221. })
  222. }
  223. /**
  224. * Run one capture-side step with its exception contained: cordis `emit`
  225. * is stop-on-throw, so a throwing listener would starve every subscriber
  226. * registered after this plugin — nothing from the backend may escape.
  227. */
  228. private contain(step: () => void): void {
  229. try {
  230. step()
  231. } catch (error) {
  232. this.ctx.logger.warn(`telemetry: capture step failed: ${String(error)}`)
  233. }
  234. }
  235. }
  236. /**
  237. * Build the per-session clean-exit marker: emitted at the session's own
  238. * disposal edge, or at coordinator dispose for sessions still alive then.
  239. */
  240. function shutdownRecord(session: Session): SessionTelemetryRecord {
  241. return {
  242. channel: 'ops',
  243. time: Date.now(),
  244. severity: 'info',
  245. attributes: { 'telemetry.op': 'shutdown', 'session.id': String(session.id) },
  246. body: { op: 'shutdown' },
  247. }
  248. }
  249. /** Map an event's own outcome flag to the pre-baked alerting severity. */
  250. function severityOf(event: SessionEvent): SessionTelemetrySeverity {
  251. switch (event.type) {
  252. case 'tool/result':
  253. return event.data.message.content[0].isError === true ? 'error' : 'info'
  254. case 'turn/end':
  255. return event.data.reason.kind === 'error' ? 'error' : 'info'
  256. default:
  257. // Merge-extensible fall-through (no assertNever): event types this coordinator
  258. // does not depend on — including plugin-merged ones it never heard of —
  259. // pass through as info; their owners' outcome semantics stay theirs.
  260. return 'info'
  261. }
  262. }
  263. /** Normalize the live bus's arbitrary thrown value into the stable operational-record shape. */
  264. function errorDetail(error: unknown): { name: string; message: string } {
  265. const normalized = error instanceof Error ? error : new Error(String(error))
  266. return { name: normalized.name, message: normalized.message }
  267. }
  268. /** Build the minimal identity attributes: envelope plus self-contained header facts. */
  269. function identityOf(session: Session, event: SessionEvent): Record<string, string | number> {
  270. const attributes: Record<string, string | number> = {
  271. 'session.id': String(session.id),
  272. 'session.format_version': session.header.version,
  273. 'event.type': event.type,
  274. 'event.seq': event.seq,
  275. }
  276. const { cwd, parentSession, isSeeded } = session.header
  277. if (cwd !== undefined) attributes['session.cwd'] = cwd
  278. if (parentSession !== undefined) attributes['session.parent_id'] = String(parentSession)
  279. // The durable fork boundary and lineage: the child ledger is complete;
  280. // parent_id and seed_length identify which leading events were inherited.
  281. if (isSeeded) attributes['session.seed_length'] = session.inheritedEventCount
  282. return attributes
  283. }