Skip to content

Commit e6ec2c4

Browse files
carderneTrigger.dev RepoOps
authored andcommitted
fix(sdk): re-suspend chat agents after unmatched wakes
Chat agents now return to a durable wait when session activity does not deliver a matching message. Repeated wakes preserve the original turn timeout instead of leaving the run active until its maximum duration. Mono-RevId: f9bb2238bbd2ea663eaac75ce3d92106a83c1104
1 parent 7f3e607 commit e6ec2c4

4 files changed

Lines changed: 184 additions & 23 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+
Chat agents now return to a durable wait when a session wake does not deliver a matching message. This prevents resumed runs from staying active until their maximum duration and preserves the configured turn timeout across repeated wakes.

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

Lines changed: 14 additions & 23 deletions
Original file line numberDiff line numberDiff line change
@@ -66,6 +66,7 @@ import type {
6666
// ESM-only `ai@7` (see ../imports/ai-runtime.ts).
6767
import { type Attributes, trace } from "@opentelemetry/api";
6868
import { traceSessionIdle } from "./sessionTracing.js";
69+
import { waitForChatRouteAfterIdle } from "./chatRouteWait.js";
6970
import {
7071
tool as aiTool,
7172
convertToModelMessages,
@@ -1780,31 +1781,21 @@ async function waitOnChatRoute<T>(
17801781
if (options.onSuspend) await options.onSuspend();
17811782

17821783
span.setAttribute("wait.resolved", "suspended");
1783-
while (true) {
1784-
/**
1785-
* The floor doubles as the wake cursor: the server completes the
1786-
* waitpoint immediately if anything sits after this sequence, so a
1787-
* floor that has advanced past an unread record parks a waitpoint
1788-
* nothing will complete. Recorded on the span so a run that never woke
1789-
* can be diagnosed from its trace alone.
1790-
*/
1791-
const wakeFrom = router.resumeFloor();
1792-
span.setAttribute("wait.lastSeqNum", wakeFrom ?? -1);
1793-
const wake = await session.in.awaitWake({
1794-
timeout: options.timeout,
1795-
lastSeqNum: wakeFrom,
1796-
});
1797-
if (!wake.ok) {
1798-
span.recordException(wake.error);
1799-
return { ok: false as const, error: wake.error };
1800-
}
1801-
1802-
const record = await router.next(route);
1803-
if (!record) continue;
1784+
const result = await waitForChatRouteAfterIdle(router, route, {
1785+
timeout: options.timeout,
1786+
wake: async (timeout, lastSeqNum) => {
1787+
span.setAttribute("wait.lastSeqNum", lastSeqNum ?? -1);
1788+
return session.in.awaitWake({ timeout, lastSeqNum });
1789+
},
1790+
});
18041791

1805-
if (options.onResume) await options.onResume();
1806-
return { ok: true as const, output: record.data as T, record };
1792+
if (!result.ok) {
1793+
span.recordException(result.error);
1794+
return result;
18071795
}
1796+
1797+
if (options.onResume) await options.onResume();
1798+
return { ok: true as const, output: result.record.data as T, record: result.record };
18081799
},
18091800
{
18101801
attributes: {
Lines changed: 101 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,101 @@
1+
import { SessionChannelRouter, WaitpointTimeoutError } from "@trigger.dev/core/v3";
2+
import { describe, expect, it } from "vitest";
3+
import { waitForChatRouteAfterIdle } from "./chatRouteWait.js";
4+
5+
function inputRouter() {
6+
return new SessionChannelRouter({
7+
kindOf: (data) => (data as { kind?: string }).kind,
8+
routes: [
9+
{ name: "messages", delivery: "queue", replayable: true, kinds: ["message"] },
10+
{ name: "control", delivery: "queue", replayable: false, kinds: ["control"] },
11+
],
12+
});
13+
}
14+
15+
describe("chat route waits after suspension", () => {
16+
it("re-suspends after a wake that has no record for the requested route", async () => {
17+
const router = inputRouter();
18+
const timeouts: Array<string | undefined> = [];
19+
let now = 0;
20+
21+
const result = await waitForChatRouteAfterIdle(router, "messages", {
22+
timeout: "10s",
23+
postWakeTimeoutMs: 0,
24+
now: () => now,
25+
wake: async (timeout) => {
26+
timeouts.push(timeout);
27+
if (timeouts.length === 1) {
28+
router.ingest({ id: "control-1", seqNum: 1, data: { kind: "control" } });
29+
now = 4_000;
30+
return { ok: true, waitpointId: "waitpoint-1" };
31+
}
32+
return { ok: false, error: new WaitpointTimeoutError("Timed out") };
33+
},
34+
});
35+
36+
expect(result.ok).toBe(false);
37+
expect(timeouts).toEqual(["10s", "6s"]);
38+
});
39+
40+
it("does not re-wake on an unrelated queued message while waiting for another route", async () => {
41+
const router = inputRouter();
42+
const cursors: Array<number | undefined> = [];
43+
44+
const result = await waitForChatRouteAfterIdle(router, "control", {
45+
postWakeTimeoutMs: 0,
46+
wake: async (_timeout, lastSeqNum) => {
47+
cursors.push(lastSeqNum);
48+
if (cursors.length === 1) {
49+
router.ingest({ id: "message-1", seqNum: 1, data: { kind: "message" } });
50+
return { ok: true, waitpointId: "waitpoint-1" };
51+
}
52+
if (lastSeqNum === undefined || lastSeqNum < 1) {
53+
return { ok: true, waitpointId: "waitpoint-repeated" };
54+
}
55+
return { ok: false, error: new WaitpointTimeoutError("Timed out") };
56+
},
57+
});
58+
59+
expect(result.ok).toBe(false);
60+
expect(cursors).toEqual([undefined, 1]);
61+
expect(router.resumeFloor()).toBe(0);
62+
expect(router.hasPending("messages")).toBe(true);
63+
});
64+
65+
it("expires an absolute timeout after an unmatched wake", async () => {
66+
const router = inputRouter();
67+
const deadline = Date.parse("2026-09-28T10:15:00Z");
68+
let now = deadline - 1_000;
69+
const timeouts: Array<string | undefined> = [];
70+
71+
const result = await waitForChatRouteAfterIdle(router, "messages", {
72+
timeout: "2026-09-28T10:15:00Z",
73+
postWakeTimeoutMs: 0,
74+
now: () => now,
75+
wake: async (timeout) => {
76+
timeouts.push(timeout);
77+
now = deadline + 1;
78+
return { ok: true, waitpointId: "waitpoint-1" };
79+
},
80+
});
81+
82+
expect(result.ok).toBe(false);
83+
expect(timeouts).toEqual(["1s"]);
84+
});
85+
86+
it("returns a matching record delivered after the channel wakes", async () => {
87+
const router = inputRouter();
88+
const message = { id: "message-1", seqNum: 1, data: { kind: "message", text: "hello" } };
89+
90+
const result = await waitForChatRouteAfterIdle(router, "messages", {
91+
timeout: "10s",
92+
postWakeTimeoutMs: 0,
93+
wake: async () => {
94+
router.ingest(message);
95+
return { ok: true, waitpointId: "waitpoint-1" };
96+
},
97+
});
98+
99+
expect(result).toEqual({ ok: true, record: message });
100+
});
101+
});
Lines changed: 64 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,64 @@
1+
import {
2+
type SessionChannelRouter,
3+
type SessionStreamRecord,
4+
WaitpointTimeoutError,
5+
} from "@trigger.dev/core/v3";
6+
import { parseNaturalLanguageDurationInMs } from "@trigger.dev/core/v3/isomorphic";
7+
8+
const POST_WAKE_TIMEOUT_MS = 5_000;
9+
10+
type ChatRouteWakeResult = { ok: true; waitpointId: string } | { ok: false; error: Error };
11+
12+
export async function waitForChatRouteAfterIdle(
13+
router: SessionChannelRouter,
14+
route: string,
15+
options: {
16+
timeout?: string;
17+
postWakeTimeoutMs?: number;
18+
now?: () => number;
19+
wake: (
20+
timeout: string | undefined,
21+
lastSeqNum: number | undefined
22+
) => Promise<ChatRouteWakeResult>;
23+
}
24+
): Promise<{ ok: true; record: SessionStreamRecord } | { ok: false; error: Error }> {
25+
const now = options.now ?? Date.now;
26+
const timeoutMs = options.timeout ? parseNaturalLanguageDurationInMs(options.timeout) : undefined;
27+
const parsedDeadline =
28+
timeoutMs !== undefined
29+
? now() + timeoutMs
30+
: options.timeout === undefined
31+
? undefined
32+
: Date.parse(options.timeout);
33+
const deadline =
34+
parsedDeadline !== undefined && Number.isFinite(parsedDeadline) ? parsedDeadline : undefined;
35+
const postWakeTimeoutMs = options.postWakeTimeoutMs ?? POST_WAKE_TIMEOUT_MS;
36+
37+
while (true) {
38+
const remainingMs = deadline === undefined ? undefined : deadline - now();
39+
if (remainingMs !== undefined && remainingMs <= 0) {
40+
return { ok: false, error: new WaitpointTimeoutError("Timed out") };
41+
}
42+
43+
const wakeTimeout =
44+
remainingMs === undefined
45+
? options.timeout
46+
: `${Math.max(1, Math.ceil(remainingMs / 1000))}s`;
47+
// A queued record holds the replay floor back, but should not wake this
48+
// waitpoint again after we have already seen it. Snapshot before checking
49+
// the route so a record arriving during registration still wakes us.
50+
const wakeFrom = router.appliedThrough();
51+
const buffered = await router.next(route, { timeoutMs: 0 });
52+
if (buffered) return { ok: true, record: buffered };
53+
54+
const wake = await options.wake(wakeTimeout, wakeFrom);
55+
if (!wake.ok) return wake;
56+
57+
const deliveryBudget =
58+
deadline === undefined
59+
? postWakeTimeoutMs
60+
: Math.max(0, Math.min(postWakeTimeoutMs, deadline - now()));
61+
const record = await router.next(route, { timeoutMs: deliveryBudget });
62+
if (record) return { ok: true, record };
63+
}
64+
}

0 commit comments

Comments
 (0)