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

Lesson 7: Streaming & Real-time — SSE, WebSocket, Voice Agent & Multimodal

Server-Sent Events streaming, WebSocket real-time, Whisper STT, TTS (ElevenLabs/OpenAI), multimodal input (image/audio/video), voice agent architecture, latency optimization.

🏗️ Architecture — Lesson 7 Lesson 7: Streaming & Real-time — SSE, WebSocket, Voice Agent & Multimodal

Enterprise AI Chatbot Platform Architecture — From Prototype to Production

Part 2: Core Chatbot Engine

xdev.asia

1. Streaming Architecture Overview

In chatbot production, streaming response is mandatory — users do not accept waiting 5-10 seconds to receive a full response. Streaming allowed display each token As soon as LLM generates, reduces perceived latency from ~5s to ~200ms.


┌──────────┐     WebSocket/SSE      ┌──────────────┐     SSE Stream      ┌───────────┐
│  Client  │ ◀──────────────────── │  API Gateway  │ ◀────────────────── │ LLM       │
│ (Browser │     Token-by-token     │  (Streaming   │     Token-by-token  │ Provider  │
│  /Mobile)│                        │   Proxy)      │                     │           │
└──────────┘                        └──────┬───────┘                     └───────────┘
                                           │
                                    ┌──────▼───────┐
                                    │  Event Store │
                                    │  (Kafka)     │
                                    └──────────────┘

2. Server-Sent Events (SSE) Implementation


// SSE Streaming Controller
class StreamingController {
  async streamChat(req: Request, res: Response): Promise<void> {
    // Set SSE headers
    res.writeHead(200, {
      'Content-Type': 'text/event-stream',
      'Cache-Control': 'no-cache',
      'Connection': 'keep-alive',
      'X-Accel-Buffering': 'no', // Disable Nginx buffering
    });

    const { conversationId, message } = req.body;

    try {
      // 1. Process input (RAG, context, etc.)
      const context = await this.chatService.prepareContext(conversationId, message);

      // 2. Stream from LLM
      const stream = await this.llmGateway.chatStream(context);
      let fullContent = '';

      for await (const chunk of stream) {
        fullContent += chunk.content;

        // Send SSE event
        const event: StreamEvent = {
          type: 'token',
          data: {
            content: chunk.content,
            tokenIndex: chunk.index,
          },
        };
        res.write(`event: token\ndata: ${JSON.stringify(event.data)}\n\n`);
      }

      // 3. Send citations after content
      if (context.citations.length > 0) {
        res.write(`event: citations\ndata: ${JSON.stringify(context.citations)}\n\n`);
      }

      // 4. Send completion event
      res.write(`event: done\ndata: ${JSON.stringify({
        messageId: crypto.randomUUID(),
        totalTokens: stream.usage?.totalTokens,
      })}\n\n`);

      // 5. Persist message async (don't block stream)
      this.messageQueue.publish('message.completed', {
        conversationId,
        content: fullContent,
        citations: context.citations,
        usage: stream.usage,
      });

    } catch (error) {
      res.write(`event: error\ndata: ${JSON.stringify({
        code: 'STREAM_ERROR',
        message: 'An error occurred while generating response',
      })}\n\n`);
    } finally {
      res.end();
    }
  }
}

3. WebSocket Real-time Communication


// WebSocket Gateway with Socket.IO
class ChatWebSocketGateway {
  private io: SocketIOServer;
  private connections = new Map<string, SocketConnection>();

  initialize(server: HTTPServer): void {
    this.io = new SocketIOServer(server, {
      cors: { origin: process.env.ALLOWED_ORIGINS?.split(',') },
      transports: ['websocket', 'polling'],
      pingInterval: 25000,
      pingTimeout: 60000,
    });

    this.io.use(this.authMiddleware.bind(this));
    this.io.on('connection', this.handleConnection.bind(this));
  }

  private async handleConnection(socket: Socket): Promise<void> {
    const { userId, tenantId } = socket.data;

    // Track connection
    this.connections.set(socket.id, {
      userId,
      tenantId,
      connectedAt: Date.now(),
    });

    // Join tenant room
    socket.join(`tenant:${tenantId}`);
    socket.join(`user:${userId}`);

    // Handle events
    socket.on('chat:message', (data) => this.handleMessage(socket, data));
    socket.on('chat:typing', (data) => this.handleTyping(socket, data));
    socket.on('chat:stop', (data) => this.handleStopGeneration(socket, data));
    socket.on('disconnect', () => this.handleDisconnect(socket));
  }

  private async handleMessage(socket: Socket, data: ChatMessageEvent): Promise<void> {
    const { conversationId, content, attachments } = data;

    // Emit typing indicator
    socket.emit('chat:bot_typing', { conversationId, isTyping: true });

    try {
      const stream = await this.chatService.streamResponse(
        socket.data.tenantId,
        conversationId,
        content,
        attachments,
      );

      const abortController = new AbortController();
      this.activeStreams.set(`${socket.id}:${conversationId}`, abortController);

      for await (const chunk of stream) {
        if (abortController.signal.aborted) break;

        socket.emit('chat:token', {
          conversationId,
          content: chunk.content,
          index: chunk.index,
        });
      }

      socket.emit('chat:complete', { conversationId });
    } catch (error) {
      socket.emit('chat:error', { conversationId, error: 'Failed to generate response' });
    } finally {
      socket.emit('chat:bot_typing', { conversationId, isTyping: false });
      this.activeStreams.delete(`${socket.id}:${conversationId}`);
    }
  }

  private handleStopGeneration(socket: Socket, data: { conversationId: string }): void {
    const key = `${socket.id}:${data.conversationId}`;
    this.activeStreams.get(key)?.abort();
  }
}

4. Streaming Resilience — Backpressure & Reconnection


// Client-side: Auto-reconnect with EventSource
class ChatStreamClient {
  private eventSource: EventSource | null = null;
  private retryCount = 0;
  private maxRetries = 3;

  connect(conversationId: string, messageId: string): void {
    const url = `/api/chat/stream?conversationId=${conversationId}&messageId=${messageId}`;

    this.eventSource = new EventSource(url);

    this.eventSource.addEventListener('token', (e) => {
      const data = JSON.parse(e.data);
      this.onToken(data.content);
      this.retryCount = 0;
    });

    this.eventSource.addEventListener('done', (e) => {
      const data = JSON.parse(e.data);
      this.onComplete(data);
      this.eventSource?.close();
    });

    this.eventSource.onerror = () => {
      if (this.retryCount < this.maxRetries) {
        this.retryCount++;
        setTimeout(() => this.reconnect(conversationId, messageId), 
          1000 * Math.pow(2, this.retryCount)); // Exponential backoff
      } else {
        this.onError('Connection lost');
      }
    };
  }

  stop(): void {
    this.eventSource?.close();
  }
}

// Server-side: Backpressure handling
class BackpressureStream {
  private buffer: string[] = [];
  private highWaterMark = 100; // Max buffered chunks
  private paused = false;

  async pipe(
    source: AsyncIterable<StreamChunk>,
    sink: Response,
  ): Promise<void> {
    for await (const chunk of source) {
      if (this.buffer.length >= this.highWaterMark) {
        // Backpressure: wait until drain
        await new Promise<void>(resolve => {
          sink.once('drain', resolve);
        });
      }

      const ok = sink.write(
        `event: token\ndata: ${JSON.stringify({ content: chunk.content })}\n\n`
      );

      if (!ok) {
        this.paused = true;
        await new Promise<void>(resolve => sink.once('drain', resolve));
        this.paused = false;
      }
    }
  }
}

5. Voice Agent Architecture — STT + TTS


┌────────────────── VOICE AGENT PIPELINE ──────────────────┐
│                                                           │
│  ┌─────────┐   ┌─────────┐   ┌─────────┐   ┌─────────┐  │
│  │  Audio  │──▶│ Whisper │──▶│ Chat    │──▶│  TTS    │  │
│  │  Input  │   │  STT    │   │ Engine  │   │ Engine  │  │
│  │ (WebRTC)│   │         │   │         │   │         │  │
│  └─────────┘   └─────────┘   └────┬────┘   └────┬────┘  │
│                                    │             │       │
│                              ┌─────▼─────┐  ┌───▼─────┐  │
│                              │  Text     │  │  Audio  │  │
│                              │  Response │  │  Output │  │
│                              └───────────┘  └─────────┘  │
│                                                           │
│  Latency Target: STT ~500ms, LLM ~200ms, TTS ~300ms     │
│  Total: < 1.5s end-to-end                                │
└───────────────────────────────────────────────────────────┘

class VoiceAgent {
  constructor(
    private sttEngine: STTEngine,
    private chatEngine: ChatEngine,
    private ttsEngine: TTSEngine,
  ) {}

  async processVoiceStream(
    audioStream: ReadableStream<Uint8Array>,
    context: VoiceContext,
  ): Promise<ReadableStream<Uint8Array>> {
    // 1. Speech-to-Text (streaming)
    const transcript = await this.sttEngine.transcribe(audioStream);
    console.log(`STT result: "${transcript.text}" (${transcript.language})`);

    // 2. Chat response (streaming)
    const textStream = this.chatEngine.streamResponse(
      context.conversationId,
      transcript.text,
    );

    // 3. Text-to-Speech (streaming — process sentence by sentence)
    return this.ttsEngine.synthesizeStream(textStream, {
      voice: context.voiceId ?? 'vi-VN-female-1',
      speed: context.speed ?? 1.0,
      format: 'opus', // Low latency codec
    });
  }
}

class WhisperSTTEngine implements STTEngine {
  async transcribe(audio: ReadableStream<Uint8Array>): Promise<TranscriptResult> {
    const audioBuffer = await this.collectStream(audio);

    const response = await this.openai.audio.transcriptions.create({
      file: new File([audioBuffer], 'audio.webm', { type: 'audio/webm' }),
      model: 'whisper-1',
      language: 'vi',
      response_format: 'verbose_json',
      timestamp_granularities: ['segment'],
    });

    return {
      text: response.text,
      language: response.language,
      segments: response.segments,
      duration: response.duration,
    };
  }
}

class OpenAITTSEngine implements TTSEngine {
  async synthesizeStream(
    textStream: AsyncIterable<string>,
    options: TTSOptions,
  ): Promise<ReadableStream<Uint8Array>> {
    // Buffer text by sentence for natural speech
    const sentenceBuffer = new SentenceBuffer();

    return new ReadableStream({
      start: async (controller) => {
        for await (const chunk of textStream) {
          sentenceBuffer.append(chunk);

          while (sentenceBuffer.hasSentence()) {
            const sentence = sentenceBuffer.flush();
            const audioChunk = await this.synthesize(sentence, options);
            controller.enqueue(audioChunk);
          }
        }

        // Flush remaining text
        const remaining = sentenceBuffer.flushAll();
        if (remaining) {
          const audioChunk = await this.synthesize(remaining, options);
          controller.enqueue(audioChunk);
        }
        controller.close();
      },
    });
  }

  private async synthesize(text: string, options: TTSOptions): Promise<Uint8Array> {
    const response = await this.openai.audio.speech.create({
      model: 'tts-1',
      voice: options.voice as 'alloy' | 'echo' | 'fable' | 'onyx' | 'nova' | 'shimmer',
      input: text,
      speed: options.speed,
      response_format: options.format,
    });

    return new Uint8Array(await response.arrayBuffer());
  }
}

6. Multimodal Input Processing


interface MultimodalMessage {
  text?: string;
  images?: ImageAttachment[];
  audio?: AudioAttachment;
  files?: FileAttachment[];
}

class MultimodalProcessor {
  async processInput(
    input: MultimodalMessage,
    modelCapabilities: ModelCapabilities,
  ): Promise<LLMMessage> {
    const contentParts: ContentPart[] = [];

    // Text
    if (input.text) {
      contentParts.push({ type: 'text', text: input.text });
    }

    // Images — resize & optimize for vision models
    if (input.images?.length) {
      for (const img of input.images) {
        if (!modelCapabilities.vision) {
          // Fallback: describe image via a vision model
          const description = await this.describeImage(img);
          contentParts.push({
            type: 'text',
            text: `[Image description: ${description}]`,
          });
        } else {
          const optimized = await this.optimizeImage(img, {
            maxWidth: 1024,
            maxHeight: 1024,
            quality: 85,
          });
          contentParts.push({
            type: 'image_url',
            image_url: {
              url: `data:${optimized.mimeType};base64,${optimized.base64}`,
              detail: img.size > 500_000 ? 'high' : 'low', // Cost optimization
            },
          });
        }
      }
    }

    // Audio — transcribe via STT
    if (input.audio) {
      const transcript = await this.sttEngine.transcribe(input.audio.data);
      contentParts.push({
        type: 'text',
        text: `[Audio transcript]: ${transcript.text}`,
      });
    }

    // Files — extract text
    if (input.files?.length) {
      for (const file of input.files) {
        const text = await this.extractText(file);
        contentParts.push({
          type: 'text',
          text: `[File: ${file.name}]\n${text}`,
        });
      }
    }

    return { role: 'user', content: contentParts };
  }
}

7. Latency Optimization Techniques

Technique Impact Implementation
Prompt caching -40% TTFT Cache system prompt prefix (Anthropic/OpenAI)
Speculative decoding -30% latency Use small model to draft, large model to verify
Edge inference -80ms network Deploy small models at edge (Cloudflare Workers AI)
Connection pooling -100ms Keep-alive connections to LLM providers
Response caching -95% latency Semantic cache for identical/similar queries
Parallel RAG -50% prep time Run retrieval + reranking in parallel

// Semantic Cache — cache responses for similar queries
class SemanticCache {
  constructor(
    private vectorStore: VectorStore,
    private similarityThreshold = 0.95,
    private ttlSeconds = 3600,
  ) {}

  async get(query: string, tenantId: string): Promise<CachedResponse | null> {
    const embedding = await this.embedder.embed(query);

    const results = await this.vectorStore.search({
      vector: embedding,
      filter: { tenantId },
      topK: 1,
    });

    if (results.length === 0) return null;
    if (results[0].score < this.similarityThreshold) return null;

    // Check TTL
    const cached = results[0].metadata as CachedResponse;
    if (Date.now() - cached.cachedAt > this.ttlSeconds * 1000) {
      return null;
    }

    return cached;
  }

  async set(query: string, response: string, tenantId: string): Promise<void> {
    const embedding = await this.embedder.embed(query);

    await this.vectorStore.upsert({
      id: `cache:${tenantId}:${this.hashQuery(query)}`,
      vector: embedding,
      metadata: {
        tenantId,
        query,
        response,
        cachedAt: Date.now(),
      },
    });
  }
}

Summary of Lesson 7

  • SSE (Server-Sent Events) is suitable for one-way streaming from server → client, simpler than WebSocket
  • WebSockets for real-time bidirectional: typing indicator, stop generation, multi-user presence
  • Voice Agent = STT (Whisper) + Chat Engine + TTS (OpenAI/ElevenLabs), target < 1.5s end-to-end
  • Multimodal: process images (vision model), audio (STT), files (text extraction) before sending LLM
  • Semantic Cache 95% reduction in latency for similar queries — very high ROI

Next article: Function Calling & Tool Use — design tool registry, safe execution sandbox, output validation, and how LLM calls external APIs.