Repository navigation
feat(telemetry): nest event loop stalls under the blocked span and sample their stacks #2541
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
852e0cf
31e8561
db9f1af
a361e7d
6ae474e
c22f2d1
51924fa
b1458ed
189f30c
0fdc80d
158e162
8cb02fe
d212698
f6d6d6a
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| 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. |
Large diffs are not rendered by default.
| 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); | ||
| 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); | ||
|
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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, Learn moreA 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. Example: A background HTTP request opens span Recommended fix: Pass the fallback span context's trace ID into Was this helpful? React with 👍 or 👎 to provide feedback.
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. also from Codex:
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; | ||
|
Member
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. from Codex:
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
Contributor
There was a problem hiding this comment. Choose a reason for hiding this commentThe 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, Learn moreThe 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. 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(); | ||
Uh oh!
There was an error while loading. Please reload this page.
There was a problem hiding this comment.
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,
blockedSpancan 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.
Was this helpful? React with 👍 or 👎 to provide feedback.