> Niveau 6 de Agent Harness Patterns, un parcours de patterns sur le fonctionnement des agents d'IA. Version web : https://harnesspatterns.dev/fr/patterns/streaming · Tous les patterns (en anglais) : https://harnesspatterns.dev/llms.txt

# Streaming et événements

Un modèle écrit sa réponse quelques tokens à la fois. Le streaming vous remet chaque morceau dès qu'il est écrit, et les événements de la boucle racontent au reste de votre app ce que fait l'agent, pendant qu'il le fait.

## Le problème

Appeler un modèle et attendre la réponse entière, ça marche dans un script. Devant une personne, non : une longue réponse peut prendre dix ou vingt secondes, et pendant tout ce temps l'écran n'affiche rien. Un agent aggrave les choses, car une requête peut enchaîner plusieurs tours et outils avant que la réponse finale existe.

Et il n'y a pas que la personne qui veut suivre. L'interface veut le texte, vos logs veulent chaque étape, la facturation veut les tokens de chaque tour. S'il faut écrire chacune de ces choses dans la boucle, la boucle finit par connaître chaque écran, fichier et base de données de votre app.

## La solution

Deux pièces qui vont ensemble. **Le streaming** : demandez la réponse au fournisseur sous forme de flux, et il envoie chaque morceau de texte dès que le modèle l'écrit. Le temps total est le même, mais le premier mot apparaît en une fraction de seconde au lieu d'arriver à la fin.

**Les événements** : la boucle annonce chaque étape sur un bus, sans se soucier de qui écoute. Votre app abonne les parties intéressées : l'écran écoute le texte, le log écoute tout, le compteur écoute la fin de chaque tour. Ajouter un auditeur ne touche pas la boucle. Voici les événements qu'astorlm émet lors d'une exécution normale :

- `turn_start`
   Un tour commence : la boucle s'apprête à appeler le modèle.
- `text_delta`
   Un morceau de la réponse est arrivé. Il y en a beaucoup par tour.
- `assistant_message`
   Le message du modèle est complet, texte et appels d'outils compris.
- `tool_execution_start`
   Un outil s'apprête à s'exécuter, avec son entrée.
- `tool_execution_end`
   L'outil a fini : sa sortie, s'il a échoué, combien de temps il a pris.
- `turn_end`
   Le tour est fini, avec les tokens de son appel au modèle.
- `session_end`
   L'exécution est finie : terminée, annulée ou en erreur.

Les appels d'outils arrivent aussi du fournisseur en streaming, quelques caractères de leurs arguments à la fois. La boucle les reconstitue avant de les exécuter : sur le bus, un appel d'outil arrive donc entier, dans `assistant_message`.

## Les personnages

Les mêmes personnages que d'habitude, sur une scène rock.

- **L'Oracle** (le modèle): Le chanteur, en haut sur l'estrade. Sa réponse descend la piste en notes, une par mot.
- **La piste** (le stream): Chaque text_delta est une poignée de notes. Sans streaming, rien ne descend avant la fin.
- **La personne à la guitare** (votre app): Joue chaque note quand elle atteint les frettes, et les paroles apparaissent pour le public.
- **Le câble** (l'EventBus): Tous les événements le parcourent. Seuls s'allument les auditeurs abonnés à cet événement.
- **Les auditeurs** (agent.on(…)): L'écran des paroles (`'text'`), un magnétophone (`'event'`) et un compteur de tokens (`turn_end`).
- **Astor** (la boucle): Fait tourner la boucle comme d'habitude, et annonce chaque étape sur le bus au fur et à mesure.
- **La scène** (les outils): Le bord de scène, la console lumière et la boîte pyro : `crowd_mood`, `set_lights` et `pyro`.

## Le code

**Avec astorlm :** le fournisseur fait toujours du streaming, il n'y a donc rien à activer. Abonnez-vous avec `agent.on()` avant d'appeler `run()` : `'text'` pour chaque morceau de texte, `'tool-start'` et `'tool-end'` pour les outils, `'event'` pour tout. Chaque appel renvoie une fonction qui annule l'abonnement.

**À partir de zéro :** la boucle du niveau 2 avec deux changements. La requête utilise `stream: true` et lit la réponse comme des server-sent events, en recollant les morceaux ; et une liste d'auditeurs reçoit chaque étape via un seul `emit()`.

**Avec astorlm**

```ts
import { OpenAIProvider, createLocalAgent, tool } from 'astorlm'
import { z } from 'zod'

// The stage's functions, wrapped as tools.
const crowdMood = tool({
  name: 'crowd_mood',
  description: 'Look at the crowd from the front of the stage: how they feel, and what they came to hear.',
  schema: z.object({}),
  execute: async () => stage.readCrowd(), // your code
})

const setLights = tool({
  name: 'set_lights',
  description: 'Set the color of every light on the stage.',
  schema: z.object({ color: z.enum(['amber', 'red', 'blue', 'white']) }),
  execute: async ({ color }) => stage.lights(color), // your code
})

const pyro = tool({
  name: 'pyro',
  description: 'Fire the pyrotechnics on both sides of the stage. Once per show.',
  schema: z.object({}),
  execute: async () => stage.firePyro(), // your code
})

const agent = await createLocalAgent({
  // Any OpenAI-compatible endpoint: OpenAI, Ollama, LM Studio, vLLM, a proxy…
  provider: new OpenAIProvider({
    baseURL: 'http://localhost:11434/v1', // e.g. Ollama's default address
    model: 'your-model', // e.g. 'llama3.1', 'gpt-4o-mini'
    apiKey: 'YOUR_API_KEY', // local servers usually ignore it
  }),
  tools: [crowdMood, setLights, pyro],
  maxTurns: 10,
})

// Streaming is always on: subscribe before run(). Each listener only hears what it asked for.
// The lyrics screen: every piece of text, the moment it arrives.
agent.on('text', (piece) => lyrics.append(piece))

// The tape deck: every event the loop emits, in order.
agent.on('event', (event) => tape.write(event))

// The meter: the tokens each turn reported.
agent.on('event', (event) => {
  if (event.type === 'turn_end' && event.usage) meter.add(event.usage.inputTokens + event.usage.outputTokens)
})

// Each on() returns a function that unsubscribes.
const stopTape = agent.on('tool-start', ({ name }) => console.log('→', name))

const last = await agent.run('Open the show: read the crowd, set the mood, and greet them.')
stopTape()
console.log(last.content) // the same text, whole, once the run is over
```

**TypeScript**

```ts
// Streaming and events from scratch. Plain fetch, no SDK.

// Any OpenAI-compatible endpoint: OpenAI, Ollama, LM Studio, vLLM, a proxy…
const LLM = {
  baseURL: 'http://localhost:11434/v1', // e.g. Ollama's default address
  model: 'your-model', // e.g. 'llama3.1', 'gpt-4o-mini'
  apiKey: 'YOUR_API_KEY', // local servers usually ignore it
}

type ToolCall = { id: string; type: 'function'; function: { name: string; arguments: string } }
type Message =
  | { role: 'user'; content: string }
  | { role: 'assistant'; content: string | null; tool_calls?: ToolCall[] }
  | { role: 'tool'; tool_call_id: string; content: string }

const tools: Record<string, (args: Record<string, string>) => Promise<string>> = { crowd_mood: crowdMood, set_lights: setLights, pyro }
const toolSchemas = [/* one JSON Schema per tool */]

// 1. The bus: a list of listeners. The loop emits; it never knows who is listening.
type AgentEvent =
  | { type: 'turn_start'; turn: number }
  | { type: 'text_delta'; text: string }
  | { type: 'tool_execution_start'; name: string }
  | { type: 'tool_execution_end'; name: string; output: string }
  | { type: 'turn_end'; turn: number; usage?: { prompt_tokens: number; completion_tokens: number } }
const listeners: ((event: AgentEvent) => void)[] = []
export const on = (listener: (event: AgentEvent) => void) => {
  listeners.push(listener)
  return () => listeners.splice(listeners.indexOf(listener), 1)
}
const emit = (event: AgentEvent) => listeners.forEach((listener) => listener(event))

// 2. One streamed model call: read the server-sent events and rebuild the message piece by piece.
async function streamReply(messages: Message[]) {
  const res = await fetch(`${LLM.baseURL}/chat/completions`, {
    method: 'POST',
    headers: { 'content-type': 'application/json', authorization: `Bearer ${LLM.apiKey}` },
    body: JSON.stringify({ model: LLM.model, messages, tools: toolSchemas, stream: true, stream_options: { include_usage: true } }),
  })
  const reader = res.body!.pipeThrough(new TextDecoderStream()).getReader()
  let text = ''
  let buffer = ''
  let usage
  const calls: ToolCall[] = []
  for (let read = await reader.read(); !read.done; read = await reader.read()) {
    buffer += read.value
    const lines = buffer.split('\n')
    buffer = lines.pop()! // a line can arrive cut in half: keep it for the next chunk
    for (const line of lines) {
      if (!line.startsWith('data: ') || line === 'data: [DONE]') continue
      const chunk = JSON.parse(line.slice(6))
      usage = chunk.usage ?? usage
      const delta = chunk.choices?.[0]?.delta
      if (delta?.content) {
        text += delta.content
        emit({ type: 'text_delta', text: delta.content }) // show it now, not at the end
      }
      // Tool calls arrive in pieces too: glue each one's arguments together by index.
      for (const part of delta?.tool_calls ?? []) {
        const call = (calls[part.index] ??= { id: '', type: 'function', function: { name: '', arguments: '' } })
        call.id ||= part.id ?? ''
        call.function.name += part.function?.name ?? ''
        call.function.arguments += part.function?.arguments ?? ''
      }
    }
  }
  const reply: Message = { role: 'assistant', content: text || null, ...(calls.length ? { tool_calls: calls } : {}) }
  return { reply, usage }
}

// 3. The loop from level 2, emitting as it goes.
export async function runAgent(prompt: string, maxTurns = 10): Promise<string> {
  const messages: Message[] = [{ role: 'user', content: prompt }]
  for (let turn = 1; turn <= maxTurns; turn++) {
    emit({ type: 'turn_start', turn })
    const { reply, usage } = await streamReply(messages)
    messages.push(reply)
    if (reply.role !== 'assistant' || !reply.tool_calls) {
      emit({ type: 'turn_end', turn, usage })
      return reply.content ?? ''
    }
    for (const call of reply.tool_calls) {
      emit({ type: 'tool_execution_start', name: call.function.name })
      const output = await tools[call.function.name]!(JSON.parse(call.function.arguments || '{}'))
      emit({ type: 'tool_execution_end', name: call.function.name, output })
      messages.push({ role: 'tool', tool_call_id: call.id, content: output })
    }
    emit({ type: 'turn_end', turn, usage })
  }
  throw new Error(`No answer after ${maxTurns} turns`)
}

// 4. Your app subscribes. The loop above doesn't change to add one.
on((e) => e.type === 'text_delta' && lyrics.append(e.text))
on((e) => tape.write(e))
on((e) => e.type === 'turn_end' && e.usage && meter.add(e.usage.prompt_tokens + e.usage.completion_tokens))
```

**Python**

```python
# Streaming and events from scratch. Standard library only, no SDK.
import json
import urllib.request

# Any OpenAI-compatible endpoint: OpenAI, Ollama, LM Studio, vLLM, a proxy...
LLM = {
    "base_url": "http://localhost:11434/v1",  # e.g. Ollama's default address
    "model": "your-model",  # e.g. "llama3.1", "gpt-4o-mini"
    "api_key": "YOUR_API_KEY",  # local servers usually ignore it
}

TOOLS = {"crowd_mood": crowd_mood, "set_lights": set_lights, "pyro": pyro}
TOOL_SCHEMAS = [...]  # one JSON Schema per tool

# 1. The bus: a list of listeners. The loop emits; it never knows who is listening.
listeners = []

def on(listener):
    listeners.append(listener)
    return lambda: listeners.remove(listener)

def emit(event):
    for listener in listeners:
        listener(event)

# 2. One streamed model call: read the server-sent events and rebuild the message piece by piece.
def stream_reply(messages):
    body = {
        "model": LLM["model"],
        "messages": messages,
        "tools": TOOL_SCHEMAS,
        "stream": True,
        "stream_options": {"include_usage": True},
    }
    request = urllib.request.Request(
        f"{LLM['base_url']}/chat/completions",
        data=json.dumps(body).encode(),
        headers={"Content-Type": "application/json", "Authorization": f"Bearer {LLM['api_key']}"},
    )
    text, calls, usage = "", {}, None
    with urllib.request.urlopen(request) as response:
        for raw in response:  # one server-sent event per line
            line = raw.decode().strip()
            if not line.startswith("data: ") or line == "data: [DONE]":
                continue
            chunk = json.loads(line[6:])
            usage = chunk.get("usage") or usage
            delta = (chunk.get("choices") or [{}])[0].get("delta", {})
            if delta.get("content"):
                text += delta["content"]
                emit({"type": "text_delta", "text": delta["content"]})  # show it now, not at the end
            # Tool calls arrive in pieces too: glue each one's arguments together by index.
            for part in delta.get("tool_calls") or []:
                call = calls.setdefault(part["index"], {"id": "", "type": "function", "function": {"name": "", "arguments": ""}})
                call["id"] = call["id"] or part.get("id") or ""
                call["function"]["name"] += part.get("function", {}).get("name") or ""
                call["function"]["arguments"] += part.get("function", {}).get("arguments") or ""
    reply = {"role": "assistant", "content": text or None}
    if calls:
        reply["tool_calls"] = [calls[i] for i in sorted(calls)]
    return reply, usage

# 3. The loop from level 2, emitting as it goes.
def run_agent(prompt, max_turns=10):
    messages = [{"role": "user", "content": prompt}]
    for turn in range(1, max_turns + 1):
        emit({"type": "turn_start", "turn": turn})
        reply, usage = stream_reply(messages)
        messages.append(reply)
        if not reply.get("tool_calls"):
            emit({"type": "turn_end", "turn": turn, "usage": usage})
            return reply["content"] or ""
        for call in reply["tool_calls"]:
            name = call["function"]["name"]
            emit({"type": "tool_execution_start", "name": name})
            output = TOOLS[name](**json.loads(call["function"]["arguments"] or "{}"))
            emit({"type": "tool_execution_end", "name": name, "output": output})
            messages.append({"role": "tool", "tool_call_id": call["id"], "content": output})
        emit({"type": "turn_end", "turn": turn, "usage": usage})
    raise RuntimeError(f"No answer after {max_turns} turns")

# 4. Your app subscribes. The loop above doesn't change to add one.
on(lambda e: e["type"] == "text_delta" and print(e["text"], end="", flush=True))
on(lambda e: tape.write(e))
on(lambda e: e["type"] == "turn_end" and e["usage"] and meter.add(e["usage"]["total_tokens"]))
```

## Points de vigilance

- **Le streaming ne rend pas le modèle plus rapide.** La réponse prend autant de temps qu'avant ; ce qui change, c'est le moment où la personne voit le premier mot. Mesurez le temps jusqu'au premier token, pas seulement le temps total.
- **Les auditeurs regardent ; ils ne pilotent pas.** Ce qu'un auditeur renvoie est ignoré, et la boucle ne l'attend pas. Pour bloquer un outil ou changer ce que lit le modèle, il vous faut un hook : le niveau suivant.
- **Gardez les auditeurs rapides.** Le bus d'astorlm les appelle l'un après l'autre, dans la boucle. Le travail lent, comme écrire dans une base de données, va dans une file où l'auditeur se contente de pousser.
- **Un auditeur cassé échoue en silence.** astorlm attrape son erreur pour que l'exécution continue, ce qui veut aussi dire que personne ne le sait. Journalisez dans vos auditeurs.
- **Les morceaux ne respectent pas les mots.** Un text_delta peut s'arrêter au milieu d'un mot ou d'un tableau Markdown. Affichez ce que vous avez et redessinez à mesure que la suite arrive.
- **Laissez la personne arrêter.** Dès qu'elle voit la réponse arriver, elle voudra couper court à une réponse qui part mal. Branchez un bouton stop sur `agent.abort()`.

## Patterns liés

- [2 · La boucle de l'agent](https://harnesspatterns.dev/fr/patterns/agent-loop.md)
- [7 · Les hooks](https://harnesspatterns.dev/fr/patterns/hooks.md)
- [15 · Observabilité et évaluations](https://harnesspatterns.dev/fr/patterns/observability.md)
- [17 · L'humain dans la boucle](https://harnesspatterns.dev/fr/patterns/human-in-the-loop.md)
