Skip to content

Commit aaf70fd

Browse files
MariefayTrigger.dev RepoOps
authored andcommitted
feat(webapp): progressive chunked loading for the run trace view
**Progressive trace loading in the run view** The run trace view now loads large traces progressively. Previously it fetched a run's entire trace in a single query and built the whole tree before anything appeared, and traces past a fixed size were refused with a message. Now the first slice of the trace renders almost immediately and the rest streams in the background as you watch, so the view feels responsive regardless of trace size. Much larger traces are viewable than before, and a trace too large to display shows a clear message instead of failing. Mono-RevId: 36800dc83b61a8d6c0a45c6f3090fd1a09df76ee
1 parent de97106 commit aaf70fd

26 files changed

Lines changed: 3560 additions & 108 deletions
Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,6 @@
1+
---
2+
area: webapp
3+
type: improvement
4+
---
5+
6+
Run traces now appear almost immediately and load progressively as you watch, and much larger traces are viewable than before.

‎apps/webapp/app/components/primitives/TreeView/TreeView.tsx‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -24,6 +24,7 @@ export type TreeViewProps<TData> = {
2424
parentRef?: MutableRefObject<HTMLElement | null>;
2525
scrollRef?: MutableRefObject<HTMLElement | null>;
2626
onScroll?: (scrollTop: number) => void;
27+
staticRowHeight?: boolean;
2728
} & Pick<UseTreeStateOutput, "getTreeProps" | "getNodeProps">;
2829

2930
export function TreeView<TData>({
@@ -37,6 +38,7 @@ export function TreeView<TData>({
3738
virtualizer,
3839
parentRef,
3940
scrollRef,
41+
staticRowHeight = false,
4042
onScroll,
4143
}: TreeViewProps<TData>) {
4244
useEffect(() => {
@@ -121,7 +123,7 @@ export function TreeView<TData>({
121123
<div
122124
key={node.id}
123125
data-index={virtualItem.index}
124-
ref={virtualizer.measureElement}
126+
ref={staticRowHeight ? undefined : virtualizer.measureElement}
125127
className="overflow-clip"
126128
{...getNodeProps(node.id)}
127129
>

‎apps/webapp/app/env.server.ts‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2494,6 +2494,8 @@ const EnvironmentSchema = z
24942494
EVENTS_CLICKHOUSE_MAX_TRACE_SUMMARY_VIEW_COUNT: z.coerce.number().int().default(25_000),
24952495
EVENTS_CLICKHOUSE_MAX_TRACE_DETAILED_SUMMARY_VIEW_COUNT: z.coerce.number().int().default(5_000),
24962496
EVENTS_CLICKHOUSE_MAX_LIVE_RELOADING_SETTING: z.coerce.number().int().default(2000),
2497+
EVENTS_CLICKHOUSE_TRACE_CHUNK_SIZE: z.coerce.number().int().positive().default(1_000),
2498+
EVENTS_CLICKHOUSE_MAX_TRACE_VIEW_COUNT: z.coerce.number().int().positive().default(250_000),
24972499

24982500
// OTLP ingest transform worker pool (opt-in). When enabled, decode/convert/enrich run in a
24992501
// worker_threads pool instead of the request event loop; the single consolidated insert path
Lines changed: 281 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,281 @@
1+
import { useEffect, useMemo, useRef, useState } from "react";
2+
import type { SpanOverride, TraceChunkCursor } from "~/v3/eventRepository/eventRepository.types";
3+
import { TraceChunkAssembler } from "~/v3/eventRepository/traceChunkAssembler";
4+
import {
5+
applyAncestorOverrides,
6+
buildTraceView,
7+
type BuildTraceViewOptions,
8+
type TraceViewEvent,
9+
} from "~/v3/eventRepository/traceViewBuilder";
10+
11+
const MAX_CHUNK_ATTEMPTS = 3;
12+
const RETRY_BACKOFF_MS = 400;
13+
const MAX_PROGRESSIVE_SPANS = 250_000;
14+
const BACKGROUND_CHUNK_SIZE = 10_000;
15+
const REBUILD_COALESCE_MS = 100;
16+
17+
type WireChunkEvent = {
18+
spanId: string;
19+
parentSpanId: string;
20+
runId: string;
21+
startTime: string;
22+
startTimeNano: string;
23+
duration: number;
24+
status: string;
25+
kind: string;
26+
message: string;
27+
metadata: string;
28+
};
29+
30+
type WireChunkResponse = {
31+
events: WireChunkEvent[];
32+
nextCursor: TraceChunkCursor | null;
33+
hasMore: boolean;
34+
};
35+
36+
type ProgressiveMeta = {
37+
firstEvents: WireChunkEvent[];
38+
supplementaryFirstEvents?: WireChunkEvent[];
39+
nextCursor: TraceChunkCursor | null;
40+
hasMore: boolean;
41+
buildOptions: BuildTraceViewOptions;
42+
showDebug: boolean;
43+
totalSpans?: number;
44+
maxSpans?: number;
45+
};
46+
47+
export type ProgressiveTraceInput = {
48+
events: TraceViewEvent[];
49+
duration: number;
50+
rootStartedAt: Date | string | undefined;
51+
rootSpanStatus: "executing" | "completed" | "failed";
52+
overridesBySpanId?: Record<string, SpanOverride>;
53+
linkedRunIdBySpanId?: Record<string, string>;
54+
progressive?: ProgressiveMeta | null;
55+
};
56+
57+
export type ProgressiveTraceState = {
58+
events: TraceViewEvent[];
59+
duration: number;
60+
rootStartedAt: Date | string | undefined;
61+
rootSpanStatus: "executing" | "completed" | "failed";
62+
overridesBySpanId: Record<string, SpanOverride>;
63+
linkedRunIdBySpanId: Record<string, string>;
64+
isComplete: boolean;
65+
isTruncated: boolean;
66+
};
67+
68+
function toChunkEvent(event: WireChunkEvent) {
69+
return { ...event, startTime: new Date(event.startTime) };
70+
}
71+
72+
function initialState(trace: ProgressiveTraceInput): ProgressiveTraceState {
73+
return {
74+
events: trace.events,
75+
duration: trace.duration,
76+
rootStartedAt: trace.rootStartedAt,
77+
rootSpanStatus: trace.rootSpanStatus,
78+
overridesBySpanId: trace.overridesBySpanId ?? {},
79+
linkedRunIdBySpanId: trace.linkedRunIdBySpanId ?? {},
80+
isComplete: !trace.progressive?.hasMore,
81+
isTruncated: false,
82+
};
83+
}
84+
85+
export function useProgressiveTrace(
86+
trace: ProgressiveTraceInput,
87+
chunkPath: string,
88+
errorsOnly: boolean
89+
): ProgressiveTraceState {
90+
const [state, setState] = useState<ProgressiveTraceState>(() => initialState(trace));
91+
const staticState = useMemo(() => (trace.progressive ? null : initialState(trace)), [trace]);
92+
93+
const progressive = trace.progressive ?? null;
94+
const identity = progressive
95+
? `${progressive.buildOptions.rootSpanId}:${progressive.firstEvents.length}:${
96+
progressive.nextCursor?.spanId ?? ""
97+
}:${progressive.hasMore}:${progressive.totalSpans ?? ""}`
98+
: `static:${trace.events.length}`;
99+
100+
const latestTraceRef = useRef(trace);
101+
latestTraceRef.current = trace;
102+
const assemblerRef = useRef<TraceChunkAssembler | null>(null);
103+
const errorsFetchedRef = useRef(false);
104+
105+
function rebuild(meta: ProgressiveMeta) {
106+
const assembler = assemblerRef.current;
107+
if (!assembler) return;
108+
const { spans, overridesBySpanId } = applyAncestorOverrides(assembler.spans);
109+
const view = buildTraceView(spans, meta.buildOptions);
110+
setState((prev) => ({
111+
events: view.events,
112+
duration: view.duration,
113+
rootStartedAt: view.rootStartedAt,
114+
rootSpanStatus: view.rootSpanStatus,
115+
overridesBySpanId,
116+
linkedRunIdBySpanId: view.linkedRunIdBySpanId,
117+
isComplete: prev.isComplete,
118+
isTruncated: prev.isTruncated,
119+
}));
120+
}
121+
122+
useEffect(() => {
123+
const current = latestTraceRef.current;
124+
setState(initialState(current));
125+
errorsFetchedRef.current = false;
126+
127+
const meta = current.progressive;
128+
if (!meta) {
129+
assemblerRef.current = null;
130+
return;
131+
}
132+
133+
const assembler = new TraceChunkAssembler();
134+
assembler.mergeChunk(meta.firstEvents.map(toChunkEvent));
135+
if (meta.supplementaryFirstEvents?.length) {
136+
assembler.mergeChunk(meta.supplementaryFirstEvents.map(toChunkEvent), {
137+
source: "deeplink",
138+
});
139+
}
140+
assemblerRef.current = assembler;
141+
142+
if (!meta.hasMore || !meta.nextCursor) {
143+
return;
144+
}
145+
146+
let cancelled = false;
147+
let cursor: TraceChunkCursor | null = meta.nextCursor;
148+
const maxSpans = meta.maxSpans ?? MAX_PROGRESSIVE_SPANS;
149+
150+
let rebuildTimer: ReturnType<typeof setTimeout> | null = null;
151+
let rebuildPending = false;
152+
const requestRebuild = () => {
153+
rebuildPending = true;
154+
if (rebuildTimer === null) {
155+
rebuildTimer = setTimeout(() => {
156+
rebuildTimer = null;
157+
rebuildPending = false;
158+
rebuild(meta!);
159+
}, REBUILD_COALESCE_MS);
160+
}
161+
};
162+
const flushRebuild = () => {
163+
if (rebuildTimer !== null) {
164+
clearTimeout(rebuildTimer);
165+
rebuildTimer = null;
166+
}
167+
if (rebuildPending) {
168+
rebuildPending = false;
169+
rebuild(meta!);
170+
}
171+
};
172+
173+
const markComplete = (truncated = false) => {
174+
flushRebuild();
175+
setState((prev) => ({
176+
...prev,
177+
isComplete: true,
178+
isTruncated: prev.isTruncated || truncated,
179+
}));
180+
};
181+
182+
async function loadRemaining() {
183+
while (!cancelled && cursor) {
184+
let data: WireChunkResponse | null = null;
185+
for (let attempt = 0; attempt < MAX_CHUNK_ATTEMPTS && !cancelled; attempt++) {
186+
data = await fetchChunk(chunkPath, meta!.showDebug, {
187+
cursor,
188+
limit: BACKGROUND_CHUNK_SIZE,
189+
});
190+
if (data) break;
191+
if (attempt < MAX_CHUNK_ATTEMPTS - 1) {
192+
await new Promise((r) => setTimeout(r, RETRY_BACKOFF_MS * (attempt + 1)));
193+
}
194+
}
195+
if (cancelled) {
196+
return;
197+
}
198+
if (!data) {
199+
markComplete();
200+
return;
201+
}
202+
203+
assembler.mergeChunk(data.events.map(toChunkEvent));
204+
requestRebuild();
205+
cursor = data.hasMore ? data.nextCursor : null;
206+
207+
if (assembler.size > maxSpans) {
208+
markComplete(true);
209+
return;
210+
}
211+
if (!cursor) {
212+
markComplete();
213+
return;
214+
}
215+
}
216+
}
217+
218+
void loadRemaining();
219+
220+
return () => {
221+
cancelled = true;
222+
if (rebuildTimer !== null) {
223+
clearTimeout(rebuildTimer);
224+
rebuildTimer = null;
225+
}
226+
};
227+
// eslint-disable-next-line react-hooks/exhaustive-deps -- re-run only on a new loader payload
228+
}, [identity, chunkPath]);
229+
230+
useEffect(() => {
231+
const meta = latestTraceRef.current.progressive;
232+
if (!errorsOnly || !meta || errorsFetchedRef.current) {
233+
return;
234+
}
235+
236+
let cancelled = false;
237+
(async () => {
238+
const data = await fetchChunk(chunkPath, meta.showDebug, { filter: "errors" });
239+
if (cancelled || !data || !assemblerRef.current) return;
240+
errorsFetchedRef.current = true;
241+
assemblerRef.current.mergeChunk(data.events.map(toChunkEvent), { source: "errors" });
242+
rebuild(meta);
243+
})();
244+
245+
return () => {
246+
cancelled = true;
247+
};
248+
// eslint-disable-next-line react-hooks/exhaustive-deps -- reacts to errorsOnly + payload
249+
}, [errorsOnly, identity, chunkPath]);
250+
251+
return staticState ?? state;
252+
}
253+
254+
async function fetchChunk(
255+
chunkPath: string,
256+
showDebug: boolean,
257+
params: { cursor?: TraceChunkCursor; filter?: "errors"; limit?: number }
258+
): Promise<WireChunkResponse | null> {
259+
const url = new URL(chunkPath, window.location.origin);
260+
if (params.cursor) {
261+
url.searchParams.set("cursorStartTime", params.cursor.startTime);
262+
url.searchParams.set("cursorSpanId", params.cursor.spanId);
263+
}
264+
if (params.limit) {
265+
url.searchParams.set("limit", String(params.limit));
266+
}
267+
if (params.filter) {
268+
url.searchParams.set("filter", params.filter);
269+
}
270+
if (showDebug) {
271+
url.searchParams.set("debug", "1");
272+
}
273+
274+
try {
275+
const response = await fetch(url.toString(), { headers: { accept: "application/json" } });
276+
if (!response.ok) return null;
277+
return (await response.json()) as WireChunkResponse;
278+
} catch {
279+
return null;
280+
}
281+
}

0 commit comments

Comments
 (0)