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/loop-stall-stacks.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,5 @@
---
'@livekit/agents': patch
---

Event loop stalls now nest under the span that was running when the loop blocked (`function_tool`, `rpc_handler`, `on_user_turn_completed`, ...) and carry `lk.blocking.stack`, the loop thread's call stack sampled by the monitor's watchdog thread through an inspector session (at the warn threshold and again at ten times it). Sampling starts after a process's first stall, or from the start with `LIVEKIT_AGENTS_LOOP_BLOCK_STACKS=1`, never with `0`; it costs a one-time enumeration of the loaded scripts on the loop when it starts and about 1% of throughput afterwards. Not enabled while an inspector is attached to the process.
14 changes: 13 additions & 1 deletion agents/api-extractor.json
Original file line number Diff line number Diff line change
Expand Up @@ -16,5 +16,17 @@
* DEFAULT VALUE: ""
*/
"extends": "../api-extractor-shared.json",
"mainEntryPointFilePath": "./dist/index.d.ts"
"mainEntryPointFilePath": "./dist/index.d.ts",
"messages": {
"extractorMessageReporting": {
/**
* Exports are not release-tagged in this package; the warning for every untagged symbol
* only buries the real changes in the API report.
*/
"ae-missing-release-tag": {
"logLevel": "none",
"addToApiReportFile": false
}
}
}
}
1,936 changes: 403 additions & 1,533 deletions agents/etc/agents.api.md

Large diffs are not rendered by default.

180 changes: 180 additions & 0 deletions agents/src/telemetry/blocked_span_tracker.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,180 @@
// SPDX-FileCopyrightText: 2026 LiveKit, Inc.
//
// SPDX-License-Identifier: Apache-2.0
import { type Context, type Span, trace } from '@opentelemetry/api';
import type { ReadableSpan, SpanProcessor } from '@opentelemetry/sdk-trace-base';

/**
* Which span was running when the event loop stalled, without sampling a stack.
*
* Every framework span is created and ended on the main thread, so a span created before a
* stall and ended after it (or at its end, when the blocking call returned) was current while
* the loop was blocked. The innermost such span, the one created last, is where the stall hurt:
* `function_tool` for a slow tool, `rpc_handler` for a slow RPC, `on_user_turn_completed` for a
* slow hook. The Python monitor reads the same span off the blocked task's context; Node has no
* cross-thread view of a task's context, so the span processor keeps the bookkeeping instead.
*
* Creation and end are timed by the wall clock at the processor call, not by the span's own
* timestamps: a span back-dated at creation (`eou_wait` starts at the user's last speech) must
* not claim a stall that predates it.
*/
export class BlockedSpanTracker implements SpanProcessor {
/** Spans not yet ended, by span id. */
readonly #open = new Map<string, { span: Span; createdAt: number }>();
/** Spans ended recently: a stall's report runs a heartbeat after the blocking call returned. */
readonly #ended: { span: Span; createdAt: number; endedAt: number }[] = [];
readonly #retention: number;
readonly #maxEnded: number;
readonly #minRetention: number;

constructor(options: { retention?: number; maxEnded?: number; minRetention?: number } = {}) {
this.#retention = options.retention ?? 5_000;
// the count bounds the memory; it never evicts a span that ended within `minRetention`,
// long enough for the late heartbeat to report a stall the span was current for, however
// many spans ended after it before the loop yielded
this.#maxEnded = options.maxEnded ?? 256;
this.#minRetention = options.minRetention ?? 1_000;
}

onStart(span: Span): void {
this.#open.set(span.spanContext().spanId, { span, createdAt: Date.now() });
}

onEnd(span: ReadableSpan): void {
const id = span.spanContext().spanId;
const entry = this.#open.get(id);
if (!entry) return;
this.#open.delete(id);
const now = Date.now();
this.#ended.push({ span: entry.span, createdAt: entry.createdAt, endedAt: now });
this.#prune(now);
}

async forceFlush(): Promise<void> {}

async shutdown(): Promise<void> {
this.#open.clear();
this.#ended.length = 0;
}

/**
* The span that was current when the loop began blocking within `[startedAt, endedAt]` (epoch
* ms): created before the window opened and still open, or ended after its first heartbeat.
* Spans of the given names are skipped (the stall span itself). `slack` accounts for heartbeat
* timing at both ends. The late heartbeat can run well after the blocking call and its span
* ended, so requiring a span to cover the report's end would discard the actual blocker.
*
* Only spans of `traceId` qualify when given: the processor sees every span of the provider,
* and a span of an unrelated trace (an application's own HTTP request, say) that happened to
* be open must not pull the stall out of the job's trace.
*
* Among the spans that qualify, the innermost of one ancestry is the answer. When two
* operations of the same kind were both in flight (two tools, two RPC handlers), timing alone
* cannot say which one blocked: the answer is then their nearest common ancestor, or nothing,
* rather than a guess at one of them. Operations of different kinds overlap all the time (a
* user turn is open while an RPC handler runs, an audio wait while an RPC blocks) and the
* newest of them is taken to be the one that blocked. That is a guess too: a tool awaiting
* I/O can resume and block while a newer RPC handler is itself awaiting I/O, and the RPC gets
* the blame. The alternative, their common ancestor, would file every stall during a wait
* under the session, so the guess is kept.
*/
blockedSpan(
startedAt: number,
endedAt: number,
exclude: ReadonlySet<string>,
slack = 2,
traceId?: string,
): Span | undefined {
const opened = startedAt + slack;
const closed = Math.min(endedAt - slack, startedAt + slack);

@devin-ai-integration devin-ai-integration Bot Sep 25, 2026 •

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Earlier operation claims later stall

When operations run consecutively without a heartbeat, blockedSpan can select an earlier operation that already ended. Its span survives the early cutoff and outranks an older span still doing the long block, misparenting the stall.

Learn more

The monitor dates a stall from the first missed heartbeat, not from the exact start of the synchronous call. This cutoff admits every span that ends after that early point, even when another operation subsequently keeps the loop blocked for most of the report. An ended span with a later creation time can then beat the operation still running when the heartbeat finally fires. The report builder derives that approximate start from the timer lag, so it cannot distinguish consecutive blocking operations by their start times alone.

Example: A turn span starts at 0 ms. An RPC span starts at 15 ms and ends at 50 ms; the turn then blocks until 700 ms without letting the 20 ms heartbeat run. The report starts near 20 ms. Both spans qualify at the 40 ms cutoff, so the RPC claims the entire 680 ms stall even though the turn caused nearly all of it.

Recommended fix: Avoid attributing an entire late heartbeat to a span that only survived the beginning of its interval when another eligible span continued through its end. Preserve the ended-RPC case where it is the only credible blocker; add a consecutive-operations regression test.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

const candidates = new Map<string, { span: Span; createdAt: number }>();
const consider = (entry: { span: Span; createdAt: number }) => {
if (entry.createdAt > opened) return;
if (exclude.has(spanName(entry.span))) return;
if (traceId !== undefined && entry.span.spanContext().traceId !== traceId) return;
candidates.set(entry.span.spanContext().spanId, entry);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Unrelated spans capture stall parenting

With another provider span open during a job stall, blockedSpan can select it as the parent. Candidates are never restricted to the fallback job trace. The stall enters an unrelated trace instead of the job timeline.

Learn more

A span processor receives every recording span created by its provider, not only LiveKit framework spans. Custom providers can therefore contribute application or auto-instrumentation spans from unrelated traces. blockedContext receives a base context carrying the current session or job root, but blockedSpan does not use that trace identity when building candidates. A newer unrelated leaf can become the parent, and trace.setSpan then replaces the base span with that unrelated span.

Example: A background HTTP request opens span http.request in trace A while job trace B runs a blocking tool. Since http.request was created later and overlaps the stall window, the emitted event_loop_blocked span becomes its child in trace A. It was expected under the blocked tool or the session root in trace B.

Recommended fix: Pass the fallback span context's trace ID into blockedSpan and discard candidates whose spanContext().traceId differs. Apply the filter before leaf and ambiguity selection so unrelated spans cannot suppress a valid same-trace candidate.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

also from Codex:

Can we filter candidates by the session/job context’s trace ID before selecting a parent? This processor receives spans from the entire provider. A newer background HTTP span from another trace can win here, moving event_loop_blocked out of the agent’s trace.

Comment thread
davidzhao marked this conversation as resolved.
};
for (const entry of this.#open.values()) consider(entry);
for (const entry of this.#ended) {
if (entry.endedAt >= closed) consider(entry);
}
if (!candidates.size) return undefined;

// the leaves: candidates no other candidate descends from
const hasChild = new Set<string>();
for (const entry of candidates.values()) {
let parent = parentSpanId(entry.span);
while (parent !== undefined && candidates.has(parent) && !hasChild.has(parent)) {
hasChild.add(parent);
parent = parentSpanId(candidates.get(parent)!.span);
}
}
const leaves = [...candidates.entries()].filter(([id]) => !hasChild.has(id));
// ties (same millisecond) go to the later-created span, which the maps yield last
let newest = leaves[0]!;
for (const leaf of leaves) if (leaf[1].createdAt >= newest[1].createdAt) newest = leaf;
const sameKind = leaves.filter(
([, entry]) => spanName(entry.span) === spanName(newest[1].span),
);
if (sameKind.length <= 1) return newest[1].span;

@chenghao-mou chenghao-mou Sep 23, 2026 •

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

from Codex:

Differently named operations have the same ambiguity. A tool can await I/O, a newer RPC can start and await I/O, then the tool can resume and block. Both spans remain open, but this selects the idle RPC.

Should we use the common ancestor across all overlapping leaves, regardless of span name?


// indistinguishable: the nearest ancestor (among the candidates) common to all of them
const chains = sameKind.map(([id]) => {
const chain: string[] = [];
let current: string | undefined = id;
while (current !== undefined && candidates.has(current)) {
chain.push(current);
current = parentSpanId(candidates.get(current)!.span);
}
return chain;
});
const shared = chains[0]!.find((id) => chains.every((chain) => chain.includes(id)));
if (shared === undefined) return undefined;
// the common ancestor is the deepest one of that kind that is not itself ambiguous
return candidates.get(shared)!.span;
}

/**
* The context carrying {@link blockedSpan}, or undefined. Candidates are confined to `base`'s
* trace (the session's or the job's) when it carries one.
*/
blockedContext(
startedAt: number,
endedAt: number,
exclude: ReadonlySet<string>,
base: Context,
slack = 2,
): Context | undefined {
const traceId = trace.getSpanContext(base)?.traceId;
const span = this.blockedSpan(startedAt, endedAt, exclude, slack, traceId);
return span ? trace.setSpan(base, span) : undefined;
}

/** @internal test hook */
get openCount(): number {
return this.#open.size;
}

#prune(now: number): void {
const cutoff = now - this.#retention;
const settled = now - this.#minRetention;
while (
this.#ended.length &&
(this.#ended[0]!.endedAt < cutoff ||
(this.#ended.length > this.#maxEnded && this.#ended[0]!.endedAt < settled))
) {
this.#ended.shift();
Comment on lines +161 to +166

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Busy turn drops the blocking parent

If a blocking tool ends before 256 newer spans in one turn, #prune evicts its span before the late heartbeat. The stall loses its tool parent and falls back to an ancestor.

Learn more

The tracker retains ended spans because the monitor emits a stall one heartbeat after the blocking call returns. Each onEnd appends to the queue and enforces a fixed count, regardless of whether the heartbeat has reported a pending stall. Synchronous work or microtasks can end many spans before the next timer turn. When that happens, the blocking span is evicted before blockedContext looks it up.

Example: A tool blocks for 200 ms and ends. Its callback then ends 256 short spans before yielding to timers. The pending stall is eventually emitted under the agent turn instead of the blocking tool.

Recommended fix: Keep ended spans long enough for the monitor's pending stall to be reported, or snapshot the eligible parent at the late heartbeat before count-based eviction can discard it. Keep a bounded policy for spans not needed by pending reports.

Devin Review


Was this helpful? React with 👍 or 👎 to provide feedback.

}
}
}

function spanName(span: Span): string {
return (span as { name?: string }).name ?? '';
}

function parentSpanId(span: Span): string | undefined {
return (span as { parentSpanContext?: { spanId: string } }).parentSpanContext?.spanId;
}

/** The one tracker the loop monitor consults; installed on every provider the framework owns. */
export const blockedSpanTracker = new BlockedSpanTracker();
3 changes: 3 additions & 0 deletions agents/src/telemetry/loop_monitor.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -99,6 +99,7 @@ describe.sequential('event loop monitor', () => {
warnThreshold: WARN,
errorThreshold: ERROR,
tickInterval: TICK,
stacks: 'never', // sampled stacks have their own tests; reports stay synchronous here
});
monitor.onReport = (report) => reports.push(report);
sessionRoot = tracer.startSpan({ name: 'agent_session' });
Expand Down Expand Up @@ -303,6 +304,7 @@ describe.sequential('event loop monitor', () => {
errorThreshold: ERROR,
tickInterval: TICK,
emitSpans: false,
stacks: 'never',
});
const workerReports: BlockedReport[] = [];
workerMonitor.onReport = (report) => workerReports.push(report);
Expand Down Expand Up @@ -416,6 +418,7 @@ describe.sequential('event loop monitor', () => {
errorThreshold: ERROR,
tickInterval: TICK,
watchdog: false,
stacks: 'never',
});
bare.start();
expect(bare.watchdogActive).toBe(false);
Expand Down
Loading
Loading