Stream AI responses
Show a model's answer in your UI as it's generated, then publish the final result without repeats on retry.
Users wait less when they can watch an answer appear. In this guide, an Inngest function streams model output to the browser token by token, reports progress, and publishes the final result as a durable step. The same pattern works for document processing, agent tool calls, and any multi-step pipeline.
1. Define the channel
Give each kind of update its own topic: status for progress, tokens for streamed output, and result for the final answer.
import { realtime, staticSchema } from "inngest";
import { z } from "zod";
export const aiChannel = realtime.channel({
// `threadId` is your own conversation identifier, not an Inngest run ID.
// The client passes the same value when it subscribes.
name: ({ threadId }: { threadId: string }) => `ai-thread:${threadId}`,
topics: {
status: {
schema: z.object({ message: z.string(), progress: z.number() }),
},
// Token payloads are tiny and very high volume, so skip runtime
// validation and keep the types only.
tokens: {
schema: staticSchema<{ token: string }>(),
},
result: {
schema: staticSchema<{
output: string;
model: string;
outputTokens: number;
}>(),
},
},
});
The tokens topic uses staticSchema<T>() because token payloads are small and frequent. The other topics validate at runtime.
2. Publish from your function
The function uses both publish methods:
inngest.realtime.publish()sends each token immediately. If the step retries, some tokens can repeat, which costs less than making every token a durable step.step.realtime.publish()sends the status and the final result as durable steps, so a retry never sends them twice.
import { OpenAI } from "openai";
import { inngest } from "../client";
import { aiChannel } from "../channels";
const openai = new OpenAI();
const MODEL = "gpt-5";
export default inngest.createFunction(
{ id: "generate-response", triggers: [{ event: "app/prompt.submitted" }] },
async ({ event, step }) => {
// One channel instance per conversation thread.
const ch = aiChannel({ threadId: event.data.threadId });
// Durable: a retry won't send this status again.
await step.realtime.publish("start", ch.status, {
message: "Generating response...",
progress: 0,
});
const generated = await step.run("stream-model", async () => {
const stream = await openai.responses.create({
model: MODEL,
input: [{ role: "user", content: event.data.prompt }],
stream: true,
});
let text = "";
let outputTokens = 0;
// The Responses API emits semantic events. Branch on `chunk.type` and
// handle only the ones your UI needs.
for await (const chunk of stream) {
if (chunk.type === "response.output_text.delta") {
text += chunk.delta;
// Non-durable on purpose: one publish per token, and a replay on
// retry is cheaper than making each token a durable step.
await inngest.realtime.publish(ch.tokens, { token: chunk.delta });
}
if (chunk.type === "response.completed") {
outputTokens = chunk.response.usage?.output_tokens ?? 0;
}
}
return { text, outputTokens };
});
// Durable: memoized as a step, so retrying past this point will not
// publish the result a second time.
await step.realtime.publish("send-result", ch.result, {
output: generated.text,
model: MODEL,
outputTokens: generated.outputTokens,
});
}
);
Call inngest.realtime.publish() inside step.run() for streamed output.
Inngest then retries the stream as one unit. Your UI should tolerate
repeated tokens if the step retries.
Each token counts as one Realtime message. Check your plan's daily message allowance in Limits before you stream long responses at high volume.
3. Mint a subscription token
Mint a token on your server after you check that the user may read the thread. See Subscription tokens for refresh behavior and more frameworks.
// app/actions.ts
"use server";
import { getClientSubscriptionToken } from "inngest/react";
import { inngest } from "@/inngest/client";
import { aiChannel } from "@/inngest/channels";
import { getSession } from "@/lib/session";
export async function fetchAIToken(threadId: string) {
// A token is a capability. Check the caller may read this thread first.
const { userId } = await getSession();
await assertUserOwnsThread(userId, threadId);
return getClientSubscriptionToken(inngest, {
channel: aiChannel({ threadId }),
topics: ["status", "tokens", "result"],
});
}
4. Render the stream
useRealtime connects to the channel and returns every message it receives. The component rebuilds the answer from the tokens messages and swaps in the final result when it arrives.
app/components/AIStream.tsx"use client";
import { useRealtime } from "inngest/react";
import { aiChannel } from "@/inngest/channels";
import { fetchAIToken } from "../actions";
export function AIStream({ threadId }: { threadId: string }) {
const { messages, connectionStatus } = useRealtime({
channel: aiChannel({ threadId }),
topics: ["status", "tokens", "result"] as const,
token: () => fetchAIToken(threadId),
enabled: !!threadId,
// `messages.all` keeps only the last 100 messages by default, which would
// silently drop the start of a long response. Disable the cap so every
// token is retained.
historyLimit: null,
// Batch re-renders instead of rendering once per token.
bufferInterval: 50,
});
// Rebuild the streamed text from the retained token messages. Skip `run`
// messages first: they are run lifecycle updates with untyped data.
let streamed = "";
for (const message of messages.all) {
if (message.kind === "run") continue;
if (message.topic === "tokens") {
streamed += message.data.token;
}
}
const status = messages.byTopic.status?.data;
const result = messages.byTopic.result?.data;
return (
<div>
<p>Connection: {connectionStatus}</p>
{status && !result && (
<p>{status.message} ({status.progress}%)</p>
)}
{streamed && !result && (
<div className="whitespace-pre-wrap">{streamed}</div>
)}
{result && (
<div>
<div className="whitespace-pre-wrap">{result.output}</div>
<p className="text-sm text-subtle">
Model: {result.model} | Output tokens: {result.outputTokens}
</p>
</div>
)}
</div>
);
}
Check message.kind === "run" before you check message.topic. Run
lifecycle messages share the same list and carry an untyped payload, so
skipping them lets TypeScript narrow message.data.
historyLimit: null keeps every token. The default limit of 100 would drop the start of a long answer. bufferInterval: 50 batches renders so React doesn't re-render once per token.
5. Start the workflow
The threadId links the function's publishes to the component's subscription. Create it before you send the event, render the component with the same value, and send the event once the connection is open so the first tokens aren't missed.
const threadId = crypto.randomUUID();
await inngest.send({
name: "app/prompt.submitted",
data: { threadId, prompt: "Summarize the key points of this document..." },
});
Next steps
- React hooks lists every
useRealtimeoption and return value. - Publishing reference documents both publish methods.
- Limits covers connection and message allowances.