Chuyển đến nội dung chính

第 8 課:Streaming 和 EventBus — 即時回應

實作串流聊天:SSE(伺服器傳送事件)、AsyncGenerator、EventBus 模式。逐一令牌交付、工具執行事件、進度追蹤。背壓處理。

🧠 人工智慧與機器學習 — 第 7 課 第 8 課:串流媒體與 EventBus — 即時 回應

從零開始搭建AI代理平台-與xClaw實戰

第 2 部分:LLM 引擎和代理核心

亞洲開發網

簡介

用戶不想等待 10 秒才能收到完整的回應。串流傳輸允許逐個令牌傳送-使用者立即看到結果。本文使用AsyncGenerator + SSE實現串流。


1. 代理 chatStream()

// packages/core/src/agent/agent.ts
async *chatStream(
  userMessage: string,
  context: ToolContext,
  additionalTools?: AdditionalTool[],
): AsyncGenerator<StreamEvent> {
  const messages = await this.buildMessages(userMessage, context);
  const tools = this.gatherTools(additionalTools);

  let iterations = 0;
  const maxIterations = this.config.maxToolIterations || 10;

  while (iterations < maxIterations) {
    let fullContent = '';
    const toolCalls: ToolCall[] = [];
    const toolCallArgs = new Map<string, string>();

    // Stream from LLM
    for await (const event of this.llm.chatStream(messages, tools)) {
      switch (event.type) {
        case 'text-delta':
          fullContent += event.delta;
          yield event; // Forward to client immediately
          break;

        case 'tool-call-start':
          yield event;
          break;

        case 'tool-call-args':
          const existing = toolCallArgs.get(event.toolCallId) || '';
          toolCallArgs.set(event.toolCallId, existing + event.args);
          yield event;
          break;

        case 'tool-call-end':
          const args = toolCallArgs.get(event.toolCallId) || '{}';
          toolCalls.push({
            id: event.toolCallId,
            name: '', // set from tool-call-start
            arguments: JSON.parse(args),
          });
          yield event;
          break;

        case 'finish':
          yield event;
          break;
      }
    }

    iterations++;

    // No tool calls → done
    if (toolCalls.length === 0) {
      await this.memory.save(context.sessionId, userMessage, fullContent);
      return;
    }

    // Execute tools & stream results
    messages.push({
      role: 'assistant',
      content: fullContent,
      toolCalls,
    });

    const results = await this.executeToolCalls(toolCalls, context);
    for (const result of results) {
      yield { type: 'tool-result', toolCallId: result.toolCallId, result };
      messages.push({
        role: 'tool',
        content: JSON.stringify(result.result),
        toolCallId: result.toolCallId,
      });
    }
    // Loop back — LLM receives tool results
  }
}

2.SSE端點

// packages/gateway/src/routes/chat.ts
app.post('/api/chat/stream', async (c) => {
  const { message, sessionId, model } = await c.req.json();
  const user = c.get('user');

  const context: ToolContext = {
    tenantId: user.tenantId,
    userId: user.sub,
    sessionId,
  };

  return streamSSE(c, async (stream) => {
    const generator = agent.chatStream(message, context);

    for await (const event of generator) {
      await stream.writeSSE({
        event: event.type,
        data: JSON.stringify(event),
      });
    }
  });
});

3.EventBus 模式

// packages/core/src/events/event-bus.ts
type EventHandler<T = unknown> = (data: T) => void | Promise<void>;

export class EventBus {
  private handlers = new Map<string, Set<EventHandler>>();

  on<T>(event: string, handler: EventHandler<T>) {
    if (!this.handlers.has(event)) {
      this.handlers.set(event, new Set());
    }
    this.handlers.get(event)!.add(handler as EventHandler);

    // Return unsubscribe function
    return () => this.handlers.get(event)?.delete(handler as EventHandler);
  }

  async emit<T>(event: string, data: T) {
    const handlers = this.handlers.get(event);
    if (!handlers) return;

    await Promise.all(
      Array.from(handlers).map(h => h(data)),
    );
  }
}

// Usage
const bus = new EventBus();
bus.on('chat:start', ({ sessionId }) => auditLog.log('chat_started', sessionId));
bus.on('tool:executed', ({ name, duration }) => metrics.recordToolExecution(name, duration));
bus.on('chat:complete', ({ usage }) => billing.recordUsage(usage));

4. 客戶端消費

// Frontend: consuming SSE stream
async function streamChat(message: string) {
  const response = await fetch('/api/chat/stream', {
    method: 'POST',
    headers: { 'Content-Type': 'application/json', Authorization: `Bearer ${token}` },
    body: JSON.stringify({ message, sessionId }),
  });

  const reader = response.body!.getReader();
  const decoder = new TextDecoder();
  let buffer = '';

  while (true) {
    const { done, value } = await reader.read();
    if (done) break;

    buffer += decoder.decode(value, { stream: true });
    const lines = buffer.split('\n');
    buffer = lines.pop() || '';

    for (const line of lines) {
      if (line.startsWith('data: ')) {
        const event = JSON.parse(line.slice(6));
        handleStreamEvent(event);
      }
    }
  }
}

function handleStreamEvent(event: StreamEvent) {
  switch (event.type) {
    case 'text-delta':
      appendToChat(event.delta); // Append token to UI
      break;
    case 'tool-call-start':
      showToolIndicator(event.toolName); // Show "Searching..."
      break;
    case 'tool-result':
      hideToolIndicator();
      break;
  }
}

5. 總結

  • AsyncGenerator — 優雅的串流 API,可組合
  • SSE — HTTP 原生,無需 WebSocket,自動重新連接
  • EventBus — 將副作用(稽核、計費、指標)與主流程分離
  • 背壓 — reader.read() 自行處理

下一篇文章: 文件處理器 - RAG Pipeline 的文檔處理。