diff --git a/.changeset/rpc-tracing-spans.md b/.changeset/rpc-tracing-spans.md new file mode 100644 index 000000000..6cd137aee --- /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.test.ts b/agents/src/job.test.ts index 089931cd7..5ce6adcca 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 73ee4c439..651076b43 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'; @@ -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, @@ -248,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, diff --git a/agents/src/telemetry/index.ts b/agents/src/telemetry/index.ts index 1fd1f1ca7..b1bca8442 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 000000000..bf1115894 --- /dev/null +++ b/agents/src/telemetry/rpc.test.ts @@ -0,0 +1,301 @@ +// 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 } 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; + const installed = install(fakeParticipant(vi.fn()), job); + + let seenJob: JobContext | undefined; + await installed.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 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); + 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, 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 () => { + // 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 new file mode 100644 index 000000000..51ed7fe1a --- /dev/null +++ b/agents/src/telemetry/rpc.ts @@ -0,0 +1,184 @@ +// 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 + * 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`); + * - `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. 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, + [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 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)); + } + 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, + }, + ); + } +} + +/** 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 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` 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), 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)); + if (typeof localParticipant?.addRpcInterceptor === 'function') { + localParticipant.addRpcInterceptor(interceptor); + } + return interceptor; +} diff --git a/agents/src/telemetry/startup_spans.test.ts b/agents/src/telemetry/startup_spans.test.ts index e36fea67e..9b853e4ab 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 5490d5f7c..5573535f7 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 494eddaee..4469a3dc1 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 3af76f7c3..6f7e00054 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/pnpm-lock.yaml b/pnpm-lock.yaml index 1333b95c2..0b0a017c1 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 985da184d..bce6271e1 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