index.ts 9.8 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213
  1. /**
  2. * OpenTelemetry backend for the DeepSeek Harness telemetry seam.
  3. *
  4. * Composes the OTel JS SDK as-is — a `LoggerProvider` with a
  5. * `BatchLogRecordProcessor` and an OTLP/HTTP log exporter — and maps each
  6. * record handed over by the seam onto `logger.emit()`. Per the seam's
  7. * boundary axiom, everything downstream of that call (batching, retry,
  8. * queueing, loss policy) is the SDK's documented behavior, configured
  9. * verbatim through the `exporter`/`processor` passthroughs. The one
  10. * backend-owned policy is an outer shutdown deadline: the SDK's export
  11. * timeout does not bound its preceding `forceFlush()` wait.
  12. *
  13. * @module @deepseek-ai/dsh-session-telemetry-otel
  14. */
  15. import { createRequire } from 'node:module'
  16. import z from 'schemastery'
  17. import type { Context } from 'cordis'
  18. import { Telemetry, TelemetryCoordinator, type TelemetryRecord, type TelemetrySeverity } from '@deepseek-ai/dsh-session-telemetry'
  19. import { APP_IDENTITY } from '@deepseek-ai/dsh-llm'
  20. import { getOrCreateAnonymousUserId } from './user-id.ts'
  21. import {
  22. BatchLogRecordProcessor,
  23. LoggerProvider,
  24. type BatchLogRecordProcessorOptions,
  25. } from '@opentelemetry/sdk-logs'
  26. import { OTLPLogExporter } from '@opentelemetry/exporter-logs-otlp-http'
  27. import type { OTLPExporterNodeConfigBase } from '@opentelemetry/otlp-exporter-base'
  28. import { SeverityNumber, type AnyValue, type Logger } from '@opentelemetry/api-logs'
  29. import { resourceFromAttributes } from '@opentelemetry/resources'
  30. // The package's own manifest is the single source of the instrumentation-scope
  31. // version (same pattern as dsh-llm's attribution identity).
  32. const { version } = createRequire(import.meta.url)('../package.json') as { version: string }
  33. /**
  34. * Plugin configuration: two verbatim SDK option shapes plus one DSH-owned
  35. * shutdown bound. The package validates its endpoint and shutdown deadline
  36. * because both must fail at plugin load rather than at first export or exit.
  37. */
  38. export interface Config {
  39. /**
  40. * Passed verbatim to the SDK's OTLP/HTTP log exporter — the complete
  41. * `OTLPExporterNodeConfigBase` shape (`headers`, `timeoutMillis`,
  42. * `compression`, `keepAlive`, …), owned and documented by the SDK. `url`
  43. * is the one field this package requires and validates itself.
  44. */
  45. exporter?: OTLPExporterNodeConfigBase & {
  46. /** Full logs endpoint (e.g. `https://collector.example.com/v1/logs`). Required; validated at plugin load. */
  47. url?: string
  48. }
  49. /**
  50. * Passed verbatim to `BatchLogRecordProcessor` (minus the exporter slot,
  51. * which this plugin fills); the SDK owns and documents these knobs.
  52. */
  53. processor?: Omit<BatchLogRecordProcessorOptions, 'exporter'>
  54. /** Maximum time spent awaiting the SDK provider's complete shutdown path. */
  55. shutdownTimeoutMillis?: number
  56. }
  57. /**
  58. * Schemastery validator for {@link Config}; cordis runs it before the plugin
  59. * starts. Shape-level only — load-bearing value checks live in the constructor
  60. * so their errors name the fields. Both SDK slots are opaque passthroughs:
  61. * the SDK owns their shapes and validates its own options;
  62. * re-declaring them field-by-field here would violate the boundary axiom
  63. * (and silently drop every field not re-declared).
  64. */
  65. export const Config: z<Config> = z.object({
  66. exporter: z.any(),
  67. processor: z.any(),
  68. shutdownTimeoutMillis: z.number(),
  69. })
  70. /** Default outer allowance for the SDK's complete shutdown sequence. */
  71. export const DEFAULT_SHUTDOWN_TIMEOUT_MILLIS = 3_000
  72. // Node clamps larger timer delays to one millisecond. This is a runtime
  73. // protocol limit, not a deployment default.
  74. const MAX_TIMER_DELAY_MILLIS = 2_147_483_647
  75. /** Severity mapping from the seam's three-level vocabulary to OTel severity numbers. */
  76. const SEVERITY: Record<TelemetrySeverity, { severityNumber: SeverityNumber; severityText: string }> = {
  77. info: { severityNumber: SeverityNumber.INFO, severityText: 'INFO' },
  78. warn: { severityNumber: SeverityNumber.WARN, severityText: 'WARN' },
  79. error: { severityNumber: SeverityNumber.ERROR, severityText: 'ERROR' },
  80. }
  81. /**
  82. * The backend plugin — the only entry a deployment loads. Constructing it
  83. * wires the SDK pipeline, registers the `telemetry` service (duplicate load
  84. * throws, cordis' standard duplicate-service behavior), and composes the
  85. * seam's {@link TelemetryCoordinator}, which installs the capture side onto
  86. * this fiber.
  87. */
  88. export class TelemetryOtel extends Telemetry {
  89. static inject = ['sessions']
  90. static Config = Config
  91. private readonly provider: LoggerProvider
  92. private readonly ledger: Logger
  93. private readonly ops: Logger
  94. private readonly shutdownTimeoutMillis: number
  95. constructor(ctx: Context, config: Config) {
  96. super(ctx)
  97. const url = config.exporter?.url
  98. if (url === undefined || url.length === 0) {
  99. throw new Error('session-telemetry-otel: exporter.url is required (the full OTLP logs endpoint)')
  100. }
  101. let parsed: URL
  102. try {
  103. parsed = new URL(url)
  104. } catch {
  105. // Re-thrown as a config error: the only way here is a malformed url string.
  106. throw new Error(`session-telemetry-otel: exporter.url is not a valid URL: ${JSON.stringify(url)}`)
  107. }
  108. if (parsed.protocol !== 'http:' && parsed.protocol !== 'https:') {
  109. throw new Error(`session-telemetry-otel: exporter.url must be http(s), got ${parsed.protocol}`)
  110. }
  111. // The one processor field checked beyond the SDK's own validation: the
  112. // SDK accepts a non-positive batch size, but its shutdown drain then
  113. // splices empty batches without consuming the queue — dispose would hang
  114. // forever with records queued. Misconfiguration fails at load instead.
  115. const batchSize = config.processor?.maxExportBatchSize
  116. if (batchSize !== undefined && (!Number.isInteger(batchSize) || batchSize < 1)) {
  117. throw new Error(`session-telemetry-otel: processor.maxExportBatchSize must be a positive integer, got ${String(batchSize)}`)
  118. }
  119. const shutdownTimeoutMillis = config.shutdownTimeoutMillis ?? DEFAULT_SHUTDOWN_TIMEOUT_MILLIS
  120. if (!Number.isFinite(shutdownTimeoutMillis) || shutdownTimeoutMillis <= 0 || shutdownTimeoutMillis > MAX_TIMER_DELAY_MILLIS) {
  121. throw new Error(`session-telemetry-otel: shutdownTimeoutMillis must be a positive finite number no greater than ${MAX_TIMER_DELAY_MILLIS}, got ${String(shutdownTimeoutMillis)}`)
  122. }
  123. this.shutdownTimeoutMillis = shutdownTimeoutMillis
  124. this.provider = new LoggerProvider({
  125. resource: resourceFromAttributes({
  126. 'service.name': APP_IDENTITY.product,
  127. 'service.version': APP_IDENTITY.version,
  128. // OTel semconv's standard user attribute, carried once per export
  129. // batch on the Resource rather than per record: the collector
  130. // aggregates by Resource, and the id is process-stable anyway.
  131. 'user.id': getOrCreateAnonymousUserId(),
  132. }),
  133. processors: [
  134. new BatchLogRecordProcessor({
  135. ...config.processor,
  136. // The complete validated exporter object, verbatim: every SDK
  137. // option (`timeoutMillis`, `compression`, `keepAlive`, …) reaches
  138. // the exporter — rebuilding selected fields here would silently
  139. // ignore the rest. App identity travels in the Resource
  140. // (service.name/version); the transport-level user-agent is the
  141. // SDK's own, per the axiom.
  142. exporter: new OTLPLogExporter(config.exporter),
  143. }),
  144. ],
  145. })
  146. this.ledger = this.provider.getLogger('@deepseek-ai/dsh-session-telemetry-otel', version)
  147. this.ops = this.provider.getLogger('@deepseek-ai/dsh-session-telemetry-otel/ops', version)
  148. new TelemetryCoordinator(ctx, this)
  149. }
  150. /**
  151. * Map one seam record onto the SDK logger for its channel — a synchronous
  152. * enqueue into the batch processor's queue.
  153. * @param record - the logical record handed over by the coordinator.
  154. */
  155. emit(record: TelemetryRecord): void {
  156. const logger = record.channel === 'ops' ? this.ops : this.ledger
  157. logger.emit({
  158. timestamp: record.time,
  159. observedTimestamp: record.time,
  160. ...SEVERITY[record.severity],
  161. // JSON-serializable by the seam's contract (validated at Session.append),
  162. // which is exactly the AnyValue subset.
  163. body: record.body as AnyValue,
  164. attributes: record.attributes,
  165. })
  166. }
  167. // The seam's optional flush() hint is deliberately NOT implemented. The
  168. // batch processor exports on its own cadence (`processor.scheduledDelayMillis`,
  169. // the SDK's documented knob), and this backend is the SDK pipeline's only
  170. // caller — forwarding the hint to `forceFlush()` was the sole source of
  171. // concurrent flushes, whose undocumented interactions with shutdown's
  172. // internal drain (concurrent-flush guard, provider-level flush timeout)
  173. // silently dropped tail records. Removal history and the revival trigger:
  174. // the revival Agent Note.
  175. /**
  176. * Ask the SDK to drain and quiesce, but reject after the backend-owned
  177. * deadline. OTel's processor export timeout wraps `exportCompleted` only;
  178. * shutdown awaits `exporter.forceFlush()` first, which can remain pending
  179. * when the transport never obtains a socket. The provider promise remains
  180. * observed after the deadline so a later rejection cannot become unhandled.
  181. * @returns resolves when the SDK pipeline quiesces, or rejects at the configured deadline.
  182. */
  183. async shutdown(): Promise<void> {
  184. const providerShutdown = this.provider.shutdown()
  185. let timer: ReturnType<typeof setTimeout> | undefined
  186. const deadline = new Promise<never>((_resolve, reject) => {
  187. timer = setTimeout(() => {
  188. reject(new Error(`session-telemetry-otel: provider shutdown exceeded ${this.shutdownTimeoutMillis}ms`))
  189. }, this.shutdownTimeoutMillis)
  190. })
  191. try {
  192. await Promise.race([providerShutdown, deadline])
  193. } finally {
  194. /* v8 ignore else -- the Promise executor assigns timer synchronously before this race starts. */
  195. if (timer !== undefined) clearTimeout(timer)
  196. }
  197. }
  198. }
  199. export default TelemetryOtel