vignette/apps/api/app/routes/sessions.py
2026-06-27 16:08:41 +09:00

1516 lines
54 KiB
Python

"""Counseling session routes.
The DB-backed source of truth is still pending, so this route uses the existing
in-process session store when DB is degraded. Unlike the previous dev fallback,
all browser calls now require a verified server-side auth session and every
session operation checks learner ownership.
"""
from __future__ import annotations
import asyncio
import json
from collections import Counter
from datetime import datetime
from typing import Literal, Optional
from fastapi import APIRouter, HTTPException, status
from pydantic import BaseModel, Field
from sse_starlette.sse import EventSourceResponse
from .. import db, session_persistence, turn_runtime
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
from ..services import evaluator, memory, orchestrator, rag, state_machine
from ..store import InProcSession, TurnRecord, store
router = APIRouter(prefix="/sessions", tags=["sessions"])
TheoryMode = Literal["humanistic", "cbt", "integrative"]
EndStateValue = str | int | float | bool | None | dict[str, float]
class SessionStartRequest(BaseModel):
persona_code: str = Field(..., examples=["P1"])
theory_mode: TheoryMode = "humanistic"
class SessionStartResponse(BaseModel):
session_id: str
case_id: str
session_no: int
stage: str
effective_openness: float
recall_summary: Optional[str] = None
degraded: bool = False
class TurnRequest(BaseModel):
text: str = Field(..., min_length=1)
class CrisisResourceResponse(BaseModel):
title: str
number: str
message: str
class TurnResponse(BaseModel):
turn_seq: int
stage: str
effective_openness: float
client_reply: Optional[str] = None
safety_flagged: bool = False
crisis_kind: str = "none"
crisis_resource: Optional[CrisisResourceResponse] = None
conversation_stopped: bool = False
class SessionEndResponse(BaseModel):
session_id: str
session_no: int
digest_pending: bool
end_state: dict[str, EndStateValue]
class LearnerSessionSummary(BaseModel):
session_id: str
persona_code: str
persona_name: str
session_no: int
status: Literal["active", "ended"]
stage: str
turn_count: int
learner_turn_count: int
client_turn_count: int
started_at: str
ended_at: str | None = None
review_ready: bool = False
class LearnerSessionsResponse(BaseModel):
source: str = "runtime"
sessions: list[LearnerSessionSummary] = Field(default_factory=list)
class SessionDetailTurn(BaseModel):
turn_seq: int
speaker: Literal["learner", "client"]
stage: str
text: str
created_at: str
class SessionDetailResponse(BaseModel):
session_id: str
case_id: str
persona_code: str
persona_name: str
theory_mode: str
status: Literal["active", "ended"]
stage: str
effective_openness: float
started_at: str
ended_at: str | None = None
turns: list[SessionDetailTurn] = Field(default_factory=list)
review_ready: bool = False
class ReviewClient(BaseModel):
name: str
initial: str
persona: str
class ReviewTechnique(BaseModel):
kind: str
label: str
class ReviewNonverbalEvent(BaseModel):
kind: Literal["audio", "silence", "pace", "barge_in"]
label: str
detail: str
class ReviewNote(BaseModel):
author: str
tone: str
title: str
body: str
quote: Optional[str] = None
class ReviewTurn(BaseModel):
id: str
ts: str
speaker: Literal["learner", "client"]
who: str
text: str
techniques: list[ReviewTechnique] = Field(default_factory=list)
nonverbal: list[ReviewNonverbalEvent] = Field(default_factory=list)
note: Optional[ReviewNote] = None
class ReviewPhaseSegment(BaseModel):
key: str
label: str
weight: float
class ReviewValencePoint(BaseModel):
t: float
v: float
class ReviewRubricRow(BaseModel):
name: str
cluster: str
ratio: float
quality: Literal["good", "watch"]
freq: str
class ReviewPoint(BaseModel):
title: str
body: str
jumpTo: Optional[str] = None
class ReviewWorksheetEvidence(BaseModel):
turnId: str
speaker: Literal["learner", "client"]
quote: str
class ReviewWorksheetItem(BaseModel):
key: str
label: str
value: Optional[str] = None
evidence: list[ReviewWorksheetEvidence] = Field(default_factory=list)
confidence: Literal["none", "low", "medium"] = "none"
emptyReason: Optional[str] = None
class ReviewWorksheetSection(BaseModel):
key: str
title: str
items: list[ReviewWorksheetItem] = Field(default_factory=list)
class ReviewCaseWorksheet(BaseModel):
status: Literal["empty", "draft_from_transcript"] = "empty"
generatedBy: str = "rule-based transcript extractor"
sections: list[ReviewWorksheetSection] = Field(default_factory=list)
limitations: list[str] = Field(default_factory=list)
class SessionReviewResponse(BaseModel):
session_id: str
client: ReviewClient
date: str
durationLabel: str
durationSeconds: int
reachedPhase: str
sessionSignal: str
supervisorState: str
supervisorName: str
summary: str
phases: list[ReviewPhaseSegment] = Field(default_factory=list)
phaseAxis: list[str] = Field(default_factory=list)
valenceAxis: list[str] = Field(default_factory=list)
clientValence: list[ReviewValencePoint] = Field(default_factory=list)
counselorBaseline: list[ReviewValencePoint] = Field(default_factory=list)
turns: list[ReviewTurn] = Field(default_factory=list)
rubric: list[ReviewRubricRow] = Field(default_factory=list)
goodMoments: list[ReviewPoint] = Field(default_factory=list)
growthPoints: list[ReviewPoint] = Field(default_factory=list)
caseWorksheet: ReviewCaseWorksheet = Field(default_factory=ReviewCaseWorksheet)
nextLine: Optional[str] = None
clientFeedback: Optional[str] = None
audioUrl: Optional[str] = None
pdfExportUrl: Optional[str] = None
degraded: bool = True
reviewReady: bool = False
_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 _build_seed_recall(*, case_id: str | None) -> memory.RecallContext:
if not case_id:
return memory.build_recall_context()
try:
db.get_pool()
except RuntimeError:
return memory.build_recall_context()
prev_summary = await _load_prev_case_summary(case_id)
pinned = list((prev_summary or {}).get("pinned_facts") or [])
return memory.build_recall_context(prev_summary=prev_summary, pinned_facts=pinned)
async def ensure_recall_context(sess: InProcSession) -> memory.RecallContext:
cached = _RECALL_CACHE.get(sess.session_id)
if cached is not None:
return cached
recall = await _build_seed_recall(case_id=sess.case_id)
_RECALL_CACHE[sess.session_id] = recall
return recall
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",
"탐색": "explore",
"개입": "intervene",
"정리": "closing",
}
def _stage_label(stage: object) -> str:
return turn_runtime.stage_label(stage)
def _ensure_learner(principal: Principal) -> None:
if principal.role != Role.LEARNER:
raise HTTPException(status.HTTP_403_FORBIDDEN, detail="only learners can use sessions")
async def _load_session_or_404(
session_id: str,
principal: Principal,
*,
allow_ended: bool = False,
include_turn_evaluation: bool = False,
) -> InProcSession:
sess, err = await turn_runtime.load_owned_session(
session_id,
principal,
allow_ended=allow_ended,
include_turn_evaluation=include_turn_evaluation,
)
if err == turn_runtime.SessionAccessError.NOT_FOUND:
raise HTTPException(status.HTTP_404_NOT_FOUND, detail="session not found")
if err == turn_runtime.SessionAccessError.FORBIDDEN:
raise HTTPException(status.HTTP_403_FORBIDDEN, detail="session does not belong to user")
if err == turn_runtime.SessionAccessError.ENDED:
raise HTTPException(status.HTTP_409_CONFLICT, detail="session already ended")
assert sess is not None
return sess
async def _end_persisted_session(sess: InProcSession, carry: memory.CarryOver) -> None:
if await session_persistence.end_session(sess, carry):
sess.ended = True
sess.ended_at = datetime.now().timestamp()
store.put(sess)
return
require_runtime_fallback_allowed("session end")
store.end(sess.session_id)
def _offset_label(seconds: float) -> str:
whole = max(0, int(round(seconds)))
minutes, sec = divmod(whole, 60)
return f"{minutes}:{sec:02d}"
def _iso(ts: float | None) -> str | None:
if ts is None:
return None
return datetime.fromtimestamp(ts).isoformat(timespec="seconds")
def _duration_label(seconds: int) -> str:
if seconds < 60:
return f"{seconds}"
minutes, sec = divmod(seconds, 60)
return f"{minutes}{sec}"
def _client_name(raw: str) -> str:
name = raw.split("(", 1)[0].strip()
return name or raw.strip() or "내담자"
def _review_summary(*, client_name: str, reached_phase: str, turns: list[ReviewTurn]) -> str:
if not turns:
return (
"아직 실제 발화가 없어 리뷰를 만들 수 없습니다. 회기를 진행한 뒤 종료하면 "
"저장된 축어록을 기준으로 리뷰가 표시됩니다."
)
learner_count = sum(1 for turn in turns if turn.speaker == "learner")
client_count = sum(1 for turn in turns if turn.speaker == "client")
return (
f"이 리뷰는 현재 세션에 저장된 실제 축어록 {len(turns)}개를 기반으로 합니다. "
f"{client_name}와의 회기는 {reached_phase} 단계까지 진행되었고, "
f"학습자 발화 {learner_count}개와 내담자 응답 {client_count}개가 기록되었습니다. "
"평가 AI 또는 교수자 코멘트가 아직 생성되지 않은 항목은 빈 상태로 남겨 둡니다."
)
def _phase_segments(stage_labels: list[str]) -> list[ReviewPhaseSegment]:
counts = Counter(stage_labels)
return [
ReviewPhaseSegment(
key=_PHASE_KEY_BY_LABEL.get(label, label),
label=label,
weight=float(count),
)
for label, count in counts.items()
if count > 0
]
def _clamp_ratio(value: float) -> float:
return round(max(0.0, min(1.0, value)), 3)
def _compact_text(text: str) -> str:
return " ".join(text.split())
def _clip_text(text: str, limit: int = 180) -> str:
compact = _compact_text(text)
if len(compact) <= limit:
return compact
return f"{compact[: max(0, limit - 1)].rstrip()}..."
def _point_title(text: str, fallback: str) -> str:
compact = _clip_text(text, 72)
for sep in (".", "", "!", "?", "\n"):
if sep in compact:
first = compact.split(sep, 1)[0].strip()
if first:
return _clip_text(first, 44)
return _clip_text(compact, 44) or fallback
def _ai_review_points(values: object, *, fallback_prefix: str) -> list[ReviewPoint]:
if not isinstance(values, list):
return []
points: list[ReviewPoint] = []
for index, value in enumerate(values, start=1):
body = _compact_text(str(value or ""))
if not body:
continue
points.append(
ReviewPoint(
title=_point_title(body, f"{fallback_prefix} {index}"),
body=body,
jumpTo=None,
)
)
return points[:3]
def _intent_deviation_points(values: object) -> list[ReviewPoint]:
if not isinstance(values, list):
return []
points: list[ReviewPoint] = []
for index, value in enumerate(values, start=1):
if not isinstance(value, dict):
continue
dimension = _compact_text(str(value.get("dimension") or f"의도 이탈 {index}"))
expected = _compact_text(str(value.get("expected") or ""))
actual = _compact_text(str(value.get("actual") or ""))
severity = _compact_text(str(value.get("severity") or "minor"))
body_parts = []
if expected:
body_parts.append(f"기대: {expected}")
if actual:
body_parts.append(f"실제: {actual}")
if severity:
body_parts.append(f"심각도: {severity}")
if body_parts:
points.append(
ReviewPoint(
title=dimension,
body=" · ".join(body_parts),
jumpTo=None,
)
)
return points[:3]
def _rubric_from_evaluation(payload: dict[str, object]) -> list[ReviewRubricRow]:
distribution = payload.get("distribution")
if not isinstance(distribution, dict):
return []
by_category = distribution.get("by_category")
if not isinstance(by_category, dict):
return []
total = int(distribution.get("total") or 0)
if total <= 0:
return []
overused = {str(item) for item in distribution.get("overused") or []}
underused = {str(item) for item in distribution.get("underused") or []}
rows: list[ReviewRubricRow] = []
for category, raw_count in sorted(by_category.items(), key=lambda item: str(item[0])):
try:
count = int(raw_count)
except (TypeError, ValueError):
continue
code = str(category)
watch = code in overused or code in underused
rows.append(
ReviewRubricRow(
name=code.replace("_", " ").title(),
cluster="평가 AI 기법 분포",
ratio=_clamp_ratio(count / max(1, total)),
quality="watch" if watch else "good",
freq=f"{count}/{total} labels",
)
)
return rows
def _review_summary_from_evaluation(
*,
fallback: str,
evaluation_record: dict[str, object] | None,
payload: dict[str, object],
) -> str:
if not evaluation_record:
return fallback
status = str(evaluation_record.get("status") or "")
if status != "ready":
error = _compact_text(str(evaluation_record.get("error") or payload.get("error") or ""))
return (
"저장된 축어록은 확인했지만 평가 AI 산출물이 아직 준비되지 않았습니다. "
+ (f"사유: {error}" if error else "평가가 완료되면 코칭 항목이 갱신됩니다.")
)
rationale = _compact_text(str(payload.get("supervisor_rationale") or ""))
critique = _compact_text(str(payload.get("supervisor_critique") or ""))
evaluated = payload.get("turns_evaluated")
prefix = f"평가 AI가 학습자 발화 {evaluated}개를 deep-loop로 분석했습니다. "
details = " ".join(part for part in [rationale, critique] if part)
return prefix + (details if details else "아래 코칭 항목은 저장된 축어록과 평가 AI 결과를 기준으로 합니다.")
def _next_line_from_evaluation(payload: dict[str, object]) -> str | None:
alternatives = payload.get("alternative_utterances")
if not isinstance(alternatives, list):
return None
for value in alternatives:
line = _compact_text(str(value or ""))
if line:
return line
return None
def _latest_client_feedback(turns: list[ReviewTurn]) -> str | None:
for turn in reversed(turns):
if turn.speaker == "client":
return _clip_text(turn.text)
return None
def _worksheet_evidence(turn: ReviewTurn) -> ReviewWorksheetEvidence:
return ReviewWorksheetEvidence(
turnId=turn.id,
speaker=turn.speaker,
quote=_clip_text(turn.text, 120),
)
def _worksheet_item(
*,
key: str,
label: str,
turns: list[ReviewTurn],
keywords: list[str],
preferred_speaker: Literal["learner", "client"] | None = None,
fallback_turn: ReviewTurn | None = None,
) -> ReviewWorksheetItem:
lowered_keywords = [keyword.lower() for keyword in keywords if keyword]
candidates = turns
if preferred_speaker:
preferred = [turn for turn in turns if turn.speaker == preferred_speaker]
candidates = preferred + [turn for turn in turns if turn.speaker != preferred_speaker]
for turn in candidates:
text = _compact_text(turn.text)
lower_text = text.lower()
if lowered_keywords and any(keyword in lower_text for keyword in lowered_keywords):
return ReviewWorksheetItem(
key=key,
label=label,
value=_clip_text(text, 140),
evidence=[_worksheet_evidence(turn)],
confidence="medium",
)
if fallback_turn is not None:
return ReviewWorksheetItem(
key=key,
label=label,
value=_clip_text(fallback_turn.text, 140),
evidence=[_worksheet_evidence(fallback_turn)],
confidence="low",
)
return ReviewWorksheetItem(
key=key,
label=label,
value=None,
evidence=[],
confidence="none",
emptyReason="저장된 축어록에서 명시 근거를 찾지 못했습니다.",
)
def _worksheet_section(
key: str,
title: str,
specs: list[tuple[str, str, list[str], Literal["learner", "client"] | None]],
turns: list[ReviewTurn],
fallback_client: ReviewTurn | None,
fallback_learner: ReviewTurn | None,
) -> ReviewWorksheetSection:
items: list[ReviewWorksheetItem] = []
for item_key, label, keywords, speaker in specs:
fallback = fallback_client if speaker == "client" else fallback_learner if speaker == "learner" else None
items.append(
_worksheet_item(
key=item_key,
label=label,
turns=turns,
keywords=keywords,
preferred_speaker=speaker,
fallback_turn=fallback if item_key in {"presenting_complaint", "first_goal"} else None,
)
)
return ReviewWorksheetSection(key=key, title=title, items=items)
def _case_worksheet_from_turns(turns: list[ReviewTurn]) -> ReviewCaseWorksheet:
if not turns:
return ReviewCaseWorksheet(
status="empty",
sections=[],
limitations=["저장된 축어록이 없어 사례개념화 워크시트를 생성하지 않았습니다."],
)
fallback_client = next((turn for turn in turns if turn.speaker == "client"), None)
fallback_learner = next((turn for turn in turns if turn.speaker == "learner"), None)
section_specs: list[
tuple[str, str, list[tuple[str, str, list[str], Literal["learner", "client"] | None]]]
] = [
(
"exploration_11",
"탐색 11항목",
[
("presenting_complaint", "주호소", ["힘들", "문제", "걱정", "불안", "우울", "스트레스", "관계"], "client"),
("trigger_context", "계기·상황", ["언제", "상황", "최근", "계기", ""], "client"),
("emotion", "정서", ["불안", "우울", "", "슬프", "답답", "무섭", "외롭", "걱정"], "client"),
("cognition", "생각", ["생각", "느낌", "해야", "", "실패", "의미"], "client"),
("behavior", "행동", ["피하", "", "", "", "", "연락", "공부", ""], "client"),
("body", "신체·수면", ["", "식욕", "", "두통", "심장", "", "피곤"], "client"),
("relationship", "관계", ["친구", "가족", "부모", "엄마", "아빠", "교수", "사람", "관계"], "client"),
("resources", "자원", ["도움", "지지", "친구", "상담", "선생님", "가족"], "client"),
("risk", "위험 신호", ["", "자살", "해치", "사라지고", "끝내", "위험"], "client"),
("motivation", "변화동기", ["", "바라", "변화", "해보고", ""], None),
("first_goal", "상담 목표 초안", ["목표", "계획", "다음", "해볼", "원하"], "learner"),
],
),
(
"five_domains",
"호소 5영역",
[
("domain_emotion", "정서", ["불안", "우울", "", "슬프", "답답", "외롭"], "client"),
("domain_cognition", "인지", ["생각", "걱정", "실패", "", "의미"], "client"),
("domain_behavior", "행동", ["피하", "연락", "공부", "", ""], "client"),
("domain_relationship", "대인관계", ["친구", "가족", "사람", "관계", "부모"], "client"),
("domain_body", "신체", ["", "식욕", "", "두통", "피곤", ""], "client"),
],
),
(
"cognitive_triad_emotions",
"인지삼제·1/2차 감정",
[
("triad_self", "자기", ["나는", "내가", "나 자신", "스스로"], "client"),
("triad_world", "타인·세계", ["사람", "세상", "학교", "가족", "친구"], "client"),
("triad_future", "미래", ["앞으로", "미래", "계속", "나중"], "client"),
("primary_emotion", "1차 감정", ["불안", "슬프", "무섭", "외롭", "걱정"], "client"),
("secondary_emotion", "2차 감정", ["", "짜증", "수치", "죄책", "부끄"], "client"),
],
),
(
"protective_barrier_quadrants",
"보호·방해 4사분면",
[
("internal_protective", "내적 보호요인", ["해보고", "버텼", "노력", "", "견뎠"], None),
("internal_barrier", "내적 방해요인", ["", "두려", "불안", "회피", "걱정"], "client"),
("external_protective", "외적 보호요인", ["친구", "가족", "상담", "교수", "도움"], "client"),
("external_barrier", "외적 방해요인", ["갈등", "압박", "비난", "스트레스", "혼자"], "client"),
],
),
(
"biopsychosocial_goals",
"생물·심리·사회 목표",
[
("bio_goal", "생물", ["", "식사", "운동", "", "피곤"], "client"),
("psy_goal", "심리", ["생각", "감정", "불안", "연습", "조절"], None),
("social_goal", "사회", ["관계", "대화", "연락", "도움", "친구"], None),
],
),
]
sections = [
_worksheet_section(
key,
title,
specs,
turns,
fallback_client,
fallback_learner,
)
for key, title, specs in section_specs
]
return ReviewCaseWorksheet(
status="draft_from_transcript",
sections=sections,
limitations=[
"저장된 축어록에서 키워드 근거를 추출한 1차 초안입니다.",
"임상팀 루브릭, 교수자 검수, 학습자 수정 입력 전에는 확정 사례개념화로 보지 않습니다.",
],
)
def _evaluation_payload(record: dict[str, object] | None) -> dict[str, object]:
if not record:
return {}
payload = record.get("payload")
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
def _seconds_label(milliseconds: int) -> str:
seconds = max(0, milliseconds) / 1000.0
if seconds >= 10:
return f"{seconds:.0f}"
return f"{seconds:.1f}"
def _review_nonverbal_events(turn: TurnRecord) -> list[ReviewNonverbalEvent]:
events: list[ReviewNonverbalEvent] = []
if turn.silence_ms is not None and turn.silence_ms >= 1000:
events.append(
ReviewNonverbalEvent(
kind="silence",
label="침묵",
detail=_seconds_label(turn.silence_ms),
)
)
if turn.speech_rate is not None:
events.append(
ReviewNonverbalEvent(
kind="pace",
label="발화 속도",
detail=f"분당 {turn.speech_rate:.0f}",
)
)
if turn.barge_in is True:
events.append(
ReviewNonverbalEvent(
kind="barge_in",
label="끼어듦",
detail="내담자 발화 중 시작",
)
)
if turn.audio_ref:
events.append(
ReviewNonverbalEvent(
kind="audio",
label="음성 입력",
detail="음성으로 기록됨",
)
)
return events
async def _evaluate_stream_turn(ctx: orchestrator.TurnContext, final_reply: str) -> Optional[dict]:
"""stream 경로 완료 후 fast-loop 평가를 계산한다. 실패는 턴 저장을 막지 않는다."""
if not final_reply:
return None
try:
hook = evaluator.make_eval_hook(
engine_client,
audit_hook=session_persistence.record_llm_call_audit,
)
return await hook(ctx, final_reply)
except Exception:
return None
def _stream_result_from_done(
ctx: orchestrator.TurnContext,
final_reply: str,
data: dict[str, object],
evaluation: Optional[dict],
) -> orchestrator.TurnResult:
assert ctx.state_after is not None
return orchestrator.TurnResult(
turn_seq=ctx.state_after.turn_seq,
stage=_stage_label(ctx.state_after.stage),
effective_openness=ctx.state_after.effective_openness,
client_reply=final_reply or None,
safety_flagged=bool(data.get("safety_flagged")),
state_after=ctx.state_after,
evaluation=evaluation,
crisis_kind=ctx.crisis.kind.value if ctx.crisis else "none",
crisis_resource=data.get("crisis_resource") if isinstance(data.get("crisis_resource"), dict) else None,
conversation_stopped=bool(data.get("conversation_stopped")),
llm_provider=str(data.get("llm_provider") or "") or None,
model=str(data.get("model") or "") or None,
tokens_in=int(data.get("tokens_in") or 0),
tokens_out=int(data.get("tokens_out") or 0),
cost_usd=float(data.get("cost_usd") or 0.0),
)
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
enriched: list[dict[str, object]] = []
for index, turn in enumerate(sess.masked_turns(), start=1):
item: dict[str, object] = dict(turn)
item["seq"] = index
enriched.append(item)
try:
result = await asyncio.wait_for(
evaluator.evaluate_session(
session_id=sess.session_id,
stage=_stage_label(sess.state.stage),
masked_turns=enriched,
engine=engine_client,
technique_codes=[],
theory_mode=sess.theory_mode,
scope="session_end",
audit_hook=session_persistence.record_llm_call_audit,
),
timeout=min(float(settings.engine_timeout), 45.0),
)
status_value = "error" if result.error else "ready"
await session_persistence.save_session_evaluation(
session_id=sess.session_id,
learner_id=sess.learner_id,
status=status_value,
source="engine",
scope=result.scope,
stage=result.stage,
payload=result.to_dict(),
error=result.error,
)
except Exception as exc:
await session_persistence.save_session_evaluation(
session_id=sess.session_id,
learner_id=sess.learner_id,
status="error",
source="engine",
scope="session_end",
stage=_stage_label(sess.state.stage),
payload={},
error=str(exc),
)
def _schedule_session_evaluation(sess: InProcSession) -> None:
if not sess.turns:
return
asyncio.create_task(_generate_and_save_session_evaluation(sess))
def _learner_summary(sess: InProcSession, *, review_ready: bool = False) -> LearnerSessionSummary:
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,
persona_name=sess.persona.display_name,
session_no=sess.session_no,
status="ended" if sess.ended else "active",
stage=_stage_label(sess.state.stage),
turn_count=len(turns),
learner_turn_count=learner_turns,
client_turn_count=client_turns,
started_at=_iso(sess.created_at) or "",
ended_at=_iso(sess.ended_at),
review_ready=review_ready,
)
async def _review_ready(sess: InProcSession, principal: Principal) -> bool:
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,
principal,
)
return bool(evaluation_record and evaluation_record.get("status") == "ready")
def _session_detail(
sess: InProcSession,
*,
review_ready: bool = False,
) -> SessionDetailResponse:
turns = _learner_visible_turns(sess)
return SessionDetailResponse(
session_id=sess.session_id,
case_id=sess.case_id,
persona_code=sess.persona_code,
persona_name=sess.persona.display_name,
theory_mode=sess.theory_mode,
status="ended" if sess.ended else "active",
stage=_stage_label(sess.state.stage),
effective_openness=round(sess.state.effective_openness, 4),
started_at=_iso(sess.created_at) or "",
ended_at=_iso(sess.ended_at),
turns=[
SessionDetailTurn(
turn_seq=turn.turn_seq,
speaker="learner" if turn.speaker == "counselor" else "client",
stage=turn.stage,
text=turn.text_masked,
created_at=_iso(turn.created_at) or "",
)
for turn in turns
],
review_ready=review_ready,
)
@router.get("", response_model=LearnerSessionsResponse)
async def list_learner_sessions(principal: CurrentPrincipal) -> LearnerSessionsResponse:
"""Return the current learner's real practice sessions."""
_ensure_learner(principal)
sessions, durable = await session_persistence.list_sessions(principal)
if not durable:
require_runtime_fallback_allowed("session list")
sessions = [
sess
for sess in store.list()
if sess.learner_id == principal.user_id
]
sessions.sort(key=lambda sess: sess.created_at, reverse=True)
summaries: list[LearnerSessionSummary] = []
for sess in sessions[:20]:
summaries.append(_learner_summary(sess, review_ready=await _review_ready(sess, principal)))
return LearnerSessionsResponse(
source="database" if durable else "runtime",
sessions=summaries,
)
@router.get("/{session_id}", response_model=SessionDetailResponse)
async def get_session_detail(
session_id: str,
principal: CurrentPrincipal,
) -> SessionDetailResponse:
"""Return a learner-owned session with transcript for resume/history."""
_ensure_learner(principal)
sess = await _load_session_or_404(
session_id,
principal,
allow_ended=True,
)
return _session_detail(sess, review_ready=await _review_ready(sess, principal))
@router.post("", response_model=SessionStartResponse, status_code=status.HTTP_201_CREATED)
async def start_session(
body: SessionStartRequest,
principal: CurrentPrincipal,
) -> SessionStartResponse:
"""Start a learner-owned practice session."""
_ensure_learner(principal)
try:
catalog_persona = await get_catalog_persona(body.persona_code)
except Exception as exc:
raise HTTPException(
status.HTTP_503_SERVICE_UNAVAILABLE,
detail="persona catalog database unavailable",
) from exc
if catalog_persona is None:
raise HTTPException(status.HTTP_404_NOT_FOUND, detail=f"unknown persona {body.persona_code}")
card = catalog_persona.card
case_context = await session_persistence.get_case_context(
learner_id=principal.user_id,
persona_id=catalog_persona.persona_id,
)
recall = await _build_seed_recall(case_id=case_context.case_id if case_context else None)
session_no = (case_context.last_session_no + 1) if case_context else 1
st = state_machine.init_state(
params=card.openness_params(),
carry=recall.carry,
)
carry_rapport = st.rapport_credit
sess = await session_persistence.create_session(
learner_id=principal.user_id,
card=card,
theory_mode=body.theory_mode,
state=st,
session_no=session_no,
carry_rapport=carry_rapport,
persona_id=catalog_persona.persona_id,
persona_version=catalog_persona.version,
case_id=case_context.case_id if case_context else None,
)
degraded = catalog_persona.degraded or sess is None
if sess is None:
require_runtime_fallback_allowed("session creation")
sess = store.create(
learner_id=principal.user_id,
persona=card,
theory_mode=body.theory_mode,
state=st,
session_no=session_no,
carry_rapport=carry_rapport,
)
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,
case_id=sess.case_id,
session_no=sess.session_no,
stage=_stage_label(st.stage),
effective_openness=round(st.effective_openness, 4),
recall_summary=recall.recall_summary,
degraded=degraded,
)
@router.get("/{session_id}/review", response_model=SessionReviewResponse)
async def get_session_review(
session_id: str,
principal: CurrentPrincipal,
) -> SessionReviewResponse:
"""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,
include_turn_evaluation=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)))
client_name = _client_name(sess.persona.display_name)
client_initial = client_name[:1] or ""
reached_phase = _stage_label(sess.state.stage)
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))
evaluation_record, evaluation_durable = await session_persistence.load_session_evaluation(
session_id,
principal,
)
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 = visible_turns[0].created_at if visible_turns else sess.created_at
turns: list[ReviewTurn] = []
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}",
ts=_offset_label(turn.created_at - first_turn_ts),
speaker=speaker,
who="학습자" if speaker == "learner" else client_name,
text=turn.text_masked,
techniques=_review_techniques_from_turn_eval(turn_eval),
nonverbal=_review_nonverbal_events(turn) if speaker == "learner" else [],
note=_review_note_from_turn_eval(turn_eval),
)
)
if not turns:
session_signal = "기록 없음"
elif sess.ended:
session_signal = "종료됨"
else:
session_signal = "진행 중"
transcript_summary = _review_summary(
client_name=client_name,
reached_phase=reached_phase,
turns=turns,
)
rubric: list[ReviewRubricRow] = []
good_moments: list[ReviewPoint] = []
growth_points: list[ReviewPoint] = []
next_line: str | None = None
if evaluation_ready:
rubric = _rubric_from_evaluation(evaluation_payload)
good_moments = _ai_review_points(
evaluation_payload.get("strengths"),
fallback_prefix="강점",
)
growth_points = _ai_review_points(
evaluation_payload.get("improvements"),
fallback_prefix="개선점",
)
if not growth_points:
growth_points = _intent_deviation_points(evaluation_payload.get("intent_deviations"))
next_line = _next_line_from_evaluation(evaluation_payload)
client_feedback = _latest_client_feedback(turns)
review_degraded = bool(turns) and not evaluation_ready
if evaluation_ready:
supervisor_state = "평가 완료"
elif evaluation_status == "error":
supervisor_state = "평가 실패"
elif turns:
supervisor_state = "평가 대기"
else:
supervisor_state = "기록 대기"
summary = _review_summary_from_evaluation(
fallback=transcript_summary,
evaluation_record=None if hidden_turns else evaluation_record,
payload=evaluation_payload,
)
if evaluation_record and not hidden_turns and not evaluation_durable:
summary += " 현재 평가는 런타임 캐시에서 복원되었습니다."
return SessionReviewResponse(
session_id=session_id,
client=ReviewClient(
name=client_name,
initial=client_initial,
persona=f"{sess.persona_code} · {sess.persona.difficulty}",
),
date=datetime.fromtimestamp(sess.created_at).strftime("%Y-%m-%d"),
durationLabel=_duration_label(duration_seconds),
durationSeconds=duration_seconds,
reachedPhase=reached_phase,
sessionSignal=session_signal,
supervisorState=supervisor_state,
supervisorName="AI",
summary=summary,
phases=_phase_segments(stage_labels),
phaseAxis=axis,
valenceAxis=axis,
clientValence=[],
counselorBaseline=[],
turns=turns,
rubric=rubric,
goodMoments=good_moments,
growthPoints=growth_points,
caseWorksheet=_case_worksheet_from_turns(turns),
nextLine=next_line,
clientFeedback=client_feedback,
audioUrl=None,
pdfExportUrl=None,
degraded=review_degraded,
reviewReady=evaluation_ready,
)
@router.post("/{session_id}/turn", response_model=TurnResponse)
async def submit_turn(
session_id: str,
body: TurnRequest,
principal: CurrentPrincipal,
) -> TurnResponse:
"""Submit one trainee utterance and return the generated client reply."""
_ensure_learner(principal)
sess = await _load_session_or_404(session_id, principal)
recall = await ensure_recall_context(sess)
kb_cues = _KB_CUES_CACHE.get(session_id) or [] # 비차단: warm 전이면 빈 단서(graceful)
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(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,
eval_hook=evaluator.make_eval_hook(
engine_client,
audit_hook=session_persistence.record_llm_call_audit,
),
audit_hook=session_persistence.record_llm_call_audit,
)
except EngineError as exc:
raise HTTPException(
status.HTTP_503_SERVICE_UNAVAILABLE,
detail=f"engine unavailable: {exc}",
) from exc
await turn_runtime.record_completed_turn(
sess,
ctx,
result,
context_prefix="session",
)
await turn_runtime.record_safety_event(sess, ctx, result)
return TurnResponse(
turn_seq=result.turn_seq,
stage=_stage_label(result.state_after.stage),
effective_openness=round(result.effective_openness, 4),
client_reply=result.client_reply,
safety_flagged=result.safety_flagged,
crisis_kind=result.crisis_kind,
crisis_resource=result.crisis_resource,
conversation_stopped=result.conversation_stopped,
)
@router.post("/{session_id}/stream")
async def stream_turn(
session_id: str,
body: TurnRequest,
principal: CurrentPrincipal,
):
"""Stream a generated client reply for one trainee utterance."""
_ensure_learner(principal)
sess = await _load_session_or_404(session_id, principal)
recall = await ensure_recall_context(sess)
kb_cues = _KB_CUES_CACHE.get(session_id) or [] # 비차단: warm 전이면 빈 단서(graceful)
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(visible_to="client"),
kb_behavior_cues=kb_cues,
theory_mode=sess.theory_mode,
)
assert ctx.state_after is not None
async def event_generator():
last_beat = asyncio.get_running_loop().time()
final_reply = ""
try:
async for ev in orchestrator.run_turn_stream(
ctx,
engine_client,
audit_hook=session_persistence.record_llm_call_audit,
):
if ev.event == "token":
text = str(ev.data.get("text", ""))
final_reply += text
yield {"event": "token", "data": text}
elif ev.event == "done":
data = {**ev.data, "stage": _stage_label(ctx.state_after.stage)}
evaluation = await _evaluate_stream_turn(ctx, final_reply)
result = _stream_result_from_done(ctx, final_reply, data, evaluation)
await turn_runtime.record_completed_turn(
sess,
ctx,
result,
context_prefix="session",
)
await turn_runtime.record_safety_event(sess, ctx, result)
yield {"event": "done", "data": json.dumps(data, ensure_ascii=False)}
else:
yield {"event": ev.event, "data": json.dumps(ev.data, ensure_ascii=False)}
now = asyncio.get_running_loop().time()
if now - last_beat >= settings.sse_heartbeat_seconds:
yield {"event": "ping", "data": "{}"}
last_beat = now
except Exception as exc:
yield {"event": "error", "data": json.dumps({"detail": str(exc)}, ensure_ascii=False)}
return
return EventSourceResponse(event_generator())
@router.post("/{session_id}/end", response_model=SessionEndResponse)
async def end_session(
session_id: str,
principal: CurrentPrincipal,
) -> SessionEndResponse:
"""End a learner-owned session and prepare carry-over state."""
_ensure_learner(principal)
sess = await _load_session_or_404(session_id, principal, allow_ended=True)
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,
)
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(
session_id=session_id,
session_no=sess.session_no,
digest_pending=carry.compression_job is not None,
end_state=carry.end_state,
)