Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 5 additions & 0 deletions .changeset/rpc-tracing-spans.md
Original file line number Diff line number Diff line change
@@ -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`.
38 changes: 37 additions & 1 deletion agents/src/job.test.ts
Original file line number Diff line number Diff line change
@@ -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 () => {}),
Expand Down Expand Up @@ -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<typeof rpcTracing.install>);
try {
const ctx = createJobContext();
const room = ctx.room as unknown as {
on: ReturnType<typeof vi.fn>;
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();
}
});
});
11 changes: 10 additions & 1 deletion agents/src/job.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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';
Expand All @@ -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,
Expand Down Expand Up @@ -248,6 +249,14 @@ export class JobContext<ProcessUserData = Record<string, unknown>> {
// 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<unknown> as JobContext);
}
});
this.#logger = log().child({
jobId: this.#info.job.id,
'lk.pii.room_name': this.#info.job.room?.name,
Expand Down
1 change: 1 addition & 0 deletions agents/src/telemetry/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Loading
Loading