// SSE 端点(P5:ai.StreamChat gRPC stream → BFF → 前端 EventSource) // 路由:GET /ai/chat/stream // 裁决依据:02-architecture-design.md §10 import { Controller, Get, Query, Res, Req } from "@nestjs/common"; import type { Request, Response } from "express"; import { TeacherService } from "./teacher.service.js"; import type { CallContext } from "../clients/types.js"; import type { ChatMessage } from "../clients/ai/ai.types.js"; @Controller("/ai/chat") export class AiChatSseController { constructor(private readonly teacherService: TeacherService) {} /** * SSE 流式聊天端点 * 前端使用 EventSource 连接,BFF 透传 ai.StreamChat gRPC stream */ @Get("stream") async streamChat( @Req() req: Request, @Res() res: Response, @Query("prompt") prompt: string, ): Promise { // 设置 SSE 响应头 res.setHeader("Content-Type", "text/event-stream"); res.setHeader("Cache-Control", "no-cache"); res.setHeader("Connection", "keep-alive"); res.setHeader("X-Accel-Buffering", "no"); // 禁用 Nginx 缓冲 // 构造 CallContext(从 req.headers 提取 x-user-id 和 x-request-id) const ctx: CallContext = { userId: (req.headers["x-user-id"] as string) ?? "dev-user", traceId: (req.headers["x-request-id"] as string) ?? "", }; // 构造聊天消息 const messages: ChatMessage[] = [{ role: "user", content: prompt }]; try { const stream = await this.teacherService.streamChat(ctx, messages); for await (const chunk of stream) { const data = JSON.stringify({ content: chunk.content, done: chunk.done, }); res.write(`data: ${data}\n\n`); if (chunk.done) { break; } } res.write("event: done\ndata: {}\n\n"); } catch (err) { const errorData = JSON.stringify({ error: err instanceof Error ? err.message : "Unknown error", }); res.write(`event: error\ndata: ${errorData}\n\n`); } finally { res.end(); } } }