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.