From f0735cbaaae2d1df1ab95a49492e01aa86388139 Mon Sep 17 00:00:00 2001 From: "rosetta-livekit-bot[bot]" <282703043+rosetta-livekit-bot[bot]@users.noreply.github.com> Date: Wed, 26 Aug 2026 12:07:47 +0000 Subject: [PATCH] fix(observability): retry recording connection failures --- .../retry-recording-connection-failures.md | 6 + agents/src/telemetry/traces.test.ts | 240 +++++++++++++++++- agents/src/telemetry/traces.ts | 163 ++++++++++-- 3 files changed, 392 insertions(+), 17 deletions(-) create mode 100644 .changeset/retry-recording-connection-failures.md diff --git a/.changeset/retry-recording-connection-failures.md b/.changeset/retry-recording-connection-failures.md new file mode 100644 index 0000000000..1587f47468 --- /dev/null +++ b/.changeset/retry-recording-connection-failures.md @@ -0,0 +1,6 @@ +--- +'@livekit/agents': patch +--- + +Retry session recording uploads after DNS, connection, and connection-timeout failures while +keeping total and connection timeouts scoped to the upload request. diff --git a/agents/src/telemetry/traces.test.ts b/agents/src/telemetry/traces.test.ts index 807c9aff3a..41c881a9dc 100644 --- a/agents/src/telemetry/traces.test.ts +++ b/agents/src/telemetry/traces.test.ts @@ -4,8 +4,33 @@ import { context as otelContext, trace } from '@opentelemetry/api'; import { InMemorySpanExporter, SimpleSpanProcessor } from '@opentelemetry/sdk-trace-base'; import { NodeTracerProvider } from '@opentelemetry/sdk-trace-node'; -import { afterEach, beforeEach, describe, expect, it } from 'vitest'; -import { setTracerProvider, tracer } from './traces.js'; +import { EventEmitter } from 'node:events'; +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import type { SessionReport } from '../voice/report.js'; +import { setTracerProvider, tracer, uploadSessionReport } from './traces.js'; + +const telemetryMocks = vi.hoisted(() => ({ + exportLogs: vi.fn(async () => {}), + submit: vi.fn(), +})); + +vi.mock('form-data', () => ({ + default: class MockFormData { + append(): void {} + + submit(...args: unknown[]): unknown { + return telemetryMocks.submit(...args); + } + }, +})); + +vi.mock('./otel_http_exporter.js', () => ({ + SimpleOTLPHttpLogExporter: class MockLogExporter { + export(...args: unknown[]): Promise { + return telemetryMocks.exportLogs(...args); + } + }, +})); /** Helper: extract parentSpanId across OTel SDK v1/v2 */ function parentSpanId(span: unknown): string | undefined { @@ -140,3 +165,214 @@ describe('register() set-once semantics', () => { expect(parentSpanId(lkSession)).toBe(userParent.spanContext().spanId); }); }); + +type SubmitCallback = (error: Error | null, response?: MockResponse) => void; + +class MockRequest extends EventEmitter { + callback?: SubmitCallback; + + destroy(error?: Error): this { + if (error) queueMicrotask(() => this.callback?.(error)); + return this; + } +} + +class MockResponse extends EventEmitter { + statusCode: number; + statusMessage: string; + + constructor(statusCode = 200, statusMessage = 'OK') { + super(); + this.statusCode = statusCode; + this.statusMessage = statusMessage; + } + + destroy(error?: Error): this { + if (error) queueMicrotask(() => this.emit('error', error)); + return this; + } +} + +function mockReport(): SessionReport { + return { + roomId: 'room-id', + jobId: 'job-id', + room: 'room-name', + timestamp: 1, + startedAt: 1, + options: {}, + chatHistory: { + items: [], + toJSON: () => ({ items: [] }), + }, + } as unknown as SessionReport; +} + +function submission( + callback: SubmitCallback, + action: (request: MockRequest, callback: SubmitCallback) => void, +): MockRequest { + const request = new MockRequest(); + request.callback = callback; + queueMicrotask(() => action(request, callback)); + return request; +} + +function respond(callback: SubmitCallback, statusCode = 200, body = Buffer.alloc(0)): void { + const response = new MockResponse(statusCode, statusCode < 400 ? 'OK' : 'Service Unavailable'); + callback(null, response); + queueMicrotask(() => { + if (body.length > 0) response.emit('data', body); + response.emit('end'); + }); +} + +function retryInfoBody(): Buffer { + const typeUrl = Buffer.from('type.googleapis.com/google.rpc.RetryInfo'); + const retryInfo = Buffer.from([0x0a, 0x00]); + const any = Buffer.concat([ + Buffer.from([0x0a, typeUrl.length]), + typeUrl, + Buffer.from([0x12, retryInfo.length]), + retryInfo, + ]); + return Buffer.concat([Buffer.from([0x1a, any.length]), any]); +} + +function networkError(code: string, message = code): NodeJS.ErrnoException { + return Object.assign(new Error(message), { code }); +} + +async function upload(): Promise { + await uploadSessionReport({ + agentName: 'test-agent', + cloudHostname: 'cloud.example.com', + report: mockReport(), + }); +} + +describe('uploadSessionReport retries', () => { + beforeEach(() => { + vi.useFakeTimers(); + vi.stubEnv('LIVEKIT_API_KEY', 'api-key'); + vi.stubEnv('LIVEKIT_API_SECRET', 'api-secret'); + telemetryMocks.submit.mockReset(); + }); + + afterEach(() => { + vi.useRealTimers(); + vi.unstubAllEnvs(); + }); + + it('uses upload-scoped extended total and connection timeouts', async () => { + const timeoutSpy = vi.spyOn(globalThis, 'setTimeout'); + telemetryMocks.submit.mockImplementation((_options: unknown, callback: SubmitCallback) => + submission(callback, (request, cb) => { + const socket = Object.assign(new EventEmitter(), { + connecting: true, + secureConnecting: true, + }); + request.emit('socket', socket); + socket.emit('secureConnect'); + respond(cb); + }), + ); + + await upload(); + + expect(timeoutSpy.mock.calls.some(([, delay]) => delay === 30_000)).toBe(true); + expect(timeoutSpy.mock.calls.some(([, delay]) => delay === 900_000)).toBe(true); + }); + + it.each([ + ['connection timeout', 'ETIMEDOUT'], + ['connection failure', 'ECONNREFUSED'], + ])('retries %s', async (_name, code) => { + telemetryMocks.submit + .mockImplementationOnce((_options: unknown, callback: SubmitCallback) => + submission(callback, (_request, cb) => cb(networkError(code))), + ) + .mockImplementationOnce((_options: unknown, callback: SubmitCallback) => + submission(callback, (_request, cb) => respond(cb)), + ); + + const uploading = upload(); + await vi.runAllTimersAsync(); + await uploading; + + expect(telemetryMocks.submit).toHaveBeenCalledTimes(2); + }); + + it('retries a response with RetryInfo', async () => { + telemetryMocks.submit + .mockImplementationOnce((_options: unknown, callback: SubmitCallback) => + submission(callback, (_request, cb) => respond(cb, 503, retryInfoBody())), + ) + .mockImplementationOnce((_options: unknown, callback: SubmitCallback) => + submission(callback, (_request, cb) => respond(cb)), + ); + + const uploading = upload(); + await vi.runAllTimersAsync(); + await uploading; + + expect(telemetryMocks.submit).toHaveBeenCalledTimes(2); + }); + + it('does not retry a response without RetryInfo', async () => { + telemetryMocks.submit.mockImplementation((_options: unknown, callback: SubmitCallback) => + submission(callback, (_request, cb) => respond(cb, 503)), + ); + + await expect(upload()).rejects.toThrow('503 Service Unavailable'); + expect(telemetryMocks.submit).toHaveBeenCalledTimes(1); + }); + + it.each(['total timeout', 'post-send disconnect', 'TLS failure'])( + 'does not retry %s', + async (failure) => { + telemetryMocks.submit.mockImplementation((_options: unknown, callback: SubmitCallback) => + submission(callback, (request, cb) => { + if (failure === 'total timeout') { + const socket = Object.assign(new EventEmitter(), { + connecting: false, + secureConnecting: false, + }); + request.emit('socket', socket); + return; + } + if (failure === 'post-send disconnect') { + const socket = Object.assign(new EventEmitter(), { + connecting: false, + secureConnecting: false, + }); + request.emit('socket', socket); + cb(networkError('ECONNRESET', 'response lost')); + return; + } + cb(networkError('CERT_HAS_EXPIRED', 'TLS failed')); + }), + ); + + const uploading = upload(); + const rejection = expect(uploading).rejects.toThrow(); + if (failure === 'total timeout') await vi.advanceTimersByTimeAsync(900_000); + + await rejection; + expect(telemetryMocks.submit).toHaveBeenCalledTimes(1); + }, + ); + + it('stops after connection retries are exhausted', async () => { + telemetryMocks.submit.mockImplementation((_options: unknown, callback: SubmitCallback) => + submission(callback, (_request, cb) => cb(networkError('ETIMEDOUT', 'connection timed out'))), + ); + + const uploading = upload(); + const rejection = expect(uploading).rejects.toThrow('connection timed out'); + await vi.runAllTimersAsync(); + await rejection; + + expect(telemetryMocks.submit).toHaveBeenCalledTimes(4); + }); +}); diff --git a/agents/src/telemetry/traces.ts b/agents/src/telemetry/traces.ts index 682682759e..e2b12e7b2c 100644 --- a/agents/src/telemetry/traces.ts +++ b/agents/src/telemetry/traces.ts @@ -23,6 +23,8 @@ import { ATTR_SERVICE_NAME } from '@opentelemetry/semantic-conventions'; import FormData from 'form-data'; import { AccessToken } from 'livekit-server-sdk'; import fs from 'node:fs/promises'; +import type { IncomingMessage } from 'node:http'; +import type { TLSSocket } from 'node:tls'; import type { ChatContent, ChatItem, ChatRole } from '../llm/index.js'; import { enableOtelLogging, log } from '../log.js'; import { filterZeroValues } from '../metrics/model_usage.js'; @@ -631,7 +633,34 @@ export async function uploadSessionReport(options: { const submitOnce = (): Promise<{ statusCode: number; statusMessage: string; body: Buffer }> => new ThrowsPromise<{ statusCode: number; statusMessage: string; body: Buffer }, Error>( (resolve, reject) => { - buildFormData().submit( + let connected = false; + let response: IncomingMessage | undefined; + let connectTimer: NodeJS.Timeout | undefined; + let timeoutError: Error | undefined; + let settled = false; + + const clearTimers = () => { + clearTimeout(connectTimer); + clearTimeout(totalTimer); + }; + const resolveOnce = (value: { + statusCode: number; + statusMessage: string; + body: Buffer; + }) => { + if (settled) return; + settled = true; + clearTimers(); + resolve(value); + }; + const rejectOnce = (error: Error, retryableConnection = false) => { + if (settled) return; + settled = true; + clearTimers(); + reject(new RecordingUploadError(error, retryableConnection)); + }; + + const request = buildFormData().submit( { protocol: 'https:', host: cloudHostname, @@ -643,17 +672,19 @@ export async function uploadSessionReport(options: { }, (err, res) => { if (err) { - reject(new Error(`Failed to upload session report: ${err.message}`)); + const error = timeoutError ?? err; + rejectOnce(error, !connected && isRetryableConnectionError(error)); return; } + connected = true; + clearTimeout(connectTimer); + response = res; const chunks: Buffer[] = []; res.on('data', (chunk: Buffer) => chunks.push(chunk)); - res.on('error', (readErr) => - reject(new Error(`Response read error: ${readErr.message}`)), - ); + res.on('error', (readErr) => rejectOnce(timeoutError ?? readErr)); res.on('end', () => - resolve({ + resolveOnce({ statusCode: res.statusCode ?? 0, statusMessage: res.statusMessage ?? '', body: Buffer.concat(chunks), @@ -661,34 +692,136 @@ export async function uploadSessionReport(options: { ); }, ); + + request.on('socket', (socket) => { + const tlsSocket = socket as TLSSocket; + const secureConnecting = (tlsSocket as TLSSocket & { secureConnecting?: boolean }) + .secureConnecting; + if (!socket.connecting && !secureConnecting) { + connected = true; + return; + } + + tlsSocket.once('secureConnect', () => { + connected = true; + clearTimeout(connectTimer); + }); + connectTimer = setTimeout(() => { + timeoutError = new RecordingUploadConnectTimeoutError(); + request.destroy(timeoutError); + }, RECORDING_UPLOAD_CONNECT_TIMEOUT_MS); + connectTimer.unref(); + }); + + const totalTimer = setTimeout(() => { + timeoutError = new RecordingUploadTotalTimeoutError(); + response?.destroy(timeoutError); + request.destroy(timeoutError); + }, RECORDING_UPLOAD_TOTAL_TIMEOUT_MS); + totalTimer.unref(); }, ); - const maxRetries = 3; - for (let attempt = 0; attempt <= maxRetries; attempt++) { + for (let attempt = 0; attempt <= RECORDING_UPLOAD_MAX_RETRIES; attempt++) { log().debug('uploading session report to LiveKit Cloud'); - const { statusCode, statusMessage, body } = await submitOnce(); + let result: { statusCode: number; statusMessage: string; body: Buffer }; + let retry: { delayMs: number; failure: string } | undefined; + try { + result = await submitOnce(); + } catch (error) { + if ( + !(error instanceof RecordingUploadError) || + !error.retryableConnection || + attempt === RECORDING_UPLOAD_MAX_RETRIES + ) { + throw error; + } + retry = { + delayMs: recordingUploadRetryDelay(attempt), + failure: error.cause.name, + }; + } + + if (retry) { + logRecordingUploadRetry(retry.failure, attempt, retry.delayMs); + await new Promise((resolve) => setTimeout(resolve, retry.delayMs)); + continue; + } + + const { statusCode, statusMessage, body } = result!; if (statusCode > 0 && statusCode < 400) { log().debug('finished uploading'); return; } const retryDelayMs = parseRetryDelayMs(body); - if (retryDelayMs === null || attempt === maxRetries) { + if (retryDelayMs === null || attempt === RECORDING_UPLOAD_MAX_RETRIES) { throw new Error( `Failed to upload session report: ${statusCode} ${statusMessage} - ${body.toString('utf-8')}`, ); } - log().warn( - `recording upload failed (attempt ${attempt + 1}/${maxRetries + 1}), retrying in ${( - retryDelayMs / 1000 - ).toFixed(1)}s`, - ); + logRecordingUploadRetry(`status ${statusCode}`, attempt, retryDelayMs); await new Promise((resolve) => setTimeout(resolve, retryDelayMs)); } } +const RECORDING_UPLOAD_TOTAL_TIMEOUT_MS = 900_000; +const RECORDING_UPLOAD_CONNECT_TIMEOUT_MS = 30_000; +const RECORDING_UPLOAD_MAX_RETRIES = 3; + +class RecordingUploadConnectTimeoutError extends Error { + constructor() { + super('Session report upload connection timed out'); + this.name = 'RecordingUploadConnectTimeoutError'; + } +} + +class RecordingUploadTotalTimeoutError extends Error { + constructor() { + super('Session report upload timed out'); + this.name = 'RecordingUploadTotalTimeoutError'; + } +} + +class RecordingUploadError extends Error { + constructor( + override readonly cause: Error, + readonly retryableConnection: boolean, + ) { + super(`Failed to upload session report: ${cause.message}`, { cause }); + this.name = 'RecordingUploadError'; + } +} + +function isRetryableConnectionError(error: Error): boolean { + if (error instanceof RecordingUploadConnectTimeoutError) return true; + if (error instanceof RecordingUploadTotalTimeoutError) return false; + const code = (error as NodeJS.ErrnoException).code; + return !( + code === undefined || + code.startsWith('ERR_TLS_') || + code.startsWith('ERR_SSL_') || + code.startsWith('ERR_OSSL_') || + code.startsWith('CERT_') || + code.startsWith('DEPTH_ZERO_') || + code.startsWith('SELF_SIGNED_') || + code.startsWith('UNABLE_TO_') + ); +} + +function recordingUploadRetryDelay(attempt: number): number { + return Math.random() * Math.min(2 ** attempt * 1000, 8000); +} + +function logRecordingUploadRetry(failure: string, attempt: number, retryDelayMs: number): void { + log().warn( + `recording upload failed (${failure}, attempt ${attempt + 1}/${ + RECORDING_UPLOAD_MAX_RETRIES + 1 + }), retrying in ${(retryDelayMs / 1000).toFixed(1)}s`, + ); +} + const RETRY_INFO_TYPE_NAME = 'google.rpc.RetryInfo'; // google.protobuf.Any.type_url only requires the segment after the last "/" to