From 9b54e5d73a20eaa0d4bf3e055571e50c43763088 Mon Sep 17 00:00:00 2001 From: David Zhao Date: Mon, 14 Sep 2026 22:45:18 -0700 Subject: [PATCH 1/4] feat(telemetry): trace RPC calls made and handled Port of livekit/agents#7134. `telemetry.rpc` installs one tracing interceptor on the job's local participant and turns each RPC into a span following the OpenTelemetry RPC semantic conventions: `rpc_call` (CLIENT, under the current span, so an RPC issued from a tool nests under `function_tool`) and `rpc_handler` (SERVER, under the primary session's root span). Attributes cover method, ids, identities, payload and response sizes, the response timeout in seconds, the RpcError code and whether a handler was registered; request and response bodies are recorded truncated to 1 KiB under `lk.pii.rpc.*` keys. Installed from JobContext.connect() after the room connects and from RoomIO on every connected transition; both are idempotent, the SDK dedups the singleton interceptor by identity. `StartSpanOptions` gains `kind`. Adaptation: the installed @livekit/rtc-node predates LocalParticipant.addRpcInterceptor (it is in an open SDK PR), so the interceptor and call shapes are declared structurally here and `install` feature-detects the hook, degrading to one debug log and a no-op like the Python layer on an older livekit-rtc. Not ported: the Python handling of asyncio cancellation of the handler chain (promises are not cancelled in Node; the SDK maps timeouts on the caller's side). Incoming invocations are dispatched by the SDK from its FFI event path, outside the job's AsyncLocalStorage, so the session and job roots did not resolve there and rpc_handler spans came out as roots of their own (seen in a live run). install() now captures the job (passed explicitly by JobContext.connect() and RoomIO, which itself runs on the SDK's event path) and the interceptor runs the handler chain inside it, which also gives the user's handler getJobContext(), as a Python handler has. Co-Authored-By: Claude Fable 5.1 --- .changeset/rpc-tracing-spans.md | 5 + agents/src/job.ts | 2 + agents/src/telemetry/index.ts | 1 + agents/src/telemetry/rpc.test.ts | 284 +++++++++++++++++++++ agents/src/telemetry/rpc.ts | 161 ++++++++++++ agents/src/telemetry/startup_spans.test.ts | 2 +- agents/src/telemetry/trace_types.test.ts | 9 + agents/src/telemetry/trace_types.ts | 18 ++ agents/src/telemetry/traces.ts | 16 +- agents/src/voice/room_io/room_io.test.ts | 47 +++- agents/src/voice/room_io/room_io.ts | 15 +- pnpm-lock.yaml | 101 ++++---- pnpm-workspace.yaml | 2 +- 13 files changed, 608 insertions(+), 55 deletions(-) create mode 100644 .changeset/rpc-tracing-spans.md create mode 100644 agents/src/telemetry/rpc.test.ts create mode 100644 agents/src/telemetry/rpc.ts diff --git a/.changeset/rpc-tracing-spans.md b/.changeset/rpc-tracing-spans.md new file mode 100644 index 0000000000..6cd137aeea --- /dev/null +++ b/.changeset/rpc-tracing-spans.md @@ -0,0 +1,5 @@ +--- +'@livekit/agents': patch +--- + +Trace RPC calls the agent performs (`rpc_call`) and handles (`rpc_handler`) through the room SDK's `RpcInterceptor` hook. `@livekit/rtc-node` `^1.1.0` is now required, the first release with `LocalParticipant.addRpcInterceptor`. diff --git a/agents/src/job.ts b/agents/src/job.ts index 73ee4c4396..7b0bb3c8b0 100644 --- a/agents/src/job.ts +++ b/agents/src/job.ts @@ -31,6 +31,7 @@ import { traceTypes, uploadSessionReport, } from './telemetry/index.js'; +import * as rpcTracing from './telemetry/rpc.js'; import { sessionSpan } from './telemetry/session_context.js'; import { ATTRIBUTE_REDACTION_ENABLED, @@ -513,6 +514,7 @@ export class JobContext> { jobCtx: this as JobContext as JobContext, }, ); + rpcTracing.install(this.#room.localParticipant, this as JobContext as JobContext); this.#onConnect(); this.#room.remoteParticipants.forEach(this.onParticipantConnected); diff --git a/agents/src/telemetry/index.ts b/agents/src/telemetry/index.ts index 1fd1f1ca71..b1bca84425 100644 --- a/agents/src/telemetry/index.ts +++ b/agents/src/telemetry/index.ts @@ -26,6 +26,7 @@ export { export * as genAI from './gen_ai.js'; export { REDACTED_EXCEPTION_MESSAGE } from './redaction.js'; export * as loopMonitor from './loop_monitor.js'; +export * as rpc from './rpc.js'; export * as traceTypes from './trace_types.js'; export { discardPreparedCloudTracer, diff --git a/agents/src/telemetry/rpc.test.ts b/agents/src/telemetry/rpc.test.ts new file mode 100644 index 0000000000..deae9916f1 --- /dev/null +++ b/agents/src/telemetry/rpc.test.ts @@ -0,0 +1,284 @@ +// SPDX-FileCopyrightText: 2026 LiveKit, Inc. +// +// SPDX-License-Identifier: Apache-2.0 + +/** + * RPC tracing (`telemetry.rpc`). + * + * The interceptor is exercised directly with fake `next` continuations, so these tests run + * against any `@livekit/rtc-node` version. With an SDK that has `RpcInterceptor` support the + * same interceptor is what `install` registers on the local participant; without it, `install` + * is a no-op, which is also covered. + */ +import { + type LocalParticipant, + type RpcCallInfo, + RpcError, + type RpcInvocationData, +} from '@livekit/rtc-node'; +import { + ROOT_CONTEXT, + SpanKind, + SpanStatusCode, + 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, vi } from 'vitest'; +import { type JobContext, getJobContext, runWithJobContext } from '../job.js'; +import { + MAX_PAYLOAD_ATTR_LEN, + TracingRpcInterceptor, + install, + interceptor as singleton, +} from './rpc.js'; +import * as traceTypes from './trace_types.js'; +import { setTracerProvider, tracer } from './traces.js'; + +function call(overrides: Partial = {}): RpcCallInfo { + return { + destinationIdentity: 'avatar-1', + method: 'playback.start', + payload: '{"id": 7}', + responseTimeout: 5000, + ...overrides, + }; +} + +function invocation(overrides: Partial = {}): RpcInvocationData { + return { + requestId: 'req-42', + callerIdentity: 'client-9', + payload: '{"q": "x"}', + responseTimeout: 10_000, + method: 'agent.lookup', + ...overrides, + }; +} + +function fakeParticipant(addRpcInterceptor: (i: unknown) => void): LocalParticipant { + return { identity: 'agent', addRpcInterceptor } as unknown as LocalParticipant; +} + +describe.sequential('rpc tracing', () => { + let exporter: InMemorySpanExporter; + let provider: NodeTracerProvider; + let originalProvider: ReturnType; + let interceptor: TracingRpcInterceptor; + + beforeEach(() => { + originalProvider = tracer.getProvider(); + exporter = new InMemorySpanExporter(); + provider = new NodeTracerProvider({ spanProcessors: [new SimpleSpanProcessor(exporter)] }); + // the context manager, so a span made active in the test is the parent of the RPC span + provider.register(); + setTracerProvider(provider); + interceptor = new TracingRpcInterceptor(); + }); + + afterEach(async () => { + setTracerProvider(originalProvider); + await provider.shutdown(); + otelContext.disable(); + trace.disable(); + vi.restoreAllMocks(); + }); + + function spans(name: string) { + return exporter.getFinishedSpans().filter((span) => span.name === name); + } + + it('nests the outgoing call span under the caller', async () => { + let parentSpanId = ''; + const result = await tracer.startActiveSpan( + async (parent) => { + parentSpanId = parent.spanContext().spanId; + return interceptor.interceptOutgoing(call(), async () => 'ok!'); + }, + { name: 'function_tool' }, + ); + + expect(result).toBe('ok!'); + const [span] = spans('rpc_call'); + expect(span).toBeDefined(); + expect(span!.kind).toBe(SpanKind.CLIENT); + expect(span!.parentSpanContext?.spanId).toBe(parentSpanId); + expect(span!.attributes).toMatchObject({ + [traceTypes.ATTR_RPC_METHOD]: 'playback.start', + [traceTypes.ATTR_RPC_DESTINATION_IDENTITY]: 'avatar-1', + [traceTypes.ATTR_RPC_PAYLOAD]: '{"id": 7}', + [traceTypes.ATTR_RPC_PAYLOAD_SIZE]: 9, + [traceTypes.ATTR_RPC_RESPONSE]: 'ok!', + [traceTypes.ATTR_RPC_RESPONSE_SIZE]: 3, + [traceTypes.ATTR_RPC_RESPONSE_TIMEOUT]: 5, // seconds on the span, ms in the SDK + }); + expect(span!.status.code).toBe(SpanStatusCode.UNSET); + }); + + it('records the code and error status of a failed outgoing call', async () => { + const error = RpcError.builtIn('RECIPIENT_NOT_FOUND'); + await expect( + interceptor.interceptOutgoing(call(), async () => { + throw error; + }), + ).rejects.toBe(error); + + const [span] = spans('rpc_call'); + expect(span!.attributes[traceTypes.ATTR_RPC_ERROR_CODE]).toBe( + RpcError.ErrorCode.RECIPIENT_NOT_FOUND, + ); + expect(span!.status.code).toBe(SpanStatusCode.ERROR); + expect(span!.events.some((event) => event.name === 'exception')).toBe(true); + }); + + it('truncates the payload and keeps its size', async () => { + const payload = 'x'.repeat(MAX_PAYLOAD_ATTR_LEN + 500); + await interceptor.interceptOutgoing( + call({ payload, responseTimeout: undefined }), + async () => '', + ); + + const [span] = spans('rpc_call'); + const attrs = span!.attributes; + expect((attrs[traceTypes.ATTR_RPC_PAYLOAD] as string).length).toBe(MAX_PAYLOAD_ATTR_LEN); + expect(attrs[traceTypes.ATTR_RPC_PAYLOAD_SIZE]).toBe(payload.length); + expect(attrs).not.toHaveProperty(traceTypes.ATTR_RPC_RESPONSE_TIMEOUT); + expect(attrs[traceTypes.ATTR_RPC_RESPONSE_SIZE]).toBe(0); + expect(attrs).not.toHaveProperty(traceTypes.ATTR_RPC_RESPONSE); // empty: not recorded + }); + + it('truncates the response like the request', async () => { + const response = 'y'.repeat(MAX_PAYLOAD_ATTR_LEN + 10); + await interceptor.interceptIncoming(invocation(), async () => response); + + const [span] = spans('rpc_handler'); + expect((span!.attributes[traceTypes.ATTR_RPC_RESPONSE] as string).length).toBe( + MAX_PAYLOAD_ATTR_LEN, + ); + expect(span!.attributes[traceTypes.ATTR_RPC_RESPONSE_SIZE]).toBe(response.length); + }); + + it('measures payload sizes in bytes, not characters', async () => { + await interceptor.interceptOutgoing(call({ payload: 'héllo' }), async () => ''); + const [span] = spans('rpc_call'); + expect(span!.attributes[traceTypes.ATTR_RPC_PAYLOAD_SIZE]).toBe(6); + }); + + it('traces an incoming invocation as a server span', async () => { + const result = await interceptor.interceptIncoming(invocation(), async () => 'found'); + + expect(result).toBe('found'); + const [span] = spans('rpc_handler'); + expect(span!.kind).toBe(SpanKind.SERVER); + expect(span!.attributes).toMatchObject({ + [traceTypes.ATTR_RPC_METHOD]: 'agent.lookup', + [traceTypes.ATTR_RPC_REQUEST_ID]: 'req-42', + [traceTypes.ATTR_RPC_CALLER_IDENTITY]: 'client-9', + [traceTypes.ATTR_RPC_RESPONSE_TIMEOUT]: 10, + [traceTypes.ATTR_RPC_RESPONSE]: 'found', + [traceTypes.ATTR_RPC_RESPONSE_SIZE]: 5, + }); + }); + + it('parents the handler span to the primary session root, not the dispatching task', async () => { + const sessionRoot = tracer.startSpan({ name: 'agent_session' }); + const job = { + _primaryAgentSession: { rootSpanContext: trace.setSpan(ROOT_CONTEXT, sessionRoot) }, + job: { id: 'AJ_test' }, + } as unknown as JobContext; + + await tracer.startActiveSpan( + () => + runWithJobContext(job, () => interceptor.interceptIncoming(invocation(), async () => '')), + { name: 'sdk_dispatch_task' }, + ); + sessionRoot.end(); + + const [span] = spans('rpc_handler'); + expect(span!.parentSpanContext?.spanId).toBe(sessionRoot.spanContext().spanId); + }); + + it('parents the handler span to the job root when no session is running', async () => { + const jobRoot = tracer.startSpan({ name: 'job_entrypoint' }); + const job = { + _jobSpanContext: trace.setSpan(ROOT_CONTEXT, jobRoot), + job: { id: 'AJ_test' }, + } as unknown as JobContext; + + await tracer.startActiveSpan( + () => + runWithJobContext(job, () => interceptor.interceptIncoming(invocation(), async () => '')), + { name: 'sdk_dispatch_task' }, + ); + jobRoot.end(); + + const [span] = spans('rpc_handler'); + expect(span!.parentSpanContext?.spanId).toBe(jobRoot.spanContext().spanId); + }); + + it('records the error code of a handler that throws an RpcError', async () => { + const error = RpcError.builtIn('APPLICATION_ERROR'); + await expect( + interceptor.interceptIncoming(invocation(), async () => { + throw error; + }), + ).rejects.toBe(error); + + const [span] = spans('rpc_handler'); + expect(span!.attributes[traceTypes.ATTR_RPC_ERROR_CODE]).toBe( + RpcError.ErrorCode.APPLICATION_ERROR, + ); + expect(span!.status.code).toBe(SpanStatusCode.ERROR); + }); + + it('records a handler exception', async () => { + const boom = new Error('bad request body'); + await expect( + interceptor.interceptIncoming(invocation(), async () => { + throw boom; + }), + ).rejects.toBe(boom); + + const [span] = spans('rpc_handler'); + expect(span!.status.code).toBe(SpanStatusCode.ERROR); + expect(span!.attributes).not.toHaveProperty(traceTypes.ATTR_RPC_ERROR_CODE); + }); + + it('restores the installed job for invocations dispatched outside its context', async () => { + // the SDK's FFI event path carries no AsyncLocalStorage: without the job captured at + // install, neither the session root nor the job root would resolve and the handler span + // would be a root of its own + const sessionRoot = tracer.startSpan({ name: 'agent_session' }); + const job = { + _primaryAgentSession: { rootSpanContext: trace.setSpan(ROOT_CONTEXT, sessionRoot) }, + job: { id: 'AJ_test' }, + } as unknown as JobContext; + install(fakeParticipant(vi.fn()), job); + + let seenJob: JobContext | undefined; + await interceptor.interceptIncoming(invocation(), async () => { + seenJob = getJobContext(false); + return ''; + }); + sessionRoot.end(); + + const [span] = spans('rpc_handler'); + expect(span!.parentSpanContext?.spanId).toBe(sessionRoot.spanContext().spanId); + expect(seenJob).toBe(job); // the handler itself runs inside the job too + }); + + it('installs the one interceptor once per participant', () => { + const addRpcInterceptor = vi.fn(); + const participant = fakeParticipant(addRpcInterceptor); + install(participant); + install(participant); + // the same singleton each time: the SDK keeps one registration per instance + expect(addRpcInterceptor).toHaveBeenCalledTimes(2); + expect(addRpcInterceptor.mock.calls[0]![0]).toBe(singleton); + expect(addRpcInterceptor.mock.calls[1]![0]).toBe(singleton); + + // no participant yet (the room is not connected): nothing to install on + expect(() => install(undefined)).not.toThrow(); + }); +}); diff --git a/agents/src/telemetry/rpc.ts b/agents/src/telemetry/rpc.ts new file mode 100644 index 0000000000..5bfe0c8f9e --- /dev/null +++ b/agents/src/telemetry/rpc.ts @@ -0,0 +1,161 @@ +// SPDX-FileCopyrightText: 2026 LiveKit, Inc. +// +// SPDX-License-Identifier: Apache-2.0 + +/** + * Trace RPCs the agent performs and handles. + * + * The room SDK exposes an `RpcInterceptor` hook (`@livekit/rtc-node` with + * `LocalParticipant.addRpcInterceptor`) that wraps every call made through `performRpc` and + * every invocation dispatched to a registered handler. This module installs one interceptor per + * local participant that turns each call into a span following the OpenTelemetry RPC semantic + * conventions: + * + * - `rpc_call` (`SpanKind.CLIENT`) for outgoing calls, under whatever span is current where the + * call is made (an RPC issued from a tool nests under `function_tool`); + * - `rpc_handler` (`SpanKind.SERVER`) for incoming invocations, under the primary agent + * session's root span, else the job's root (`job_entrypoint`): the SDK dispatches them on a + * task whose context carries neither. + * + * Payloads are recorded truncated under `lk.pii` keys. Participant identities are application + * identifiers, not end-user data, and are recorded as is. + */ +import { + type IncomingRpcNext, + type LocalParticipant, + type OutgoingRpcNext, + type RpcCallInfo, + RpcError, + type RpcInterceptor, + type RpcInvocationData, +} from '@livekit/rtc-node'; +import { type Attributes, SpanKind } from '@opentelemetry/api'; +import { type JobContext, getJobContext, runWithJobContext } from '../job.js'; +import { jobRootContext, sessionRootContext } from './session_context.js'; +import * as traceTypes from './trace_types.js'; +import { tracer } from './traces.js'; +import { recordException } from './utils.js'; + +/** Request and response payloads longer than this many characters are truncated in span attributes. */ +export const MAX_PAYLOAD_ATTR_LEN = 1024; + +function truncate(payload: string): string { + return payload.slice(0, MAX_PAYLOAD_ATTR_LEN); +} + +function payloadAttributes(payload: string): Attributes { + const attrs: Attributes = { + [traceTypes.ATTR_RPC_PAYLOAD_SIZE]: Buffer.byteLength(payload, 'utf8'), + }; + if (payload) attrs[traceTypes.ATTR_RPC_PAYLOAD] = truncate(payload); + return attrs; +} + +function responseAttributes(response: string | undefined | null): Attributes { + const text = response ?? ''; + const attrs: Attributes = { + [traceTypes.ATTR_RPC_RESPONSE_SIZE]: Buffer.byteLength(text, 'utf8'), + }; + if (text) attrs[traceTypes.ATTR_RPC_RESPONSE] = truncate(text); + return attrs; +} + +function toError(error: unknown): Error { + return error instanceof Error ? error : new Error(String(error)); +} + +/** An `RpcInterceptor` emitting `rpc_call` / `rpc_handler` spans. */ +export class TracingRpcInterceptor implements RpcInterceptor { + async interceptOutgoing(call: RpcCallInfo, next: OutgoingRpcNext): Promise { + const attributes: Attributes = { + [traceTypes.ATTR_RPC_METHOD]: call.method, + [traceTypes.ATTR_RPC_DESTINATION_IDENTITY]: call.destinationIdentity, + ...payloadAttributes(call.payload), + }; + if (call.responseTimeout !== undefined) { + attributes[traceTypes.ATTR_RPC_RESPONSE_TIMEOUT] = call.responseTimeout / 1000; + } + + return tracer.startActiveSpan( + async (span) => { + let response: string; + try { + response = await next(call); + } catch (error) { + if (error instanceof RpcError) { + span.setAttribute(traceTypes.ATTR_RPC_ERROR_CODE, Number(error.code)); + } + recordException(span, toError(error)); + throw error; + } + span.setAttributes(responseAttributes(response)); + return response; + }, + { name: 'rpc_call', kind: SpanKind.CLIENT, attributes }, + ); + } + + async interceptIncoming(invocation: RpcInvocationData, next: IncomingRpcNext): Promise { + // the SDK dispatches invocations from its FFI event path, outside the job's + // AsyncLocalStorage: restore the job captured at install so the session and job roots + // resolve (and the handler sees getJobContext(), as a Python handler does) + const job = getJobContext(false) ?? installedJob; + if (job !== undefined && getJobContext(false) === undefined) { + return runWithJobContext(job, () => this.handleIncoming(invocation, next)); + } + return this.handleIncoming(invocation, next); + } + + private async handleIncoming( + invocation: RpcInvocationData, + next: IncomingRpcNext, + ): Promise { + const attributes: Attributes = { + [traceTypes.ATTR_RPC_METHOD]: invocation.method, + [traceTypes.ATTR_RPC_REQUEST_ID]: invocation.requestId, + [traceTypes.ATTR_RPC_CALLER_IDENTITY]: invocation.callerIdentity, + [traceTypes.ATTR_RPC_RESPONSE_TIMEOUT]: invocation.responseTimeout / 1000, + ...payloadAttributes(invocation.payload), + }; + + return tracer.startActiveSpan( + async (span) => { + let response: string; + try { + response = await next(invocation); + } catch (error) { + if (error instanceof RpcError) { + span.setAttribute(traceTypes.ATTR_RPC_ERROR_CODE, Number(error.code)); + } + recordException(span, toError(error)); + throw error; + } + span.setAttributes(responseAttributes(response)); + return response; + }, + { + name: 'rpc_handler', + // the session timeline, from whatever task the SDK dispatches on; before or after the + // session, the job's own timeline + context: sessionRootContext() ?? jobRootContext(), + kind: SpanKind.SERVER, + attributes, + }, + ); + } +} + +/** The one interceptor every participant gets; the SDK dedups registrations by identity. */ +export const interceptor = new TracingRpcInterceptor(); + +/** The job the interceptor was installed for: incoming invocations arrive outside its context. */ +let installedJob: JobContext | undefined; + +/** + * Trace RPCs on `localParticipant`. Idempotent: the SDK keeps one registration per interceptor + * instance. A missing participant (the room is not connected yet) installs nothing. + */ +export function install(localParticipant: LocalParticipant | undefined, jobCtx?: JobContext): void { + installedJob = jobCtx ?? getJobContext(false) ?? installedJob; + localParticipant?.addRpcInterceptor(interceptor); +} diff --git a/agents/src/telemetry/startup_spans.test.ts b/agents/src/telemetry/startup_spans.test.ts index e36fea67ee..9b853e4abf 100644 --- a/agents/src/telemetry/startup_spans.test.ts +++ b/agents/src/telemetry/startup_spans.test.ts @@ -47,7 +47,7 @@ function mockRoom(overrides: Record = {}): Room { connect: vi.fn(async () => undefined), isConnected: false, remoteParticipants: new Map(), - localParticipant: { sid: 'PA_agent', identity: 'agent-1' }, + localParticipant: { sid: 'PA_agent', identity: 'agent-1', addRpcInterceptor: vi.fn() }, ...overrides, } as unknown as Room; } diff --git a/agents/src/telemetry/trace_types.test.ts b/agents/src/telemetry/trace_types.test.ts index 5490d5f7c7..5573535f77 100644 --- a/agents/src/telemetry/trace_types.test.ts +++ b/agents/src/telemetry/trace_types.test.ts @@ -186,6 +186,15 @@ const SAFE_KEYS = new Set([ 'lk.disconnect_reason', 'lk.old_state', 'lk.new_state', + // rpc (semconv names, ids, sizes, codes; the payload keys are tagged) + 'rpc.method', + 'lk.rpc.request_id', + 'lk.rpc.caller_identity', + 'lk.rpc.destination_identity', + 'lk.rpc.payload_size', + 'lk.rpc.response_size', + 'lk.rpc.response_timeout', + 'lk.rpc.error_code', 'lk.close_reason', 'lk.close.drain', 'lk.shutdown.reason', diff --git a/agents/src/telemetry/trace_types.ts b/agents/src/telemetry/trace_types.ts index 494eddaee9..4469a3dc18 100644 --- a/agents/src/telemetry/trace_types.ts +++ b/agents/src/telemetry/trace_types.ts @@ -86,6 +86,24 @@ export const ATTR_DISCONNECT_REASON = 'lk.disconnect_reason'; export const ATTR_OLD_STATE = 'lk.old_state'; export const ATTR_NEW_STATE = 'lk.new_state'; +// rpc (`rpc.method` from the OpenTelemetry RPC semantic conventions, plus lk.rpc.* details) +export const ATTR_RPC_METHOD = 'rpc.method'; +export const ATTR_RPC_REQUEST_ID = 'lk.rpc.request_id'; +export const ATTR_RPC_CALLER_IDENTITY = 'lk.rpc.caller_identity'; +export const ATTR_RPC_DESTINATION_IDENTITY = 'lk.rpc.destination_identity'; +/** Request payload, truncated to `telemetry.rpc.MAX_PAYLOAD_ATTR_LEN` characters. */ +export const ATTR_RPC_PAYLOAD = 'lk.pii.rpc.payload'; +/** Request payload size in bytes, before truncation. */ +export const ATTR_RPC_PAYLOAD_SIZE = 'lk.rpc.payload_size'; +/** Response payload, truncated like the request. */ +export const ATTR_RPC_RESPONSE = 'lk.pii.rpc.response'; +/** Response payload size in bytes, before truncation. */ +export const ATTR_RPC_RESPONSE_SIZE = 'lk.rpc.response_size'; +/** Seconds the caller waits for a response. */ +export const ATTR_RPC_RESPONSE_TIMEOUT = 'lk.rpc.response_timeout'; +/** The `RpcError` code the call failed with. */ +export const ATTR_RPC_ERROR_CODE = 'lk.rpc.error_code'; + // session close / job shutdown export const ATTR_CLOSE_REASON = 'lk.close_reason'; export const ATTR_CLOSE_DRAIN = 'lk.close.drain'; diff --git a/agents/src/telemetry/traces.ts b/agents/src/telemetry/traces.ts index 3af76f7c3e..6f7e00054c 100644 --- a/agents/src/telemetry/traces.ts +++ b/agents/src/telemetry/traces.ts @@ -8,6 +8,7 @@ import { type Context, ProxyTracerProvider, type Span, + type SpanKind, type SpanOptions, type Tracer, type TracerProvider, @@ -73,6 +74,8 @@ export interface StartSpanOptions { endOnExit?: boolean; /** Optional start time for the span in milliseconds (Date.now() format) */ startTime?: number; + /** The span's kind (client, server, ...); defaults to INTERNAL */ + kind?: SpanKind; } /** @@ -147,6 +150,7 @@ class DynamicTracer { { attributes: options.attributes, startTime: options.startTime, + kind: options.kind, }, ctx, ); @@ -165,7 +169,11 @@ class DynamicTracer { async startActiveSpan(fn: (span: Span) => Promise, options: StartSpanOptions): Promise { const ctx = options.context || otelContext.active(); const endOnExit = options.endOnExit === undefined ? true : options.endOnExit; // default true - const opts: SpanOptions = { attributes: options.attributes, startTime: options.startTime }; + const opts: SpanOptions = { + attributes: options.attributes, + startTime: options.startTime, + kind: options.kind, + }; // Directly return the tracer's startActiveSpan result - it handles async correctly return await this.tracer.startActiveSpan(options.name, opts, ctx, async (span) => { @@ -210,7 +218,11 @@ class DynamicTracer { startActiveSpanSync(fn: (span: Span) => T, options: StartSpanOptions): T { const ctx = options.context || otelContext.active(); const endOnExit = options.endOnExit === undefined ? true : options.endOnExit; // default true - const opts: SpanOptions = { attributes: options.attributes, startTime: options.startTime }; + const opts: SpanOptions = { + attributes: options.attributes, + startTime: options.startTime, + kind: options.kind, + }; return this.tracer.startActiveSpan(options.name, opts, ctx, (span) => { try { diff --git a/agents/src/voice/room_io/room_io.test.ts b/agents/src/voice/room_io/room_io.test.ts index 3c3bfb466a..050d48031c 100644 --- a/agents/src/voice/room_io/room_io.test.ts +++ b/agents/src/voice/room_io/room_io.test.ts @@ -1,7 +1,7 @@ // SPDX-FileCopyrightText: 2026 LiveKit, Inc. // // SPDX-License-Identifier: Apache-2.0 -import { AudioFrame } from '@livekit/rtc-node'; +import { AudioFrame, ConnectionState, RoomEvent } from '@livekit/rtc-node'; import { ROOT_CONTEXT, trace } from '@opentelemetry/api'; import { InMemorySpanExporter, SimpleSpanProcessor } from '@opentelemetry/sdk-trace-base'; import { NodeTracerProvider } from '@opentelemetry/sdk-trace-node'; @@ -12,6 +12,7 @@ import { RealtimeModel } from '../../llm/index.js'; import { log } from '../../log.js'; import { IdentityTransform } from '../../stream/identity_transform.js'; import { setTracerProvider, tracer } from '../../telemetry/index.js'; +import * as rpcTracing from '../../telemetry/rpc.js'; import { DEFAULT_API_CONNECT_OPTIONS } from '../../types.js'; import { AgentSessionEventTypes, CloseReason, createCloseEvent } from '../events.js'; import { AudioInput, AudioOutput, TextOutput } from '../io.js'; @@ -106,6 +107,50 @@ function createFakeSession(llm?: RealtimeModel): FakeSession { }; } +describe('RoomIO rpc tracing', () => { + afterEach(() => { + vi.restoreAllMocks(); + }); + + /** + * A session may start on a room that connects later (ctx.connect() after session.start(), or + * a room the user connects). Tracing goes in on the connected transition, not only when the + * room is already up at start(); install is idempotent so a reconnect is harmless. + */ + it('installs rpc tracing when the room connects', () => { + const install = vi.spyOn(rpcTracing, 'install').mockReturnValue(true); + const room = createFakeRoom(); + const session = createFakeSession(); + const roomIO = new RoomIO({ + agentSession: session as unknown as RoomIOArgs['agentSession'], + room: room as unknown as RoomIOArgs['room'], + inputOptions: { audioEnabled: false, textEnabled: false }, + outputOptions: { audioEnabled: false, transcriptionEnabled: false }, + }); + + try { + roomIO.start(); + const onConnectionState = room.on.mock.calls.find( + ([event]) => event === RoomEvent.ConnectionStateChanged, + )?.[1] as ((state: ConnectionState) => void) | undefined; + expect(onConnectionState).toBeDefined(); + + onConnectionState!(ConnectionState.CONN_DISCONNECTED); + expect(install).not.toHaveBeenCalled(); + + room.isConnected = true; + onConnectionState!(ConnectionState.CONN_CONNECTED); + expect(install).toHaveBeenCalledTimes(1); + expect(install).toHaveBeenCalledWith(room.localParticipant, undefined); + + onConnectionState!(ConnectionState.CONN_CONNECTED); // reconnected + expect(install).toHaveBeenCalledTimes(2); // same singleton each time; the SDK dedups + } finally { + void roomIO.close(); + } + }); +}); + describe('RoomIO agent state attributes', () => { it('handles a failed update and publishes later state changes', async () => { const error = new Error('attribute update failed'); diff --git a/agents/src/voice/room_io/room_io.ts b/agents/src/voice/room_io/room_io.ts index 3d9df3df24..e85c26347a 100644 --- a/agents/src/voice/room_io/room_io.ts +++ b/agents/src/voice/room_io/room_io.ts @@ -24,6 +24,7 @@ import { RealtimeModel } from '../../llm/index.js'; import { log } from '../../log.js'; import { IdentityTransform } from '../../stream/identity_transform.js'; import { participantAttributes, traceTypes, tracer } from '../../telemetry/index.js'; +import * as rpcTracing from '../../telemetry/rpc.js'; import { DEFAULT_API_CONNECT_OPTIONS } from '../../types.js'; import { Future, IdleTimeoutError, Task, waitForAbort, waitUntilTimeout } from '../../utils.js'; import { type AgentSession } from '../agent_session.js'; @@ -258,12 +259,14 @@ export class RoomIO { this.agentSession._addSessionEvent?.('connection_state_changed', { [traceTypes.ATTR_CONNECTION_STATE]: ConnectionState[state] ?? String(state), }); - if ( - state === ConnectionState.CONN_CONNECTED && - this.room.isConnected && - !this.roomConnectedFuture.done - ) { - this.roomConnectedFuture.resolve(); + if (state === ConnectionState.CONN_CONNECTED && this.room.isConnected) { + // on every connect and reconnect; install is idempotent (one interceptor instance, + // deduped by the SDK), so JobContext.connect() installing too is fine. The job is passed + // explicitly: this runs on the SDK's event path, outside the job's AsyncLocalStorage + rpcTracing.install(this.room.localParticipant, this.jobContext); + if (!this.roomConnectedFuture.done) { + this.roomConnectedFuture.resolve(); + } } }; diff --git a/pnpm-lock.yaml b/pnpm-lock.yaml index 1333b95c25..0b0a017c1b 100644 --- a/pnpm-lock.yaml +++ b/pnpm-lock.yaml @@ -7,8 +7,8 @@ settings: catalogs: default: '@livekit/rtc-node': - specifier: ^0.13.34 - version: 0.13.34 + specifier: ^1.1.0 + version: 1.1.0 '@types/ws': specifier: ^8.5.10 version: 8.18.1 @@ -237,7 +237,7 @@ importers: devDependencies: '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@22.19.1) @@ -360,7 +360,7 @@ importers: version: 0.1.9 '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@opentelemetry/api': specifier: ^1.9.0 version: 1.9.0 @@ -434,7 +434,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -462,7 +462,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -493,7 +493,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -518,7 +518,7 @@ importers: version: link:../../agents '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -552,7 +552,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@22.19.1) @@ -580,7 +580,7 @@ importers: version: link:../../agents '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -617,7 +617,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -651,7 +651,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -685,7 +685,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -710,7 +710,7 @@ importers: version: link:../../agents '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -744,7 +744,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -778,7 +778,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -824,7 +824,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -851,7 +851,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -879,7 +879,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -897,7 +897,7 @@ importers: dependencies: '@livekit/plugins-krisp-viva-internal': specifier: 0.2.0 - version: 0.2.0(@livekit/agents@agents)(@livekit/rtc-node@0.13.34) + version: 0.2.0(@livekit/agents@agents)(@livekit/rtc-node@1.1.0) devDependencies: '@livekit/agents': specifier: workspace:* @@ -910,7 +910,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -941,7 +941,7 @@ importers: version: link:../../agents '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -972,7 +972,7 @@ importers: version: link:../../agents '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -1028,7 +1028,7 @@ importers: version: link:../../agents '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -1056,7 +1056,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -1081,7 +1081,7 @@ importers: version: link:../../agents '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -1109,7 +1109,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -1140,7 +1140,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -1180,7 +1180,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -1211,7 +1211,7 @@ importers: version: link:../openai '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -1233,7 +1233,7 @@ importers: version: link:../../agents '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -1255,7 +1255,7 @@ importers: version: link:../../agents '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -1286,7 +1286,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -1329,7 +1329,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -1354,7 +1354,7 @@ importers: version: link:../../agents '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -1388,7 +1388,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -1419,7 +1419,7 @@ importers: version: 0.2.7 '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -1447,7 +1447,7 @@ importers: version: link:../../agents '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -1472,7 +1472,7 @@ importers: version: link:../../agents '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -1503,7 +1503,7 @@ importers: version: link:../../agents '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@types/node': specifier: ^22.5.5 version: 22.19.1 @@ -1528,7 +1528,7 @@ importers: version: link:../../agents '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -1565,7 +1565,7 @@ importers: version: link:../test '@livekit/rtc-node': specifier: 'catalog:' - version: 0.13.34 + version: 1.1.0 '@microsoft/api-extractor': specifier: ^7.58.12 version: 7.58.12(@types/node@25.6.0) @@ -2482,6 +2482,10 @@ packages: resolution: {integrity: sha512-22EWklyNUGJ5EamOcC/eUnLy0Gu3SpzxzUkuXCixCequwNDmj8PGjKgtpWThlqTWkzvh9xsAKmQoDWhlXhkrzA==} engines: {node: '>= 18'} + '@livekit/rtc-node@1.1.0': + resolution: {integrity: sha512-KFNvPz7qQSJBTttDukanhMJ+bCNxj6O3A3fXl6GJvOEamJJk7GqKCWgAdLW3JLjuT1pPQDyQPJJnYdR25N+MXA==} + engines: {node: '>= 18'} + '@livekit/throws-transformer@0.1.8': resolution: {integrity: sha512-AaSwQfIaG6YArKOCO+5/DgI8HfL19oHdrkyI0LJJVCqGeZXzPsZIO/Nm/SKUHbT1ERGhF5c34X8CcVWeOePpIQ==} hasBin: true @@ -6260,10 +6264,10 @@ snapshots: '@livekit/plugins-krisp-viva-aarch64-unknown-linux-gnu@0.2.0': optional: true - '@livekit/plugins-krisp-viva-internal@0.2.0(@livekit/agents@agents)(@livekit/rtc-node@0.13.34)': + '@livekit/plugins-krisp-viva-internal@0.2.0(@livekit/agents@agents)(@livekit/rtc-node@1.1.0)': dependencies: '@livekit/agents': link:agents - '@livekit/rtc-node': 0.13.34 + '@livekit/rtc-node': 1.1.0 ffi-rs: 1.3.2 pino: 9.6.0 pino-pretty: 13.0.0 @@ -6326,6 +6330,15 @@ snapshots: pino: 9.6.0 pino-pretty: 13.0.0 + '@livekit/rtc-node@1.1.0': + dependencies: + '@datastructures-js/deque': 1.0.8 + '@livekit/mutex': 1.1.1 + '@livekit/rtc-ffi-bindings': 0.12.73 + '@livekit/typed-emitter': 3.0.0 + pino: 9.6.0 + pino-pretty: 13.0.0 + '@livekit/throws-transformer@0.1.8(typescript@5.9.3)': dependencies: glob: 13.0.6 diff --git a/pnpm-workspace.yaml b/pnpm-workspace.yaml index 985da184df..bce6271e17 100644 --- a/pnpm-workspace.yaml +++ b/pnpm-workspace.yaml @@ -28,7 +28,7 @@ minimumReleaseAgeExclude: - sharp@0.35.4 catalog: - '@livekit/rtc-node': ^0.13.34 + '@livekit/rtc-node': ^1.1.0 ws: ^8.21.0 '@types/ws': ^8.5.10 From f92c389939f0a4fefae8466d3b6690814358a7b5 Mon Sep 17 00:00:00 2001 From: David Zhao Date: Fri, 25 Sep 2026 19:41:05 -0700 Subject: [PATCH 2/4] review: one RPC interceptor per job, installed by the job whenever its room connects - an interceptor is bound to the job it is installed for, one per job, so a process hosting several jobs routes each participant's invocations to its own job rather than the last one installed - the job installs on its room's connected transition, whoever connected it: ctx.connect(), a session, or the entrypoint connecting ctx.room itself. RoomIO no longer installs; the job covers it Co-Authored-By: Claude Fable 5.1 --- agents/src/job.test.ts | 38 ++++++++++++++++++- agents/src/job.ts | 11 +++++- agents/src/telemetry/rpc.test.ts | 46 +++++++++++++++-------- agents/src/telemetry/rpc.ts | 48 +++++++++++++++++------- agents/src/voice/room_io/room_io.test.ts | 47 +---------------------- agents/src/voice/room_io/room_io.ts | 5 --- 6 files changed, 111 insertions(+), 84 deletions(-) diff --git a/agents/src/job.test.ts b/agents/src/job.test.ts index 089931cd77..5ce6adccaa 100644 --- a/agents/src/job.test.ts +++ b/agents/src/job.test.ts @@ -1,13 +1,14 @@ // SPDX-FileCopyrightText: 2026 LiveKit, Inc. // // SPDX-License-Identifier: Apache-2.0 -import { RoomEvent } from '@livekit/rtc-node'; +import { ConnectionState, RoomEvent } from '@livekit/rtc-node'; import type { Room } from '@livekit/rtc-node'; import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; import type { InferenceExecutor } from './ipc/inference_executor.js'; import { JobContext, type JobProcess, type RunningJobInfo } from './job.js'; import { log } from './log.js'; import { SimulationContext, parseSimulationDispatch } from './simulation.js'; +import * as rpcTracing from './telemetry/rpc.js'; const { deleteRoomMock, roomServiceClientMock, setupCloudTracerMock } = vi.hoisted(() => ({ deleteRoomMock: vi.fn(async () => {}), @@ -470,3 +471,38 @@ describe('JobContext observability URL', () => { expect(setupCloudTracerMock).not.toHaveBeenCalled(); }); }); + +describe('JobContext rpc tracing', () => { + it('installs rpc tracing whenever its room connects, however it was connected', () => { + // the entrypoint may connect ctx.room itself, with no session and no ctx.connect(): the + // job listens to its room, so the hook is there as long as the room is connected + const install = vi + .spyOn(rpcTracing, 'install') + .mockReturnValue({} as unknown as ReturnType); + try { + const ctx = createJobContext(); + const room = ctx.room as unknown as { + on: ReturnType; + isConnected: boolean; + localParticipant?: unknown; + }; + const onConnectionState = room.on.mock.calls.find( + ([event]) => event === RoomEvent.ConnectionStateChanged, + )?.[1] as ((state: ConnectionState) => void) | undefined; + expect(onConnectionState).toBeDefined(); + + onConnectionState!(ConnectionState.CONN_DISCONNECTED); + expect(install).not.toHaveBeenCalled(); + + room.isConnected = true; + room.localParticipant = { identity: 'agent' }; + onConnectionState!(ConnectionState.CONN_CONNECTED); + expect(install).toHaveBeenCalledWith(room.localParticipant, ctx); + + onConnectionState!(ConnectionState.CONN_CONNECTED); // a reconnect installs again + expect(install).toHaveBeenCalledTimes(2); + } finally { + install.mockRestore(); + } + }); +}); diff --git a/agents/src/job.ts b/agents/src/job.ts index 7b0bb3c8b0..651076b439 100644 --- a/agents/src/job.ts +++ b/agents/src/job.ts @@ -10,7 +10,7 @@ import type { Room, RtcConfiguration, } from '@livekit/rtc-node'; -import { ParticipantKind, RoomEvent, TrackKind } from '@livekit/rtc-node'; +import { ConnectionState, ParticipantKind, RoomEvent, TrackKind } from '@livekit/rtc-node'; import { ThrowsPromise } from '@livekit/throws-transformer/throws'; import type { Context } from '@opentelemetry/api'; import { RoomServiceClient } from 'livekit-server-sdk'; @@ -249,6 +249,14 @@ export class JobContext> { // its lk.simulator attribute, and gating on simulationContext() here would // miss user-token runs where no dispatch rides the job. this.#room.on(RoomEvent.ParticipantDisconnected, this.onParticipantDisconnected); + // RPC tracing goes in whenever the room is connected, whether by connect() below, by a + // session, or by the entrypoint connecting ctx.room itself. On the SDK's event path the job + // is passed explicitly: there is no AsyncLocalStorage there + this.#room.on(RoomEvent.ConnectionStateChanged, (state: ConnectionState) => { + if (state === ConnectionState.CONN_CONNECTED && this.#room.isConnected) { + rpcTracing.install(this.#room.localParticipant, this as JobContext as JobContext); + } + }); this.#logger = log().child({ jobId: this.#info.job.id, 'lk.pii.room_name': this.#info.job.room?.name, @@ -514,7 +522,6 @@ export class JobContext> { jobCtx: this as JobContext as JobContext, }, ); - rpcTracing.install(this.#room.localParticipant, this as JobContext as JobContext); this.#onConnect(); this.#room.remoteParticipants.forEach(this.onParticipantConnected); diff --git a/agents/src/telemetry/rpc.test.ts b/agents/src/telemetry/rpc.test.ts index deae9916f1..778666bad2 100644 --- a/agents/src/telemetry/rpc.test.ts +++ b/agents/src/telemetry/rpc.test.ts @@ -27,12 +27,7 @@ import { InMemorySpanExporter, SimpleSpanProcessor } from '@opentelemetry/sdk-tr import { NodeTracerProvider } from '@opentelemetry/sdk-trace-node'; import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; import { type JobContext, getJobContext, runWithJobContext } from '../job.js'; -import { - MAX_PAYLOAD_ATTR_LEN, - TracingRpcInterceptor, - install, - interceptor as singleton, -} from './rpc.js'; +import { MAX_PAYLOAD_ATTR_LEN, TracingRpcInterceptor, install } from './rpc.js'; import * as traceTypes from './trace_types.js'; import { setTracerProvider, tracer } from './traces.js'; @@ -254,10 +249,10 @@ describe.sequential('rpc tracing', () => { _primaryAgentSession: { rootSpanContext: trace.setSpan(ROOT_CONTEXT, sessionRoot) }, job: { id: 'AJ_test' }, } as unknown as JobContext; - install(fakeParticipant(vi.fn()), job); + const installed = install(fakeParticipant(vi.fn()), job); let seenJob: JobContext | undefined; - await interceptor.interceptIncoming(invocation(), async () => { + await installed.interceptIncoming(invocation(), async () => { seenJob = getJobContext(false); return ''; }); @@ -268,17 +263,36 @@ describe.sequential('rpc tracing', () => { expect(seenJob).toBe(job); // the handler itself runs inside the job too }); - it('installs the one interceptor once per participant', () => { + it('installs one interceptor per job, the same one on every connect', () => { + const jobA = { job: { id: 'AJ_a' } } as unknown as JobContext; + const jobB = { job: { id: 'AJ_b' } } as unknown as JobContext; const addRpcInterceptor = vi.fn(); const participant = fakeParticipant(addRpcInterceptor); - install(participant); - install(participant); - // the same singleton each time: the SDK keeps one registration per instance - expect(addRpcInterceptor).toHaveBeenCalledTimes(2); - expect(addRpcInterceptor.mock.calls[0]![0]).toBe(singleton); - expect(addRpcInterceptor.mock.calls[1]![0]).toBe(singleton); + const first = install(participant, jobA); + const again = install(participant, jobA); // a reconnect + const other = install(participant, jobB); + // the same instance for a job each time: the SDK keeps one registration per instance + expect(again).toBe(first); + expect(other).not.toBe(first); + expect(addRpcInterceptor.mock.calls.map(([i]) => i)).toEqual([first, first, other]); // no participant yet (the room is not connected): nothing to install on - expect(() => install(undefined)).not.toThrow(); + expect(() => install(undefined, jobA)).not.toThrow(); + }); + + it('routes an invocation to the job its participant was installed for', async () => { + // two jobs in one process: each participant's invocations run under their own job, not + // the last one installed + const jobA = { job: { id: 'AJ_a' } } as unknown as JobContext; + const jobB = { job: { id: 'AJ_b' } } as unknown as JobContext; + const forA = install(fakeParticipant(vi.fn()), jobA); + install(fakeParticipant(vi.fn()), jobB); + + let seenJob: JobContext | undefined; + await forA.interceptIncoming(invocation(), async () => { + seenJob = getJobContext(false); + return ''; + }); + expect(seenJob).toBe(jobA); }); }); diff --git a/agents/src/telemetry/rpc.ts b/agents/src/telemetry/rpc.ts index 5bfe0c8f9e..c8e47e4ff8 100644 --- a/agents/src/telemetry/rpc.ts +++ b/agents/src/telemetry/rpc.ts @@ -8,8 +8,8 @@ * The room SDK exposes an `RpcInterceptor` hook (`@livekit/rtc-node` with * `LocalParticipant.addRpcInterceptor`) that wraps every call made through `performRpc` and * every invocation dispatched to a registered handler. This module installs one interceptor per - * local participant that turns each call into a span following the OpenTelemetry RPC semantic - * conventions: + * job on its local participant that turns each call into a span following the OpenTelemetry RPC + * semantic conventions: * * - `rpc_call` (`SpanKind.CLIENT`) for outgoing calls, under whatever span is current where the * call is made (an RPC issued from a tool nests under `function_tool`); @@ -64,8 +64,13 @@ function toError(error: unknown): Error { return error instanceof Error ? error : new Error(String(error)); } -/** An `RpcInterceptor` emitting `rpc_call` / `rpc_handler` spans. */ +/** + * An `RpcInterceptor` emitting `rpc_call` / `rpc_handler` spans. Bound to the job it was + * installed for: incoming invocations arrive outside that job's context (see interceptIncoming). + */ export class TracingRpcInterceptor implements RpcInterceptor { + constructor(private readonly job?: JobContext) {} + async interceptOutgoing(call: RpcCallInfo, next: OutgoingRpcNext): Promise { const attributes: Attributes = { [traceTypes.ATTR_RPC_METHOD]: call.method, @@ -97,9 +102,9 @@ export class TracingRpcInterceptor implements RpcInterceptor { async interceptIncoming(invocation: RpcInvocationData, next: IncomingRpcNext): Promise { // the SDK dispatches invocations from its FFI event path, outside the job's - // AsyncLocalStorage: restore the job captured at install so the session and job roots - // resolve (and the handler sees getJobContext(), as a Python handler does) - const job = getJobContext(false) ?? installedJob; + // AsyncLocalStorage: restore the job this interceptor was installed for so the session and + // job roots resolve (and the handler sees getJobContext(), as a Python handler does) + const job = getJobContext(false) ?? this.job; if (job !== undefined && getJobContext(false) === undefined) { return runWithJobContext(job, () => this.handleIncoming(invocation, next)); } @@ -145,17 +150,32 @@ export class TracingRpcInterceptor implements RpcInterceptor { } } -/** The one interceptor every participant gets; the SDK dedups registrations by identity. */ -export const interceptor = new TracingRpcInterceptor(); +/** One interceptor per job, so a process hosting several jobs routes each invocation to its own. */ +const interceptors = new WeakMap(); +let joblessInterceptor: TracingRpcInterceptor | undefined; -/** The job the interceptor was installed for: incoming invocations arrive outside its context. */ -let installedJob: JobContext | undefined; +/** The interceptor for `job` (or the one shared outside any job), created on first use. */ +function interceptorFor(job: JobContext | undefined): TracingRpcInterceptor { + if (job === undefined) return (joblessInterceptor ??= new TracingRpcInterceptor()); + let interceptor = interceptors.get(job); + if (!interceptor) { + interceptor = new TracingRpcInterceptor(job); + interceptors.set(job, interceptor); + } + return interceptor; +} /** - * Trace RPCs on `localParticipant`. Idempotent: the SDK keeps one registration per interceptor - * instance. A missing participant (the room is not connected yet) installs nothing. + * Trace RPCs on `localParticipant` for `jobCtx` (else the current job). Idempotent: one + * interceptor per job, and the SDK keeps one registration per interceptor instance, so every + * connect and reconnect may install again. A missing participant (the room is not connected + * yet) installs nothing. Returns the interceptor installed. */ -export function install(localParticipant: LocalParticipant | undefined, jobCtx?: JobContext): void { - installedJob = jobCtx ?? getJobContext(false) ?? installedJob; +export function install( + localParticipant: LocalParticipant | undefined, + jobCtx?: JobContext, +): TracingRpcInterceptor { + const interceptor = interceptorFor(jobCtx ?? getJobContext(false)); localParticipant?.addRpcInterceptor(interceptor); + return interceptor; } diff --git a/agents/src/voice/room_io/room_io.test.ts b/agents/src/voice/room_io/room_io.test.ts index 050d48031c..3c3bfb466a 100644 --- a/agents/src/voice/room_io/room_io.test.ts +++ b/agents/src/voice/room_io/room_io.test.ts @@ -1,7 +1,7 @@ // SPDX-FileCopyrightText: 2026 LiveKit, Inc. // // SPDX-License-Identifier: Apache-2.0 -import { AudioFrame, ConnectionState, RoomEvent } from '@livekit/rtc-node'; +import { AudioFrame } from '@livekit/rtc-node'; import { ROOT_CONTEXT, trace } from '@opentelemetry/api'; import { InMemorySpanExporter, SimpleSpanProcessor } from '@opentelemetry/sdk-trace-base'; import { NodeTracerProvider } from '@opentelemetry/sdk-trace-node'; @@ -12,7 +12,6 @@ import { RealtimeModel } from '../../llm/index.js'; import { log } from '../../log.js'; import { IdentityTransform } from '../../stream/identity_transform.js'; import { setTracerProvider, tracer } from '../../telemetry/index.js'; -import * as rpcTracing from '../../telemetry/rpc.js'; import { DEFAULT_API_CONNECT_OPTIONS } from '../../types.js'; import { AgentSessionEventTypes, CloseReason, createCloseEvent } from '../events.js'; import { AudioInput, AudioOutput, TextOutput } from '../io.js'; @@ -107,50 +106,6 @@ function createFakeSession(llm?: RealtimeModel): FakeSession { }; } -describe('RoomIO rpc tracing', () => { - afterEach(() => { - vi.restoreAllMocks(); - }); - - /** - * A session may start on a room that connects later (ctx.connect() after session.start(), or - * a room the user connects). Tracing goes in on the connected transition, not only when the - * room is already up at start(); install is idempotent so a reconnect is harmless. - */ - it('installs rpc tracing when the room connects', () => { - const install = vi.spyOn(rpcTracing, 'install').mockReturnValue(true); - const room = createFakeRoom(); - const session = createFakeSession(); - const roomIO = new RoomIO({ - agentSession: session as unknown as RoomIOArgs['agentSession'], - room: room as unknown as RoomIOArgs['room'], - inputOptions: { audioEnabled: false, textEnabled: false }, - outputOptions: { audioEnabled: false, transcriptionEnabled: false }, - }); - - try { - roomIO.start(); - const onConnectionState = room.on.mock.calls.find( - ([event]) => event === RoomEvent.ConnectionStateChanged, - )?.[1] as ((state: ConnectionState) => void) | undefined; - expect(onConnectionState).toBeDefined(); - - onConnectionState!(ConnectionState.CONN_DISCONNECTED); - expect(install).not.toHaveBeenCalled(); - - room.isConnected = true; - onConnectionState!(ConnectionState.CONN_CONNECTED); - expect(install).toHaveBeenCalledTimes(1); - expect(install).toHaveBeenCalledWith(room.localParticipant, undefined); - - onConnectionState!(ConnectionState.CONN_CONNECTED); // reconnected - expect(install).toHaveBeenCalledTimes(2); // same singleton each time; the SDK dedups - } finally { - void roomIO.close(); - } - }); -}); - describe('RoomIO agent state attributes', () => { it('handles a failed update and publishes later state changes', async () => { const error = new Error('attribute update failed'); diff --git a/agents/src/voice/room_io/room_io.ts b/agents/src/voice/room_io/room_io.ts index e85c26347a..d545569554 100644 --- a/agents/src/voice/room_io/room_io.ts +++ b/agents/src/voice/room_io/room_io.ts @@ -24,7 +24,6 @@ import { RealtimeModel } from '../../llm/index.js'; import { log } from '../../log.js'; import { IdentityTransform } from '../../stream/identity_transform.js'; import { participantAttributes, traceTypes, tracer } from '../../telemetry/index.js'; -import * as rpcTracing from '../../telemetry/rpc.js'; import { DEFAULT_API_CONNECT_OPTIONS } from '../../types.js'; import { Future, IdleTimeoutError, Task, waitForAbort, waitUntilTimeout } from '../../utils.js'; import { type AgentSession } from '../agent_session.js'; @@ -260,10 +259,6 @@ export class RoomIO { [traceTypes.ATTR_CONNECTION_STATE]: ConnectionState[state] ?? String(state), }); if (state === ConnectionState.CONN_CONNECTED && this.room.isConnected) { - // on every connect and reconnect; install is idempotent (one interceptor instance, - // deduped by the SDK), so JobContext.connect() installing too is fine. The job is passed - // explicitly: this runs on the SDK's event path, outside the job's AsyncLocalStorage - rpcTracing.install(this.room.localParticipant, this.jobContext); if (!this.roomConnectedFuture.done) { this.roomConnectedFuture.resolve(); } From cd106047ad90ce5f59aa6f07505a90e57655d152 Mon Sep 17 00:00:00 2001 From: David Zhao Date: Fri, 25 Sep 2026 20:15:07 -0700 Subject: [PATCH 3/4] review: a participant without the interceptor hook installs no rpc tracing The warm transfer lifecycle tests stand in a local participant without addRpcInterceptor; the install threw inside the room's connection handler and stalled the session. Tracing steps aside instead: an older room SDK or a stand-in participant gets no spans, never a broken connection. Co-Authored-By: Claude Fable 5.1 --- agents/src/telemetry/rpc.test.ts | 3 +++ agents/src/telemetry/rpc.ts | 7 +++++-- 2 files changed, 8 insertions(+), 2 deletions(-) diff --git a/agents/src/telemetry/rpc.test.ts b/agents/src/telemetry/rpc.test.ts index 778666bad2..bf1115894d 100644 --- a/agents/src/telemetry/rpc.test.ts +++ b/agents/src/telemetry/rpc.test.ts @@ -278,6 +278,9 @@ describe.sequential('rpc tracing', () => { // no participant yet (the room is not connected): nothing to install on expect(() => install(undefined, jobA)).not.toThrow(); + // a participant without the hook (an older SDK, a stand-in): tracing steps aside rather + // than break the connection it is installed from + expect(() => install({ identity: 'agent' } as unknown as LocalParticipant, jobA)).not.toThrow(); }); it('routes an invocation to the job its participant was installed for', async () => { diff --git a/agents/src/telemetry/rpc.ts b/agents/src/telemetry/rpc.ts index c8e47e4ff8..51ed7fe1a2 100644 --- a/agents/src/telemetry/rpc.ts +++ b/agents/src/telemetry/rpc.ts @@ -169,13 +169,16 @@ function interceptorFor(job: JobContext | undefined): TracingRpcInterceptor { * Trace RPCs on `localParticipant` for `jobCtx` (else the current job). Idempotent: one * interceptor per job, and the SDK keeps one registration per interceptor instance, so every * connect and reconnect may install again. A missing participant (the room is not connected - * yet) installs nothing. Returns the interceptor installed. + * yet), or one without the interceptor hook (an older room SDK, a stand-in in tests), installs + * nothing: tracing must never be what breaks a connection. Returns the interceptor. */ export function install( localParticipant: LocalParticipant | undefined, jobCtx?: JobContext, ): TracingRpcInterceptor { const interceptor = interceptorFor(jobCtx ?? getJobContext(false)); - localParticipant?.addRpcInterceptor(interceptor); + if (typeof localParticipant?.addRpcInterceptor === 'function') { + localParticipant.addRpcInterceptor(interceptor); + } return interceptor; } From 47f0ad01b4f1fb70f57b95de153f856d0cb62a37 Mon Sep 17 00:00:00 2001 From: David Zhao Date: Fri, 25 Sep 2026 22:38:59 -0700 Subject: [PATCH 4/4] room_io: drop a no-op reshuffle left from the removed interceptor install Co-Authored-By: Claude Fable 5.1 --- agents/src/voice/room_io/room_io.ts | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) diff --git a/agents/src/voice/room_io/room_io.ts b/agents/src/voice/room_io/room_io.ts index d545569554..3d9df3df24 100644 --- a/agents/src/voice/room_io/room_io.ts +++ b/agents/src/voice/room_io/room_io.ts @@ -258,10 +258,12 @@ export class RoomIO { this.agentSession._addSessionEvent?.('connection_state_changed', { [traceTypes.ATTR_CONNECTION_STATE]: ConnectionState[state] ?? String(state), }); - if (state === ConnectionState.CONN_CONNECTED && this.room.isConnected) { - if (!this.roomConnectedFuture.done) { - this.roomConnectedFuture.resolve(); - } + if ( + state === ConnectionState.CONN_CONNECTED && + this.room.isConnected && + !this.roomConnectedFuture.done + ) { + this.roomConnectedFuture.resolve(); } };