Skip to content

Commit 0754931

Browse files
d-csTrigger.dev RepoOps
authored andcommitted
fix(sdk): separate chat idle and durable wait traces
Chat agent traces now show warm idle time and durable waits separately, with shorter message labels and fewer repeated session IDs. Durable wait spans open the waitpoint inspector while waiting and after completion. This improves trace visibility without changing idle settings or compute billing behavior. Mono-RevId: 1a071a992f017a495f9b0becd8511a3a1ee9a9a6
1 parent 4132259 commit 0754931

5 files changed

Lines changed: 312 additions & 52 deletions

File tree

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,5 @@
1+
---
2+
"@trigger.dev/sdk": patch
3+
---
4+
5+
Show warm idle time and durable waits separately in chat agent traces. Durable waits now open the waitpoint inspector while waiting and after completion. Message span names are shorter, and repeated session IDs no longer crowd message-wait and default output-stream spans.

‎packages/trigger-sdk/src/v3/ai.ts‎

Lines changed: 11 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -65,6 +65,7 @@ import type {
6565
// Runtime VALUES go through the ESM/CJS shim so the CJS build can `require`
6666
// ESM-only `ai@7` (see ../imports/ai-runtime.ts).
6767
import { type Attributes, trace } from "@opentelemetry/api";
68+
import { traceSessionIdle } from "./sessionTracing.js";
6869
import {
6970
tool as aiTool,
7071
convertToModelMessages,
@@ -1786,7 +1787,9 @@ async function waitOnChatRoute<T>(
17861787
async (span) => {
17871788
const idleMs = (options.idleTimeoutInSeconds ?? 0) * 1000;
17881789
if (idleMs > 0) {
1789-
const warm = await router.next(route, { timeoutMs: idleMs });
1790+
const warm = await traceSessionIdle(session.id, idleMs / 1000, () =>
1791+
router.next(route, { timeoutMs: idleMs })
1792+
);
17901793
if (warm) {
17911794
span.setAttribute("wait.resolved", "idle");
17921795
return { ok: true as const, output: warm.data as T, record: warm };
@@ -1842,10 +1845,6 @@ async function waitOnChatRoute<T>(
18421845
session: session.id,
18431846
io: "in",
18441847
route,
1845-
...accessoryAttributes({
1846-
items: [{ text: `${session.id}.in:${route}`, variant: "normal" }],
1847-
style: "codepath",
1848-
}),
18491848
},
18501849
}
18511850
);
@@ -8116,7 +8115,7 @@ function chatAgent<
81168115
const preloadResult = await messagesInput.waitWithIdleTimeout({
81178116
idleTimeoutInSeconds: effectivePreloadIdleTimeout,
81188117
timeout: effectivePreloadTimeout,
8119-
spanName: "waiting for first message",
8118+
spanName: "first message",
81208119
skipSuspend: exitAfterPreloadIdle,
81218120
onSuspend: onChatSuspend
81228121
? async () => {
@@ -8298,7 +8297,7 @@ function chatAgent<
82988297
const continuationResult = await messagesInput.waitWithIdleTimeout({
82998298
idleTimeoutInSeconds: effectiveIdleTimeout,
83008299
timeout: effectiveTurnTimeout,
8301-
spanName: "waiting for first message (continuation)",
8300+
spanName: "first message (continuation)",
83028301
onSuspend: onChatSuspend
83038302
? async () => {
83048303
await tracer.startActiveSpan(
@@ -9948,7 +9947,7 @@ function chatAgent<
99489947
const next = await messagesInput.waitWithIdleTimeout({
99499948
idleTimeoutInSeconds: effectiveIdleTimeout,
99509949
timeout: effectiveTurnTimeout,
9951-
spanName: "waiting for next message",
9950+
spanName: "next message",
99529951
onSuspend: onChatSuspend
99539952
? async () => {
99549953
await tracer.startActiveSpan(
@@ -10347,7 +10346,7 @@ function chatAgent<
1034710346
const next = await messagesInput.waitWithIdleTimeout({
1034810347
idleTimeoutInSeconds: effectiveIdleTimeout,
1034910348
timeout: effectiveTurnTimeout,
10350-
spanName: "waiting for next message (after error)",
10349+
spanName: "next message (after error)",
1035110350
});
1035210351

1035310352
if (!next.ok) {
@@ -12356,8 +12355,8 @@ function createChatSession<TClientData = unknown>(
1235612355
timeout,
1235712356
spanName:
1235812357
currentPayload.trigger === "preload"
12359-
? "waiting for first message"
12360-
: "waiting for first message (continuation)",
12358+
? "first message"
12359+
: "first message (continuation)",
1236112360
});
1236212361
if (!result.ok || runSignal.aborted) {
1236312362
stop.cleanup();
@@ -12397,7 +12396,7 @@ function createChatSession<TClientData = unknown>(
1239712396
const next = await messagesInput.waitWithIdleTimeout({
1239812397
idleTimeoutInSeconds,
1239912398
timeout,
12400-
spanName: "waiting for next message",
12399+
spanName: "next message",
1240112400
});
1240212401
if (!next.ok || runSignal.aborted) {
1240312402
stop.cleanup();
Lines changed: 204 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,204 @@
1+
import { createServer } from "node:http";
2+
import { setTimeout as sleep } from "node:timers/promises";
3+
import { SpanStatusCode, trace } from "@opentelemetry/api";
4+
import {
5+
SemanticInternalAttributes as Attr,
6+
SessionChannelRouter,
7+
usage,
8+
WaitpointTimeoutError,
9+
} from "@trigger.dev/core/v3";
10+
import { TracingSDK, type TracingSDKConfig } from "@trigger.dev/core/v3/otel";
11+
import { DevUsageManager } from "@trigger.dev/core/v3/workers";
12+
import { afterAll, assert, beforeAll, beforeEach, describe, expect, it } from "vitest";
13+
import { traceSessionIdle, traceSessionWait } from "./sessionTracing.js";
14+
import { tracer } from "./tracer.js";
15+
16+
type Exporter = NonNullable<TracingSDKConfig["exporters"]>[number];
17+
type ExportedSpan = Parameters<Exporter["export"]>[0][number];
18+
type WireSpan = {
19+
name: string;
20+
attributes: Array<{ key: string; value: { stringValue?: string; boolValue?: boolean } }>;
21+
};
22+
23+
const spans: ExportedSpan[] = [];
24+
const wireSpans: WireSpan[] = [];
25+
// A real local OTLP receiver also captures partial spans, which external
26+
// exporters intentionally omit. No tracer, router, or usage mocks are used.
27+
const receiver = createServer(async (request, response) => {
28+
const chunks: Buffer[] = [];
29+
for await (const chunk of request) chunks.push(Buffer.from(chunk));
30+
if (request.url === "/v1/traces") {
31+
const body = JSON.parse(Buffer.concat(chunks).toString());
32+
for (const resource of body.resourceSpans ?? []) {
33+
for (const scope of resource.scopeSpans ?? []) wireSpans.push(...scope.spans);
34+
}
35+
}
36+
response.writeHead(200, { "Content-Type": "application/json" }).end("{}");
37+
});
38+
let tracing: TracingSDK;
39+
40+
beforeAll(async () => {
41+
await new Promise<void>((resolve) => receiver.listen(0, "127.0.0.1", resolve));
42+
const address = receiver.address();
43+
if (!address || typeof address === "string") throw new Error("Missing OTLP receiver address");
44+
tracing = new TracingSDK({
45+
url: `http://127.0.0.1:${address.port}`,
46+
forceFlushTimeoutMillis: 5_000,
47+
exporters: [
48+
{
49+
export(batch, done) {
50+
spans.push(...batch);
51+
done({ code: 0 });
52+
},
53+
async shutdown() {},
54+
},
55+
],
56+
});
57+
});
58+
59+
beforeEach(async () => {
60+
await tracing.flush();
61+
spans.length = 0;
62+
wireSpans.length = 0;
63+
usage.reset();
64+
usage.setGlobalUsageManager(new DevUsageManager());
65+
});
66+
67+
afterAll(async () => {
68+
usage.reset();
69+
await tracing?.shutdown();
70+
await new Promise<void>((resolve, reject) =>
71+
receiver.close((error) => (error ? reject(error) : resolve()))
72+
);
73+
});
74+
75+
function messageRouter() {
76+
return new SessionChannelRouter({
77+
kindOf: () => "message",
78+
routes: [{ name: "messages", delivery: "queue", replayable: true, kinds: ["message"] }],
79+
});
80+
}
81+
82+
describe("session trace phases", () => {
83+
it.each([30, 10])(
84+
"shows a message arriving within a %ss idle window without a waitpoint",
85+
async (seconds) => {
86+
const router = messageRouter();
87+
const message = { id: "m1", seqNum: 1, data: { text: "hello" } };
88+
const measurement = usage.start();
89+
const result = await tracer.startActiveSpan("next message", async () => {
90+
const pending = traceSessionIdle("session_test", seconds, () =>
91+
router.next("messages", { timeoutMs: seconds * 1000 })
92+
);
93+
router.ingest(message);
94+
return pending;
95+
});
96+
const sample = usage.stop(measurement);
97+
98+
expect(result).toBe(message);
99+
expect(sample.cpuTime).toBe(sample.wallTime);
100+
expect(spans.map((span) => span.name)).toEqual(["idle", "next message"]);
101+
const [idle, parent] = spans;
102+
assert(idle);
103+
assert(parent);
104+
expect(idle.parentSpanContext?.spanId).toBe(parent.spanContext().spanId);
105+
expect(idle.attributes["wait.idleTimeoutInSeconds"]).toBe(seconds);
106+
expect(idle.attributes[Attr.ENTITY_TYPE]).toBeUndefined();
107+
expect(idle.attributes.session).toBe("session_test");
108+
}
109+
);
110+
111+
it("exports a linked waitpoint while waiting and keeps idle and durable usage separate", async () => {
112+
const router = messageRouter();
113+
let finish!: () => void;
114+
const completed = new Promise<void>((resolve) => {
115+
finish = resolve;
116+
});
117+
let entered!: () => void;
118+
const waiting = new Promise<void>((resolve) => {
119+
entered = resolve;
120+
});
121+
const measurement = usage.start();
122+
const result = { ok: true as const, waitpointId: "waitpoint_test" };
123+
const pending = tracer.startActiveSpan("next message", async () => {
124+
expect(
125+
await traceSessionIdle("session_test", 0.01, () =>
126+
router.next("messages", { timeoutMs: 10 })
127+
)
128+
).toBeUndefined();
129+
return traceSessionWait("session_test", result.waitpointId, () =>
130+
usage.pauseAsync(async () => {
131+
expect(trace.getActiveSpan()).toBeDefined();
132+
entered();
133+
await completed;
134+
return result;
135+
})
136+
);
137+
});
138+
await waiting;
139+
try {
140+
const pausedUsage = measurement.sample().cpuTime;
141+
await tracing.flush();
142+
await sleep(10);
143+
// Usage rounds to milliseconds, so two samples can differ by 1ms.
144+
expect(Math.abs(measurement.sample().cpuTime - pausedUsage)).toBeLessThanOrEqual(1);
145+
const partial = wireSpans.find(
146+
(span) =>
147+
span.name === "wait.forToken()" &&
148+
span.attributes.some((attr) => attr.key === Attr.SPAN_PARTIAL && attr.value.boolValue)
149+
);
150+
expect(partial?.attributes).toEqual(
151+
expect.arrayContaining([
152+
{ key: Attr.ENTITY_TYPE, value: { stringValue: "waitpoint" } },
153+
{ key: Attr.ENTITY_ID, value: { stringValue: "waitpoint_test" } },
154+
])
155+
);
156+
} finally {
157+
finish();
158+
await pending;
159+
}
160+
expect(await pending).toBe(result);
161+
const sample = usage.stop(measurement);
162+
expect(sample.wallTime - sample.cpuTime).toBeGreaterThanOrEqual(10);
163+
const idle = spans.find((span) => span.name === "idle")!;
164+
const wait = spans.find((span) => span.name === "wait.forToken()")!;
165+
const parent = spans.find((span) => span.name === "next message")!;
166+
expect(idle.parentSpanContext?.spanId).toBe(parent.spanContext().spanId);
167+
expect(wait.parentSpanContext?.spanId).toBe(parent.spanContext().spanId);
168+
expect(wait.attributes[Attr.ENTITY_ID]).toBe(result.waitpointId);
169+
expect(wait.status.code).toBe(SpanStatusCode.UNSET);
170+
});
171+
172+
it("retains the waitpoint identity and timeout result on a failed wait", async () => {
173+
const result = { ok: false as const, error: new WaitpointTimeoutError("Timed out") };
174+
expect(await traceSessionWait("session_test", "waitpoint_timeout", async () => result)).toBe(
175+
result
176+
);
177+
expect(spans).toHaveLength(1);
178+
const span = spans[0];
179+
assert(span);
180+
expect(span.attributes[Attr.ENTITY_ID]).toBe("waitpoint_timeout");
181+
expect(span.status.code).toBe(SpanStatusCode.ERROR);
182+
expect(span.events[0]?.attributes?.["exception.message"]).toBe("Timed out");
183+
});
184+
185+
it.each(["idle", "wait"])(
186+
"ends the %s span and propagates an aborted operation",
187+
async (phase) => {
188+
const error = new Error("Operation aborted");
189+
const fail = async (): Promise<never> => {
190+
throw error;
191+
};
192+
const pending =
193+
phase === "idle"
194+
? traceSessionIdle("session_test", 30, fail)
195+
: traceSessionWait("session_test", "waitpoint_aborted", fail);
196+
await expect(pending).rejects.toBe(error);
197+
expect(spans).toHaveLength(1);
198+
const span = spans[0];
199+
assert(span);
200+
expect(span.status.code).toBe(SpanStatusCode.ERROR);
201+
expect(span.endTime[0]).toBeGreaterThan(0);
202+
}
203+
);
204+
});
Lines changed: 53 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,53 @@
1+
import { accessoryAttributes, SemanticInternalAttributes } from "@trigger.dev/core/v3";
2+
import { SpanStatusCode } from "@opentelemetry/api";
3+
import { tracer } from "./tracer.js";
4+
5+
/** The warm window uses compute; it must not be presented as a waitpoint. */
6+
export function traceSessionIdle<T>(
7+
sessionId: string,
8+
idleTimeoutInSeconds: number,
9+
read: () => Promise<T>
10+
): Promise<T> {
11+
return tracer.startActiveSpan("idle", read, {
12+
attributes: {
13+
[SemanticInternalAttributes.STYLE_ICON]: "sessions",
14+
session: sessionId,
15+
io: "in",
16+
"wait.phase": "idle",
17+
"wait.idleTimeoutInSeconds": idleTimeoutInSeconds,
18+
...accessoryAttributes({
19+
items: [{ text: `up to ${idleTimeoutInSeconds}s`, variant: "normal" }],
20+
style: "codepath",
21+
}),
22+
},
23+
});
24+
}
25+
26+
/** Attach the waitpoint before starting the wait, including for partial spans. */
27+
export function traceSessionWait<T extends { ok: boolean; error?: Error }>(
28+
sessionId: string,
29+
waitpointId: string,
30+
wait: () => Promise<T>
31+
): Promise<T> {
32+
return tracer.startActiveSpan(
33+
"wait.forToken()",
34+
async (span) => {
35+
const result = await wait();
36+
if (!result.ok && result.error) {
37+
span.recordException(result.error);
38+
span.setStatus({ code: SpanStatusCode.ERROR });
39+
}
40+
return result;
41+
},
42+
{
43+
attributes: {
44+
[SemanticInternalAttributes.STYLE_ICON]: "wait",
45+
[SemanticInternalAttributes.ENTITY_TYPE]: "waitpoint",
46+
[SemanticInternalAttributes.ENTITY_ID]: waitpointId,
47+
session: sessionId,
48+
io: "in",
49+
"wait.phase": "suspended",
50+
},
51+
}
52+
);
53+
}

0 commit comments

Comments
 (0)