import { attachOpenTelemetryBridge, TelemetryRow, WhisperTelemetryCollector, } from "../src/services/transcription/metrics"; const baseStatus = { connectionState: "connecting" as const, websocketState: "connecting" as const, browserOnline: true, isRecovering: false, isBufferingAudio: false, isDrainingAudio: false, transportBufferedAmountBytes: 0, pendingBufferedAudioBytes: 0, persistedBufferedAudioBytes: 0, persistedBufferedAudioSegments: 0, totalBufferedAudioBytes: 0, bufferingStartedAt: null, drainStartedAt: null, lastDrainProgressAt: null, drainCycles: 0, drainExitReason: null, reconnectAttempt: null, nextReconnectDelayMs: null, finishDeliveryState: null, maxReconnectAttempts: 3, remainingReconnectAttempts: 3, lastDisconnect: null, }; describe("telemetry v2 core", () => { it("caps retained rows with a ring buffer", () => { const collector = new WhisperTelemetryCollector({ provider: "sofya_as_service", onEmit: jest.fn(), rowLimit: 5, }); collector.markSessionStart(1_000); collector.clearTelemetryRows(); for (let index = 0; index < 10; index += 1) { collector.recordTranscriptDiarization(1_001 + index); } const rows = collector.getTelemetryRows(); expect(rows).toHaveLength(5); expect(rows.every((row) => row.metric === "tx.diarization.received")).toBe(true); }); it("memoizes snapshots while revision is unchanged", () => { const collector = new WhisperTelemetryCollector({ provider: "sofya_as_service", onEmit: jest.fn(), }); collector.markSessionStart(1_000); const first = collector.getTelemetrySnapshot(1_500); const second = collector.getTelemetrySnapshot(1_500); expect(first).toBe(second); collector.recordTranscriptDiarization(1_600); const third = collector.getTelemetrySnapshot(1_600); expect(third).not.toBe(second); }); it("computes derived continuity, delivery, and reconnect metrics", () => { const collector = new WhisperTelemetryCollector({ provider: "sofya_as_service", onEmit: jest.fn(), }); collector.markSessionStart(1_000); collector.markAudioChunkProduced(120, 1_010); collector.markAudioChunkSent(120, 1_011); collector.markAudioChunkDropped(80, "terminal_disconnected", 1_012); collector.recordResilienceStatus( { ...baseStatus, connectionState: "reconnecting", websocketState: "closed", isRecovering: true, reconnectAttempt: 1, nextReconnectDelayMs: 100, remainingReconnectAttempts: 2, }, 1_100 ); collector.recordResilienceStatus( { ...baseStatus, connectionState: "connected", websocketState: "open", reconnectAttempt: null, nextReconnectDelayMs: null, remainingReconnectAttempts: 3, }, 1_300 ); const snapshot = collector.getTelemetrySnapshot(1_500); expect(snapshot.derived.continuity_loss_rate).toBe(1); expect(snapshot.derived.delivery_rate).toBe(1); expect(snapshot.derived.reconnect_success_rate).toBe(1); expect(snapshot.derived.mean_recovery_time_ms).toBeGreaterThanOrEqual(200); }); }); describe("open telemetry bridge", () => { it("maps telemetry rows into OTel instruments and unsubscribes cleanly", () => { const listeners = new Set<(row: TelemetryRow) => void>(); const source = { on: (_event: "telemetry_row", listener: (row: TelemetryRow) => void) => { listeners.add(listener); }, off: (_event: "telemetry_row", listener: (row: TelemetryRow) => void) => { listeners.delete(listener); }, getTelemetrySnapshot: () => ({ schemaVersion: 2 as const, updatedAt: 1, provider: "sofya_as_service" as const, capabilities: { transcriptUi: true, session: true, connection: true, recovery: true, buffering: true, browserNetwork: true, audioCapture: true, }, window: { sessionId: "session-1", startedAt: 1, endedAt: null, durationMs: null, }, status: { connectionState: "connected" as const, websocketState: "open" as const, browserOnline: true, isRecovering: false, isBufferingAudio: false, isDrainingAudio: false, transportBufferedAmountBytes: 0, pendingBufferedAudioBytes: 0, persistedBufferedAudioBytes: 0, persistedBufferedAudioSegments: 0, totalBufferedAudioBytes: 0, lastDrainProgressAt: null, drainCycles: 0, drainExitReason: null, reconnectAttempt: null, nextReconnectDelayMs: null, remainingReconnectAttempts: null, finishDeliveryState: null, lastDisconnectCode: null, lastDisconnectReason: null, networkEffectiveType: "4g", }, counters: {}, gauges: { "network.rtt_ms": 20, }, histograms: {}, derived: { continuity_loss_rate: null, delivery_rate: null, reconnect_success_rate: null, mean_recovery_time_ms: null, stall_ratio: null, buffer_pressure_ratio: null, }, }), }; const counterAdd = jest.fn(); const histogramRecord = jest.fn(); const gaugeAdd = jest.fn(); const meter = { createCounter: () => ({ add: counterAdd, }), createHistogram: () => ({ record: histogramRecord, }), createUpDownCounter: () => ({ add: gaugeAdd, }), } as any; const detach = attachOpenTelemetryBridge(source, { meter }); const rowListener = Array.from(listeners)[0]; expect(typeof rowListener).toBe("function"); rowListener({ ts: 2, sessionId: "session-1", metric: "tx.partial.received", value: 2, type: "counter", }); rowListener({ ts: 3, sessionId: "session-1", metric: "audio.chunk.size_bytes", value: 640, type: "histogram", }); rowListener({ ts: 4, sessionId: "session-1", metric: "network.rtt_ms", value: 25, type: "gauge", }); rowListener({ ts: 5, sessionId: "session-1", metric: "network.effective_type_changed", value: 1, type: "event", tags: { value: "4g" }, }); expect(counterAdd).toHaveBeenCalled(); expect(histogramRecord).toHaveBeenCalledWith(640, undefined); expect(gaugeAdd).toHaveBeenCalled(); detach(); expect(listeners.size).toBe(0); }); });