Skip to content

Commit de97106

Browse files
MariefayTrigger.dev RepoOps
authored andcommitted
fix(webapp): bound trace API ClickHouse reads by run completion
Bound the time window the run trace API reads for finished runs, so retrieving the trace of an old run no longer scans every day of stored events since the run was created. For a finished run, events stored more than 7 days after the run completed are no longer read. In rare cases, a span that finishes more than 7 days after the run, or a cancellation the run would inherit from a parent recorded that late, is not reflected: the span is returned in its last recorded state. Runs that are still executing are not affected. Mono-RevId: 549b041b446c38ab61e8ce68c809a8d376df3dd6
1 parent e9f56b8 commit de97106

8 files changed

Lines changed: 249 additions & 6 deletions

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

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -4,6 +4,7 @@ import { z } from "zod";
44
import { $replica } from "~/db.server";
55
import { anyResource, createLoaderApiRoute } from "~/services/routeBuilders/apiBuilder.server";
66
import { getEventRepositoryForStore } from "~/v3/eventRepository/index.server";
7+
import { getTraceInsertedAtEnd } from "~/v3/eventRepository/traceInsertedAtBound";
78
import { getTaskEventStoreTableForRun } from "~/v3/taskEventStore.server";
89
import { findRunByIdWithMollifierFallback } from "~/v3/mollifier/readFallback.server";
910
import { buildSyntheticTraceBody } from "~/v3/mollifier/syntheticApiResponses.server";
@@ -99,7 +100,8 @@ export const loader = createLoaderApiRoute(
99100
run.traceId,
100101
run.spanId,
101102
run.createdAt,
102-
run.completedAt ?? undefined
103+
run.completedAt ?? undefined,
104+
{ insertedAtEnd: getTraceInsertedAtEnd(run) }
103105
);
104106

105107
if (!traceSummary) {

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

Lines changed: 9 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -1765,7 +1765,7 @@ export class ClickhouseEventRepository implements IEventRepository {
17651765
anchorSpanId: string;
17661766
startCreatedAt: Date;
17671767
endCreatedAt?: Date;
1768-
options?: { includeDebugLogs?: boolean };
1768+
options?: { includeDebugLogs?: boolean; insertedAtEnd?: Date };
17691769
limit?: number;
17701770
}): Promise<{
17711771
records: TaskEventDetailedSummaryV1Result[];
@@ -2546,7 +2546,7 @@ export class ClickhouseEventRepository implements IEventRepository {
25462546
traceId: string;
25472547
startCreatedAt?: Date;
25482548
endCreatedAt?: Date;
2549-
options?: { includeDebugLogs?: boolean };
2549+
options?: { includeDebugLogs?: boolean; insertedAtEnd?: Date };
25502550
spanIds?: string[];
25512551
parentSpanIds?: string[];
25522552
limit?: number;
@@ -2578,6 +2578,12 @@ export class ClickhouseEventRepository implements IEventRepository {
25782578
queryBuilder.where("inserted_at >= {insertedAtStart: DateTime64(3)}", {
25792579
insertedAtStart: convertDateToClickhouseDateTime(startCreatedAtWithBuffer),
25802580
});
2581+
2582+
if (options?.insertedAtEnd) {
2583+
queryBuilder.where("inserted_at <= {insertedAtEnd: DateTime64(3)}", {
2584+
insertedAtEnd: convertDateToClickhouseDateTime(options.insertedAtEnd),
2585+
});
2586+
}
25812587
}
25822588
}
25832589

@@ -2724,7 +2730,7 @@ export class ClickhouseEventRepository implements IEventRepository {
27242730
anchorSpanId: string,
27252731
startCreatedAt: Date,
27262732
endCreatedAt?: Date,
2727-
options?: { includeDebugLogs?: boolean }
2733+
options?: { includeDebugLogs?: boolean; insertedAtEnd?: Date }
27282734
): Promise<TraceDetailedSummary | undefined> {
27292735
const limit = this._config.maximumTraceDetailedSummaryViewCount;
27302736

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -500,7 +500,7 @@ export class EventRepository implements IEventRepository {
500500
anchorSpanId: string,
501501
startCreatedAt: Date,
502502
endCreatedAt?: Date,
503-
options?: { includeDebugLogs?: boolean }
503+
options?: { includeDebugLogs?: boolean; insertedAtEnd?: Date }
504504
): Promise<TraceDetailedSummary | undefined> {
505505
const events = await this.taskEventStore.findDetailedTraceEvents(
506506
storeTable,

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -430,7 +430,7 @@ export interface IEventRepository {
430430
anchorSpanId: string,
431431
startCreatedAt: Date,
432432
endCreatedAt?: Date,
433-
options?: { includeDebugLogs?: boolean }
433+
options?: { includeDebugLogs?: boolean; insertedAtEnd?: Date }
434434
): Promise<TraceDetailedSummary | undefined>;
435435

436436
// Streams a trace's events in start_time order, one at a time, without ever
Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,18 @@
1+
import type { TaskRunStatus } from "@trigger.dev/database";
2+
import { isFinalRunStatus } from "~/v3/taskStatus";
3+
4+
// Spans in a run's subtree can finish after the run, so their final rows land later.
5+
const TRACE_INSERTED_AT_END_BUFFER_MS = 7 * 24 * 60 * 60 * 1000;
6+
7+
export function getTraceInsertedAtEnd(run: {
8+
status: TaskRunStatus;
9+
completedAt: Date | null;
10+
updatedAt: Date;
11+
}): Date | undefined {
12+
if (!isFinalRunStatus(run.status)) {
13+
return undefined;
14+
}
15+
16+
const finishedAt = run.completedAt ?? run.updatedAt;
17+
return new Date(finishedAt.getTime() + TRACE_INSERTED_AT_END_BUFFER_MS);
18+
}

‎apps/webapp/test/getTraceDetailedSubtreeSummary.integration.test.ts‎

Lines changed: 167 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -426,3 +426,170 @@ describe("getTraceDetailedSubtreeSummary", () => {
426426
INTEGRATION_TIMEOUT_MS
427427
);
428428
});
429+
430+
describe("getTraceDetailedSubtreeSummary insertedAtEnd bound", () => {
431+
const DAY_MS = 24 * 60 * 60 * 1000;
432+
433+
clickhouseTest(
434+
"bounds full and descendant fetches by insertedAtEnd but not ancestor fetches",
435+
async ({ clickhouseContainer }) => {
436+
const clickhouse = new ClickHouse({
437+
url: clickhouseContainer.getConnectionUrl(),
438+
logLevel: "warn",
439+
});
440+
const repository = new ClickhouseEventRepository({
441+
clickhouse,
442+
version: "v2",
443+
maximumTraceDetailedSummaryViewCount: TRACE_ROW_LIMIT,
444+
});
445+
446+
const environmentId = "env_trace_bound";
447+
const baseMs = Date.now() - 30 * DAY_MS;
448+
const expiresAt = convertDateToClickhouseDateTime(new Date(baseMs + 400 * DAY_MS));
449+
450+
function row(
451+
traceId: string,
452+
spanId: string,
453+
parentSpanId: string,
454+
startOffsetMs: number,
455+
insertedOffsetMs: number,
456+
status: string
457+
): TaskEventV2Input {
458+
return {
459+
environment_id: environmentId,
460+
organization_id: "org_trace_bound",
461+
project_id: "proj_trace_bound",
462+
task_identifier: "trace-bound-test-task",
463+
run_id: `run_${spanId}`,
464+
start_time: formatClickhouseStartTime(baseMs, startOffsetMs),
465+
inserted_at: convertDateToClickhouseDateTime(new Date(baseMs + insertedOffsetMs)),
466+
duration: "1000000000",
467+
trace_id: traceId,
468+
span_id: spanId,
469+
parent_span_id: parentSpanId,
470+
message: spanId,
471+
kind: "SPAN",
472+
status,
473+
attributes: {},
474+
metadata: "{}",
475+
expires_at: expiresAt,
476+
};
477+
}
478+
479+
// Full fetch: the requested run is the trace root and completes at +10s.
480+
const fullTrace = "c".repeat(32);
481+
const fullCompletedMs = 10_000;
482+
// Subtree walk: the requested run starts after its parent, so the full fetch can't re-root.
483+
const walkTrace = "d".repeat(32);
484+
const walkCreatedMs = 100_000;
485+
const walkCompletedMs = 110_000;
486+
// Ancestor: the parent's cancellation is stored long after the bound.
487+
const ancestorTrace = "e".repeat(32);
488+
489+
const [insertError] = await clickhouse.taskEventsV2.insert(
490+
[
491+
row(fullTrace, "fullroot", "", 0, 0, "PARTIAL"),
492+
row(fullTrace, "fullroot", "", 0, fullCompletedMs, "OK"),
493+
row(fullTrace, "fullinside", "fullroot", 1_000, 1_000, "PARTIAL"),
494+
row(fullTrace, "fullinside", "fullroot", 1_000, fullCompletedMs + 6 * DAY_MS, "OK"),
495+
row(fullTrace, "fulllate", "fullroot", 2_000, 2_000, "PARTIAL"),
496+
row(fullTrace, "fulllate", "fullroot", 2_000, fullCompletedMs + 8 * DAY_MS, "OK"),
497+
498+
row(walkTrace, "walkroot", "", 0, 0, "PARTIAL"),
499+
row(walkTrace, "walkanchor", "walkroot", walkCreatedMs, walkCreatedMs, "PARTIAL"),
500+
row(walkTrace, "walkanchor", "walkroot", walkCreatedMs, walkCompletedMs, "OK"),
501+
row(
502+
walkTrace,
503+
"walkinside",
504+
"walkanchor",
505+
walkCreatedMs + 1_000,
506+
walkCreatedMs + 1_000,
507+
"PARTIAL"
508+
),
509+
row(
510+
walkTrace,
511+
"walkinside",
512+
"walkanchor",
513+
walkCreatedMs + 1_000,
514+
walkCompletedMs + DAY_MS,
515+
"OK"
516+
),
517+
row(
518+
walkTrace,
519+
"walklate",
520+
"walkanchor",
521+
walkCreatedMs + 2_000,
522+
walkCreatedMs + 2_000,
523+
"PARTIAL"
524+
),
525+
row(
526+
walkTrace,
527+
"walklate",
528+
"walkanchor",
529+
walkCreatedMs + 2_000,
530+
walkCompletedMs + 9 * DAY_MS,
531+
"OK"
532+
),
533+
534+
row(ancestorTrace, "ancroot", "", 0, 0, "PARTIAL"),
535+
row(ancestorTrace, "ancroot", "", 0, 20 * DAY_MS, "CANCELLED"),
536+
row(ancestorTrace, "ancanchor", "ancroot", walkCreatedMs, walkCreatedMs, "PARTIAL"),
537+
],
538+
{ clickhouse_settings: { async_insert: 0 } }
539+
);
540+
expect(insertError).toBeNull();
541+
542+
const fullCompletedAt = new Date(baseMs + fullCompletedMs);
543+
const fullBounded = await repository.getTraceDetailedSubtreeSummary(
544+
"taskEvent",
545+
environmentId,
546+
fullTrace,
547+
"fullroot",
548+
new Date(baseMs),
549+
fullCompletedAt,
550+
{ insertedAtEnd: new Date(fullCompletedAt.getTime() + 7 * DAY_MS) }
551+
);
552+
expect(findSpan(fullBounded?.rootSpan, "fullinside")?.data.isPartial).toBe(false);
553+
expect(findSpan(fullBounded?.rootSpan, "fulllate")?.data.isPartial).toBe(true);
554+
555+
const fullUnbounded = await repository.getTraceDetailedSubtreeSummary(
556+
"taskEvent",
557+
environmentId,
558+
fullTrace,
559+
"fullroot",
560+
new Date(baseMs),
561+
fullCompletedAt
562+
);
563+
expect(findSpan(fullUnbounded?.rootSpan, "fullinside")?.data.isPartial).toBe(false);
564+
expect(findSpan(fullUnbounded?.rootSpan, "fulllate")?.data.isPartial).toBe(false);
565+
566+
const walkCompletedAt = new Date(baseMs + walkCompletedMs);
567+
const walk = await repository.getTraceDetailedSubtreeSummary(
568+
"taskEvent",
569+
environmentId,
570+
walkTrace,
571+
"walkanchor",
572+
new Date(baseMs + walkCreatedMs),
573+
walkCompletedAt,
574+
{ insertedAtEnd: new Date(walkCompletedAt.getTime() + 7 * DAY_MS) }
575+
);
576+
expect(walk?.rootSpan.id).toBe("walkanchor");
577+
expect(walk?.rootSpan.parentId).toBe("walkroot");
578+
expect(findSpan(walk?.rootSpan, "walkinside")?.data.isPartial).toBe(false);
579+
expect(findSpan(walk?.rootSpan, "walklate")?.data.isPartial).toBe(true);
580+
581+
const ancestor = await repository.getTraceDetailedSubtreeSummary(
582+
"taskEvent",
583+
environmentId,
584+
ancestorTrace,
585+
"ancanchor",
586+
new Date(baseMs + walkCreatedMs),
587+
undefined,
588+
{ insertedAtEnd: new Date(baseMs + 7 * DAY_MS) }
589+
);
590+
expect(ancestor?.rootSpan.id).toBe("ancanchor");
591+
expect(ancestor?.rootSpan.data.isCancelled).toBe(true);
592+
},
593+
INTEGRATION_TIMEOUT_MS
594+
);
595+
});
Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,48 @@
1+
import { describe, expect, it } from "vitest";
2+
import { getTraceInsertedAtEnd } from "~/v3/eventRepository/traceInsertedAtBound";
3+
4+
const DAY_MS = 24 * 60 * 60 * 1000;
5+
const COMPLETED_AT = new Date("2026-05-13T20:13:30.310Z");
6+
const UPDATED_AT = new Date("2026-05-14T09:00:00.000Z");
7+
8+
describe("getTraceInsertedAtEnd", () => {
9+
it("bounds a completed run at completedAt + 7 days", () => {
10+
const end = getTraceInsertedAtEnd({
11+
status: "COMPLETED_WITH_ERRORS",
12+
completedAt: COMPLETED_AT,
13+
updatedAt: UPDATED_AT,
14+
});
15+
16+
expect(end).toEqual(new Date(COMPLETED_AT.getTime() + 7 * DAY_MS));
17+
});
18+
19+
it("falls back to updatedAt + 7 days for a final run with no completedAt", () => {
20+
const end = getTraceInsertedAtEnd({
21+
status: "CANCELED",
22+
completedAt: null,
23+
updatedAt: UPDATED_AT,
24+
});
25+
26+
expect(end).toEqual(new Date(UPDATED_AT.getTime() + 7 * DAY_MS));
27+
});
28+
29+
it("does not bound a run that is not in a final status", () => {
30+
const end = getTraceInsertedAtEnd({
31+
status: "EXECUTING",
32+
completedAt: null,
33+
updatedAt: UPDATED_AT,
34+
});
35+
36+
expect(end).toBeUndefined();
37+
});
38+
39+
it("does not bound a non-final run even if completedAt is set", () => {
40+
const end = getTraceInsertedAtEnd({
41+
status: "WAITING_TO_RESUME",
42+
completedAt: COMPLETED_AT,
43+
updatedAt: UPDATED_AT,
44+
});
45+
46+
expect(end).toBeUndefined();
47+
});
48+
});

‎docs/management/runs/retrieve-trace.mdx‎

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,3 +6,5 @@ openapi: "v3-openapi GET /api/v1/runs/{runId}/trace"
66
Returns the OpenTelemetry trace subtree for the run you request. The response `trace.rootSpan` is that run's span — not necessarily the trace-wide root — with its descendant spans nested under `children`.
77

88
For a child or nested run inside a large trace, this endpoint scopes the tree to that run so you still get a useful subtree even when the full trace has more spans than the platform can return in one response.
9+
10+
For a finished run, spans that finish more than 7 days after the run are returned in their last recorded state.

0 commit comments

Comments
 (0)