# 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.

```ts {{ title: "TypeScript", filename: "inngest/channels.ts" }}
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;
      }>(),
    },
  },
});
```

Channels and topics are plain strings in Python. Pass them to `realtime.publish()` and `realtime.get_subscription_token()`.

Channels and topics are plain strings in Go. Pass them to `realtime.Publish()` and encode payloads yourself, for example with `json.Marshal`.

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.

```ts {{ title: "TypeScript", filename: "inngest/functions/generate.ts" }}
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,
    });
  }
);
```

```python {{ title: "Python", filename: "generate.py" }}
import inngest
from inngest.experimental import realtime

from .client import inngest_client
from .model import stream_completion

MODEL = "gpt-5"

@inngest_client.create_function(
    fn_id="generate-response",
    trigger=inngest.TriggerEvent(event="app/prompt.submitted"),
)
async def generate_response(ctx: inngest.Context) -> None:
    # One channel per conversation thread.
    channel = f"ai-thread:{ctx.event.data['threadId']}"
    prompt = str(ctx.event.data["prompt"])

    # Python has no durable publish step, so wrap one-time publishes in
    # ctx.step.run(). A retry won't send this status again.
    async def publish_start() -> None:
        await realtime.publish(
            client=inngest_client,
            channel=channel,
            topic="status",
            data={"message": "Generating response...", "progress": 0},
        )

    await ctx.step.run("start", publish_start)

    async def stream_model() -> dict[str, object]:
        text = ""
        output_tokens = 0

        # stream_completion() wraps your model provider's streaming API, such as
        # the OpenAI Responses API with stream=True, and yields text deltas
        # followed by a completion chunk with usage.
        async for chunk in stream_completion(model=MODEL, prompt=prompt):
            if chunk.type == "delta":
                text += chunk.delta

                # Immediate on purpose: one publish per token, and a replay on
                # retry is cheaper than making each token a durable step.
                await realtime.publish(
                    client=inngest_client,
                    channel=channel,
                    topic="tokens",
                    data={"token": chunk.delta},
                )

            if chunk.type == "completed":
                output_tokens = chunk.output_tokens

        return {"text": text, "output_tokens": output_tokens}

    generated = await ctx.step.run("stream-model", stream_model)

    # Memoized as a step, so retrying past this point will not publish the
    # result a second time.
    async def publish_result() -> None:
        await realtime.publish(
            client=inngest_client,
            channel=channel,
            topic="result",
            data={
                "output": generated["text"],
                "model": MODEL,
                "outputTokens": generated["output_tokens"],
            },
        )

    await ctx.step.run("send-result", publish_result)
```

```go {{ title: "Go", filename: "generate.go" }}
import (
	"context"
	"encoding/json"

	"github.com/inngest/inngestgo"
	"github.com/inngest/inngestgo/realtime"
	"github.com/inngest/inngestgo/step"
)

const model = "gpt-5"

type PromptSubmitted struct {
	ThreadID string `json:"threadId"`
	Prompt   string `json:"prompt"`
}

type Generated struct {
	Text         string `json:"text"`
	OutputTokens int    `json:"outputTokens"`
}

func publishJSON(ctx context.Context, channel, topic string, v any) error {
	data, err := json.Marshal(v)
	if err != nil {
		return err
	}
	return realtime.Publish(ctx, channel, topic, data)
}

func RegisterGenerate(client inngestgo.Client) (inngestgo.ServableFunction, error) {
	return inngestgo.CreateFunction(
		client,
		inngestgo.FunctionOpts{ID: "generate-response"},
		inngestgo.EventTrigger("app/prompt.submitted", nil),
		func(ctx context.Context, input inngestgo.Input[PromptSubmitted]) (any, error) {
			// One channel per conversation thread.
			ch := "ai-thread:" + input.Event.Data.ThreadID

			// Outside a step: sent once, not repeated on replay or retry.
			if err := publishJSON(ctx, ch, "status", map[string]any{
				"message":  "Generating response...",
				"progress": 0,
			}); err != nil {
				return nil, err
			}

			generated, err := step.Run(ctx, "stream-model", func(ctx context.Context) (Generated, error) {
				text := ""
				// streamCompletion calls your model provider's streaming API
				// and invokes the callback for each text delta.
				outputTokens, err := streamCompletion(ctx, model, input.Event.Data.Prompt, func(delta string) error {
					text += delta
					// Inside a step: one publish per token. A retry of the
					// step can repeat tokens.
					return publishJSON(ctx, ch, "tokens", map[string]string{"token": delta})
				})
				return Generated{Text: text, OutputTokens: outputTokens}, err
			})
			if err != nil {
				return nil, err
			}

			// Outside a step, after the last step: sent once.
			return nil, publishJSON(ctx, ch, "result", map[string]any{
				"output":       generated.Text,
				"model":        model,
				"outputTokens": generated.OutputTokens,
			})
		},
	)
}
```

> **Info:** 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](/docs-markdown/realtime/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](/docs-markdown/realtime/guides/subscription-tokens) for refresh behavior and more frameworks.

```ts {{ title: "Next.js Server Action" }}
// 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"],
  });
}
```

```ts {{ title: "Express" }}
import { getClientSubscriptionToken } from "inngest/react";
import { inngest } from "./inngest/client";
import { aiChannel } from "./inngest/channels";

app.post("/api/ai-token", async (req, res) => {
  const { threadId } = req.body;

  const { userId } = getSession(req);
  await assertUserOwnsThread(userId, threadId);

  const token = await getClientSubscriptionToken(inngest, {
    channel: aiChannel({ threadId }),
    topics: ["status", "tokens", "result"],
  });

  res.json(token);
});
```

```python {{ title: "Python" }}
# FastAPI
import typing

import fastapi
from inngest.experimental import realtime

from .client import inngest_client
from .stubs import assert_user_owns_thread, get_session

app = fastapi.FastAPI()

class TokenRequest(typing.TypedDict):
    threadId: str

@app.post("/api/ai-token")
async def ai_token(
    request: fastapi.Request, body: TokenRequest
) -> typing.Mapping[str, object]:
    thread_id = body["threadId"]

    # A token is a capability. Check the caller may read this thread first.
    session = get_session(request)
    await assert_user_owns_thread(session.user_id, thread_id)

    return await realtime.get_subscription_token(
        client=inngest_client,
        channel=f"ai-thread:{thread_id}",
        topics=["status", "tokens", "result"],
    )
```

Mint tokens with the TypeScript or Python SDK.

## 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.

```tsx {{ filename: "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>
  );
}
```

> **Callout:** 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.

```ts {{ title: "TypeScript" }}
const threadId = crypto.randomUUID();

await inngest.send({
  name: "app/prompt.submitted",
  data: { threadId, prompt: "Summarize the key points of this document..." },
});
```

```python {{ title: "Python" }}
import uuid

import inngest

from .client import inngest_client

async def start_workflow() -> str:
    thread_id = str(uuid.uuid4())

    await inngest_client.send(
        inngest.Event(
            name="app/prompt.submitted",
            data={
                "threadId": thread_id,
                "prompt": "Summarize the key points of this document...",
            },
        )
    )
    return thread_id
```

```go {{ title: "Go" }}
import (
	"context"
	"crypto/rand"

	"github.com/inngest/inngestgo"
)

func startWorkflow(ctx context.Context, client inngestgo.Client) (string, error) {
	threadID := rand.Text()

	_, err := client.Send(ctx, inngestgo.Event{
		Name: "app/prompt.submitted",
		Data: map[string]any{
			"threadId": threadID,
			"prompt":   "Summarize the key points of this document...",
		},
	})
	return threadID, err
}
```

## Next steps

- [React hooks](/docs-markdown/realtime/guides/react-hooks) lists every `useRealtime` option and return value.
- [Publishing reference](/docs-markdown/reference/typescript/v4/realtime/publishing) documents both publish methods.
- [Limits](/docs-markdown/realtime/limits) covers connection and message allowances.