Advanced and operational chat.agent capabilities for Trigger.dev, loaded on demand. Load this when working on the raw Sessions primitive (sessions / SessionHand
复制下面这句话,粘贴给 Claude Code、Codex、Cursor 等 AI 编程工具,它会读取安装说明并在你确认后完成安装。
请阅读 https://ai.atlankj.com/install/asset/gh-trigger-chat-agent-advanced-35d5573ca89b ,按照其中的说明把「trigger-chat-agent-advanced」安装到你(当前 AI 工具)中。执行前先告诉我将运行的命令和写入的位置,等我确认。
查看 AI 将读取的安装说明正在读取 GitHub 原文…
内容来自 GitHub 原始文件,由原作者维护。在 GitHub 查看
chat.agent is built on Sessions: a durable, task-bound, bi-directional I/O channel pair keyed
on a stable externalId (e.g. chatId) that outlives any single run. This skill covers the layers
beneath and around the everyday agent: the raw sessions API, server-side AgentChat, durable
sub-agents, actions / background injection, fast starts, compaction and recovery, and the wire
protocol for custom transports.
Two chat namespaces are easy to confuse: the agent definition imports chat from
@trigger.dev/sdk/ai; Head Start / Node-listener server entries import chat from
@trigger.dev/sdk/chat-server.
Happy path: drive an agent from server-side code (task, webhook, or script) with AgentChat.
import { AgentChat } from "@trigger.dev/sdk/chat";
import type { myAgent } from "./trigger/my-agent";
const chat = new AgentChat<typeof myAgent>({ agent: "my-chat", clientData: { userId: "user_123" } });
const stream = await chat.sendMessage("Review PR #42");
const text = await stream.text();
await chat.close();
sendMessage() triggers a run on the first call, then reuses it via input streams. ChatStream
exposes text(), result() ({ text, toolCalls, toolResults }), messages() (UIMessage
snapshots), and the raw .stream. Other methods: steer(text), stop(), sendRaw(uiMessages),
sendAction(action), preload(), reconnect().
Reach for sessions directly when the chat abstraction does not fit: agent inboxes, approval flows,
server-to-server pipelines. sessions.start is idempotent on (env, externalId); externalId
cannot start with session_.
import { sessions } from "@trigger.dev/sdk";
const { id, publicAccessToken } = await sessions.start({
type: "chat.agent",
externalId: chatId,
taskIdentifier: "my-chat",
triggerConfig: { tags: [`chat:${chatId}`], basePayload: { chatId, trigger: "preload" } },
});
const session = sessions.open(chatId); // no network call; methods are lazy
await session.out.append({ kind: "message", text: "hello" });
const next = await session.in.once<MyEvent>({ timeoutMs: 30_000 });
sessions.open(id).in also has send, on(handler), peek, wait (suspends the run, only inside
task.run()), and waitWithIdleTimeout. .out has append, pipe, writer, read,
writeControl, and trimTo. List with sessions.list({ type, tag, status, ... }) (for await),
mutate with sessions.update, end with sessions.close (terminal, idempotent).
AgentChat inside an AI SDK tool() delegates to a durable sub-agent; its response streams as
preliminary tool results. Give the tool a toModelOutput so the model sees a compact summary.
import { tool } from "ai";
import { AgentChat } from "@trigger.dev/sdk/chat";
import { z } from "zod";
const researchTool = tool({
description: "Delegate research to a specialist agent.",
inputSchema: z.object({ topic: z.string() }),
execute: async function* ({ topic }, { abortSignal }) {
const chat = new AgentChat({ agent: "research-agent" });
const stream = await chat.sendMessage(topic, { abortSignal });
yield* stream.messages(); // UIMessage snapshots become preliminary tool results
await chat.close();
},
toModelOutput: ({ output: message }) => {
const lastText = message?.parts?.findLast((p: { type: string }) => p.type === "text") as
| { text?: string }
| undefined;
return { type: "text", value: lastText?.text ?? "Done." };
},
});
For a subtask exposed via execute: ai.toolExecute(task), stream progress to the agent's run with
chat.stream.writer({ target: "root" }). target accepts "self" | "parent" | "root" | <runId>.
Inside the subtask, read context with ai.toolCallId() and ai.chatContextOrThrow<typeof myChat>()
({ chatId, turn, continuation, clientData }).
import { chat, ai } from "@trigger.dev/sdk/ai";
const { waitUntilComplete } = chat.stream.writer({
target: "root",
execute: ({ write }) =>
write({ type: "data-research-status", id: partId, data: { query, status: "in-progress" } }),
});
await waitUntilComplete();
chat.defer(promise) runs work in parallel with streaming (all deferred promises are awaited, with a
5s timeout, before onTurnComplete). chat.inject(messages) queues ModelMessage[] that drain at
the next turn start or prepareStep boundary.
Two lanes, decided by role. A role: "system" message goes to the model's instructions, where it is
trusted like the system prompt, and applies to the next turn only. Any other role joins the
conversation and is untrusted by construction, so put checkable facts there and directives in the
system lane. The instructions lane reaches the model only through the managed streamText (or a
chat.toStreamTextOptions() spread), since that is where the SDK can set instructions.
export const myChat = chat.agent({
id: "my-chat",
onTurnComplete: async ({ messages }) => {
chat.defer(
(async () => {
const analysis = await analyzeConversation(messages);
chat.inject([{ role: "system", content: `[Analysis]\n\n${analysis}` }]);
})()
);
},
registry,
run: async ({ messages, signal, streamText }) =>
streamText({ messages, abortSignal: signal, stopWhen: stepCountIs(15) }),
});
compaction.shouldCompact decides when, summarize produces the summary that replaces the model
messages. UI messages are preserved by default (customize via compactUIMessages). The prepareStep
that performs inner-loop compaction rides the managed streamText, which composes a prepareStep
you pass after it. Spreading chat.toStreamTextOptions() and then passing your own replaces it,
switching compaction off.
compaction: {
shouldCompact: ({ totalTokens }) => (totalTokens ?? 0) > 80_000,
summarize: async ({ messages }) =>
(await generateText({
model: anthropic("claude-haiku-4-5"),
messages: [...messages, { role: "user", content: "Summarize concisely." }],
})).text,
},
actionSchema validates; onAction edits via chat.history (slice, replace, rollbackTo,
remove, getPendingToolCalls, extractNewToolResults). An action fires hydrateMessages and
onAction only. Return nothing for an edit-only action: no model call, and the turn counter does not
advance. Return chat.turn() to answer after the edit: a turn runs on the edited history with
everything a turn has (the agent's system prompt and tools, steering, compaction, injected
instructions, onTurnStart and onTurnComplete, persistence), and run() receives it with
trigger: "action-turn". onAction has no streamText argument, and returning a
StreamTextResult, string or UIMessage from it throws.
Persistence splits by model. With transcript storage (storage on chat.agent; the platform
snapshot by default), the runtime hands storage a changeset with reason: "action" after an action
that changed the conversation: an undo is one truncateAfter, a regenerate is a truncateAfter
followed by the new answer's put when the turn completes, and an edit is a put for the edited
id. With the deprecated hydrateMessages your store is the source of truth and the runtime does
not write, so mirror every mutation yourself: a regenerate is a delete and an insert, and the
answer that follows chat.turn() arrives through onTurnComplete like any turn's answer.
export const myChat = chat.agent({
id: "my-chat",
actionSchema: z.discriminatedUnion("type", [
z.object({ type: z.literal("undo") }),
z.object({ type: z.literal("regenerate") }),
z.object({ type: z.literal("rollback"), targetMessageId: z.string() }),
]),
onAction: async ({ action }) => {
if (action.type === "undo") chat.history.slice(0, -2); // edit only
if (action.type === "rollback") chat.history.rollbackTo(action.targetMessageId);
if (action.type === "regenerate") {
chat.history.slice(0, -1);
return chat.turn(); // answer the edited history
}
},
run: async ({ messages, signal, streamText }) =>
streamText({ model: anthropic("claude-sonnet-4-5"), messages, abortSignal: signal }),
});
Send from the browser through useChat, so the answer a turn produces renders like any turn:
sendMessage(undefined, { body: { action: { type: "regenerate" } } }), regenerate({ body: { action } }),
or useChatActions({ sendMessage }) from @trigger.dev/sdk/chat/react. For regeneration, use
regenerate() from useChat (Vercel AI SDK) which removes the last assistant message before
streaming the new one; calling sendMessage with a regenerate action appends without removal,
leaving both answers visible. transport.sendAction(chatId, action) returns a raw stream the
caller must read and apply itself. Server-side, use
agentChat.sendAction({ type: "rollback", targetMessageId: "msg-3" }).
chat.headStart (from @trigger.dev/sdk/chat-server, NOT /ai) returns a Web Fetch handler that
serves turn 1 from your own warm process, then hands off to the agent on turn 2+. Tools passed here
must be schema-only (a module importing ai + zod only); heavy executes stay in the task.
import { chat } from "@trigger.dev/sdk/chat-server";
import { anthropic } from "@ai-sdk/anthropic";
import { headStartTools } from "@/lib/chat-tools/schemas";
export const chatHandler = chat.headStart({
agentId: "my-chat",
// `streamText` from the run argument owns `messages`, `prompt`, `stopWhen`
// and `abortSignal`: the handover needs `stopWhen: stepCountIs(1)` so the agent,
// not this handler, runs step 2 onward. Passing any of them is a type error.
run: async ({ streamText }) =>
streamText({
model: anthropic("claude-sonnet-4-6"),
system: "You are helpful.",
tools: headStartTools,
}),
});
// Next.js: export const POST = chatHandler; Transport: headStart: "/api/chat"
Node-only frameworks wrap a Web Fetch handler with chat.toNodeListener(handler). Use the same
model on both sides to avoid a tone shift between turn 1 and turn 2+.
chat.local<T>({ id }) is module-level, shallow-proxy, run-scoped state. Initialize it in onBoot
(fires on every fresh worker, including continuation runs), never onChatStart.
const userContext = chat.local<{ name: string; plan: "free" | "pro" }>({ id: "userContext" });
export const myChat = chat.agent({
id: "my-chat",
onBoot: async ({ clientData }) => userContext.init({ name: "Alice", plan: "pro" }),
run: async ({ messages, signal }) => streamText({ /* ... */ }),
});
A message sent while a turn is streaming should NOT cancel the stream. Configure
pendingMessages (shouldInject, prepare, onReceived, onInjected) on the agent so the managed
streamText's prepareStep folds them in at the next boundary. An injected steering message is part
of the conversation your hooks see, so it arrives in uiMessages and newUIMessages at
onTurnComplete and an app persisting from there stores it without extra work. On the frontend, usePendingMessages
returns pending, steer(text), queue(text), and promoteToSteering(id); send via
transport.sendPendingMessage(chatId, uiMessage, metadata?).
onRecoveryBoot fires only when a partial assistant message exists on the tail (interrupted
deploy, crash, OOM retry). It does NOT fire on chat.requestUpgrade(), which is a graceful exit with
no partial. chat.requestUpgrade() (called in onTurnStart / onValidateMessages to skip run(),
or in run() / chat.defer() to exit after the turn) rotates the Session's currentRunId to a run
on the latest deployment without a client reconnect. Pair it with a contract version on clientData.
const SUPPORTED_VERSIONS = new Set(["v2", "v3"]);
onTurnStart: async ({ clientData }) => {
if (clientData?.protocolVersion && !SUPPORTED_VERSIONS.has(clientData.protocolVersion)) {
chat.requestUpgrade();
}
},
For OOM resilience, set oomMachine (and machine) on the agent so retries land on a larger preset.
@trigger.dev/sdk/ai/test runs the real turn loop in-memory. Import it before the agent module
so the resource catalog is installed. Drive with sendMessage, sendRegenerate, sendAction,
sendStop, sendHeadStart, sendHandover; seed state with seedSnapshot / seedSessionOutTail /
seedSessionOutPartial / seedSessionInTail; assert against turn.chunks and harness.allChunks.
import { mockChatAgent } from "@trigger.dev/sdk/ai/test"; // BEFORE the agent module
import { myChatAgent } from "./my-chat.js";
const harness = mockChatAgent(myChatAgent, { chatId: "test-1", clientData: { model } });
try {
const turn = await harness.sendMessage({ id: "u1", role: "user", parts: [{ type: "text", text: "hi" }] });
// assert against turn.chunks
} finally {
await harness.close();
}
Options include mode ("preload" | "submit-message" | "handover-prepare" | "continuation"),
preload, continuation, previousRunId, snapshot, taskContext, and setupLocals. Set
taskContext.ctx.attempt.number > 1 to simulate an OOM-retry attempt. runInMockTaskContext drives a
non-chat task offline.
Endpoints: POST /api/v1/sessions (create), GET /realtime/v1/sessions/{id}/out (SSE),
POST /realtime/v1/sessions/{id}/in/append, POST /api/v1/sessions/{id}/close. ChatInputChunk is
{ kind: "message"; payload: ChatTaskWirePayload } | { kind: "stop"; message? }. The
ChatTaskWirePayload carries chatId, trigger (submit-message | regenerate-message | preload | close | action | handover-prepare), message?, metadata?, action?, continuation?,
previousRunId?, and more. Control records are header-form: trigger-control: turn-complete (with
optional public-access-token, session-in-event-id) and trigger-control: upgrade-required. The
TS helpers SSEStreamSubscription and controlSubtype(headers) (documented in
docs/ai-chat/client-protocol.mdx) handle batch decoding and control-record filtering for you.
CRITICAL: sending a follow-up by re-POSTing POST /api/v1/sessions.
// Wrong - a cached re-POST silently drops basePayload.message; basePayload is trigger config, not a channel
await fetch("/api/v1/sessions", { method: "POST", body: JSON.stringify({ ...createBody }) });
// Correct - append to the session's input channel
await fetch(`/realtime/v1/sessions/${id}/in/append`, { method: "POST", body: JSON.stringify({ kind: "message", payload }) });
Using the wrong token for .in / .out. Use publicAccessToken from the create response
body (session-scoped). The x-trigger-jwt response header is run-scoped and cannot subscribe.
Initializing chat.local in onChatStart. It is skipped on continuation runs, so run()
crashes with chat.local can only be modified after initialization. Init in onBoot.
chat.defer for the message-history write. A mid-stream refresh would read []. await that
write inline before the model streams; reserve chat.defer for analytics, audit, cache warming.
Giving the HITL tool an execute. streamText calls it immediately. Leave it execute-less;
the frontend supplies the answer via addToolOutput + sendAutomaticallyWhen.
Declaring sub-agent / heavy tools only on streamText. Also declare them on
chat.agent({ tools }) (or pass to convertToModelMessages(uiMessages, { tools }) in a custom
agent) so toModelOutput re-applies on every turn.
Importing heavy-execute tools into the Head Start route module. This is a build-time import
chain problem; runtime strip helpers do not fix it. Keep schemas in an ai + zod-only module.
Returning a megabyte tool output on the stream. One tool-output-available record over ~1 MiB
throws ChatChunkTooLargeError. Persist to your store, write the row first, then emit only an id.
Setting X-Peek-Settled: 1 on the active-send path. It races the new turn's first chunk and
closes the stream early. Use it only on reconnect-on-reload paths.
Note on docs vocabulary: agent-side examples in some docs still use the legacy
trigger:turn-completechunk type. That is the agent-emit vocabulary. A custom reader must filter on thetrigger-controlheader, not onchunk.type.MCP-driven agent chats (
list_agents,start_agent_chat,send_agent_message,close_agent_chat) are MCP server tools used from Claude Code / Cursor, not importable SDK functions. See/mcp-tools#agent-chat-tools.
trigger-authoring-chat-agent skill - the everyday chat.agent({...}) definition, lifecycle hooks, and
the useTriggerChatTransport happy path. Start there before reaching for this skill.trigger-realtime-and-frontend skill - Realtime hooks and frontend streaming beyond the chat transport.trigger-authoring-tasks skill - base task() semantics, ctx, and standard lifecycle hooks.Reference docs ship beside this skill in the same package, read them locally (no network), pinned to your installed version. The sources: frontmatter above lists every doc this skill draws from, all under @trigger.dev/sdk/docs/ai-chat/ (including patterns/). For HITL, sessions, and sub-agents start with sessions.mdx, server-chat.mdx, client-protocol.mdx, patterns/human-in-the-loop.mdx, patterns/sub-agents.mdx.
For trigger.config.ts and build extensions a chat-agent task may need (Prisma, Playwright, Python, etc.), read the bundled config docs under @trigger.dev/sdk/docs/config/ (config/extensions/ for the per-extension setup).
This skill is bundled inside @trigger.dev/sdk and read directly from node_modules, so it always matches your installed SDK version (see the adjacent package.json). The full documentation for these APIs ships alongside it under @trigger.dev/sdk/docs/.