Skip to content

Commit 8178ba7

Browse files
MariefayTrigger.dev RepoOps
authored andcommitted
feat(webapp,core): page large run traces from the trace API
The run trace API (`GET /api/v1/runs/:runId/trace`) can page large traces. Pass `page[size]` and `page[after]` to read every span of a run at the root of its trace. Requests without these parameters are unchanged. Mono-RevId: 3a811d1a61118cfe089ba692b4edd56282674c65
1 parent b43672d commit 8178ba7

27 files changed

Lines changed: 1861 additions & 173 deletions

‎apps/webapp/app/routes/api.v1.runs.$runId.trace.ts‎

Lines changed: 119 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -1,10 +1,23 @@
1+
import { trace } from "@opentelemetry/api";
12
import { json } from "@remix-run/server-runtime";
3+
import type { RetrieveRunTracePageResponseBody } from "@trigger.dev/core/v3";
24
import { BatchId } from "@trigger.dev/core/v3/isomorphic";
35
import { z } from "zod";
46
import { $replica } from "~/db.server";
7+
import { env } from "~/env.server";
58
import { anyResource, createLoaderApiRoute } from "~/services/routeBuilders/apiBuilder.server";
9+
import { clampToEmergencySpanCap } from "~/v3/eventRepository/emergencySpanCap.server";
610
import { getEventRepositoryForStore } from "~/v3/eventRepository/index.server";
711
import { getTraceInsertedAtEnd } from "~/v3/eventRepository/traceInsertedAtBound";
12+
import { getRunTracePage } from "~/v3/eventRepository/runTracePage.server";
13+
import { FEATURE_FLAG } from "~/v3/featureFlags";
14+
import { makeFlag } from "~/v3/featureFlags.server";
15+
import { toPageSpan } from "~/v3/eventRepository/tracePage";
16+
import {
17+
DEFAULT_TRACE_PAGE_SIZE,
18+
parseTracePageRequest,
19+
type TracePageRequest,
20+
} from "~/v3/eventRepository/tracePageRequest";
821
import { getTaskEventStoreTableForRun } from "~/v3/taskEventStore.server";
922
import { findRunByIdWithMollifierFallback } from "~/v3/mollifier/readFallback.server";
1023
import { buildSyntheticTraceBody } from "~/v3/mollifier/syntheticApiResponses.server";
@@ -78,7 +91,22 @@ export const loader = createLoaderApiRoute(
7891
},
7992
},
8093
},
81-
async ({ resource: resolved, authentication }) => {
94+
async ({ resource: resolved, authentication, request }) => {
95+
const searchParams = new URL(request.url).searchParams;
96+
97+
if (
98+
hasPageParams(searchParams) &&
99+
(await isTracePagingEnabled(authentication.environment.organization.featureFlags))
100+
) {
101+
const pageRequest = parseTracePageRequest(searchParams, DEFAULT_TRACE_PAGE_SIZE);
102+
if (pageRequest.isErr()) {
103+
return json({ error: PAGE_REQUEST_ERRORS[pageRequest.error] }, { status: 400 });
104+
}
105+
if (pageRequest.value) {
106+
return tracePageResponse(resolved, authentication.environment, pageRequest.value);
107+
}
108+
}
109+
82110
if (resolved.source === "buffer") {
83111
// Buffered runs have no events ingested yet — the drainer hasn't
84112
// materialised the PG row and the worker hasn't started executing.
@@ -104,6 +132,11 @@ export const loader = createLoaderApiRoute(
104132
{ insertedAtEnd: getTraceInsertedAtEnd(run) }
105133
);
106134

135+
trace.getActiveSpan()?.setAttributes({
136+
"run.depth": run.depth,
137+
"trace.truncated": traceSummary?.isTruncated ?? false,
138+
});
139+
107140
if (!traceSummary) {
108141
return json({ error: "Trace not found" }, { status: 404 });
109142
}
@@ -116,3 +149,88 @@ export const loader = createLoaderApiRoute(
116149
);
117150
}
118151
);
152+
153+
function hasPageParams(searchParams: URLSearchParams): boolean {
154+
return searchParams.has("page[size]") || searchParams.has("page[after]");
155+
}
156+
157+
// Off: page parameters are ignored and the request gets today's tree response.
158+
function isTracePagingEnabled(orgFeatureFlags: unknown): Promise<boolean> {
159+
return makeFlag()({
160+
key: FEATURE_FLAG.publicTracePagingEnabled,
161+
defaultValue: false,
162+
overrides: (orgFeatureFlags as Record<string, unknown>) ?? {},
163+
});
164+
}
165+
166+
const PAGE_REQUEST_ERRORS = {
167+
invalid_page_size: "page[size] must be a positive integer",
168+
invalid_cursor: "page[after] is not a valid cursor; pass back pagination.next unchanged",
169+
} as const;
170+
171+
const CHILD_RUN_ERROR =
172+
"Paging is only supported for runs at the root of their trace. Request this run's trace without page parameters.";
173+
174+
// The same synthetic root span the unpaged response returns for a run still in the trigger buffer.
175+
function bufferedRunPage(
176+
run: Parameters<typeof buildSyntheticTraceBody>[0]
177+
): RetrieveRunTracePageResponseBody {
178+
if (!run.spanId) {
179+
return { data: [], attemptFailures: [], pagination: {} };
180+
}
181+
182+
const { rootSpan } = buildSyntheticTraceBody(run).trace;
183+
return { data: [toPageSpan(rootSpan)], attemptFailures: [], pagination: {} };
184+
}
185+
186+
async function tracePageResponse(
187+
resolved: ResolvedRun,
188+
environment: { id: string; organization: { id: string } },
189+
pageRequest: TracePageRequest
190+
) {
191+
// Paging reads the whole trace, so the run's span must be the trace root (no parent run or span).
192+
if (resolved.run.parentTaskRunId || resolved.run.parentSpanId) {
193+
return json({ error: CHILD_RUN_ERROR }, { status: 400 });
194+
}
195+
196+
if (resolved.source === "buffer") {
197+
return json(bufferedRunPage(resolved.run), { status: 200 });
198+
}
199+
200+
const capped = env.TRACE_VIEW_EMERGENCY_SPAN_CAP !== undefined;
201+
if (capped && pageRequest.after) {
202+
return json(
203+
{ error: "Trace paging is temporarily limited to the first page" },
204+
{ status: 503, headers: { "Retry-After": "60", "x-should-retry": "false" } }
205+
);
206+
}
207+
208+
const run = resolved.run;
209+
const page = await getRunTracePage({
210+
repository: await getEventRepositoryForStore(run.taskEventStore, environment.organization.id),
211+
storeTable: getTaskEventStoreTableForRun(run),
212+
environmentId: environment.id,
213+
run,
214+
pageRequest: { ...pageRequest, limit: clampToEmergencySpanCap(pageRequest.limit) },
215+
// Like the dashboard, stop paging under the emergency cap instead of reading the whole trace.
216+
stopAfterPage: capped,
217+
});
218+
219+
if (page.isOk()) {
220+
return json(page.value, { status: 200 });
221+
}
222+
223+
switch (page.error) {
224+
case "store_unavailable":
225+
// The SDK retries 5xx; tell it not to, so a struggling store isn't hit harder.
226+
return json(
227+
{ error: "The trace store is temporarily unavailable" },
228+
{ status: 503, headers: { "Retry-After": "5", "x-should-retry": "false" } }
229+
);
230+
case "paging_unsupported":
231+
return json(
232+
{ error: "Trace paging isn't supported by this deployment's trace store" },
233+
{ status: 501 }
234+
);
235+
}
236+
}

‎apps/webapp/app/v3/eventRepository/clickhouseEventRepository.server.ts‎

Lines changed: 24 additions & 42 deletions
Original file line numberDiff line numberDiff line change
@@ -75,6 +75,7 @@ import type {
7575
TraceChunk,
7676
TraceChunkCursor,
7777
TraceChunkEvent,
78+
TraceChunkScopeOptions,
7879
TraceDetailedSummary,
7980
TraceErrorEvents,
8081
TraceEventOptions,
@@ -1757,11 +1758,11 @@ export class ClickhouseEventRepository implements IEventRepository {
17571758
startCreatedAt: Date,
17581759
endCreatedAt: Date | undefined,
17591760
cursor: TraceChunkCursor | undefined,
1760-
options?: { includeDebugLogs?: boolean; limit?: number; tailInsertedAtSinceMs?: number }
1761+
options?: TraceChunkScopeOptions & { limit?: number }
17611762
): Promise<TraceChunk | undefined> {
17621763
const limit = options?.limit ?? this.maximumTraceChunkSize;
17631764

1764-
const { events, nextCursor, hasMore } = await this.#fetchTraceChunkRecords({
1765+
const { events, nextCursor, hasMore, droppedKeyRows } = await this.#fetchTraceChunkRecords({
17651766
environmentId,
17661767
traceId,
17671768
startCreatedAt,
@@ -1775,37 +1776,10 @@ export class ClickhouseEventRepository implements IEventRepository {
17751776
events: events.map((record) => this.#toTraceChunkEvent(record)),
17761777
nextCursor,
17771778
hasMore,
1779+
droppedKeyRows,
17781780
};
17791781
}
17801782

1781-
async getTraceSpanCount(
1782-
storeTable: TaskEventStoreTable,
1783-
environmentId: string,
1784-
traceId: string,
1785-
startCreatedAt: Date,
1786-
endCreatedAt: Date | undefined,
1787-
options?: { includeDebugLogs?: boolean }
1788-
): Promise<number | undefined> {
1789-
const queryBuilder = this.#createTraceSpanCountQueryBuilder();
1790-
this.#applyTraceScopeWhere(queryBuilder, {
1791-
environmentId,
1792-
traceId,
1793-
startCreatedAt,
1794-
endCreatedAt,
1795-
options,
1796-
});
1797-
1798-
const [queryError, records] = await queryBuilder.execute();
1799-
1800-
if (queryError) {
1801-
logger.error("getTraceSpanCount failed", { error: queryError, traceId });
1802-
return undefined;
1803-
}
1804-
1805-
const count = records?.[0]?.count;
1806-
return count === undefined ? undefined : Number(count);
1807-
}
1808-
18091783
#createTraceChunkQueryBuilder() {
18101784
return this._version === "v2"
18111785
? this._clickhouse.taskEventsV2.traceChunkQueryBuilder()
@@ -1827,11 +1801,12 @@ export class ClickhouseEventRepository implements IEventRepository {
18271801
endCreatedAt?: Date;
18281802
cursor?: TraceChunkCursor;
18291803
limit: number;
1830-
options?: { includeDebugLogs?: boolean; tailInsertedAtSinceMs?: number };
1804+
options?: TraceChunkScopeOptions;
18311805
}): Promise<{
18321806
events: TaskEventChunkV2Result[];
18331807
nextCursor: TraceChunkCursor | null;
18341808
hasMore: boolean;
1809+
droppedKeyRows: boolean;
18351810
}> {
18361811
const queryBuilder = this.#applyTraceChunkScope({
18371812
environmentId,
@@ -1869,21 +1844,28 @@ export class ClickhouseEventRepository implements IEventRepository {
18691844
groupBuilder.where(clause, params);
18701845
groupBuilder.orderBy(TRACE_CHUNK_ORDER_BY);
18711846
// The next cursor skips past this key, so rows beyond the cap are dropped.
1872-
groupBuilder.limit(this.maximumKeyRows);
1847+
groupBuilder.limit(this.maximumKeyRows + 1);
18731848

18741849
const [groupError, groupRecords] = await groupBuilder.execute();
18751850
if (groupError) {
18761851
throw groupError;
18771852
}
18781853

1854+
const rows = groupRecords ?? [];
18791855
return {
1880-
events: groupRecords ?? [],
1856+
events: rows.slice(0, this.maximumKeyRows),
18811857
nextCursor: slice.nextCursor,
18821858
hasMore: slice.hasMore,
1859+
droppedKeyRows: rows.length > this.maximumKeyRows,
18831860
};
18841861
}
18851862

1886-
return { events: slice.events, nextCursor: slice.nextCursor, hasMore: slice.hasMore };
1863+
return {
1864+
events: slice.events,
1865+
nextCursor: slice.nextCursor,
1866+
hasMore: slice.hasMore,
1867+
droppedKeyRows: false,
1868+
};
18871869
}
18881870

18891871
#applyTraceChunkScope({
@@ -1897,7 +1879,7 @@ export class ClickhouseEventRepository implements IEventRepository {
18971879
traceId: string;
18981880
startCreatedAt: Date;
18991881
endCreatedAt?: Date;
1900-
options?: { includeDebugLogs?: boolean; tailInsertedAtSinceMs?: number };
1882+
options?: TraceChunkScopeOptions;
19011883
}) {
19021884
const queryBuilder = this.#createTraceChunkQueryBuilder();
19031885
this.#applyTraceScopeWhere(queryBuilder, {
@@ -1910,12 +1892,6 @@ export class ClickhouseEventRepository implements IEventRepository {
19101892
return queryBuilder;
19111893
}
19121894

1913-
#createTraceSpanCountQueryBuilder() {
1914-
return this._version === "v2"
1915-
? this._clickhouse.taskEventsV2.traceSpanCountQueryBuilder()
1916-
: this._clickhouse.taskEvents.traceSpanCountQueryBuilder();
1917-
}
1918-
19191895
#applyTraceScopeWhere(
19201896
queryBuilder: { where: (clause: string, params?: any) => unknown },
19211897
{
@@ -1929,7 +1905,7 @@ export class ClickhouseEventRepository implements IEventRepository {
19291905
traceId: string;
19301906
startCreatedAt: Date;
19311907
endCreatedAt?: Date;
1932-
options?: { includeDebugLogs?: boolean; tailInsertedAtSinceMs?: number };
1908+
options?: TraceChunkScopeOptions;
19331909
}
19341910
) {
19351911
queryBuilder.where("environment_id = {environmentId: String}", { environmentId });
@@ -1964,6 +1940,12 @@ export class ClickhouseEventRepository implements IEventRepository {
19641940
),
19651941
});
19661942
}
1943+
1944+
if (options?.insertedAtEnd) {
1945+
queryBuilder.where("inserted_at <= {insertedAtEnd: DateTime64(3)}", {
1946+
insertedAtEnd: convertDateToClickhouseDateTime(options.insertedAtEnd),
1947+
});
1948+
}
19671949
}
19681950

19691951
if (options?.includeDebugLogs === false) {

‎apps/webapp/app/v3/eventRepository/eventRepository.server.ts‎

Lines changed: 2 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -71,6 +71,7 @@ import type {
7171
TraceChunk,
7272
TraceChunkCursor,
7373
TraceChunkEvent,
74+
TraceChunkScopeOptions,
7475
TraceDetailedSummary,
7576
TraceErrorEvents,
7677
TraceEventOptions,
@@ -487,22 +488,11 @@ export class EventRepository implements IEventRepository {
487488
_startCreatedAt: Date,
488489
_endCreatedAt: Date | undefined,
489490
_cursor: TraceChunkCursor | undefined,
490-
_options?: { includeDebugLogs?: boolean; limit?: number }
491+
_options?: TraceChunkScopeOptions & { limit?: number }
491492
): Promise<TraceChunk | undefined> {
492493
return undefined;
493494
}
494495

495-
public async getTraceSpanCount(
496-
_storeTable: TaskEventStoreTable,
497-
_environmentId: string,
498-
_traceId: string,
499-
_startCreatedAt: Date,
500-
_endCreatedAt: Date | undefined,
501-
_options?: { includeDebugLogs?: boolean }
502-
): Promise<number | undefined> {
503-
return undefined;
504-
}
505-
506496
public async getTraceErrorEvents(
507497
_storeTable: TaskEventStoreTable,
508498
_environmentId: string,

‎apps/webapp/app/v3/eventRepository/eventRepository.types.ts‎

Lines changed: 11 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -311,10 +311,20 @@ export type TraceChunkEvent = {
311311

312312
export type { TraceChunkCursor };
313313

314+
export type TraceChunkScopeOptions = {
315+
includeDebugLogs?: boolean;
316+
// Live tail: only rows written at or after this time (ms since epoch).
317+
tailInsertedAtSinceMs?: number;
318+
// Only rows written at or before this time.
319+
insertedAtEnd?: Date;
320+
};
321+
314322
export type TraceChunk = {
315323
events: Array<TraceChunkEvent>;
316324
nextCursor: TraceChunkCursor | null;
317325
hasMore: boolean;
326+
// A span had more rows at one timestamp than one read allows; the extra rows were skipped.
327+
droppedKeyRows?: boolean;
318328
};
319329

320330
export type TraceErrorEvents = {
@@ -449,18 +459,9 @@ export interface IEventRepository {
449459
startCreatedAt: Date,
450460
endCreatedAt: Date | undefined,
451461
cursor: TraceChunkCursor | undefined,
452-
options?: { includeDebugLogs?: boolean; limit?: number; tailInsertedAtSinceMs?: number }
462+
options?: TraceChunkScopeOptions & { limit?: number }
453463
): Promise<TraceChunk | undefined>;
454464

455-
getTraceSpanCount(
456-
storeTable: TaskEventStoreTable,
457-
environmentId: string,
458-
traceId: string,
459-
startCreatedAt: Date,
460-
endCreatedAt: Date | undefined,
461-
options?: { includeDebugLogs?: boolean }
462-
): Promise<number | undefined>;
463-
464465
getTraceErrorEvents(
465466
storeTable: TaskEventStoreTable,
466467
environmentId: string,

0 commit comments

Comments
 (0)