Files
resume-agent/backend/app/builder_sse.py
T
hypandClaude ae2d9b128d feat: builder 简历生成 + 轻度优化 + 简历导入交付副本
自内部仓库剥离深度优化与 RAG 知识库后的交付版本:
- Builder 对话式简历生成(FSM + 意图路由 LLM 兜底增强)
- 条目级轻度优化:事实覆盖门禁 + STAR/bullet 修复链,功能/简介/成果与技术栈同级保护
- 简历导入:DOCX/PDF 解析、结构归一、手机号脱敏
- PostgreSQL 运行时 + Alembic 迁移链

Co-Authored-By: Claude <noreply@anthropic.com>
2026-08-05 11:08:30 +08:00

75 lines
2.7 KiB
Python

"""SSE transport for observable Builder chat turns."""
from __future__ import annotations
import json
import logging
from collections.abc import Callable, Iterator
from queue import Queue
from threading import Thread
from typing import Any
from fastapi.responses import StreamingResponse
from .fsm import FSMError
from .models import ActionResponse
BuilderOperation = Callable[[], ActionResponse]
PHASE_LABELS = {
"suggesting_next": "正在判断下一步建议",
"structuring": "正在整理信息",
"checking_gaps": "正在检查可补充的信息",
"rewriting": "正在生成候选改写",
"saving": "正在写入简历",
}
def stream_builder_message(operation: BuilderOperation) -> StreamingResponse:
events: Queue[tuple[str, dict[str, Any]] | None] = Queue()
def emit(event: str, data: dict[str, Any] | None = None) -> None:
events.put((event, data or {}))
def worker() -> None:
try:
result = operation()
_emit_statuses(emit, tuple(result.builder_stream_phases))
for chunk in _chunks(str((result.turn.content if result.turn else "") or "")):
emit("delta", {"text": chunk})
emit("complete", result.model_dump(mode="json"))
except FSMError as exc:
emit("error", {"code": exc.code, "message": exc.message, "status_code": exc.status_code})
except Exception as exc: # pragma: no cover - defensive transport boundary
logging.getLogger(__name__).exception("builder SSE operation failed")
emit("error", {"code": "builder_stream_failed", "message": "Resume assistant stream failed. Please retry.", "status_code": 502, "reason_code": type(exc).__name__})
finally:
events.put(None)
def generate() -> Iterator[str]:
thread = Thread(target=worker, name="resume-builder-sse", daemon=True)
thread.start()
while True:
item = events.get()
if item is None:
break
event, data = item
yield _frame(event, data)
return StreamingResponse(generate(), media_type="text/event-stream", headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"})
def _emit_statuses(emit: Callable[[str, dict[str, Any]], None], phases: tuple[str, ...]) -> None:
for phase in phases:
emit("status", {"phase": phase, "label": PHASE_LABELS[phase]})
def _chunks(text: str, size: int = 24) -> Iterator[str]:
if not text:
return
for index in range(0, len(text), size):
yield text[index : index + size]
def _frame(event: str, data: dict[str, Any]) -> str:
return f"event: {event}\ndata: {json.dumps(data, ensure_ascii=False, separators=(',', ':'))}\n\n"