feat(sdk): ctx on chat.task hooks; ai-chat E2B sandbox; docs patterns

- Pass TaskRunContext through chat lifecycle events, CompactedEvent, and
  ChatTaskRunPayload; use ctx.run.id for chat access tokens
- Export TaskRunContext from @trigger.dev/sdk
- ai-chat reference: executeCode via E2B, code-sandbox module, warm on
  onTurnStart, dispose on token onWait and onComplete; chat.local run id
- Docs: database persistence + code sandbox pattern pages; reference and
  backend updates for ctx; chat.defer anchor; navigation
This commit is contained in:
Eric Allam
2026-03-27 17:07:25 +00:00
parent 430a81f1b2
commit cb543f3343
8 changed files with 298 additions and 31 deletions
+6
View File
@@ -0,0 +1,6 @@
---
"@trigger.dev/sdk": patch
---
Add `TaskRunContext` (`ctx`) to all `chat.task` lifecycle events, `CompactedEvent`, and `ChatTaskRunPayload`. Export `TaskRunContext` from `@trigger.dev/sdk`.
+43 -11
View File
@@ -13,6 +13,7 @@ import {
type TaskIdentifier,
type TaskOptions,
type TaskSchema,
type TaskRunContext,
type TaskWithSchema,
} from "@trigger.dev/core/v3";
import type {
@@ -44,6 +45,9 @@ import type { ResolvedPrompt } from "./prompt.js";
import { streams } from "./streams.js";
import { createTask } from "./shared.js";
import { tracer } from "./tracer.js";
/** Re-export for typing `ctx` in `chat.task` hooks without importing `@trigger.dev/core`. */
export type { TaskRunContext } from "@trigger.dev/core/v3";
import {
CHAT_STREAM_KEY as _CHAT_STREAM_KEY,
CHAT_MESSAGES_STREAM_ID,
@@ -612,6 +616,11 @@ export type ChatTaskSignals = {
*/
export type ChatTaskRunPayload<TClientData = unknown> = ChatTaskPayload<TClientData> &
ChatTaskSignals & {
/**
* Task run context — same object as the `ctx` passed to a standard `task({ run })` handlers second argument.
* Use for tags, metadata, parent run links, or any API that needs the full run record.
*/
ctx: TaskRunContext;
/** Token usage from the previous turn. Undefined on turn 0. */
previousTurnUsage?: LanguageModelUsage;
/** Cumulative token usage across all completed turns so far. */
@@ -742,6 +751,8 @@ interface CompactionState {
const chatCompactionStateKey = locals.create<CompactionState>("chat.compaction");
const chatOnCompactedKey =
locals.create<(event: CompactedEvent) => Promise<void> | void>("chat.onCompacted");
/** @internal Full task `ctx` for the active `chat.task` run (for hooks invoked from nested compaction). */
const chatTaskRunContextKey = locals.create<TaskRunContext>("chat.taskRunContext");
const chatPrepareMessagesKey =
locals.create<(event: PrepareMessagesEvent<unknown>) => ModelMessage[] | Promise<ModelMessage[]>>(
"chat.prepareMessages"
@@ -980,6 +991,8 @@ export type CompactionChunkData = {
* Event passed to the `onCompacted` callback.
*/
export type CompactedEvent = {
/** Task run context — same as `task` lifecycle hooks and `chat.task` `run({ ctx })`. */
ctx: TaskRunContext;
/** The generated summary text. */
summary: string;
/** The messages that were compacted (pre-compaction). */
@@ -1280,6 +1293,7 @@ async function chatCompact(
const onCompactedHook = locals.get(chatOnCompactedKey);
if (onCompactedHook) {
await onCompactedHook({
ctx: locals.get(chatTaskRunContextKey)!,
summary,
messages,
messageCount: messages.length,
@@ -1863,6 +1877,8 @@ async function pipeChat(
* Event passed to the `onPreload` callback.
*/
export type PreloadEvent<TClientData = unknown> = {
/** Task run context — same as `task({ run })` second-argument `ctx`. */
ctx: TaskRunContext;
/** The unique identifier for the chat session. */
chatId: string;
/** The Trigger.dev run ID for this conversation. */
@@ -1879,6 +1895,8 @@ export type PreloadEvent<TClientData = unknown> = {
* Event passed to the `onChatStart` callback.
*/
export type ChatStartEvent<TClientData = unknown> = {
/** Task run context — same as `task({ run })` second-argument `ctx`. */
ctx: TaskRunContext;
/** The unique identifier for the chat session. */
chatId: string;
/** The initial model-ready messages for this conversation. */
@@ -1903,6 +1921,8 @@ export type ChatStartEvent<TClientData = unknown> = {
* Event passed to the `onTurnStart` callback.
*/
export type TurnStartEvent<TClientData = unknown, TUIM extends UIMessage = UIMessage> = {
/** Task run context — same as `task({ run })` second-argument `ctx`. */
ctx: TaskRunContext;
/** The unique identifier for the chat session. */
chatId: string;
/** The accumulated model-ready messages (all turns so far, including new user message). */
@@ -1935,6 +1955,8 @@ export type TurnStartEvent<TClientData = unknown, TUIM extends UIMessage = UIMes
* Event passed to the `onTurnComplete` callback.
*/
export type TurnCompleteEvent<TClientData = unknown, TUIM extends UIMessage = UIMessage> = {
/** Task run context — same as `task({ run })` second-argument `ctx`. */
ctx: TaskRunContext;
/** The unique identifier for the chat session. */
chatId: string;
/** The full accumulated conversation in model format (all turns so far). */
@@ -2022,8 +2044,9 @@ export type ChatTaskOptions<
* chat.task({
* id: "my-chat",
* clientDataSchema: z.object({ model: z.string().optional(), userId: z.string() }),
* run: async ({ messages, clientData, signal }) => {
* run: async ({ messages, clientData, ctx, signal }) => {
* // clientData is typed as { model?: string; userId: string }
* // ctx is the same TaskRunContext as in task({ run: (payload, { ctx }) => ... })
* },
* });
* ```
@@ -2034,7 +2057,8 @@ export type ChatTaskOptions<
* The run function for the chat task.
*
* Receives a `ChatTaskRunPayload` with the conversation messages, chat session ID,
* trigger type, and abort signals (`signal`, `cancelSignal`, `stopSignal`).
* trigger type, task `ctx` (same as `task({ run })`s second argument), and abort signals
* (`signal`, `cancelSignal`, `stopSignal`).
*
* **Auto-piping:** If this function returns a value with `.toUIMessageStream()`,
* the stream is automatically piped to the frontend.
@@ -2049,7 +2073,7 @@ export type ChatTaskOptions<
*
* @example
* ```ts
* onPreload: async ({ chatId, clientData }) => {
* onPreload: async ({ ctx, chatId, clientData }) => {
* await db.chat.create({ data: { id: chatId } });
* userContext.init(await loadUser(clientData.userId));
* }
@@ -2064,7 +2088,7 @@ export type ChatTaskOptions<
*
* @example
* ```ts
* onChatStart: async ({ chatId, messages, clientData }) => {
* onChatStart: async ({ ctx, chatId, messages, clientData }) => {
* await db.chat.create({ data: { id: chatId, userId: clientData.userId } });
* }
* ```
@@ -2080,7 +2104,7 @@ export type ChatTaskOptions<
*
* @example
* ```ts
* onTurnStart: async ({ chatId, uiMessages }) => {
* onTurnStart: async ({ ctx, chatId, uiMessages }) => {
* await db.chat.update({ where: { id: chatId }, data: { messages: uiMessages } });
* }
* ```
@@ -2097,7 +2121,7 @@ export type ChatTaskOptions<
*
* @example
* ```ts
* onBeforeTurnComplete: async ({ writer, usage }) => {
* onBeforeTurnComplete: async ({ ctx, writer, usage }) => {
* if (usage?.inputTokens && usage.inputTokens > 5000) {
* writer.write({ type: "data-compaction", id: generateId(), data: { status: "compacting" } });
* // ... compact messages ...
@@ -2117,7 +2141,7 @@ export type ChatTaskOptions<
*
* @example
* ```ts
* onCompacted: async ({ summary, totalTokens, chatId }) => {
* onCompacted: async ({ ctx, summary, totalTokens, chatId }) => {
* logger.info("Compacted", { totalTokens, chatId });
* await db.compactionLog.create({ data: { chatId, summary } });
* }
@@ -2177,7 +2201,7 @@ export type ChatTaskOptions<
*
* @example
* ```ts
* onTurnComplete: async ({ chatId, messages }) => {
* onTurnComplete: async ({ ctx, chatId, messages }) => {
* await db.chat.update({ where: { id: chatId }, data: { messages } });
* }
* ```
@@ -2366,8 +2390,10 @@ function chatTask<
...restOptions,
run: async (
payload: ChatTaskWirePayload<TUIMessage, inferSchemaIn<TClientDataSchema>>,
{ signal: runSignal }
{ signal: runSignal, ctx }
) => {
locals.set(chatTaskRunContextKey, ctx);
// Set gen_ai.conversation.id on the run-level span for dashboard context
const activeSpan = trace.getActiveSpan();
if (activeSpan) {
@@ -2433,7 +2459,7 @@ function chatTask<
activeSpan.setAttribute("chat.preloaded", true);
}
const currentRunId = taskContext.ctx?.run.id ?? "";
const currentRunId = ctx.run.id;
let preloadAccessToken = "";
if (currentRunId) {
try {
@@ -2461,6 +2487,7 @@ function chatTask<
async () => {
await withChatWriter(async (writer) => {
await onPreload({
ctx,
chatId: payload.chatId,
runId: currentRunId,
chatAccessToken: preloadAccessToken,
@@ -2671,7 +2698,7 @@ function chatTask<
// Mint a scoped public access token once per turn, reused for
// onChatStart, onTurnStart, onTurnComplete, and the turn-complete chunk.
const currentRunId = taskContext.ctx?.run.id ?? "";
const currentRunId = ctx.run.id;
let turnAccessToken = "";
if (currentRunId) {
try {
@@ -2694,6 +2721,7 @@ function chatTask<
async () => {
await withChatWriter(async (writer) => {
await onChatStart({
ctx,
chatId: currentWirePayload.chatId,
messages: accumulatedMessages,
clientData,
@@ -2728,6 +2756,7 @@ function chatTask<
async () => {
await withChatWriter(async (writer) => {
await onTurnStart({
ctx,
chatId: currentWirePayload.chatId,
messages: accumulatedMessages,
uiMessages: accumulatedUIMessages,
@@ -2799,6 +2828,7 @@ function chatTask<
preloaded,
previousTurnUsage,
totalUsage: cumulativeUsage,
ctx,
signal: combinedSignal,
cancelSignal,
stopSignal,
@@ -3057,6 +3087,7 @@ function chatTask<
const onCompactedHook = locals.get(chatOnCompactedKey);
if (onCompactedHook) {
await onCompactedHook({
ctx,
summary,
messages: accumulatedMessages,
messageCount: accumulatedMessages.length,
@@ -3108,6 +3139,7 @@ function chatTask<
}
const turnCompleteEvent = {
ctx,
chatId: currentWirePayload.chatId,
messages: accumulatedMessages,
uiMessages: accumulatedUIMessages,
+2 -2
View File
@@ -22,9 +22,9 @@ export type { Context };
import type { Context } from "./shared.js";
import type { ApiClientConfiguration } from "@trigger.dev/core/v3";
import type { ApiClientConfiguration, TaskRunContext } from "@trigger.dev/core/v3";
export type { ApiClientConfiguration };
export type { ApiClientConfiguration, TaskRunContext };
export {
ApiError,
+110 -16
View File
@@ -2178,6 +2178,9 @@ importers:
'@ai-sdk/react':
specifier: ^3.0.0
version: 3.0.51(react@19.1.0)(zod@3.25.76)
'@e2b/code-interpreter':
specifier: ^2.4.0
version: 2.4.0
'@prisma/adapter-pg':
specifier: ^7.4.2
version: 7.4.2
@@ -2907,19 +2910,6 @@ importers:
specifier: 5.5.4
version: 5.5.4
references/secure-exec-sandbox:
dependencies:
'@trigger.dev/sdk':
specifier: workspace:*
version: link:../../packages/trigger-sdk
devDependencies:
'@trigger.dev/build':
specifier: workspace:*
version: link:../../packages/build
trigger.dev:
specifier: workspace:*
version: link:../../packages/cli-v3
references/seed:
dependencies:
'@sinclair/typebox':
@@ -4105,6 +4095,9 @@ packages:
'@bufbuild/protobuf@1.10.0':
resolution: {integrity: sha512-QDdVFLoN93Zjg36NoQPZfsVH9tZew7wKDKyV5qRdj8ntT4wQCOradQjRaTdwMhWUYsgKsvCINKKm87FdEk96Ag==}
'@bufbuild/protobuf@2.11.0':
resolution: {integrity: sha512-sBXGT13cpmPR5BMgHE6UEEfEaShh5Ror6rfN3yEK5si7QVrtZg8LEPQb0VVhiLRUslD2yLnXtnRzG035J/mZXQ==}
'@bufbuild/protobuf@2.2.5':
resolution: {integrity: sha512-/g5EzJifw5GF8aren8wZ/G5oMuPoGeS6MQD3ca8ddcvdXR5UELUfdTZITCGNhNXynY/AYl3Z4plmxdj/tRl/hQ==}
@@ -4365,6 +4358,10 @@ packages:
resolution: {integrity: sha512-T54U7WS56ou11ytoxlYllBRBM+MYBpOvVZQa1p1qE4KDZBKJd9m1kAA0PqHjy5T6f/tSv4w5wlq4oyExl4QLLA==}
engines: {node: '>=18'}
'@e2b/code-interpreter@2.4.0':
resolution: {integrity: sha512-6yLi0I2/FBhB6beqtuiFMTajB3PHQw+5+apuI3QdEv1Rx8e2Xp0h37SL0Y6tuVJQ2Frp7P3wWT77Wvy6l6OHPA==}
engines: {node: '>=20'}
'@effect/platform@0.63.2':
resolution: {integrity: sha512-b39pVFw0NGo/tXjGShW7Yg0M+kG7bRrFR6+dQ3aIu99ePTkTp6bGb/kDB7n+dXsFFdIqHsQGYESeYcOQngxdFQ==}
peerDependencies:
@@ -12060,6 +12057,10 @@ packages:
balanced-match@1.0.2:
resolution: {integrity: sha512-3oSeUO0TMV67hN1AmbXsK4yaqU7tjiHlbxRDZOpH0KW9+CeX4bRAaX0Anxt0tx2MrpRpWwQaPwIlISEJhYU5Pw==}
balanced-match@4.0.4:
resolution: {integrity: sha512-BLrgEcRTwX2o6gGxGOCNyMvGSp35YofuYzw9h1IMTRmKqttAZZVU67bdb9Pr2vUHA8+j3i2tJfjO6C6+4myGTA==}
engines: {node: 18 || 20 || >=22}
bare-events@2.8.2:
resolution: {integrity: sha512-riJjyv1/mHLIPX4RwiK+oW9/4c3TEUeORHKefKAKnZ5kyslbN+HXowtbaVEqt4IMUB7OXlfixcs6gsFeo/jhiQ==}
peerDependencies:
@@ -12177,6 +12178,10 @@ packages:
brace-expansion@2.0.1:
resolution: {integrity: sha512-XnAIvQ8eM+kC6aULx6wuQiwVsnzsi9d3WxzV3FpWTGA19F621kwdbsAcFKXgKUHZWsy+mY6iL1sHTxWEFCytDA==}
brace-expansion@5.0.4:
resolution: {integrity: sha512-h+DEnpVvxmfVefa4jFbCf5HdH5YMDXRsmKflpf1pILZWRFlTbJpxeU55nJl4Smt5HQaGzg1o6RHFPJaOqnmBDg==}
engines: {node: 18 || 20 || >=22}
braces@3.0.3:
resolution: {integrity: sha512-yQbXgO/OSZVD2IsiLlro+7Hf6Q18EJrKSEsdoMzKePKXct3gvD8oLcOQdIzGupr5Fj+EDe8gO/lxc1BzfMpxvA==}
engines: {node: '>=8'}
@@ -13317,6 +13322,9 @@ packages:
resolution: {integrity: sha512-ens7BiayssQz/uAxGzH8zGXCtiV24rRWXdjNha5V4zSOcxmAZsfGVm/PPFbwQdqEkDnhG+SyR9E3zSHUbOKXBQ==}
engines: {node: '>= 8.0'}
dockerfile-ast@0.7.1:
resolution: {integrity: sha512-oX/A4I0EhSkGqrFv0YuvPkBUSYp1XiY8O8zAKc8Djglx8ocz+JfOr8gP0ryRMC2myqvDLagmnZaU9ot1vG2ijw==}
dockerode@4.0.6:
resolution: {integrity: sha512-FbVf3Z8fY/kALB9s+P9epCpWhfi/r0N2DgYYcYpsAUlaTxPjdsitsFobnltb+lyCgAIvf9C+4PSWlTnHlJMf1w==}
engines: {node: '>= 8.0'}
@@ -13405,6 +13413,10 @@ packages:
resolution: {integrity: sha512-ii/Bw55ecxgORqkArKNbuVTwqLgVZ0rH1X3J/NOe4LMZaVETm3qNpPBjoPkpQAsQjw2ew0Ad2sd54epqm9nLCw==}
engines: {node: '>=18'}
e2b@2.10.4:
resolution: {integrity: sha512-yLJ/aBpGymfHIsgLIuFwBiNGwUNApzp5Xse86hWmPd/jH3Qjl6Gr7ibQZj9VxP+oFxDCDP1Che7wQyNRTbY2xA==}
engines: {node: '>=20'}
eastasianwidth@0.2.0:
resolution: {integrity: sha512-I88TYZWc9XiYHRQ4/3c5rjjfgkjhLyW2luGIheGERbNQ6OY7yTybanSpDXZa8y7VUP9YmDcYa+eyq4ca7iLqWA==}
@@ -14550,6 +14562,12 @@ packages:
deprecated: Old versions of glob are not supported, and contain widely publicized security vulnerabilities, which have been fixed in the current version. Please update. Support for old versions may be purchased (at exorbitant rates) by contacting i@izs.me
hasBin: true
glob@11.1.0:
resolution: {integrity: sha512-vuNwKSaKiqm7g0THUBu2x7ckSs3XJLXE+2ssL7/MfTGPLLcrJQ/4Uq1CjPTtO5cCIiRxqvN6Twy1qOwhL0Xjcw==}
engines: {node: 20 || >=22}
deprecated: Old versions of glob are not supported, and contain widely publicized security vulnerabilities, which have been fixed in the current version. Please update. Support for old versions may be purchased (at exorbitant rates) by contacting i@izs.me
hasBin: true
glob@7.2.3:
resolution: {integrity: sha512-nFR0zLpU2YCaRxwoCJvL6UvCH2JFyFVIvwTLsIf21AuHlMskA1hhTdk+LlYJtOlYt9v6dvszD2BGRqBL+iQK9Q==}
deprecated: Glob versions prior to v9 are no longer supported
@@ -15265,6 +15283,10 @@ packages:
resolution: {integrity: sha512-cub8rahkh0Q/bw1+GxP7aeSe29hHHn2V4m29nnDlvCdlgU+3UGxkZp7Z53jLUdpX3jdTO0nJZUDl3xvbWc2Xog==}
engines: {node: 20 || >=22}
jackspeak@4.1.1:
resolution: {integrity: sha512-zptv57P3GpL+O0I7VdMJNBZCu+BPHVQUk55Ft8/QCJjTVxrnJHuVuX/0Bl2A6/+2oyR/ZMEuFKwmzqqZ/U5nPQ==}
engines: {node: 20 || >=22}
javascript-stringify@2.1.0:
resolution: {integrity: sha512-JVAfqNPTvNq3sB/VHQJAFxN/sPgKnsKrCwyRt15zwNCdrMMJDdcEOdubuy+DuJYYdm0ox1J4uzEuYKkN+9yhVg==}
@@ -16313,6 +16335,10 @@ packages:
resolution: {integrity: sha512-ethXTt3SGGR+95gudmqJ1eNhRO7eGEGIgYA9vnPatK4/etz2MEVDno5GMCibdMTuBMyElzIlgxMna3K94XDIDQ==}
engines: {node: 20 || >=22}
minimatch@10.2.4:
resolution: {integrity: sha512-oRjTw/97aTBN0RHbYCdtF1MQfvusSIBQM0IZEgzl6426+8jSC0nF1a/GmnVLpfB9yyr6g6FTqWqiZVbxrtaCIg==}
engines: {node: 18 || 20 || >=22}
minimatch@3.1.2:
resolution: {integrity: sha512-J7p63hRiAjw1NDEww1W7i37+ByIrOWO5XQQAzZ3VOcL0PNybwpfmV/N05zFAzwQ9USyEcX6t3UO+K5aqBQOIHw==}
@@ -16949,9 +16975,15 @@ packages:
zod:
optional: true
openapi-fetch@0.14.1:
resolution: {integrity: sha512-l7RarRHxlEZYjMLd/PR0slfMVse2/vvIAGm75/F7J6MlQ8/b9uUQmUF2kCPrQhJqMXSxmYWObVgeYXbFYzZR+A==}
openapi-fetch@0.9.8:
resolution: {integrity: sha512-zM6elH0EZStD/gSiNlcPrzXcVQ/pZo3BDvC6CDwRDUt1dDzxlshpmQnpD6cZaJ39THaSmwVCxxRrPKNM1hHrDg==}
openapi-typescript-helpers@0.0.15:
resolution: {integrity: sha512-opyTPaunsklCBpTK8JGef6mfPhLSnyy5a0IN9vKtx3+4aExf+KxEqYwIy3hqkedXIB97u357uLMJsOnm3GVjsw==}
openapi-typescript-helpers@0.0.8:
resolution: {integrity: sha512-1eNjQtbfNi5Z/kFhagDIaIRj6qqDzhjNJKz8cmMW0CVdGwT6e1GLbAfgI0d28VTJa1A8jz82jm/4dG8qNoNS8g==}
@@ -23157,6 +23189,8 @@ snapshots:
'@bufbuild/protobuf@1.10.0': {}
'@bufbuild/protobuf@2.11.0': {}
'@bufbuild/protobuf@2.2.5': {}
'@bugsnag/cuid@3.1.1': {}
@@ -23469,6 +23503,11 @@ snapshots:
'@connectrpc/connect': 1.4.0(@bufbuild/protobuf@1.10.0)
undici: 5.29.0
'@connectrpc/connect-web@2.0.0-rc.3(@bufbuild/protobuf@2.11.0)(@connectrpc/connect@2.0.0-rc.3(@bufbuild/protobuf@2.11.0))':
dependencies:
'@bufbuild/protobuf': 2.11.0
'@connectrpc/connect': 2.0.0-rc.3(@bufbuild/protobuf@2.11.0)
'@connectrpc/connect-web@2.0.0-rc.3(@bufbuild/protobuf@2.2.5)(@connectrpc/connect@2.0.0-rc.3(@bufbuild/protobuf@2.2.5))':
dependencies:
'@bufbuild/protobuf': 2.2.5
@@ -23478,6 +23517,10 @@ snapshots:
dependencies:
'@bufbuild/protobuf': 1.10.0
'@connectrpc/connect@2.0.0-rc.3(@bufbuild/protobuf@2.11.0)':
dependencies:
'@bufbuild/protobuf': 2.11.0
'@connectrpc/connect@2.0.0-rc.3(@bufbuild/protobuf@2.2.5)':
dependencies:
'@bufbuild/protobuf': 2.2.5
@@ -23535,6 +23578,10 @@ snapshots:
dependencies:
e2b: 1.2.1
'@e2b/code-interpreter@2.4.0':
dependencies:
e2b: 2.10.4
'@effect/platform@0.63.2(@effect/schema@0.72.2(effect@3.7.2))(effect@3.7.2)':
dependencies:
'@effect/schema': 0.72.2(effect@3.7.2)
@@ -32917,6 +32964,8 @@ snapshots:
balanced-match@1.0.2: {}
balanced-match@4.0.4: {}
bare-events@2.8.2:
optional: true
@@ -33058,6 +33107,10 @@ snapshots:
dependencies:
balanced-match: 1.0.2
brace-expansion@5.0.4:
dependencies:
balanced-match: 4.0.4
braces@3.0.3:
dependencies:
fill-range: 7.1.1
@@ -34207,6 +34260,11 @@ snapshots:
transitivePeerDependencies:
- supports-color
dockerfile-ast@0.7.1:
dependencies:
vscode-languageserver-textdocument: 1.0.12
vscode-languageserver-types: 3.17.5
dockerode@4.0.6:
dependencies:
'@balena/dockerignore': 1.0.2
@@ -34309,6 +34367,19 @@ snapshots:
openapi-fetch: 0.9.8
platform: 1.3.6
e2b@2.10.4:
dependencies:
'@bufbuild/protobuf': 2.11.0
'@connectrpc/connect': 2.0.0-rc.3(@bufbuild/protobuf@2.11.0)
'@connectrpc/connect-web': 2.0.0-rc.3(@bufbuild/protobuf@2.11.0)(@connectrpc/connect@2.0.0-rc.3(@bufbuild/protobuf@2.11.0))
chalk: 5.3.0
compare-versions: 6.1.1
dockerfile-ast: 0.7.1
glob: 11.1.0
openapi-fetch: 0.14.1
platform: 1.3.6
tar: 7.5.6
eastasianwidth@0.2.0: {}
ecc-jsbn@0.1.2:
@@ -35879,7 +35950,7 @@ snapshots:
glob@10.3.10:
dependencies:
foreground-child: 3.1.1
foreground-child: 3.3.1
jackspeak: 2.3.6
minimatch: 9.0.5
minipass: 7.1.2
@@ -35887,7 +35958,7 @@ snapshots:
glob@10.3.4:
dependencies:
foreground-child: 3.1.1
foreground-child: 3.3.1
jackspeak: 2.3.6
minimatch: 9.0.5
minipass: 7.1.2
@@ -35904,13 +35975,22 @@ snapshots:
glob@11.0.0:
dependencies:
foreground-child: 3.1.1
foreground-child: 3.3.1
jackspeak: 4.0.1
minimatch: 10.0.1
minipass: 7.1.2
package-json-from-dist: 1.0.0
path-scurry: 2.0.0
glob@11.1.0:
dependencies:
foreground-child: 3.3.1
jackspeak: 4.1.1
minimatch: 10.2.4
minipass: 7.1.2
package-json-from-dist: 1.0.0
path-scurry: 2.0.0
glob@7.2.3:
dependencies:
fs.realpath: 1.0.0
@@ -36711,6 +36791,10 @@ snapshots:
optionalDependencies:
'@pkgjs/parseargs': 0.11.0
jackspeak@4.1.1:
dependencies:
'@isaacs/cliui': 8.0.2
javascript-stringify@2.1.0: {}
jest-worker@27.5.1:
@@ -38073,6 +38157,10 @@ snapshots:
dependencies:
brace-expansion: 2.0.1
minimatch@10.2.4:
dependencies:
brace-expansion: 5.0.4
minimatch@3.1.2:
dependencies:
brace-expansion: 1.1.11
@@ -38788,10 +38876,16 @@ snapshots:
transitivePeerDependencies:
- encoding
openapi-fetch@0.14.1:
dependencies:
openapi-typescript-helpers: 0.0.15
openapi-fetch@0.9.8:
dependencies:
openapi-typescript-helpers: 0.0.8
openapi-typescript-helpers@0.0.15: {}
openapi-typescript-helpers@0.0.8: {}
opener@1.5.2: {}
+1
View File
@@ -17,6 +17,7 @@
"@ai-sdk/react": "^3.0.0",
"@prisma/adapter-pg": "^7.4.2",
"@prisma/client": "^7.4.2",
"@e2b/code-interpreter": "^2.4.0",
"@trigger.dev/sdk": "workspace:*",
"ai": "^6.0.0",
"next": "15.3.3",
+58
View File
@@ -5,6 +5,7 @@ import type { InferUITools, UIDataTypes, UIMessage } from "ai";
import { z } from "zod";
import os from "node:os";
import TurndownService from "turndown";
import { codeSandboxRun, runWithCodeSandbox } from "@/lib/code-sandbox";
const turndown = new TurndownService();
@@ -226,12 +227,69 @@ export const posthogQuery = tool({
},
});
export const executeCode = tool({
description:
"Run code in an isolated E2B sandbox (Python by default; other languages supported by E2B). " +
"Use for calculations, data analysis, or transforming tool outputs (e.g. PostHog query results). " +
"The sandbox persists across turns in the same run until the chat idles and suspends.",
inputSchema: z.object({
code: z.string().describe("Source code to execute in the sandbox"),
language: z
.string()
.optional()
.describe("Language id (e.g. python, javascript). Defaults to python."),
}),
execute: async function executeCodeExecute({ code, language }) {
const runId = codeSandboxRun.runId;
if (!runId?.trim()) {
return {
error:
"Code sandbox run id is not set yet (call from the chat task after onTurnStart), or this tool is not wired to that task.",
};
}
const out = await runWithCodeSandbox(runId, async function runInSandbox(sandbox) {
const execution = await sandbox.runCode(code, {
...(language?.trim() ? { language: language.trim() } : {}),
timeoutMs: 60_000,
});
if (execution.error) {
return {
error: `${execution.error.name}: ${execution.error.value}`,
traceback: execution.error.traceback,
stdout: execution.logs.stdout.join("\n"),
stderr: execution.logs.stderr.join("\n"),
};
}
const mainText = execution.text;
const resultSnippets = execution.results
.map(function mapResult(r) {
return r.text ?? r.markdown ?? r.json;
})
.filter(Boolean)
.slice(0, 5);
return {
text: mainText,
results: resultSnippets,
stdout: execution.logs.stdout.join("\n"),
stderr: execution.logs.stderr.join("\n"),
};
});
return out;
},
});
/** Tool set passed to `streamText` for the main `chat.task` run (includes PostHog). */
export const chatTools = {
inspectEnvironment,
webFetch,
deepResearch,
posthogQuery,
executeCode,
};
type ChatToolSet = typeof chatTools;
@@ -0,0 +1,57 @@
/**
* E2B sandboxes keyed by Trigger run id.
*
* - Warmed from `chat.task` `onTurnStart` (non-blocking) so the first `executeCode` tool call is faster.
* - Disposed in task `onWait` when `wait.type === "token"` (input-stream suspend, same path as `wait.for` tokens).
* - `onComplete` disposes any leftover sandbox if the run ends without hitting another token wait.
*
* No extra `chat.task` SDK hook is required for the suspend boundary — platform `onWait` is sufficient.
*/
import { chat } from "@trigger.dev/sdk/ai";
import { Sandbox } from "@e2b/code-interpreter";
const sandboxPromises = new Map<string, Promise<Sandbox>>();
/** Run id for the active chat turn — set from `onTurnStart` so tools can key the sandbox without `taskContext`. */
export const codeSandboxRun = chat.local<{ runId: string }>({ id: "codeSandboxRun" });
export function warmCodeSandbox(runId: string): void {
codeSandboxRun.init({ runId });
if (!process.env.E2B_API_KEY?.trim()) return;
if (sandboxPromises.has(runId)) return;
sandboxPromises.set(runId, Sandbox.create());
}
export async function runWithCodeSandbox<T>(
runId: string,
runner: (sandbox: Sandbox) => Promise<T>
): Promise<T | { error: string }> {
if (!process.env.E2B_API_KEY?.trim()) {
return { error: "Code sandbox not configured. Set E2B_API_KEY in the Trigger environment." };
}
let promise = sandboxPromises.get(runId);
if (!promise) {
promise = Sandbox.create();
sandboxPromises.set(runId, promise);
}
try {
const sandbox = await promise;
return await runner(sandbox);
} catch (err) {
return { error: err instanceof Error ? err.message : String(err) };
}
}
export async function disposeCodeSandboxForRun(runId: string): Promise<void> {
const promise = sandboxPromises.get(runId);
if (!promise) return;
sandboxPromises.delete(runId);
try {
const sandbox = await promise;
await sandbox.kill();
} catch {
/* best-effort cleanup */
}
}
+21 -2
View File
@@ -21,6 +21,7 @@ import {
webFetch,
type ChatUiMessage,
} from "@/lib/chat-tools";
import { disposeCodeSandboxForRun, warmCodeSandbox } from "@/lib/code-sandbox";
const adapter = new PrismaPg({ connectionString: process.env.DATABASE_URL! });
const prisma = new PrismaClient({ adapter });
@@ -77,7 +78,8 @@ const systemPrompt = prompts.define({
- If you don't know something, say so — don't make things up.
## Capabilities
You can inspect the execution environment, fetch web pages, and perform multi-URL deep research.
You can inspect the execution environment, fetch web pages, perform multi-URL deep research,
query PostHog with HogQL, and run short code snippets in an isolated sandbox (e.g. to analyze query results).
When the user asks you to research a topic, use the deep research tool with relevant URLs.
## Tone
@@ -335,7 +337,8 @@ export const aiChat = chat
// #endregion
// #region onTurnStart — persist messages + write status via writer
onTurnStart: async ({ chatId, uiMessages, writer }) => {
onTurnStart: async ({ chatId, uiMessages, writer, runId }) => {
warmCodeSandbox(runId);
writer.write({ type: "data-turn-status", data: { status: "preparing" } });
chat.defer(
prisma.chat.update({
@@ -346,6 +349,22 @@ export const aiChat = chat
},
// #endregion
onWait: async ({ wait, ctx }) => {
if (wait.type === "token") {
await disposeCodeSandboxForRun(ctx.run.id);
}
},
onResume: async ({ wait, ctx }) => {
if (wait.type === "token") {
logger.debug("Chat resumed after input-stream wait", { runId: ctx.run.id });
}
},
onComplete: async ({ ctx }) => {
await disposeCodeSandboxForRun(ctx.run.id);
},
// #region onTurnComplete — persist + background self-review via chat.inject()
onTurnComplete: async ({
chatId,