feat: P1 풀빌드 — React 프론트 7화면 + 백엔드 상담루프·평가·음성·RAG

web (Vite+React19+TS, Cloudflare Pages 배포):
- 디자인토큰(세이지틸/테라코타 SSOT), 앱셸, 공통 UI 프리미티브
- 7화면: 로그인/학습자홈/상담세션/회기리뷰/교수자/관리자/설정
- ClientAvatar: SVG 반구상 흉상 4상태 + RMS 립싱크 + 6파라미터 정서
- 회기리뷰는 외부 레퍼런스 디자인을 Vignette 토큰으로 리스킨

api (FastAPI):
- 게이트웨이 /v1/generate·/v1/stream 어댑터(상주풀/EngineSession 보존)
- services: 페르소나 L0~L6 빌더 / 결정론 상태머신 / 가드레일 /
  턴 오케스트레이터 / 회기간 메모리 / 평가AI / 음성 / RAG
- store: DB off 폴백(in-memory), sessions 실구현

검증:
- web: node22 tsc+vite build 통과(node23 segfault 회피), Pages 배포 200
- api: app.main import 통과
- 핫픽스: Topbar initials undefined-safe (undefined.trim 크래시)
- E2E: 서연(P1) 상담 1턴 — 좋은/나쁜 상담에 차등 반응 실증
This commit is contained in:
Yun Chan 2026-06-25 23:37:22 +09:00
parent 859ab26314
commit 24b1b7a6e1
84 changed files with 19645 additions and 107 deletions

205
apps/api/app/routes/eval.py Normal file
View file

@ -0,0 +1,205 @@
"""평가 라우트 — 교수자/관리자용 평가 조회 + 재평가 트리거 (Features: evaluator).
services/evaluator.py 2-loop 평가(fast/deep) 교수자(TEACHER)·관리자(ADMIN)에게 노출한다.
학습자(LEARNER)에겐 평가 결과가 직접 노출되지 않는다(설계서 §4.1 2-레이어: RBAC × AIView).
라우터는 RBAC(레이어2) 강제한다 require_role(TEACHER, ADMIN). 평가 AI 전부 봐도 되므로
(레이어1 AIView.EVALUATOR) CCD/정답 누설 걱정은 client AI 책임이고 여기선 무관.
엔드포인트:
GET /eval/health 헬스(소유 트랙 전환 확인)
POST /eval/sessions/{id}/turn 단일 fast-loop 재평가 트리거
POST /eval/sessions/{id}/reevaluate 회기 deep-loop 재평가 트리거(전체 축어록)
GET /eval/sessions/{id}/evaluation 회기 평가 조회(분포 + 최근 deep 결과)
DB(feedback_scores/supervisor_comment) SoR 적재는 Phase 2. 현재는 in-proc store + 엔진 직접 호출
(degraded). DB 붙으면 조회 경로를 turns.evaluation / supervisor_comment 조인으로 교체한다.
"""
from __future__ import annotations
from typing import Annotated, Any, Optional
from fastapi import APIRouter, Depends, HTTPException, status
from pydantic import BaseModel, Field
from ..deps import Principal, Role, require_role
from ..engine_client import EngineError, engine_client
from ..services import evaluator
from ..services.evaluator import SessionEvaluation, TurnEvaluation
from ..store import store
router = APIRouter(prefix="/eval", tags=["eval"])
# 교수자/관리자만 평가 조회·트리거 (학습자 비노출)
TeacherOrAdmin = Annotated[Principal, Depends(require_role(Role.TEACHER, Role.ADMIN))]
# ── 요청/응답 모델 ──────────────────────────────────────
class ReevaluateRequest(BaseModel):
scope: str = Field("session_end", description="'session_end' | 'stage_transition'")
class TurnReevaluateRequest(BaseModel):
turn_seq: int = Field(..., ge=0, description="재평가할 상담자 발화의 turn_seq")
class EvaluationSummary(BaseModel):
"""회기 평가 조회 응답(분포 + deep 결과 합본)."""
session_id: str
stage: str
deep: Optional[dict[str, Any]] = None
distribution: dict[str, Any] = Field(default_factory=dict)
# ── in-proc 평가 결과 캐시 (DB 적재 전 degraded 보관) ───────────────────────
# DB 가 붙으면 turns.evaluation / supervisor_comment 로 대체. 지금은 트리거 결과를 보관해
# 조회 GET 이 재호출 없이 마지막 deep 결과를 돌려주게 한다.
_DEEP_CACHE: dict[str, SessionEvaluation] = {}
def _load_session_or_404(session_id: str):
sess = store.get(session_id)
if sess is None:
raise HTTPException(status.HTTP_404_NOT_FOUND, detail="session not found")
return sess
def _theory_mode_of(sess) -> Optional[str]:
tt = getattr(sess.persona, "theory_target", None)
if isinstance(tt, (list, tuple)) and tt:
return ", ".join(str(x) for x in tt)
# store 가 theory_mode 문자열도 보유(InProcSession.theory_mode)
return getattr(sess, "theory_mode", None)
@router.get("/health")
async def eval_health() -> dict[str, str]:
"""평가 라우터 헬스 — Features:evaluator 로 전환됨."""
return {"status": "ok", "owner": "features:evaluator", "loops": "fast,deep"}
# ════════════════════════════════════════════════════════════════════════════
# 회기 deep-loop 재평가 트리거 (교수자/관리자)
# ════════════════════════════════════════════════════════════════════════════
@router.post("/sessions/{session_id}/reevaluate", response_model=SessionEvaluation)
async def reevaluate_session(
session_id: str,
body: ReevaluateRequest,
principal: TeacherOrAdmin,
) -> SessionEvaluation:
"""회기 전체 deep-loop 재평가(슈퍼바이저 rationale/critique + 개선점 + 대안발화).
in-proc store 마스킹 축어록을 evaluator.evaluate_session 으로 평가한다.
엔진 장애는 503 으로 변환(평가는 비치명적이지만 트리거는 사용자 명시 요청이라 에러 노출).
"""
sess = _load_session_or_404(session_id)
masked = sess.masked_turns()
# 발화 seq 보강(deep 프롬프트 가독성 — store 가 seq 미포함이라 인덱스로 부여)
enriched: list[dict[str, Any]] = []
for i, t in enumerate(masked):
item = dict(t)
item.setdefault("seq", i)
enriched.append(item)
# 누적 기법 코드 — DB 미가용이라 fast 결과가 없으면 빈 분포(deep LLM 정성 평가는 그대로 유효).
technique_codes: list[str] = []
try:
result = await evaluator.evaluate_session(
session_id=session_id,
stage=sess.state.stage.value,
masked_turns=enriched,
engine=engine_client,
technique_codes=technique_codes,
theory_mode=_theory_mode_of(sess),
scope=body.scope if body.scope in ("session_end", "stage_transition") else "session_end",
)
except EngineError as e:
raise HTTPException(status.HTTP_503_SERVICE_UNAVAILABLE, detail=f"engine unavailable: {e}")
if result.error and result.error.startswith("engine_error"):
raise HTTPException(status.HTTP_503_SERVICE_UNAVAILABLE, detail=result.error)
_DEEP_CACHE[session_id] = result
return result
# ════════════════════════════════════════════════════════════════════════════
# 단일 턴 fast-loop 재평가 트리거 (교수자/관리자)
# ════════════════════════════════════════════════════════════════════════════
@router.post("/sessions/{session_id}/turn", response_model=TurnEvaluation)
async def reevaluate_turn(
session_id: str,
body: TurnReevaluateRequest,
principal: TeacherOrAdmin,
) -> TurnEvaluation:
"""단일 상담자 발화 fast-loop 재평가(기법/내담자상태/적절성/의도이탈).
store 축어록에서 해당 turn_seq 상담자 발화 + 직후 내담자 응답을 재구성해
경량 TurnContext evaluator.evaluate_turn 호출한다.
"""
sess = _load_session_or_404(session_id)
# 대상 상담자 발화 + 직후 내담자 응답 찾기
target_idx: Optional[int] = None
for i, tr in enumerate(sess.turns):
if tr.speaker == "counselor" and tr.turn_seq == body.turn_seq:
target_idx = i
break
if target_idx is None:
raise HTTPException(
status.HTTP_404_NOT_FOUND, detail=f"counselor turn_seq {body.turn_seq} not found"
)
learner = sess.turns[target_idx]
client_reply = ""
if target_idx + 1 < len(sess.turns) and sess.turns[target_idx + 1].speaker == "client":
client_reply = sess.turns[target_idx + 1].text_masked
# 평가용 경량 TurnContext 재구성(prepare_turn 의 결정론 산출과 동형). 엔진 호출 없음.
from ..services.orchestrator import TurnContext # 지연 import(소유권 경계)
recent = [
{"speaker": tr.speaker, "text": tr.text_masked} for tr in sess.turns[max(0, target_idx - 4):target_idx]
]
ctx = TurnContext(
session_id=session_id,
case_id=sess.case_id,
persona=sess.persona,
state_before=sess.state,
learner_text_raw=learner.text,
learner_text_masked=learner.text_masked,
state_after=sess.state, # 조회 시점 상태(정밀 재현은 DB 스냅샷 도입 시)
recent_turns=recent,
)
result = await evaluator.evaluate_turn(ctx, client_reply, engine=engine_client)
if result.error and result.error.startswith("engine_error"):
raise HTTPException(status.HTTP_503_SERVICE_UNAVAILABLE, detail=result.error)
return result
# ════════════════════════════════════════════════════════════════════════════
# 회기 평가 조회 (교수자/관리자) — 마지막 deep 결과 + 분포
# ════════════════════════════════════════════════════════════════════════════
@router.get("/sessions/{session_id}/evaluation", response_model=EvaluationSummary)
async def get_session_evaluation(
session_id: str,
principal: TeacherOrAdmin,
) -> EvaluationSummary:
"""회기 평가 조회(읽기) — 마지막 deep 재평가 결과 + 기법 분포.
DB 적재 degraded: deep 결과는 reevaluate 트리거가 보관한 캐시에서, 분포는 결과에서.
아직 평가 트리거가 없었다면 deep=None + 분포.
"""
_load_session_or_404(session_id)
cached = _DEEP_CACHE.get(session_id)
if cached is None:
return EvaluationSummary(session_id=session_id, stage="", deep=None, distribution={})
return EvaluationSummary(
session_id=session_id,
stage=cached.stage,
deep=cached.to_dict(),
distribution=cached.distribution.model_dump(),
)

321
apps/api/app/routes/kb.py Normal file
View file

@ -0,0 +1,321 @@
"""지식베이스(KB) 라우트 — 하이브리드 검색 + 인덱싱 트리거(관리자).
설계서 §3.6·§4.3 / MASTERPLAN §3.6:
GET /kb/health 라우터 + RAG 구성요소 readiness
POST /kb/search 정적 지식 하이브리드 검색(정책 4-튜플, visible_to DB 강제)
POST /kb/eval-grounding 평가 AI 채점 근거(evaluator 정책, label_id 동봉)
POST /kb/index 문서 인덱싱 트리거(관리자, content_hash 증분, 오프라인 배치)
정보비대칭은 *DB WHERE* 강제한다(services/rag.py POLICIES). 라우트는 role
요청 컨텍스트로만 결정하고, 검색 함수가 정책 화이트리스트를 고정한다(코드경로 부재 1차방어).
DB/임베딩 모델 미가용(Docker off, rag 의존성 미설치)이면 rag.NotConfigured 503 변환.
무거운 import(FlagEmbedding/torch) services/rag.py 함수 내부로 가둔다(이식성).
"""
from __future__ import annotations
from typing import Annotated, Any, Literal, Optional
from fastapi import APIRouter, Depends, HTTPException, status
from pydantic import BaseModel, Field
from ..db import acquire, get_pool
from ..deps import AIView, Principal, Role, require_role
from ..services import rag
router = APIRouter(prefix="/kb", tags=["kb"])
# ── 요청/응답 모델 ──────────────────────────────────────
RoleLiteral = Literal["client", "counselor", "evaluator"]
class KBSearchRequest(BaseModel):
query: str = Field(..., min_length=1) # PII 마스킹된 질의(마스킹은 가드레일 책임)
role: RoleLiteral = "evaluator" # 검색 주체(정보비대칭 정책 선택)
k: int = Field(default=5, ge=1, le=50)
rerank: bool = True
# 정책 화이트리스트를 *좁히는* 추가 필터만 허용(넓히지 못함 — 정보비대칭 보존)
kb_kind: Optional[list[str]] = None
source_id: Optional[list[str]] = None
sensitivity_max: Optional[int] = Field(default=None, ge=0, le=3)
# 감사 귀속(선택)
session_id: Optional[str] = None
turn_id: Optional[str] = None
class ChunkOut(BaseModel):
chunk_id: int
score: float
kb_kind: str
heading_path: Optional[str] = None
context_prefix: Optional[str] = None
body: Optional[str] = None # expose_body=True(상담사/평가) 정책에서만
behavior_cue: Optional[str] = None # 내담자 정책: 본문 비노출, 행동단서만(M6)
label_id: Optional[int] = None # 평가 정책에서만
meta: dict[str, Any] = Field(default_factory=dict)
source_id: Optional[str] = None
class KBSearchResponse(BaseModel):
chunks: list[ChunkOut]
policy: str
top1_score: float
crag_pass: bool # top1 >= 임계(F-06: 미달 시 관찰 프레이밍)
latency_ms: int
degraded: bool = False # reranker/embed 폴백 투명성
class MemoryRecallRequest(BaseModel):
case_id: str = Field(..., min_length=1) # UUID — 학습자별 케이스 스코프(M5/T4)
query: str = Field(..., min_length=1)
k: int = Field(default=5, ge=1, le=20)
session_id: Optional[str] = None
turn_id: Optional[str] = None
class IndexChunkIn(BaseModel):
seq: int
chunk_text: str = Field(..., min_length=1)
heading_path: Optional[str] = None
context_prefix: Optional[str] = None # Contextual Retrieval 프리픽스(색인 대상)
kb_kind: Optional[str] = None
visible_to: Optional[list[str]] = None # 미지정 시 {client,counselor,evaluator}
sensitivity: Optional[int] = Field(default=None, ge=0, le=3)
label_id: Optional[int] = None # taxonomy 정답 라벨 FK
meta: Optional[dict[str, Any]] = None
token_count: Optional[int] = None
class IndexRequestIn(BaseModel):
source_id: str = Field(..., min_length=1)
doc_uri: str = Field(..., min_length=1)
version: int = 1
content_hash: Optional[str] = None
chunks: list[IndexChunkIn]
class IndexResponse(BaseModel):
doc_id: Optional[int]
chunks_indexed: int
skipped_unchanged: bool # content_hash 동일 → 증분 스킵
embedded: bool # 임베딩 적재 여부(모델 미가용 시 False)
degraded: bool = False
# ── 헬퍼: rag.NotConfigured → 503 ───────────────────────
def _to_chunk_out(c: rag.RetrievedChunk) -> ChunkOut:
return ChunkOut(
chunk_id=c.chunk_id,
score=round(c.score, 6),
kb_kind=c.kb_kind,
heading_path=c.heading_path,
context_prefix=c.context_prefix,
body=c.body,
behavior_cue=c.behavior_cue,
label_id=c.label_id,
meta=c.meta,
source_id=c.source_id,
)
# ════════════════════════════════════════════════════════════════════════════
# 헬스 — 라우터 + RAG readiness (모델/DB 미가용도 정직하게 보고)
# ════════════════════════════════════════════════════════════════════════════
@router.get("/health")
async def kb_health() -> dict[str, object]:
"""KB 라우터 + RAG 구성요소 readiness.
DB /임베딩 모델 가용 여부를 *크래시 없이* 점검(미가용=degraded). 부트/디버그용.
"""
db_ready = False
try:
get_pool()
db_ready = True
except RuntimeError:
db_ready = False
# 임베딩 모델은 무거우므로 *로드하지 않고* 설치 가능성만 가볍게 확인(import 시도 X).
return {
"status": "ok" if db_ready else "degraded",
"owner": "features:rag",
"db_pool": db_ready,
"crag_threshold": rag.CRAG_TOP1_THRESHOLD,
"policies": [r.value for r in rag.POLICIES],
}
# ════════════════════════════════════════════════════════════════════════════
# 지식 검색 — 정책 4-튜플(role)로 분기, visible_to DB 강제
# ════════════════════════════════════════════════════════════════════════════
@router.post("/search", response_model=KBSearchResponse)
async def search(body: KBSearchRequest) -> KBSearchResponse:
"""정적 지식 KB 하이브리드 검색(dense pgvector cosine + sparse tsvector + 리랭킹).
role 정책 4-튜플(사전필터·가중치·본문노출·라벨) 고정한다 호출부가 넓힌다.
AI RLS 컨텍스트(app.current_ai_view) 커넥션에 주입해 visible_to 2 강제.
"""
ai_role = rag.AIRole(body.role)
filters: dict[str, Any] = {}
if body.kb_kind:
filters["kb_kind"] = body.kb_kind
if body.source_id:
filters["source_id"] = body.source_id
if body.sensitivity_max is not None:
filters["sensitivity_max"] = body.sensitivity_max
# RLS 컨텍스트(레이어1): app.current_ai_view = role → visible_to WHERE DB 강제
try:
async with acquire(ai_view=AIView(body.role).value) as conn:
result = await rag.search_kb(
conn,
query=body.query,
role=ai_role,
k=body.k,
filters=filters or None,
rerank=body.rerank,
)
# 감사 적재(best-effort — 로그 실패가 검색을 막지 않음)
try:
await rag.log_retrieval(
conn,
result=result,
ai_role=body.role,
session_id=body.session_id,
turn_id=body.turn_id,
)
except Exception:
pass
except rag.NotConfigured as e:
raise HTTPException(
status.HTTP_503_SERVICE_UNAVAILABLE,
detail=f"RAG not configured: {e}",
)
except RuntimeError as e:
# DB 풀 미초기화(lifespan 밖) — 시연/테스트 degraded
raise HTTPException(status.HTTP_503_SERVICE_UNAVAILABLE, detail=f"DB not ready: {e}")
return KBSearchResponse(
chunks=[_to_chunk_out(c) for c in result.chunks],
policy=result.policy_name,
top1_score=round(result.top1_score, 6),
crag_pass=result.top1_score >= rag.CRAG_TOP1_THRESHOLD,
latency_ms=result.latency_ms,
degraded=result.degraded,
)
# ════════════════════════════════════════════════════════════════════════════
# 평가 근거 — evaluator 정책 래퍼(label_id 동봉, CRAG 게이트)
# ════════════════════════════════════════════════════════════════════════════
@router.post("/eval-grounding", response_model=KBSearchResponse)
async def eval_grounding(body: KBSearchRequest) -> KBSearchResponse:
"""평가 AI 채점 근거 회수(DSM/이론/taxonomy 정답라벨 + 논평).
role 무시하고 evaluator 정책 고정(평가 전용 경로). crag_pass=False 호출부가
'관찰 프레이밍'으로 다운그레이드(F-06).
"""
try:
async with acquire(ai_view=AIView.EVALUATOR.value) as conn:
result = await rag.retrieve_eval_grounding(
conn,
query=body.query,
k=body.k,
kinds=body.kb_kind,
rerank=body.rerank,
)
try:
await rag.log_retrieval(
conn,
result=result,
ai_role="evaluator",
session_id=body.session_id,
turn_id=body.turn_id,
)
except Exception:
pass
except rag.NotConfigured as e:
raise HTTPException(status.HTTP_503_SERVICE_UNAVAILABLE, detail=f"RAG not configured: {e}")
except RuntimeError as e:
raise HTTPException(status.HTTP_503_SERVICE_UNAVAILABLE, detail=f"DB not ready: {e}")
return KBSearchResponse(
chunks=[_to_chunk_out(c) for c in result.chunks],
policy=result.policy_name,
top1_score=round(result.top1_score, 6),
crag_pass=result.top1_score >= rag.CRAG_TOP1_THRESHOLD,
latency_ms=result.latency_ms,
degraded=result.degraded,
)
# ════════════════════════════════════════════════════════════════════════════
# 페르소나 메모리 회상 — 내담자 연속성(case 스코프, episodic)
# ════════════════════════════════════════════════════════════════════════════
@router.post("/persona-memory", response_model=KBSearchResponse)
async def persona_memory(body: MemoryRecallRequest) -> KBSearchResponse:
"""회기 시작 episodic recall(app.turn_embedding, case_id 스코프 강제).
CCD/정답/평가는 경로에 구조적으로 부재(코드경로 부재 1차방어). 반환은 turn_id+점수만
(본문은 호출부 memory.build_recall_context turns 조인). 내담자 RLS 주입.
"""
try:
async with acquire(ai_view=AIView.CLIENT.value) as conn:
result = await rag.retrieve_persona_memory(
conn,
case_id=body.case_id,
query=body.query,
k=body.k,
)
except rag.NotConfigured as e:
raise HTTPException(status.HTTP_503_SERVICE_UNAVAILABLE, detail=f"RAG not configured: {e}")
except RuntimeError as e:
raise HTTPException(status.HTTP_503_SERVICE_UNAVAILABLE, detail=f"DB not ready: {e}")
return KBSearchResponse(
chunks=[_to_chunk_out(c) for c in result.chunks],
policy=result.policy_name,
top1_score=round(result.top1_score, 6),
crag_pass=result.top1_score >= rag.CRAG_TOP1_THRESHOLD,
latency_ms=result.latency_ms,
degraded=result.degraded,
)
# ════════════════════════════════════════════════════════════════════════════
# 인덱싱 트리거 — 관리자 전용(content_hash 증분, 오프라인 배치)
# ════════════════════════════════════════════════════════════════════════════
@router.post("/index", response_model=IndexResponse, status_code=status.HTTP_202_ACCEPTED)
async def index_document(
body: IndexRequestIn,
principal: Annotated[Principal, Depends(require_role(Role.ADMIN))],
) -> IndexResponse:
"""문서 인덱싱(관리자, RBAC ADMIN 강제). content_hash 증분 + 청크 임베딩 적재.
임베딩은 무거운 작업 본래 BackgroundTasks/배치 워커 위임 권장(202 Accepted).
DSM verbatim 저작권(license C/D) source 등록 시점 external_llm_ok 가드 책임.
모델 미가용 embedding NULL 폴백(BM25 , degraded=True) 크래시 X.
"""
req = rag.IndexRequest(
source_id=body.source_id,
doc_uri=body.doc_uri,
version=body.version,
content_hash=body.content_hash,
chunks=[c.model_dump() for c in body.chunks],
)
try:
# 관리자 인덱싱은 RLS 미적용(쓰기 — kb 스키마 직접). role 주입 없이 acquire.
async with acquire() as conn:
result = await rag.index_document(conn, req)
except rag.NotConfigured as e:
raise HTTPException(status.HTTP_503_SERVICE_UNAVAILABLE, detail=f"RAG not configured: {e}")
except RuntimeError as e:
raise HTTPException(status.HTTP_503_SERVICE_UNAVAILABLE, detail=f"DB not ready: {e}")
return IndexResponse(
doc_id=result.doc_id,
chunks_indexed=result.chunks_indexed,
skipped_unchanged=result.skipped_unchanged,
embedded=result.embedded,
degraded=result.degraded,
)

View file

@ -1,13 +1,15 @@
"""상담 세션 라우트 — 시작 / 턴 / 종료 + SSE 스트림 스텁.
"""상담 세션 라우트 — 시작 / 턴 / 스트림 / 종료 (services 실호출).
흐름 (설계서 §2 회기 라이프사이클 + 마스터플랜 §2.2 사이클):
POST /sessions 회기 시작 (case_profile/summary 회상 + session_state 초기화)
POST /sessions/{id}/turn 수련생 발화 1 (가드레일상태머신내담자AI평가)
GET /sessions/{id}/stream 내담자 AI 응답 SSE 스트림 (Cloudflare 우회 heartbeat)
POST /sessions/{id}/end 회기 종료 (무손실 carry-over + LLM 압축 트리거)
POST /sessions 회기 시작 (페르소나 + 회상 + 상태머신 init)
POST /sessions/{id}/turn 수련생 발화 1 (가드레일상태머신내담자AI출력가드)
GET /sessions/{id}/stream 내담자 AI 응답 SSE 스트림 (heartbeat 포함)
POST /sessions/{id}/end 회기 종료 (무손실 carry-over + 압축 트리거)
DB(NAS Postgres) SoR 이지만 Docker off 에서도 엔진만 있으면 1턴이 돌도록
**store(in-memory)** 폴백을 둔다(degraded). 인증도 dev 폴백을 허용한다(개발 편의).
상태머신(라포탐색개입정리) 백엔드가 결정론적으로 소유(LLM 아님, 마스터플랜 §0).
파일은 핸들러 시그니처 + 계약 + TODO. 실제 상태머신/가드레일/압축은 Phase 1~2a 트랙 A.
"""
from __future__ import annotations
@ -15,19 +17,41 @@ from __future__ import annotations
import asyncio
import json
from typing import Annotated, Literal, Optional
from uuid import UUID, uuid4
from fastapi import APIRouter, Depends, HTTPException, status
from pydantic import BaseModel, Field
from sse_starlette.sse import EventSourceResponse
from fastapi import Cookie
from ..config import settings
from ..deps import CurrentPrincipal, HumanDB
from ..engine_client import EngineMessage, StreamRequest, engine_client, EngineError
from ..deps import Principal, Role
from ..engine_client import EngineError, engine_client
from ..services import memory, orchestrator, persona, state_machine
from ..store import TurnRecord, store
router = APIRouter(prefix="/sessions", tags=["sessions"])
Stage = Literal["라포", "탐색", "개입", "정리"]
StageLiteral = Literal["라포", "탐색", "개입", "정리"]
# ── 인증 — dev 폴백 허용 (쿠키 없으면 dev learner) ───────────────────────────
async def get_principal_dev(
session_cookie: Annotated[Optional[str], Cookie(alias="__Host-vignette_sid")] = None,
) -> Principal:
"""세션 쿠키 → Principal. 미인증(쿠키 없음)이면 dev learner 폴백.
개발/시연(쿠키 없음, DB off)에서도 상담 루프가 돌게 한다.
prod 에선 auth.py BFF + Redis 세션이 완성되면 deps.get_current_principal 교체.
TODO: Redis 세션 룩업으로 user_id/role/cohort 복원.
"""
if not session_cookie:
return Principal(user_id="dev-learner", role=Role.LEARNER, cohort_ids=[])
# TODO: Redis 세션 검증. 현재는 쿠키 존재만으로 dev learner.
return Principal(user_id="dev-user", role=Role.LEARNER, cohort_ids=[])
DevPrincipal = Annotated[Principal, Depends(get_principal_dev)]
# ── 요청/응답 모델 ──────────────────────────────────────
@ -37,12 +61,13 @@ class SessionStartRequest(BaseModel):
class SessionStartResponse(BaseModel):
session_id: UUID
case_id: UUID
session_id: str
case_id: str
session_no: int
stage: Stage
# 회기 시작 회상 요약 (큰그림→세부, UI 카드용. CCD/정답은 절대 미포함)
stage: StageLiteral
effective_openness: float
recall_summary: Optional[str] = None
degraded: bool = False # DB 미가용 in-proc 모드 여부(시연 투명성)
class TurnRequest(BaseModel):
@ -51,138 +76,269 @@ class TurnRequest(BaseModel):
class TurnResponse(BaseModel):
turn_seq: int
stage: Stage
stage: StageLiteral
effective_openness: float
# 내담자 응답은 스트림(GET /stream)으로 받는 게 기본. 동기 응답은 폴백/테스트용.
client_reply: Optional[str] = None
safety_flagged: bool = False
crisis_kind: str = "none"
class SessionEndResponse(BaseModel):
session_id: UUID
session_id: str
session_no: int
digest_pending: bool # 압축은 비동기 비블로킹 (설계서 §2-C)
end_state: dict
# ── 핸들러 ──────────────────────────────────────────────
# ════════════════════════════════════════════════════════════════════════════
# 회기 시작
# ════════════════════════════════════════════════════════════════════════════
@router.post("", response_model=SessionStartResponse, status_code=status.HTTP_201_CREATED)
async def start_session(
body: SessionStartRequest,
principal: CurrentPrincipal,
conn: HumanDB,
principal: DevPrincipal,
) -> SessionStartResponse:
"""회기 시작 — 회상 + 상태 복원 (설계서 §2-A).
"""회기 시작 — 페르소나 핀 + 회상 + 결정론 상태 init (설계서 §2-A).
절차:
1. persona_card(approved) 조회 + (persona_id, learner_id) -> case_profile upsert
2. case_digest + 직전 session_summary + episodic recall (Phase 2a, 1차는 단일회기)
3. session_state 초기화: stage='라포', carry-over (rapport×0.7, ideation 보수적 유지) [P2]
4. sessions insert
TODO: persona 조회/회상/상태머신 init 구현 (트랙 A). 현재 스텁 응답.
DB 가용 : persona_card(approved) 조회 + case_profile/직전 summary 회상.
DB 미가용(degraded): 시드 페르소나(persona.SEED) + 회상( 회기)으로 in-proc.
"""
# TODO: SELECT persona_id FROM app.persona_card WHERE code=$1 AND status='approved'
# TODO: init_session_state_from_history() — 결정론 carry-over
session_id = uuid4()
case_id = uuid4()
return SessionStartResponse(
session_id=session_id,
case_id=case_id,
card = persona.get_seed_persona(body.persona_code)
if card is None:
# TODO: DB app.persona_card WHERE code=$1 AND status='approved' 조회 경로
raise HTTPException(status.HTTP_404_NOT_FOUND, detail=f"unknown persona {body.persona_code}")
# 회상 — DB/RAG 미가용 시 빈 컨텍스트(첫 회기). 가용 시 case_digest/summary/episodic 주입.
# TODO(Phase 2a): memory.build_recall_context(case_digest=..., prev_summary=..., episodic_snippets=...)
recall = memory.build_recall_context()
# 결정론 상태 init (carry-over 가 있으면 이월; 첫 회기는 None)
st = state_machine.init_state(
base_resistance=card.base_resistance(),
unlock_rate=card.unlock_rate(),
decay_floor=card.decay_floor(),
ideation_baseline=card.ideation_baseline(),
carry=recall.carry,
)
sess = store.create(
learner_id=principal.user_id,
persona=card,
theory_mode=body.theory_mode,
state=st,
session_no=1,
stage="라포",
recall_summary=None, # Phase 2a 회상 채움
)
# 회상 핀(pinned facts)을 세션에 묶어 둔다(턴마다 재조립). store 는 간단히 state 만 보유하므로
# recall_summary/pinned 는 in-proc 캐시로 별도 보관.
_RECALL_CACHE[sess.session_id] = recall
return SessionStartResponse(
session_id=sess.session_id,
case_id=sess.case_id,
session_no=sess.session_no,
stage=st.stage.value, # type: ignore[arg-type]
effective_openness=round(st.effective_openness, 4),
recall_summary=recall.recall_summary,
degraded=True, # 현재 in-proc 경로(DB 붙으면 False 분기)
)
# 회상 컨텍스트 in-proc 캐시 (회기 내 재사용, recall_context). DB 붙으면 session_state.recall_context.
_RECALL_CACHE: dict[str, memory.RecallContext] = {}
def _load_session_or_404(session_id: str):
sess = store.get(session_id)
if sess is None:
raise HTTPException(status.HTTP_404_NOT_FOUND, detail="session not found")
if sess.ended:
raise HTTPException(status.HTTP_409_CONFLICT, detail="session already ended")
return sess
# ════════════════════════════════════════════════════════════════════════════
# 턴 (동기 폴백 — 기본 UX 는 /stream)
# ════════════════════════════════════════════════════════════════════════════
@router.post("/{session_id}/turn", response_model=TurnResponse)
async def submit_turn(
session_id: UUID,
session_id: str,
body: TurnRequest,
principal: CurrentPrincipal,
conn: HumanDB,
principal: DevPrincipal,
) -> TurnResponse:
"""수련생 발화 1턴 (마스터플랜 §2.2 / 설계서 §2-B).
파이프라인 (전부 백엔드 결정론 게이트):
1. [입력 가드레일] Presidio PII 마스킹 + 위기분류(실제위기 vs 페르소나 연기) [R7/F-03]
2. [상태머신] effective_openness = clamp(stage.openness
+ rapport_credit*unlock_rate - resistance*decay, 0, 1) [P2, 결정론]
3. [모순 검사] pinned_fact locked 모순 -> 차단·재생성 (설계서 §2-B)
4. [내담자 AI] engine_client.stream/generate (CCD 직접노출 금지, Structured Outputs)
5. [출력 가드레일] 자살수단 차단, ideation_stage <= 3 상한 [R5]
6. [평가 AI] fast-loop 4차원 태깅 (deep-loop 단계전환/회기말)
7. [working 갱신] session_state UPSERT (체크포인트)
8. [로깅] turns insert + 임베딩 (재귀학습 원천)
TODO: 1~8 구현 (트랙 A). 현재 스텁: 발화 검증만.
오케스트레이터로 1~8단계 결정론 파이프라인 실행. 내담자 응답은 동기로 번에 받는다
(기본 UX GET /stream 토큰 스트리밍; 경로는 폴백/테스트).
"""
# TODO: load session_state, run guardrail + state machine deterministically
# 동기 응답은 폴백. 기본 UX 는 GET /stream 으로 토큰 스트리밍.
sess = _load_session_or_404(session_id)
recall = _RECALL_CACHE.get(session_id) or memory.RecallContext()
ctx = orchestrator.prepare_turn(
session_id=session_id,
case_id=sess.case_id,
card=sess.persona,
state=sess.state,
learner_text=body.text,
recall_summary=recall.recall_summary,
pinned_facts=recall.pinned_facts,
recent_turns=sess.recent_turns(),
)
# 수련생 발화 로깅(② episodic 미러) — 마스킹본 저장
assert ctx.state_after is not None
store.append_turn(
session_id,
TurnRecord(
turn_seq=ctx.state_after.turn_seq,
speaker="counselor",
stage=ctx.state_after.stage.value,
text=body.text,
text_masked=ctx.learner_text_masked,
),
)
try:
result = await orchestrator.run_turn_generate(ctx, engine_client)
except EngineError as e:
raise HTTPException(status.HTTP_503_SERVICE_UNAVAILABLE, detail=f"engine unavailable: {e}")
# 내담자 응답 로깅 + 상태 체크포인트(① working UPSERT 미러)
if result.client_reply:
store.append_turn(
session_id,
TurnRecord(
turn_seq=result.turn_seq,
speaker="client",
stage=result.stage,
text=result.client_reply,
text_masked=result.client_reply, # 내담자 응답은 합성(원문PII 없음)
),
)
store.update_state(session_id, result.state_after)
return TurnResponse(
turn_seq=0,
stage="라포",
effective_openness=0.15,
client_reply=None,
safety_flagged=False,
turn_seq=result.turn_seq,
stage=result.stage, # type: ignore[arg-type]
effective_openness=round(result.effective_openness, 4),
client_reply=result.client_reply,
safety_flagged=result.safety_flagged,
crisis_kind=result.crisis_kind,
)
@router.get("/{session_id}/stream")
async def stream_client_reply(
session_id: UUID,
principal: CurrentPrincipal,
# ════════════════════════════════════════════════════════════════════════════
# 스트림 (기본 UX — SSE 토큰)
# ════════════════════════════════════════════════════════════════════════════
@router.post("/{session_id}/stream")
async def stream_turn(
session_id: str,
body: TurnRequest,
principal: DevPrincipal,
):
"""내담자 AI 응답 SSE 스트림 (마스터플랜 §1.1 SSE 분리경로).
"""수련생 발화 1턴을 받아 내담자 AI 응답을 SSE 토큰 스트림으로 흘린다.
- Cloudflare 100 timeout 회피: settings.sse_heartbeat_seconds 마다 ping 이벤트 [R2]
- 게이트웨이 SSE(engine_client.stream) 프록시해 토큰을 재방출
- 이벤트: {event: "token"|"done"|"safety"|"ping", data: ...}
TODO: 게이트웨이와 StreamRequest 바디 결합(현재 stage 회상 컨텍스트 없이 placeholder).
상태머신 컨텍스트(L3 stage/openness) + 마스킹된 최근 N턴 주입.
- Cloudflare 100 timeout 회피: settings.sse_heartbeat_seconds 마다 ping [R2]
- 오케스트레이터 run_turn_stream(가드레일·상태머신·페르소나·출력가드 적용) 프록시
- 이벤트: token | ping | safety | done | error
"""
sess = _load_session_or_404(session_id)
recall = _RECALL_CACHE.get(session_id) or memory.RecallContext()
ctx = orchestrator.prepare_turn(
session_id=session_id,
case_id=sess.case_id,
card=sess.persona,
state=sess.state,
learner_text=body.text,
recall_summary=recall.recall_summary,
pinned_facts=recall.pinned_facts,
recent_turns=sess.recent_turns(),
)
assert ctx.state_after is not None
# 수련생 발화 로깅 + 상태 체크포인트(스트림은 응답 전 상태 갱신 — 결정론이라 무방)
store.append_turn(
session_id,
TurnRecord(
turn_seq=ctx.state_after.turn_seq,
speaker="counselor",
stage=ctx.state_after.stage.value,
text=body.text,
text_masked=ctx.learner_text_masked,
),
)
store.update_state(session_id, ctx.state_after)
async def event_generator():
# heartbeat 와 엔진 스트림을 병행 (Cloudflare 버퍼링/타임아웃 회피)
last_beat = asyncio.get_event_loop().time()
# TODO: 실제 StreamRequest 조립 — session_state 에서 stage/openness/최근턴 로드
req = StreamRequest(
ai_role="client",
tier="client",
messages=[
EngineMessage(role="system", content="<persona L0~L2 cache_control 주입 TODO>", cache=True),
EngineMessage(role="user", content="<masked latest learner turn TODO>"),
],
)
final_reply = ""
try:
async for chunk in engine_client.stream(req):
yield {"event": "token", "data": chunk}
async for ev in orchestrator.run_turn_stream(ctx, engine_client):
if ev.event == "token":
final_reply += ev.data.get("text", "")
yield {"event": ev.event, "data": json.dumps(ev.data, ensure_ascii=False)}
now = asyncio.get_event_loop().time()
if now - last_beat >= settings.sse_heartbeat_seconds:
yield {"event": "ping", "data": "{}"}
last_beat = now
yield {"event": "done", "data": json.dumps({"session_id": str(session_id)})}
except EngineError as e:
yield {"event": "error", "data": json.dumps({"detail": str(e)})}
except Exception as e: # 방어 — 어떤 예외도 SSE error 프레임으로
yield {"event": "error", "data": json.dumps({"detail": str(e)}, ensure_ascii=False)}
return
# 내담자 응답 로깅(② episodic) — 스트림 종료 후
if final_reply:
store.append_turn(
session_id,
TurnRecord(
turn_seq=ctx.state_after.turn_seq,
speaker="client",
stage=ctx.state_after.stage.value,
text=final_reply,
text_masked=final_reply,
),
)
return EventSourceResponse(event_generator())
# ════════════════════════════════════════════════════════════════════════════
# 회기 종료
# ════════════════════════════════════════════════════════════════════════════
@router.post("/{session_id}/end", response_model=SessionEndResponse)
async def end_session(
session_id: UUID,
principal: CurrentPrincipal,
conn: HumanDB,
session_id: str,
principal: DevPrincipal,
) -> SessionEndResponse:
"""회기 종료 — carry-over + 압축 트리거 (설계서 §2-C, 비동기 비블로킹).
"""회기 종료 — 무손실 carry-over + 압축 트리거 (설계서 §2-C, 비동기 비블로킹).
절차:
(A) 무손실 carry-over: end_state = session_state 종료 snapshot (코드 복사, LLM 미경유) [P4]
(B) salience 산출 -> 망각/유지
(C) narrative 압축 (LLM, 상주 claude -p 재사용) 비동기
(D~G) digest 임베딩 / case_profile 병합 / pinned 모순처리 / RAG 동기화
+ 상주 프로세스 회수 (말투표류 회기경계 차단)
TODO: (A) 동기 수행 (B~G) BackgroundTasks/큐로 비블로킹. 현재 스텁.
(A) 무손실 carry-over: end_state = state.snapshot() (코드 복사, LLM 미경유) [P4]
(C) narrative 압축(LLM) CompressionJob 으로 큐잉(여기선 페이로드만; 실제 호출은 후속 워커)
"""
# TODO: UPDATE app.sessions SET ended_at=now(); copy end_state; enqueue compression
return SessionEndResponse(session_id=session_id, session_no=1, digest_pending=True)
sess = store.get(session_id)
if sess is None:
raise HTTPException(status.HTTP_404_NOT_FOUND, detail="session not found")
recall = _RECALL_CACHE.get(session_id) or memory.RecallContext()
carry = memory.make_carry_over(
state=sess.state,
session_id=session_id,
case_id=sess.case_id,
session_no=sess.session_no,
masked_turns=sess.masked_turns(),
prev_rapport_credit=sess.prev_rapport_credit,
open_threads=recall.open_threads,
)
# TODO(Phase 2a): BackgroundTasks 로 carry.compression_job 을
# engine_client.generate(GenerateRequest(ai_role='evaluator', tier='feedback',
# messages=memory.build_compression_messages(job))) 호출 → session_summary UPSERT + 임베딩.
# 현재는 큐잉만(digest_pending=True). DB 없으면 압축 결과 적재 생략.
store.end(session_id)
_RECALL_CACHE.pop(session_id, None)
return SessionEndResponse(
session_id=session_id,
session_no=sess.session_no,
digest_pending=carry.compression_job is not None,
end_state=carry.end_state,
)

View file

@ -0,0 +1,437 @@
"""음성 라우트 — OpenAI STT/TTS 캐스케이드 + WSS 실시간 턴테이킹.
한신대 요구 '음성 필수'. 학습자가 마이크로 말하면 STT orchestrator 상담 1
내담자 텍스트 TTS 오디오 + 립싱크 힌트(설계 §4.3 RMS) 역방향으로 흘린다.
캐스케이드(설계 §5.2 음성 오브 4상태 listeningthinkingspeakingidle):
[클라] audio_start(JSON) 바이너리 오디오 청크들 audio_end(JSON)
[서버] state(listening) STT transcript(JSON) state(thinking)
orchestrator.run_turn(가드레일·상태머신·페르소나·내담자AI·출력가드)
reply(JSON, 내담자 텍스트 + stage/openness) state(speaking)
[tts_chunk(JSON: seq/rms) + 바이너리 오디오] × N tts_end(JSON) state(idle)
프로토콜(JSON 제어 + 바이너리 오디오 혼합, 단일 WS):
- 클라서버 텍스트 = JSON 제어({"type": ...}); 클라서버 바이너리 = 오디오 청크
- 서버클라 텍스트 = JSON 이벤트; 서버클라 바이너리 = TTS 오디오 청크
- TTS 바이너리 청크 *직전* 메타 JSON(tts_chunk: seq, rms) 보내 프론트가 짝짓는다.
음성 미설정(OPENAI_API_KEY 없음): GET /voice/health 503 degraded,
WS 핸드셰이크 직후 degraded 이벤트 + close(1011). 절대 크래시 금지.
"""
from __future__ import annotations
import json
from typing import Optional
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
from fastapi.responses import JSONResponse
from starlette.websockets import WebSocketState
from ..engine_client import EngineError, engine_client
from ..services import memory, orchestrator, persona
from ..services import voice as voice_svc
from ..services.voice import VoicePreset, VoiceUnavailable, resolve_voice, voice_service
from ..store import TurnRecord, store
router = APIRouter(prefix="/voice", tags=["voice"])
# WS close 코드(섹션별 의미 명시)
WS_CLOSE_DEGRADED = 1011 # 서버측 음성 미설정/장애
WS_CLOSE_BAD_REQUEST = 1008 # 프로토콜 위반(세션 누락 등)
# 한 발화당 누적 오디오 상한(메모리 방어, ~10MB)
_MAX_AUDIO_BYTES = 10 * 1024 * 1024
# ════════════════════════════════════════════════════════════════════════════
# 헬스 — 음성 가용성(키 설정) 노출
# ════════════════════════════════════════════════════════════════════════════
@router.get("/health")
async def voice_health() -> JSONResponse:
"""음성 라우터 헬스. 키 미설정이면 503 degraded(시연 투명성)."""
available = voice_service.is_available()
body = {
"status": "ok" if available else "degraded",
"available": available,
"stt_model": voice_svc.STT_MODEL,
"tts_model": voice_svc.TTS_MODEL,
"reason": None if available else "OPENAI_API_KEY 미설정",
}
return JSONResponse(body, status_code=200 if available else 503)
# ════════════════════════════════════════════════════════════════════════════
# WebSocket — 실시간 음성 캐스케이드
# ════════════════════════════════════════════════════════════════════════════
@router.websocket("/ws")
async def voice_ws(websocket: WebSocket) -> None:
"""음성 실시간 턴 캐스케이드.
쿼리: ?session_id=<hex> (없으면 persona_code 일회용 in-proc 세션 생성 시연용)
오디오 in(바이너리) STT 상담 1 TTS out(바이너리) + 립싱크 힌트.
"""
await websocket.accept()
# 1) 음성 미설정 → degraded 알리고 정상 종료(크래시 금지)
if not voice_service.is_available():
await _safe_send_json(
websocket,
{"type": "degraded", "reason": "OPENAI_API_KEY 미설정 — 음성 기능 비활성"},
)
await _safe_close(websocket, WS_CLOSE_DEGRADED)
return
# 2) 세션 바인딩 — session_id 우선, 없으면 persona_code 로 시연 세션 생성
session_id, voice_preset, err = _bind_session(websocket)
if err is not None:
await _safe_send_json(websocket, {"type": "error", "detail": err})
await _safe_close(websocket, WS_CLOSE_BAD_REQUEST)
return
assert session_id is not None and voice_preset is not None
await _safe_send_json(
websocket,
{
"type": "ready",
"session_id": session_id,
"voice": voice_preset.openai_voice,
"preset": voice_preset.preset,
"state": "idle",
},
)
audio_buf = bytearray()
receiving = False
try:
while True:
msg = await websocket.receive()
mtype = msg.get("type")
if mtype == "websocket.disconnect":
break
# ── 바이너리 = 오디오 청크 누적 ──
if msg.get("bytes") is not None:
if not receiving:
# audio_start 없이 들어온 바이너리 — 관용적으로 자동 시작
receiving = True
audio_buf.clear()
await _safe_send_json(websocket, {"type": "state", "state": "listening"})
audio_buf.extend(msg["bytes"])
if len(audio_buf) > _MAX_AUDIO_BYTES:
await _safe_send_json(
websocket,
{"type": "error", "detail": "audio too large — 발화를 짧게 끊어 주세요"},
)
audio_buf.clear()
receiving = False
continue
# ── 텍스트 = JSON 제어 ──
text = msg.get("text")
if text is None:
continue
try:
ctrl = json.loads(text)
except (json.JSONDecodeError, TypeError):
await _safe_send_json(websocket, {"type": "error", "detail": "invalid control json"})
continue
ctype = ctrl.get("type")
if ctype == "audio_start":
receiving = True
audio_buf.clear()
await _safe_send_json(websocket, {"type": "state", "state": "listening"})
elif ctype == "audio_end":
receiving = False
await _handle_utterance(
websocket,
session_id=session_id,
voice_preset=voice_preset,
audio=bytes(audio_buf),
fmt=ctrl.get("format"),
)
audio_buf.clear()
elif ctype == "text_turn":
# 음성 없이 텍스트만 보내는 경로(접근성/디버그): STT 건너뛰고 바로 턴.
receiving = False
audio_buf.clear()
learner_text = (ctrl.get("text") or "").strip()
if learner_text:
await _run_turn_and_speak(
websocket,
session_id=session_id,
voice_preset=voice_preset,
learner_text=learner_text,
)
elif ctype == "ping":
await _safe_send_json(websocket, {"type": "pong"})
elif ctype == "close":
break
except WebSocketDisconnect:
pass
except Exception as e: # 어떤 예외도 WS 를 깨끗이 닫고 알린다(크래시 금지)
await _safe_send_json(websocket, {"type": "error", "detail": f"voice ws error: {e}"})
finally:
await _safe_close(websocket)
# ════════════════════════════════════════════════════════════════════════════
# 발화 1건 처리 — STT → 턴 → TTS
# ════════════════════════════════════════════════════════════════════════════
async def _handle_utterance(
websocket: WebSocket,
*,
session_id: str,
voice_preset: VoicePreset,
audio: bytes,
fmt: Optional[str],
) -> None:
"""오디오 1발화 → STT → 상담 턴 → TTS 캐스케이드."""
if not audio:
await _safe_send_json(websocket, {"type": "transcript", "text": "", "final": True})
await _safe_send_json(websocket, {"type": "state", "state": "idle"})
return
# 1) STT (thinking 진입)
await _safe_send_json(websocket, {"type": "state", "state": "thinking"})
filename, content_type = _audio_meta(fmt)
try:
stt = await voice_service.transcribe(
audio, filename=filename, content_type=content_type
)
except VoiceUnavailable as e:
await _safe_send_json(websocket, {"type": "degraded", "reason": str(e)})
await _safe_send_json(websocket, {"type": "state", "state": "idle"})
return
except Exception as e:
await _safe_send_json(websocket, {"type": "error", "detail": f"STT 실패: {e}"})
await _safe_send_json(websocket, {"type": "state", "state": "idle"})
return
learner_text = stt.text
await _safe_send_json(
websocket,
{"type": "transcript", "text": learner_text, "final": True, "speaker": "counselor"},
)
if not learner_text:
# 무음/인식 실패 — 턴 진행 안 함
await _safe_send_json(websocket, {"type": "state", "state": "idle"})
return
await _run_turn_and_speak(
websocket,
session_id=session_id,
voice_preset=voice_preset,
learner_text=learner_text,
)
async def _run_turn_and_speak(
websocket: WebSocket,
*,
session_id: str,
voice_preset: VoicePreset,
learner_text: str,
) -> None:
"""상담 1턴(orchestrator) → 내담자 텍스트 → TTS 오디오/립싱크 힌트 역방향 전송."""
sess = store.get(session_id)
if sess is None or sess.ended:
await _safe_send_json(websocket, {"type": "error", "detail": "세션 없음/종료됨"})
await _safe_send_json(websocket, {"type": "state", "state": "idle"})
return
recall = memory.RecallContext()
ctx = orchestrator.prepare_turn(
session_id=session_id,
case_id=sess.case_id,
card=sess.persona,
state=sess.state,
learner_text=learner_text,
recall_summary=recall.recall_summary,
pinned_facts=recall.pinned_facts,
recent_turns=sess.recent_turns(),
)
assert ctx.state_after is not None
# 학습자 발화 로깅(마스킹본) — sessions.py 패턴과 동일
store.append_turn(
session_id,
TurnRecord(
turn_seq=ctx.state_after.turn_seq,
speaker="counselor",
stage=ctx.state_after.stage.value,
text=learner_text,
text_masked=ctx.learner_text_masked,
),
)
# 2) 내담자 AI 1턴(동기 — 음성은 TTS 전 전체 텍스트가 필요)
try:
result = await orchestrator.run_turn_generate(ctx, engine_client)
except EngineError as e:
await _safe_send_json(websocket, {"type": "error", "detail": f"engine unavailable: {e}"})
await _safe_send_json(websocket, {"type": "state", "state": "idle"})
return
reply = result.client_reply or ""
# 내담자 응답 로깅 + 상태 체크포인트
if reply:
store.append_turn(
session_id,
TurnRecord(
turn_seq=result.turn_seq,
speaker="client",
stage=result.stage,
text=reply,
text_masked=reply,
),
)
store.update_state(session_id, result.state_after)
# 내담자 텍스트 이벤트(설계 §5.3 자막 — partial 없이 final)
await _safe_send_json(
websocket,
{
"type": "reply",
"text": reply,
"speaker": "client",
"stage": result.stage,
"effective_openness": round(result.effective_openness, 4),
"turn_seq": result.turn_seq,
"safety_flagged": result.safety_flagged,
"crisis_kind": result.crisis_kind,
},
)
if not reply:
await _safe_send_json(websocket, {"type": "state", "state": "idle"})
return
# 3) TTS (speaking) — 청크별 메타 JSON(립싱크 rms) + 바이너리 오디오
await _safe_send_json(
websocket,
{"type": "state", "state": "speaking", "voice": voice_preset.openai_voice},
)
try:
n = 0
async for ck in voice_service.synthesize_stream(reply, voice_preset):
# 메타 먼저(프론트가 직후 바이너리와 짝지음) — 설계 §4.3 RMS 1채널
await _safe_send_json(
websocket, {"type": "tts_chunk", "seq": ck.seq, "rms": round(ck.rms, 4)}
)
await _safe_send_bytes(websocket, ck.audio)
n += 1
await _safe_send_json(websocket, {"type": "tts_end", "chunks": n})
except VoiceUnavailable as e:
await _safe_send_json(websocket, {"type": "degraded", "reason": str(e)})
except Exception as e:
await _safe_send_json(websocket, {"type": "error", "detail": f"TTS 실패: {e}"})
await _safe_send_json(websocket, {"type": "state", "state": "idle"})
# ════════════════════════════════════════════════════════════════════════════
# 세션 바인딩 / 메타 헬퍼
# ════════════════════════════════════════════════════════════════════════════
def _bind_session(
websocket: WebSocket,
) -> tuple[Optional[str], Optional[VoicePreset], Optional[str]]:
"""쿼리에서 세션을 바인딩(또는 시연 세션 생성)하고 voice preset 을 해석.
우선순위:
?session_id=<hex> 기존 세션(REST 시작된) 음성 부착
?persona_code=P1[&preset=] in-proc 시연 세션 생성(DB off 폴백)
반환 (session_id, voice_preset, error).
"""
qp = websocket.query_params
explicit_preset = qp.get("preset")
session_id = qp.get("session_id")
if session_id:
sess = store.get(session_id)
if sess is None:
return None, None, f"unknown session {session_id}"
if sess.ended:
return None, None, "session already ended"
vp = resolve_voice(persona_code=sess.persona.code, preset=explicit_preset)
return session_id, vp, None
# persona_code 로 시연 세션 생성(REST 미경유 음성 단독 데모)
persona_code = qp.get("persona_code")
if not persona_code:
return None, None, "session_id 또는 persona_code 쿼리 필요"
card = persona.get_seed_persona(persona_code)
if card is None:
return None, None, f"unknown persona {persona_code}"
from ..services import state_machine
st = state_machine.init_state(
base_resistance=card.base_resistance(),
unlock_rate=card.unlock_rate(),
decay_floor=card.decay_floor(),
ideation_baseline=card.ideation_baseline(),
)
sess = store.create(
learner_id="dev-learner-voice",
persona=card,
theory_mode="humanistic",
state=st,
session_no=1,
)
vp = resolve_voice(persona_code=card.code, preset=explicit_preset)
return sess.session_id, vp, None
def _audio_meta(fmt: Optional[str]) -> tuple[str, str]:
"""클라가 알려준 포맷 → (filename, content_type). 기본 webm/opus."""
f = (fmt or "webm").lower().lstrip(".")
table = {
"webm": ("audio.webm", "audio/webm"),
"ogg": ("audio.ogg", "audio/ogg"),
"opus": ("audio.ogg", "audio/ogg"),
"wav": ("audio.wav", "audio/wav"),
"mp3": ("audio.mp3", "audio/mpeg"),
"mp4": ("audio.mp4", "audio/mp4"),
"m4a": ("audio.m4a", "audio/mp4"),
"pcm": ("audio.wav", "audio/wav"),
}
return table.get(f, ("audio.webm", "audio/webm"))
# ── 안전 송수신(연결 끊김 시 조용히 무시) ───────────────────────────────────
async def _safe_send_json(websocket: WebSocket, payload: dict) -> None:
if websocket.client_state != WebSocketState.CONNECTED:
return
try:
await websocket.send_text(json.dumps(payload, ensure_ascii=False))
except Exception:
pass
async def _safe_send_bytes(websocket: WebSocket, data: bytes) -> None:
if websocket.client_state != WebSocketState.CONNECTED:
return
try:
await websocket.send_bytes(data)
except Exception:
pass
async def _safe_close(websocket: WebSocket, code: int = 1000) -> None:
if websocket.client_state == WebSocketState.DISCONNECTED:
return
try:
await websocket.close(code=code)
except Exception:
pass
__all__ = ["router"]