대시보드 폴드아웃/드릴다운 정리 + 페르소나 역린·misconduct 반응 + 게이트웨이 격리·RAG 비차단 수정

SSOT 대시보드:
- 한신대 기술분석 PDF(19쪽) 정합성 분석 + 이번 세션 발견 섹션 추가
- 섹션 폴드아웃(접기)·상단 목차(드릴다운)·모두 펼치기/접기 — 내용 보존, 레이아웃만 정리

페르소나 반응 강화('저항·반응 조절' 핵심 차별):
- PersonaCard.triggers(역린) 필드 + CCD 핵심상처 파생 역린 블록
- L0에 무례·모욕·조롱 시 현실적 동맹 균열 반응 지침

버그·성능 수정(라이브/E2E로 포착):
- 게이트웨이 페르소나 격리: --append-system-prompt를 --system-prompt(교체)로 + --exclude-dynamic-system-prompt-sections (내담자 캐릭터 붕괴·개발맥락 누출 차단)
- RAG: 임베더 동기 로드(약 7-13초)를 _warm_rag_caches 백그라운드 warm으로(세션 생성 블로킹 회귀 수정)
- voice TTS RMS 데드힌트 제거, init_state OpennessParams 파라미터객체화
- 한국어 PII(날짜·금액·주소) 마스킹 보강
- 레이아웃 시각 게이트: 폼 컨트롤 값 스크롤 오탐 제외(7/7)

검증: 백엔드 84/84, E2E 42(데스크톱 27·모바일 11·아바타 4), 시각 게이트 7/7
This commit is contained in:
Yun Chan 2026-06-27 02:30:46 +09:00
parent cb2aebd76c
commit 085460b5e0
327 changed files with 31226 additions and 1829 deletions

View file

@ -25,6 +25,13 @@ from pydantic import BaseModel
from ..auth_sessions import InactiveUserError, SessionUser, create_session, revoke_session
from ..config import settings
from ..deps import CurrentPrincipal, Principal, Role
from ..saml import (
SamlIdentity,
acs_url_for_entity_id,
build_authn_request,
parse_fixture_response,
redirect_binding_url,
)
router = APIRouter(prefix="/auth", tags=["auth"])
@ -41,7 +48,15 @@ class OAuthState:
created_at: float
@dataclass(slots=True)
class SamlState:
request_id: str
next_path: str
created_at: float
_oauth_states: dict[str, OAuthState] = {}
_saml_states: dict[str, SamlState] = {}
class MeResponse(BaseModel):
@ -52,8 +67,17 @@ class MeResponse(BaseModel):
cohort_ids: list[str]
class AuthProviderStatus(BaseModel):
provider: Literal["google", "saml"]
configured: bool
enabled: bool
login_path: str
class AuthConfigResponse(BaseModel):
google_oauth_configured: bool
saml_configured: bool
providers: list[AuthProviderStatus]
allowed_email_domains: list[str]
redirect_uri: str
dev_login_enabled: bool
@ -84,6 +108,37 @@ def _normalize_email_set(values: list[str]) -> set[str]:
return {email for value in values if (email := _normalize_email(value))}
def _google_configured() -> bool:
return bool(settings.oauth_google_client_id and settings.oauth_google_client_secret)
def _saml_configured() -> bool:
return bool(
settings.auth_saml_enabled
and settings.saml_sp_entity_id.strip()
and settings.saml_sso_url.strip()
)
def _auth_provider_statuses() -> list[AuthProviderStatus]:
google_ready = _google_configured()
saml_ready = _saml_configured()
return [
AuthProviderStatus(
provider="google",
configured=google_ready,
enabled=google_ready,
login_path="/auth/login?provider=google",
),
AuthProviderStatus(
provider="saml",
configured=saml_ready,
enabled=saml_ready,
login_path="/auth/login?provider=saml",
),
]
def allowed_email_domains() -> set[str]:
"""Configured login email domains, normalized for claim checks."""
return {
@ -137,6 +192,15 @@ def _role_for_email(email: str) -> Role:
return Role.LEARNER
def _role_for_saml_identity(identity: SamlIdentity) -> Role:
hinted = (identity.role_hint or "").strip().lower()
if hinted in {"admin", "administrator"}:
return Role.ADMIN
if hinted in {"teacher", "instructor", "faculty"}:
return Role.TEACHER
return _role_for_email(identity.email)
def _safe_next_path(next_path: str | None) -> str:
if not next_path or not next_path.startswith("/") or next_path.startswith("//"):
return "/"
@ -211,6 +275,13 @@ def _prune_oauth_states() -> None:
_oauth_states.pop(key, None)
def _prune_saml_states() -> None:
cutoff = time.time() - OAUTH_STATE_TTL_SECONDS
stale = [key for key, value in _saml_states.items() if value.created_at < cutoff]
for key in stale:
_saml_states.pop(key, None)
def _cookie_secure() -> bool:
# The __Host- prefix requires Secure, Path=/, and no Domain. Modern Chrome
# accepts Secure cookies on localhost, which keeps dev and prod semantics
@ -291,10 +362,12 @@ def _dev_login_available(request: Request) -> bool:
@router.get("/config", response_model=AuthConfigResponse)
async def auth_config(request: Request) -> AuthConfigResponse:
"""Return non-secret login configuration for the browser login screen."""
google_ready = _google_configured()
saml_ready = _saml_configured()
return AuthConfigResponse(
google_oauth_configured=bool(
settings.oauth_google_client_id and settings.oauth_google_client_secret
),
google_oauth_configured=google_ready,
saml_configured=saml_ready,
providers=_auth_provider_statuses(),
allowed_email_domains=sorted(allowed_email_domains()),
redirect_uri=settings.oauth_redirect_uri,
dev_login_enabled=_dev_login_available(request),
@ -308,9 +381,33 @@ async def login(
next: Annotated[str | None, Query()] = None,
) -> RedirectResponse:
"""Start Google OIDC authorization code + PKCE login."""
if provider == "saml":
if not _saml_configured():
return _frontend_login_redirect("saml_not_configured", request)
_prune_saml_states()
relay_state = secrets.token_urlsafe(32)
acs_url = acs_url_for_entity_id(settings.saml_sp_entity_id)
request_id, authn_request_xml = build_authn_request(
sp_entity_id=settings.saml_sp_entity_id,
sso_url=settings.saml_sso_url,
acs_url=acs_url,
)
_saml_states[relay_state] = SamlState(
request_id=request_id,
next_path=_safe_next_path(next),
created_at=time.time(),
)
return RedirectResponse(
redirect_binding_url(
sso_url=settings.saml_sso_url,
authn_request_xml=authn_request_xml,
relay_state=relay_state,
),
status_code=302,
)
if provider != "google":
return _frontend_login_redirect("unsupported_provider", request)
if not settings.oauth_google_client_id or not settings.oauth_google_client_secret:
if not _google_configured():
return _frontend_login_redirect("not_configured", request)
_prune_oauth_states()
@ -406,6 +503,58 @@ async def callback(
return response
@router.post("/saml/acs")
async def saml_acs(request: Request) -> RedirectResponse:
"""Accept a minimal unsigned SAMLResponse for local fixture SAML proof.
Signed SAML verification is intentionally not implemented. When
SAML_X509_CERT_FINGERPRINT is configured, this endpoint refuses to trust the
response so production does not silently run unsigned SAML.
"""
if not _saml_configured():
return _frontend_login_redirect("saml_not_configured", request)
if settings.saml_x509_cert_fingerprint.strip():
return _frontend_login_redirect("saml_signature_verification_required", request)
if settings.environment != "dev":
return _frontend_login_redirect("saml_fixture_acs_dev_only", request)
form = await request.form()
relay_state = str(form.get("RelayState") or "")
encoded_response = str(form.get("SAMLResponse") or "")
if not relay_state or not encoded_response:
return _frontend_login_redirect("saml_missing_callback", request)
_prune_saml_states()
stored = _saml_states.pop(relay_state, None)
if stored is None:
return _frontend_login_redirect("saml_invalid_state", request)
try:
identity = parse_fixture_response(encoded_response)
email = validate_google_identity_domain(
email=identity.email,
email_verified=True,
hosted_domain=_email_domain(identity.email),
)
except (HTTPException, ValueError):
return _frontend_login_redirect("saml_assertion_invalid", request)
role = _role_for_saml_identity(identity)
try:
sid, _ = await create_session(
email=email,
display_name=identity.display_name or email,
role=role.value,
cohort_ids=[],
)
except InactiveUserError:
return _frontend_login_redirect("inactive_user", request)
response = RedirectResponse(_frontend_url(stored.next_path, request), status_code=302)
_set_session_cookie(response, sid)
return response
@router.post("/dev-login", response_model=MeResponse)
async def dev_login(request: Request, body: DevLoginRequest, response: Response) -> MeResponse:
"""Dev-only server login for local E2E and manual testing.

View file

@ -2,15 +2,23 @@
from __future__ import annotations
from typing import Any
from typing import Annotated, Any, Literal
from fastapi import APIRouter, HTTPException, Response, status
from fastapi import APIRouter, Depends, HTTPException, Response, status
from pydantic import BaseModel
from ..deps import CurrentPrincipal
from ..persona_repository import CatalogPersona, list_catalog_personas
from ..deps import CurrentPrincipal, Principal, Role, require_role
from ..persona_repository import (
CatalogPersona,
PersonaReviewAction,
PersonaReviewItem,
list_catalog_personas,
list_persona_review_queue,
update_persona_review_status,
)
router = APIRouter(prefix="/personas", tags=["personas"])
TeacherOrAdmin = Annotated[Principal, Depends(require_role(Role.TEACHER, Role.ADMIN))]
class PersonaSummary(BaseModel):
@ -25,6 +33,24 @@ class PersonaSummary(BaseModel):
degraded: bool = False
class PersonaReviewSummary(BaseModel):
persona_id: str
code: str
version: int
status: Literal["draft", "review", "approved", "archived"]
display_name: str
difficulty: str
theory_target: list[str]
source_provenance: str
is_synthetic: bool
created_at: str | None = None
approved_at: str | None = None
class PersonaReviewDecisionRequest(BaseModel):
action: PersonaReviewAction
def _first_text_value(data: dict[str, Any]) -> str:
for value in data.values():
if isinstance(value, str) and value.strip():
@ -46,6 +72,27 @@ def _summary(entry: CatalogPersona) -> PersonaSummary:
)
def _review_summary(entry: PersonaReviewItem) -> PersonaReviewSummary:
return PersonaReviewSummary(
persona_id=entry.persona_id,
code=entry.code,
version=entry.version,
status=entry.status,
display_name=entry.display_name,
difficulty=entry.difficulty,
theory_target=entry.theory_target,
source_provenance=entry.source_provenance,
is_synthetic=entry.is_synthetic,
created_at=entry.created_at,
approved_at=entry.approved_at,
)
def _ensure_teacher_or_admin(principal: Principal) -> None:
if principal.role not in {Role.TEACHER, Role.ADMIN}:
raise HTTPException(status.HTTP_403_FORBIDDEN, detail="only teachers and admins can review personas")
@router.get("", response_model=list[PersonaSummary])
async def list_personas(response: Response, _principal: CurrentPrincipal) -> list[PersonaSummary]:
"""Return latest approved personas from app.persona_card."""
@ -64,3 +111,47 @@ async def list_personas(response: Response, _principal: CurrentPrincipal) -> lis
response.headers["X-Vignette-Catalog-Source"] = "database"
return [_summary(entry) for entry in personas]
@router.get("/review", response_model=list[PersonaReviewSummary])
async def list_persona_reviews(principal: TeacherOrAdmin) -> list[PersonaReviewSummary]:
"""Return draft/review personas awaiting faculty approval."""
_ensure_teacher_or_admin(principal)
try:
queue = await list_persona_review_queue(role=principal.role.value)
except Exception as exc:
raise HTTPException(
status.HTTP_503_SERVICE_UNAVAILABLE,
detail="persona review queue database unavailable",
) from exc
return [_review_summary(entry) for entry in queue]
@router.post("/review/{persona_id}", response_model=PersonaReviewSummary)
async def decide_persona_review(
persona_id: str,
request: PersonaReviewDecisionRequest,
principal: TeacherOrAdmin,
) -> PersonaReviewSummary:
"""Approve a persona for learners or return it to draft for changes."""
_ensure_teacher_or_admin(principal)
try:
updated = await update_persona_review_status(
persona_id=persona_id,
action=request.action,
reviewer_id=principal.user_id,
role=principal.role.value,
)
except ValueError as exc:
raise HTTPException(status.HTTP_403_FORBIDDEN, detail=str(exc)) from exc
except Exception as exc:
raise HTTPException(
status.HTTP_503_SERVICE_UNAVAILABLE,
detail="persona review update database unavailable",
) from exc
if updated is None:
raise HTTPException(
status.HTTP_404_NOT_FOUND,
detail="persona review item not found or not pending review",
)
return _review_summary(updated)

View file

@ -18,13 +18,13 @@ from fastapi import APIRouter, HTTPException, status
from pydantic import BaseModel, Field
from sse_starlette.sse import EventSourceResponse
from .. import session_persistence
from .. import db, session_persistence
from ..config import settings
from ..deps import CurrentPrincipal, Principal, Role
from ..engine_client import EngineError, engine_client
from ..persona_repository import get_catalog_persona
from ..runtime_policy import require_runtime_fallback_allowed, runtime_fallback_allowed
from ..services import evaluator, memory, orchestrator, state_machine
from ..services import evaluator, memory, orchestrator, rag, state_machine
from ..store import InProcSession, TurnRecord, store
router = APIRouter(prefix="/sessions", tags=["sessions"])
@ -193,6 +193,165 @@ class SessionReviewResponse(BaseModel):
_RECALL_CACHE: dict[str, memory.RecallContext] = {}
# 세션별 KB 증상 행동단서(회기 1회 산출·캐시). 빈 list 캐시 = 회기 내 재시도 안 함(안정성).
_KB_CUES_CACHE: dict[str, list[str]] = {}
_LEARNER_VISIBLE_AI_ROLE = "counselor"
# ────────────────────────────────────────────────────────────────────────────
# RAG 배선 헬퍼 — 내담자(CLIENT) 뷰. 임베더/KB/DB 풀 미가용 시 빈 값으로 graceful
# degradation: 상담 루프를 절대 막지 않는다(라이브 루프 비차단이 계약). routes/kb.py가
# 같은 예외를 503으로 올리는 것과 의도적으로 다르다. 임베딩은 rag가 스레드풀로 offload.
# ────────────────────────────────────────────────────────────────────────────
_RAG_RECALL_K = 5
_KB_CUES_K = 4
def _persona_kb_query(card) -> str:
"""페르소나 증상·호소 → KB 행동단서 검색 질의(임베더/tsquery 입력 전용, LLM 미주입).
질의는 프롬프트에 들어가지 않는다. 회수된 behavior_cue만 L2로 주입되고, CLIENT 정책
(expose_body=False) 본문을 잘라 '행동단서' 돌려준다(CCD 본문 비노출 자동 보존).
"""
parts: list[str] = []
presenting = getattr(card, "presenting", None) or {}
if presenting.get("주호소"):
parts.append(str(presenting["주호소"]))
if presenting.get("표층"):
parts.append(str(presenting["표층"]))
dsm = getattr(card, "dsm5_dimensional", None) or {}
parts.extend(str(key) for key in dsm.keys() if key != "note")
return " ".join(p for p in parts if p).strip()
async def _retrieve_kb_behavior_cues(card) -> list[str]:
"""KB 증상 행동단서 회수(CLIENT 정책). 미가용 시 빈 리스트(비차단)."""
query = _persona_kb_query(card)
if not query:
return []
try:
async with db.acquire(ai_view=rag.AIRole.CLIENT.value) as conn:
result = await rag.search_kb(
conn,
query=query,
role=rag.AIRole.CLIENT,
k=_KB_CUES_K,
)
return [c.behavior_cue for c in result.chunks if c.behavior_cue]
except Exception:
# rag.NotConfigured(임베더/KB 미가용)·RuntimeError(풀 미초기화)·DB 오류 포함.
# 비치명적: 빈 단서로 진행. CancelledError는 BaseException이라 미포착.
return []
async def _ensure_kb_cues(session_id: str, card) -> list[str]:
"""세션별 KB 행동단서(회기 1회 산출·캐시, 서버 재시작/재개 시 lazy 재계산)."""
cached = _KB_CUES_CACHE.get(session_id)
if cached is not None:
return cached
cues = await _retrieve_kb_behavior_cues(card)
_KB_CUES_CACHE[session_id] = cues
return cues
async def _load_prev_case_summary(case_id: str) -> Optional[dict]:
"""직전 회기 요약(case 스코프) → build_recall_context 입력. 미존재/미가용 시 None."""
try:
async with db.acquire(ai_view=rag.AIRole.CLIENT.value) as conn:
row = await conn.fetchrow(
"""
SELECT digest, open_threads, end_state
FROM app.session_summary
WHERE case_id = $1::uuid
ORDER BY session_no DESC, created_at DESC
LIMIT 1
""",
case_id,
)
except Exception:
return None
if row is None:
return None
return {
"digest": row["digest"],
"open_threads": list(row["open_threads"] or []),
"end_state": dict(row["end_state"] or {}),
}
async def _hydrate_episodic_text(conn, result) -> list[str]:
"""retrieve_persona_memory가 돌려준 turn_id → app.turns 마스킹 본문 조인(내담자 발화)."""
turn_ids = [c.meta.get("turn_id") for c in result.chunks if c.meta.get("turn_id")]
if not turn_ids:
return []
rows = await conn.fetch(
"""
SELECT id, text_masked FROM app.turns
WHERE id = ANY($1::uuid[]) AND speaker = 'client'
""",
turn_ids,
)
by_id = {str(r["id"]): r["text_masked"] for r in rows}
return [by_id[t] for t in turn_ids if by_id.get(t)]
def _recall_query(prev_summary: Optional[dict], card) -> str:
"""episodic recall 질의: 직전 open_threads 우선, 없으면 주호소."""
if prev_summary:
threads = prev_summary.get("open_threads") or []
if threads:
return " ".join(str(t) for t in threads)
presenting = getattr(card, "presenting", None) or {}
return str(presenting.get("주호소") or "").strip()
async def _episodic_recall_snippets(case_id: str, query: str) -> list[str]:
"""case 스코프 episodic 벡터 recall → 내담자 발화 단편(마스킹본). 미가용 시 []."""
if not query:
return []
try:
async with db.acquire(ai_view=rag.AIRole.CLIENT.value) as conn:
result = await rag.retrieve_persona_memory(
conn, case_id=case_id, query=query, k=_RAG_RECALL_K,
)
return await _hydrate_episodic_text(conn, result)
except Exception:
return []
async def _build_start_recall(*, case_id: str, card) -> memory.RecallContext:
"""회기 시작 회상 조립: prev_summary(case) + episodic recall을 build_recall_context로
합본. 구간 graceful(미가용 회상).
"""
try:
db.get_pool() # 풀 미초기화 시 RuntimeError → 첫 회기와 동일한 빈 회상
except RuntimeError:
return memory.build_recall_context()
prev_summary = await _load_prev_case_summary(case_id)
query = _recall_query(prev_summary, card)
episodic = await _episodic_recall_snippets(case_id, query)
pinned = list((prev_summary or {}).get("pinned_facts") or [])
return memory.build_recall_context(
prev_summary=prev_summary,
episodic_snippets=episodic,
pinned_facts=pinned,
)
async def _warm_rag_caches(session_id: str, case_id: str, card) -> None:
"""RAG 회상·KB 행동단서를 **백그라운드**로 산출해 캐시한다(요청 경로 비차단).
BGE-M3 임베더 로드(~ ) 회기 시작/ 응답을 막지 않도록 create_task로 띄운다.
warm 완료 턴은 회상/단서로 진행(graceful), 이후 턴부터 RAG 주입. 구간 비치명적.
"""
try:
_RECALL_CACHE[session_id] = await _build_start_recall(case_id=case_id, card=card)
except Exception:
pass
try:
_KB_CUES_CACHE[session_id] = await _retrieve_kb_behavior_cues(card)
except Exception:
pass
_PHASE_KEY_BY_LABEL = {
"라포": "rapport",
@ -481,6 +640,92 @@ def _evaluation_payload(record: dict[str, object] | None) -> dict[str, object]:
return payload if isinstance(payload, dict) else {}
# fast-loop 턴 평가(TechniqueCategory) → 프론트 sr-technique--{kind} 시각 매핑.
_TECHNIQUE_KIND_BY_CATEGORY = {
"relational": "empathy",
"exploratory": "explore",
"intervention": "confront",
"stabilizing": "reflect",
"structuring": "closed",
}
def _review_techniques_from_turn_eval(ev: dict[str, object] | None) -> list[ReviewTechnique]:
"""턴 평가의 기법 태그를 리뷰 칩으로. label_ko 우선, category로 색 kind 결정."""
if not isinstance(ev, dict):
return []
out: list[ReviewTechnique] = []
for tag in ev.get("techniques") or []:
if not isinstance(tag, dict):
continue
label = str(tag.get("label_ko") or tag.get("code") or "").strip()
if not label:
continue
kind = _TECHNIQUE_KIND_BY_CATEGORY.get(str(tag.get("category") or ""), "explore")
out.append(ReviewTechnique(kind=kind, label=label))
return out
def _review_note_from_turn_eval(ev: dict[str, object] | None) -> Optional[ReviewNote]:
"""의도이탈(있으면 우선) 또는 적절성 신호를 턴 노트로. tone: good|warn(프론트 계약)."""
if not isinstance(ev, dict):
return None
dev = ev.get("intent_deviation")
if isinstance(dev, dict):
dimension = str(dev.get("dimension") or "").strip()
expected = str(dev.get("expected") or "").strip()
actual = str(dev.get("actual") or "").strip()
body = " / ".join(p for p in (f"권장: {expected}" if expected else "", f"실제: {actual}" if actual else "") if p)
return ReviewNote(
author="평가 AI",
tone="warn",
title=f"의도와 다른 부분 · {dimension}".rstrip(" ·") or "의도와 다른 부분",
body=body or "권장 반응과 실제 반응에 차이가 있었어요.",
)
appropriateness = str(ev.get("appropriateness") or "neutral")
note_text = str(ev.get("appropriateness_note") or "").strip()
if appropriateness == "pos":
return ReviewNote(author="평가 AI", tone="good", title="적절한 개입", body=note_text or "이 개입은 흐름에 적절했어요.")
if appropriateness == "warn" and note_text:
return ReviewNote(author="평가 AI", tone="warn", title="점검해볼 지점", body=note_text)
return None
async def _record_safety_event(sess: InProcSession, ctx, result) -> None:
"""위기 escalate 시 app.safety_events 적재(교수자 감사·알림 레코드). C2.
비차단: DB 미가용(degraded)·FK 미충족(in-memory 세션) graceful skip 상담 루프를
절대 막지 않는다. 실시간 교수자 push 알림은 후속( 레코드가 1 알림원).
"""
crisis = getattr(ctx, "crisis", None)
if crisis is None or not getattr(crisis, "escalate", False):
return
kind = getattr(crisis.kind, "value", None) or str(getattr(crisis, "kind", "crisis"))
try:
async with db.acquire() as conn:
await conn.execute(
"""
INSERT INTO app.safety_events
(session_id, trigger_type, ko_risk_level, escalated, detail)
VALUES ($1::uuid, $2, $3, TRUE, $4::jsonb)
""",
sess.session_id,
kind,
int(getattr(crisis, "risk_level", 0) or 0),
json.dumps({
"matched": list(getattr(crisis, "matched", []) or []),
"stage": getattr(result, "stage", None),
"turn_seq": getattr(result, "turn_seq", None),
}),
)
except Exception:
pass # 비차단(R5): 적재 실패가 위기 대응/상담을 막지 않음.
def _learner_visible_turns(sess: InProcSession) -> list[TurnRecord]:
return sess.turns_visible_to(_LEARNER_VISIBLE_AI_ROLE)
async def _generate_and_save_session_evaluation(sess: InProcSession) -> None:
if not sess.turns:
return
@ -535,8 +780,9 @@ def _schedule_session_evaluation(sess: InProcSession) -> None:
def _learner_summary(sess: InProcSession, *, review_ready: bool = False) -> LearnerSessionSummary:
learner_turns = sum(1 for turn in sess.turns if turn.speaker == "counselor")
client_turns = sum(1 for turn in sess.turns if turn.speaker == "client")
turns = _learner_visible_turns(sess)
learner_turns = sum(1 for turn in turns if turn.speaker == "counselor")
client_turns = sum(1 for turn in turns if turn.speaker == "client")
return LearnerSessionSummary(
session_id=sess.session_id,
persona_code=sess.persona_code,
@ -544,7 +790,7 @@ def _learner_summary(sess: InProcSession, *, review_ready: bool = False) -> Lear
session_no=sess.session_no,
status="ended" if sess.ended else "active",
stage=_stage_label(sess.state.stage),
turn_count=len(sess.turns),
turn_count=len(turns),
learner_turn_count=learner_turns,
client_turn_count=client_turns,
started_at=_iso(sess.created_at) or "",
@ -554,7 +800,10 @@ def _learner_summary(sess: InProcSession, *, review_ready: bool = False) -> Lear
async def _review_ready(sess: InProcSession, principal: Principal) -> bool:
if not sess.ended or not sess.turns:
turns = _learner_visible_turns(sess)
if not sess.ended or not turns:
return False
if len(turns) != len(sess.turns):
return False
evaluation_record, _ = await session_persistence.load_session_evaluation(
sess.session_id,
@ -568,6 +817,7 @@ def _session_detail(
*,
review_ready: bool = False,
) -> SessionDetailResponse:
turns = _learner_visible_turns(sess)
return SessionDetailResponse(
session_id=sess.session_id,
case_id=sess.case_id,
@ -587,7 +837,7 @@ def _session_detail(
text=turn.text_masked,
created_at=_iso(turn.created_at) or "",
)
for turn in sess.turns
for turn in turns
],
review_ready=review_ready,
)
@ -649,10 +899,7 @@ async def start_session(
recall = memory.build_recall_context()
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(),
params=card.openness_params(),
carry=recall.carry,
)
@ -680,7 +927,11 @@ async def start_session(
)
else:
store.put(sess)
# 즉시 빈/carry 회상으로 응답을 막지 않는다. RAG 회상·KB 단서(임베더 로드 수 초)는
# 백그라운드 warm으로 캐시 — 회기 시작/턴 응답이 임베더 로드에 블로킹되지 않게(성능 회귀 방지).
_RECALL_CACHE[sess.session_id] = recall
asyncio.create_task(_warm_rag_caches(sess.session_id, sess.case_id, card))
return SessionStartResponse(
session_id=sess.session_id,
@ -701,6 +952,8 @@ async def get_session_review(
"""Return a learner-safe review built only from the stored session transcript."""
_ensure_learner(principal)
sess = await _load_session_or_404(session_id, principal, allow_ended=True)
visible_turns = _learner_visible_turns(sess)
hidden_turns = len(visible_turns) != len(sess.turns)
end_ts = sess.ended_at or datetime.now().timestamp()
duration_seconds = max(0, int(round(end_ts - sess.created_at)))
@ -708,7 +961,7 @@ async def get_session_review(
client_initial = client_name[:1] or ""
reached_phase = _stage_label(sess.state.stage)
stage_labels = [turn.stage for turn in sess.turns] or [reached_phase]
stage_labels = [turn.stage for turn in visible_turns] or [reached_phase]
axis = ["0:00"]
if duration_seconds > 0:
axis.append(_offset_label(duration_seconds))
@ -717,16 +970,20 @@ async def get_session_review(
session_id,
principal,
)
evaluation_payload = _evaluation_payload(evaluation_record)
evaluation_status = str(evaluation_record.get("status") or "") if evaluation_record else ""
evaluation_ready = evaluation_status == "ready"
evaluation_payload = {} if hidden_turns else _evaluation_payload(evaluation_record)
evaluation_status = (
"" if hidden_turns else str(evaluation_record.get("status") or "") if evaluation_record else ""
)
evaluation_ready = not hidden_turns and evaluation_status == "ready"
first_turn_ts = sess.turns[0].created_at if sess.turns else sess.created_at
first_turn_ts = visible_turns[0].created_at if visible_turns else sess.created_at
turns: list[ReviewTurn] = []
for index, turn in enumerate(sess.turns):
for index, turn in enumerate(visible_turns):
speaker: Literal["learner", "client"] = (
"learner" if turn.speaker == "counselor" else "client"
)
# 턴별 fast-loop 평가는 학습자 발화에만 부착(기법 태깅·노트). hidden 시 노출 안 함.
turn_eval = turn.evaluation if (speaker == "learner" and not hidden_turns) else None
turns.append(
ReviewTurn(
id=f"t{index + 1}",
@ -734,8 +991,8 @@ async def get_session_review(
speaker=speaker,
who="학습자" if speaker == "learner" else client_name,
text=turn.text_masked,
techniques=[],
note=None,
techniques=_review_techniques_from_turn_eval(turn_eval),
note=_review_note_from_turn_eval(turn_eval),
)
)
@ -783,10 +1040,10 @@ async def get_session_review(
summary = _review_summary_from_evaluation(
fallback=transcript_summary,
evaluation_record=evaluation_record,
evaluation_record=None if hidden_turns else evaluation_record,
payload=evaluation_payload,
)
if evaluation_record and not evaluation_durable:
if evaluation_record and not hidden_turns and not evaluation_durable:
summary += " 현재 평가는 런타임 캐시에서 복원되었습니다."
return SessionReviewResponse(
@ -832,6 +1089,7 @@ async def submit_turn(
_ensure_learner(principal)
sess = await _load_session_or_404(session_id, principal)
recall = _RECALL_CACHE.get(session_id) or memory.RecallContext()
kb_cues = _KB_CUES_CACHE.get(session_id) or [] # 비차단: warm 전이면 빈 단서(graceful)
ctx = orchestrator.prepare_turn(
session_id=session_id,
@ -841,18 +1099,25 @@ async def submit_turn(
learner_text=body.text,
recall_summary=recall.recall_summary,
pinned_facts=recall.pinned_facts,
recent_turns=sess.recent_turns(),
recent_turns=sess.recent_turns(visible_to="client"),
kb_behavior_cues=kb_cues,
theory_mode=sess.theory_mode,
)
assert ctx.state_after is not None
try:
result = await orchestrator.run_turn_generate(ctx, engine_client)
result = await orchestrator.run_turn_generate(
ctx,
engine_client,
eval_hook=evaluator.make_eval_hook(engine_client),
)
except EngineError as exc:
raise HTTPException(
status.HTTP_503_SERVICE_UNAVAILABLE,
detail=f"engine unavailable: {exc}",
) from exc
# 턴별 fast-loop 평가는 학습자(상담자) 발화에 부착(기법 태깅·적절성·의도이탈).
await _append_session_turn(
sess,
TurnRecord(
@ -861,6 +1126,7 @@ async def submit_turn(
stage=_stage_label(ctx.state_after.stage),
text=body.text,
text_masked=ctx.learner_text_masked,
evaluation=result.evaluation,
),
)
@ -873,9 +1139,15 @@ async def submit_turn(
stage=_stage_label(result.state_after.stage),
text=result.client_reply,
text_masked=result.client_reply,
llm_provider=result.llm_provider,
model=result.model,
tokens_in=result.tokens_in,
tokens_out=result.tokens_out,
cost_usd=result.cost_usd,
),
)
await _update_session_state(sess, result.state_after)
await _record_safety_event(sess, ctx, result) # C2: 위기 escalate 시 safety_events 적재(비차단)
return TurnResponse(
turn_seq=result.turn_seq,
@ -897,6 +1169,7 @@ async def stream_turn(
_ensure_learner(principal)
sess = await _load_session_or_404(session_id, principal)
recall = _RECALL_CACHE.get(session_id) or memory.RecallContext()
kb_cues = _KB_CUES_CACHE.get(session_id) or [] # 비차단: warm 전이면 빈 단서(graceful)
ctx = orchestrator.prepare_turn(
session_id=session_id,
@ -906,7 +1179,9 @@ async def stream_turn(
learner_text=body.text,
recall_summary=recall.recall_summary,
pinned_facts=recall.pinned_facts,
recent_turns=sess.recent_turns(),
recent_turns=sess.recent_turns(visible_to="client"),
kb_behavior_cues=kb_cues,
theory_mode=sess.theory_mode,
)
assert ctx.state_after is not None
@ -941,6 +1216,11 @@ async def stream_turn(
stage=_stage_label(ctx.state_after.stage),
text=final_reply,
text_masked=final_reply,
llm_provider=str(ev.data.get("llm_provider") or ""),
model=str(ev.data.get("model") or ""),
tokens_in=int(ev.data.get("tokens_in") or 0),
tokens_out=int(ev.data.get("tokens_out") or 0),
cost_usd=float(ev.data.get("cost_usd") or 0.0),
),
)
yield {"event": "done", "data": json.dumps(data, ensure_ascii=False)}
@ -980,6 +1260,7 @@ async def end_session(
await _end_persisted_session(sess, carry)
_RECALL_CACHE.pop(session_id, None)
_KB_CUES_CACHE.pop(session_id, None)
_schedule_session_evaluation(sess)
return SessionEndResponse(

View file

@ -14,6 +14,8 @@ cleanly instead of crashing.
from __future__ import annotations
import json
import hashlib
import time
from typing import Optional
from fastapi import APIRouter, WebSocket, WebSocketDisconnect
@ -27,7 +29,7 @@ from ..deps import Principal, Role
from ..engine_client import EngineError, engine_client
from ..persona_repository import get_catalog_persona
from ..runtime_policy import require_runtime_fallback_allowed, runtime_fallback_allowed
from ..services import memory, orchestrator, state_machine
from ..services import evaluator, memory, orchestrator, state_machine
from ..services import voice as voice_svc
from ..services.voice import VoicePreset, VoiceUnavailable, resolve_voice, voice_service
from ..store import InProcSession, TurnRecord, store
@ -113,6 +115,8 @@ async def voice_ws(websocket: WebSocket) -> None:
audio_buf = bytearray()
receiving = False
audio_started_at: float | None = None
last_audio_end_at: float | None = None
try:
while True:
@ -126,6 +130,7 @@ async def voice_ws(websocket: WebSocket) -> None:
if not receiving:
# Be tolerant when audio arrives before audio_start.
receiving = True
audio_started_at = time.monotonic()
audio_buf.clear()
await _safe_send_json(websocket, {"type": "state", "state": "listening"})
audio_buf.extend(msg["bytes"])
@ -151,11 +156,16 @@ async def voice_ws(websocket: WebSocket) -> None:
ctype = ctrl.get("type")
if ctype == "audio_start":
receiving = True
audio_started_at = time.monotonic()
audio_buf.clear()
await _safe_send_json(websocket, {"type": "state", "state": "listening"})
elif ctype == "audio_end":
receiving = False
audio_ended_at = time.monotonic()
silence_ms = _safe_int(ctrl.get("silence_ms"))
if silence_ms is None and last_audio_end_at is not None and audio_started_at is not None:
silence_ms = max(0, int((audio_started_at - last_audio_end_at) * 1000))
await _handle_utterance(
websocket,
session_id=session_id,
@ -163,7 +173,13 @@ async def voice_ws(websocket: WebSocket) -> None:
voice_preset=voice_preset,
audio=bytes(audio_buf),
fmt=ctrl.get("format"),
audio_started_at=audio_started_at,
audio_ended_at=audio_ended_at,
silence_ms=silence_ms,
barge_in=_safe_bool(ctrl.get("barge_in")),
)
last_audio_end_at = audio_ended_at
audio_started_at = None
audio_buf.clear()
elif ctype == "text_turn":
@ -202,6 +218,10 @@ async def _handle_utterance(
voice_preset: VoicePreset,
audio: bytes,
fmt: Optional[str],
audio_started_at: float | None = None,
audio_ended_at: float | None = None,
silence_ms: int | None = None,
barge_in: bool | None = None,
) -> None:
"""Transcribe one utterance, generate the client reply, then synthesize TTS."""
if not audio:
@ -226,6 +246,9 @@ async def _handle_utterance(
return
learner_text = stt.text
audio_ref = _voice_audio_ref(audio, fmt)
duration_s = stt.duration or _elapsed_seconds(audio_started_at, audio_ended_at)
speech_rate = _estimate_speech_rate(learner_text, duration_s)
await _safe_send_json(
websocket,
{"type": "transcript", "text": learner_text, "final": True, "speaker": "counselor"},
@ -240,6 +263,10 @@ async def _handle_utterance(
principal=principal,
voice_preset=voice_preset,
learner_text=learner_text,
audio_ref=audio_ref,
silence_ms=silence_ms,
speech_rate=speech_rate,
barge_in=barge_in,
)
@ -250,6 +277,10 @@ async def _run_turn_and_speak(
principal: Principal,
voice_preset: VoicePreset,
learner_text: str,
audio_ref: str | None = None,
silence_ms: int | None = None,
speech_rate: float | None = None,
barge_in: bool | None = None,
) -> None:
"""Run one counseling turn and stream synthesized client speech."""
sess, err = await _load_voice_session(session_id, principal)
@ -267,13 +298,18 @@ async def _run_turn_and_speak(
learner_text=learner_text,
recall_summary=recall.recall_summary,
pinned_facts=recall.pinned_facts,
recent_turns=sess.recent_turns(),
recent_turns=sess.recent_turns(visible_to="client"),
theory_mode=sess.theory_mode,
)
assert ctx.state_after is not None
# Voice needs the full client reply before TTS starts.
try:
result = await orchestrator.run_turn_generate(ctx, engine_client)
result = await orchestrator.run_turn_generate(
ctx,
engine_client,
eval_hook=evaluator.make_eval_hook(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"})
@ -290,6 +326,11 @@ async def _run_turn_and_speak(
stage=ctx.state_after.stage.value,
text=learner_text,
text_masked=ctx.learner_text_masked,
audio_ref=audio_ref,
silence_ms=silence_ms,
speech_rate=speech_rate,
barge_in=barge_in,
evaluation=result.evaluation,
),
)
if reply:
@ -302,6 +343,11 @@ async def _run_turn_and_speak(
stage=result.stage,
text=reply,
text_masked=reply,
llm_provider=result.llm_provider,
model=result.model,
tokens_in=result.tokens_in,
tokens_out=result.tokens_out,
cost_usd=result.cost_usd,
),
)
await _update_voice_state(sess, result.state_after)
@ -333,10 +379,7 @@ async def _run_turn_and_speak(
try:
n = 0
async for ck in voice_service.synthesize_stream(reply, voice_preset):
# Metadata precedes the binary chunk so the client can pair them.
await _safe_send_json(
websocket, {"type": "tts_chunk", "seq": ck.seq, "rms": round(ck.rms, 4)}
)
# 바이너리 오디오 청크만 송신(프론트가 Web Audio AnalyserNode로 립싱크 자체 산출).
await _safe_send_bytes(websocket, ck.audio)
n += 1
await _safe_send_json(websocket, {"type": "tts_end", "chunks": n})
@ -451,10 +494,7 @@ async def _bind_session(
card = catalog_persona.card
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(),
params=card.openness_params(),
)
sess = await session_persistence.create_session(
learner_id=principal.user_id,
@ -511,6 +551,52 @@ def _audio_meta(fmt: Optional[str]) -> tuple[str, str]:
return table.get(f, ("audio.webm", "audio/webm"))
def _voice_audio_ref(audio: bytes, fmt: Optional[str]) -> str | None:
if not audio:
return None
f = (fmt or "webm").lower().lstrip(".") or "webm"
digest = hashlib.sha256(audio).hexdigest()[:24]
return f"voice:{f}:sha256:{digest}"
def _elapsed_seconds(started_at: float | None, ended_at: float | None) -> float | None:
if started_at is None or ended_at is None:
return None
return max(0.001, ended_at - started_at)
def _estimate_speech_rate(text: str, duration_s: float | None) -> float | None:
if not text or not duration_s or duration_s <= 0:
return None
units = sum(1 for ch in text if not ch.isspace())
if units <= 0:
return None
return round((units / duration_s) * 60.0, 2)
def _safe_int(value: object) -> int | None:
if value is None:
return None
try:
return int(value)
except (TypeError, ValueError):
return None
def _safe_bool(value: object) -> bool | None:
if value is None:
return None
if isinstance(value, bool):
return value
if isinstance(value, str):
normalized = value.strip().lower()
if normalized in {"1", "true", "yes", "y"}:
return True
if normalized in {"0", "false", "no", "n"}:
return False
return bool(value)
async def _safe_send_json(websocket: WebSocket, payload: dict) -> None:
if websocket.client_state != WebSocketState.CONNECTED:
return