diff --git a/apps/api/app/routes/sessions.py b/apps/api/app/routes/sessions.py index 6d1b498..7afab38 100644 --- a/apps/api/app/routes/sessions.py +++ b/apps/api/app/routes/sessions.py @@ -10,9 +10,7 @@ from __future__ import annotations import asyncio import json -import re import secrets -from collections import Counter from datetime import datetime from typing import Literal, Optional @@ -27,13 +25,41 @@ 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, guardrail, live_coach, memory, orchestrator, rag, session_metrics, state_machine +from ..services import evaluator, guardrail, live_coach, memory, orchestrator, rag, state_machine +from ..session_read_model import ( + LearnerDashboardResponse, + LearnerSessionsResponse, + LearnerSessionSummary, + LEARNER_VISIBLE_AI_ROLE, + ReviewCaseWorksheet, + ReviewCaseWorksheetSaveRequest, + ReviewWorksheetItem, + ReviewWorksheetSection, + SessionArchiveResponse, + SessionDetailResponse, + SessionReviewReadInput, + SessionReviewResponse, + SessionShareDeleteResponse, + SessionShareResponse, + StageLabel, + build_session_review, + dashboard_achievements as _dashboard_achievements, + dashboard_feedback as _dashboard_feedback, + dashboard_growth as _dashboard_growth, + dashboard_overview as _dashboard_overview, + dashboard_persona_progress as _dashboard_persona_progress, + iso as _iso, + learner_summary as _learner_summary, + learner_visible_turns as _learner_visible_turns, + session_detail as _session_detail, + session_share_payload as _session_share_payload, + stage_label as _stage_label, +) from ..store import InProcSession, TurnRecord, store router = APIRouter(prefix="/sessions", tags=["sessions"]) TheoryMode = Literal["humanistic", "cbt", "integrative"] -StageLabel = Literal["라포", "탐색", "개입", "정리"] EndStateValue = str | int | float | bool | None | dict[str, float] @@ -91,284 +117,6 @@ class SessionEndResponse(BaseModel): 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: StageLabel - turn_count: int - learner_turn_count: int - client_turn_count: int - started_at: str - ended_at: str | None = None - review_ready: bool = False - archived: bool = False - archived_at: str | None = None - - -class LearnerSessionsResponse(BaseModel): - source: str = "runtime" - sessions: list[LearnerSessionSummary] = Field(default_factory=list) - - -class LearnerDashboardOverview(BaseModel): - total_sessions: int = 0 - completed_sessions: int = 0 - active_sessions: int = 0 - review_ready_sessions: int = 0 - archived_sessions: int = 0 - learner_turns: int = 0 - client_turns: int = 0 - last_practiced_at: str | None = None - - -class LearnerDashboardGrowthPoint(BaseModel): - session_id: str - session_no: int - persona_code: str - stage: StageLabel - started_at: str - ended_at: str | None = None - score: float | None = None - rapport: float | None = None - technique_count: int = 0 - watch_count: int = 0 - - -class LearnerDashboardGrowth(BaseModel): - first_score: float | None = None - latest_score: float | None = None - score_delta: float | None = None - avg_score: float | None = None - avg_rapport: float | None = None - trend: str = "insufficient" - evaluated_sessions: int = 0 - top_techniques: list[str] = Field(default_factory=list) - points: list[LearnerDashboardGrowthPoint] = Field(default_factory=list) - - -class LearnerDashboardPersonaProgress(BaseModel): - persona_code: str - persona_name: str - sessions: int = 0 - completed_sessions: int = 0 - active_sessions: int = 0 - review_ready_sessions: int = 0 - latest_at: str | None = None - latest_stage: StageLabel | None = None - latest_score: float | None = None - trend: str = "insufficient" - - -class LearnerDashboardAchievement(BaseModel): - id: str - label: str - state: Literal["done", "available", "locked"] = "locked" - detail: str - - -class LearnerDashboardFeedbackItem(BaseModel): - session_id: str - persona_code: str - persona_name: str - session_no: int - stage: StageLabel - turn_seq: int - created_at: str - score: float | None = None - rapport: float | None = None - note: str - techniques: list[str] = Field(default_factory=list) - - -class LearnerDashboardResponse(BaseModel): - source: str = "runtime" - overview: LearnerDashboardOverview - growth: LearnerDashboardGrowth - persona_progress: list[LearnerDashboardPersonaProgress] = Field(default_factory=list) - achievements: list[LearnerDashboardAchievement] = Field(default_factory=list) - recent_feedback: list[LearnerDashboardFeedbackItem] = Field(default_factory=list) - message: str - - -class SessionArchiveResponse(BaseModel): - session_id: str - archived: bool - archived_at: str | None = None - source: str = "runtime" - session: LearnerSessionSummary - - -class SessionDetailTurn(BaseModel): - turn_seq: int - speaker: Literal["learner", "client"] - stage: StageLabel - 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: StageLabel - 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", "saved_by_learner"] = "empty" - generatedBy: str = "rule-based transcript extractor" - sections: list[ReviewWorksheetSection] = Field(default_factory=list) - limitations: list[str] = Field(default_factory=list) - savedAt: Optional[str] = None - - -class ReviewCaseWorksheetSaveRequest(BaseModel): - sections: list[ReviewWorksheetSection] = Field(default_factory=list) - limitations: list[str] = Field(default_factory=list) - - -class SessionTeacherReviewStatus(BaseModel): - status: Literal["pending", "viewed", "closed"] = "pending" - note: str = "" - reviewerId: str | None = None - reviewedAt: str | None = None - updatedAt: str | None = None - - -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 - teacherReview: SessionTeacherReviewStatus | None = None - - -class SessionShareResponse(BaseModel): - shareUrl: str - title: str - description: str - imageUrl: str - createdAt: str - - -class SessionShareDeleteResponse(BaseModel): - revoked: bool - - _RECALL_CACHE: dict[str, memory.RecallContext] = {} # 세션별 KB 증상 행동단서(회기 1회 산출·캐시). 빈 list 캐시 = 회기 내 재시도 안 함(안정성). _KB_CUES_CACHE: dict[str, list[str]] = {} @@ -498,9 +246,24 @@ async def _ensure_kb_cues(session_id: str, card) -> list[str]: async def _load_prev_case_summary(case_id: str) -> Optional[dict]: """직전 회기 요약(case 스코프) → build_recall_context 입력. 미존재/미가용 시 None.""" + case_memory = await _load_case_memory(case_id) + return case_memory.get("prev_summary") + + +async def _load_case_memory(case_id: str) -> dict: + """case-level 큰그림 + 직전 요약 + client-visible pinned fact를 한 번에 읽는다.""" + empty = {"case_digest": None, "prev_summary": None, "pinned_facts": []} try: async with db.acquire(ai_view=rag.AIRole.CLIENT.value) as conn: - row = await conn.fetchrow( + case_row = await conn.fetchrow( + """ + SELECT case_digest + FROM app.case_profile + WHERE case_id = $1::uuid + """, + case_id, + ) + summary_row = await conn.fetchrow( """ SELECT digest, open_threads, end_state FROM app.session_summary @@ -510,14 +273,33 @@ async def _load_prev_case_summary(case_id: str) -> Optional[dict]: """, case_id, ) + fact_rows = await conn.fetch( + """ + SELECT value + FROM app.pinned_fact + WHERE case_id = $1::uuid + AND status IN ('stable', 'evolving', 'locked') + AND $2 = ANY(visible_to) + ORDER BY updated_at DESC + LIMIT 12 + """, + case_id, + rag.AIRole.CLIENT.value, + ) except Exception: - return None - if row is None: - return None + return empty + + prev_summary = None + if summary_row is not None: + prev_summary = { + "digest": summary_row["digest"], + "open_threads": list(summary_row["open_threads"] or []), + "end_state": dict(summary_row["end_state"] or {}), + } return { - "digest": row["digest"], - "open_threads": list(row["open_threads"] or []), - "end_state": dict(row["end_state"] or {}), + "case_digest": (case_row["case_digest"] if case_row is not None else None) or None, + "prev_summary": prev_summary, + "pinned_facts": [row["value"] for row in fact_rows if row["value"]], } @@ -569,14 +351,15 @@ async def _build_start_recall(*, case_id: str, card) -> memory.RecallContext: db.get_pool() # 풀 미초기화 시 RuntimeError → 첫 회기와 동일한 빈 회상 except RuntimeError: return memory.build_recall_context() - prev_summary = await _load_prev_case_summary(case_id) + case_memory = await _load_case_memory(case_id) + prev_summary = case_memory.get("prev_summary") 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( + case_digest=case_memory.get("case_digest"), prev_summary=prev_summary, episodic_snippets=episodic, - pinned_facts=pinned, + pinned_facts=case_memory.get("pinned_facts") or [], ) @@ -587,9 +370,12 @@ async def _build_seed_recall(*, case_id: str | None) -> memory.RecallContext: 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) + case_memory = await _load_case_memory(case_id) + return memory.build_recall_context( + case_digest=case_memory.get("case_digest"), + prev_summary=case_memory.get("prev_summary"), + pinned_facts=case_memory.get("pinned_facts") or [], + ) async def ensure_recall_context(sess: InProcSession) -> memory.RecallContext: @@ -618,18 +404,6 @@ async def _warm_rag_caches(session_id: str, case_id: str, card) -> None: 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) -> Principal: if principal.role == Role.LEARNER: return principal @@ -711,610 +485,41 @@ async def _end_persisted_session(sess: InProcSession, carry: memory.CarryOver) - sess.ended = True sess.ended_at = datetime.now().timestamp() store.put(sess) + asyncio.create_task(_write_episodic_embeddings(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}" +async def _write_episodic_embeddings(sess: InProcSession) -> None: + """Best-effort M2 episodic writer. - -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 또는 교수자 코멘트가 아직 생성되지 않은 항목은 빈 상태로 남겨 둡니다." + Only masked client-visible client turns are eligible. Missing BGE-M3, pgvector, + or DB runtime should not change session-end persistence semantics. + """ + inputs = rag.episodic_turn_inputs_from_records( + session_id=sess.session_id, + case_id=sess.case_id, + turns=sess.turns, ) - - -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 _saved_case_worksheet_from_payload(payload: dict[str, object] | None) -> ReviewCaseWorksheet | None: - if not payload: - return None + if not inputs: + return try: - worksheet = ReviewCaseWorksheet.model_validate(payload) + db.get_pool() + async with db.acquire( + role="learner", + user_id=sess.learner_id, + ai_view=rag.AIRole.CLIENT.value, + ) as conn: + await rag.write_persona_turn_embeddings(conn, turns=inputs) except Exception: - return None - return worksheet.model_copy(update={"status": "saved_by_learner"}) - - -def _worksheet_share_highlights(worksheet: ReviewCaseWorksheet, *, limit: int = 4) -> list[dict[str, str]]: - highlights: list[dict[str, str]] = [] - for section in worksheet.sections: - for item in section.items: - value = _compact_text(item.value or "") - if not value: - continue - highlights.append( - { - "section": section.title, - "label": item.label, - "value": "비공개 요약 항목", - } - ) - if len(highlights) >= limit: - return highlights - return highlights - - -def _review_point_titles(points: list[ReviewPoint], *, limit: int = 3) -> list[str]: - return [_clip_text(point.title or point.body, 72) for point in points[:limit] if (point.title or point.body)] - - -def _share_image_url() -> str: - return f"{settings.frontend_base_url.rstrip('/')}/design-elements/clinical-paper-ambient.png" - - -def _session_share_payload(review: SessionReviewResponse) -> dict[str, object]: - title = f"Vignette 회기 리뷰 · {review.client.name} {review.date}" - description = _clip_text(review.summary, 156) - return { - "version": 1, - "title": title, - "description": description, - "summary": _clip_text(review.summary, 420), - "clientName": review.client.name, - "persona": review.client.persona, - "date": review.date, - "durationLabel": review.durationLabel, - "reachedPhase": review.reachedPhase, - "sessionSignal": review.sessionSignal, - "reviewReady": review.reviewReady, - "goodMoments": _review_point_titles(review.goodMoments), - "growthPoints": _review_point_titles(review.growthPoints), - "worksheetHighlights": _worksheet_share_highlights(review.caseWorksheet), - "imageUrl": _share_image_url(), - "appUrl": settings.frontend_base_url.rstrip("/"), - "privacy": "공유 카드에는 회기 원문 축어록과 학습자 식별 정보를 포함하지 않습니다.", - } + pass def _public_share_url(request: Request, token: str) -> str: base = str(request.base_url).rstrip("/") return f"{base}/share/session/{token}" - -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_body_markdown(text: str) -> str: - """평가 노트 본문을 회기 리뷰용 제한 markdown으로 정돈한다.""" - body = text.strip() - body = re.sub( - r"`?\beffective[_\s-]?openness\b`?(?!\(유효 개방도\))", - "`effective openness(유효 개방도)`", - body, - flags=re.IGNORECASE, - ) - body = re.sub( - r"(? str | None: - clean = " ".join(str(text or "").split()) - if not clean: - return None - sentences = [part.strip() for part in re.split(r"(?<=[.!?。!?])\s+", clean) if part.strip()] - for sentence in sentences: - if "가장 큰 마음" in sentence: - return sentence if len(sentence) <= limit else f"{sentence[: limit - 3].rstrip()}..." - for sentence in sentences: - if "?" in sentence: - return sentence if len(sentence) <= limit else f"{sentence[: limit - 3].rstrip()}..." - return clean if len(clean) <= limit else f"{clean[: limit - 3].rstrip()}..." - - -def _review_note_from_turn_eval( - ev: dict[str, object] | None, - learner_text: str | None = None, -) -> Optional[ReviewNote]: - """의도이탈(있으면 우선) 또는 적절성 신호를 턴 노트로. tone: good|warn(프론트 계약).""" - if not isinstance(ev, dict): - return None - quote = _review_quote_excerpt(learner_text) - 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_parts = ( - f"- 권장: {expected}" if expected else "", - f"- 실제: {actual}" if actual else "", - ) - body = "\n".join(p for p in body_parts if p) - return ReviewNote( - author="평가 AI", - tone="warn", - title=f"의도와 다른 부분 · {dimension}".rstrip(" ·") or "의도와 다른 부분", - body=_review_note_body_markdown(body or "권장 반응과 실제 반응에 차이가 있었어요."), - quote=quote, - ) - 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=_review_note_body_markdown( - note_text - or ( - "타당화·공감·탐색이 회기 흐름에 맞았습니다.\n\n" - "`effective openness(유효 개방도)`가 낮은 내담자라면 다음 질문은 " - "더 작고 구체적인 선택지로 낮춰도 좋습니다." - ) - ), - quote=quote, - ) - if appropriateness == "warn" and note_text: - return ReviewNote( - author="평가 AI", - tone="warn", - title="점검해볼 지점", - body=_review_note_body_markdown(note_text), - quote=quote, - ) - 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: @@ -1355,10 +560,6 @@ def _stream_result_from_done( ) -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 @@ -1413,34 +614,6 @@ def _schedule_session_evaluation(sess: InProcSession) -> None: asyncio.create_task(_generate_and_save_session_evaluation(sess)) -def _learner_summary( - sess: InProcSession, - *, - review_ready: bool = False, - archived: bool = False, - archived_at: str | None = None, -) -> 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, - archived=archived, - archived_at=archived_at, - ) - - async def _load_learner_sessions( principal: Principal, *, @@ -1496,187 +669,6 @@ async def _archive_map( return records -def _dashboard_growth_point(point: session_metrics.SessionGrowthPoint) -> LearnerDashboardGrowthPoint: - return LearnerDashboardGrowthPoint( - session_id=point.session_id, - session_no=point.session_no, - persona_code=point.persona_code, - stage=_stage_label(point.stage), - started_at=point.started_at, - ended_at=point.ended_at, - score=point.score, - rapport=point.rapport, - technique_count=point.technique_count, - watch_count=point.watch_count, - ) - - -def _dashboard_growth(sessions: list[InProcSession]) -> LearnerDashboardGrowth: - metrics = session_metrics.build_learner_growth( - sessions, - learner_label=lambda _learner_id: "나", - limit=1, - ) - if not metrics: - return LearnerDashboardGrowth() - item = metrics[0] - points = [_dashboard_growth_point(point) for point in item.points] - return LearnerDashboardGrowth( - first_score=item.first_score, - latest_score=item.latest_score, - score_delta=item.score_delta, - avg_score=item.avg_score, - avg_rapport=item.avg_rapport, - trend=item.trend, - evaluated_sessions=sum(1 for point in points if point.score is not None), - top_techniques=item.top_techniques, - points=points, - ) - - -def _dashboard_persona_progress( - sessions: list[InProcSession], - review_ready: dict[str, bool], -) -> list[LearnerDashboardPersonaProgress]: - grouped: dict[str, list[InProcSession]] = {} - for sess in sessions: - grouped.setdefault(sess.persona_code, []).append(sess) - - rows: list[LearnerDashboardPersonaProgress] = [] - for persona_code, items in grouped.items(): - ordered = sorted(items, key=session_metrics.session_activity_time) - latest = ordered[-1] - metrics = session_metrics.build_learner_growth( - ordered, - learner_label=lambda _learner_id: "나", - limit=1, - ) - growth = metrics[0] if metrics else None - rows.append( - LearnerDashboardPersonaProgress( - persona_code=persona_code, - persona_name=latest.persona.display_name, - sessions=len(ordered), - completed_sessions=sum(1 for sess in ordered if sess.ended), - active_sessions=sum(1 for sess in ordered if not sess.ended), - review_ready_sessions=sum( - 1 for sess in ordered if review_ready.get(sess.session_id, False) - ), - latest_at=session_metrics.iso_datetime( - session_metrics.session_activity_time(latest) - ), - latest_stage=_stage_label(latest.state.stage), - latest_score=growth.latest_score if growth else None, - trend=growth.trend if growth else "insufficient", - ) - ) - return sorted( - rows, - key=lambda row: row.latest_at or "", - reverse=True, - ) - - -def _achievement_state(done: bool, available: bool) -> Literal["done", "available", "locked"]: - if done: - return "done" - if available: - return "available" - return "locked" - - -def _dashboard_achievements( - sessions: list[InProcSession], - review_ready: dict[str, bool], -) -> list[LearnerDashboardAchievement]: - completed = sum(1 for sess in sessions if sess.ended) - active = sum(1 for sess in sessions if not sess.ended) - review_count = sum(1 for ready in review_ready.values() if ready) - by_persona = Counter(sess.persona_code for sess in sessions) - max_persona_sessions = max(by_persona.values(), default=0) - persona_coverage = len(by_persona) - return [ - LearnerDashboardAchievement( - id="first_session_complete", - label="첫 회기 완료", - state=_achievement_state(completed >= 1, active >= 1), - detail="한 회기를 종료하면 리뷰와 워크시트 흐름이 열립니다.", - ), - LearnerDashboardAchievement( - id="review_ready", - label="리뷰 확인 가능", - state=_achievement_state(review_count >= 1, completed >= 1), - detail=f"현재 리뷰 가능한 회기 {review_count}건입니다.", - ), - LearnerDashboardAchievement( - id="persona_repeat", - label="같은 내담자 반복 연습", - state=_achievement_state(max_persona_sessions >= 3, max_persona_sessions >= 1), - detail="같은 페르소나를 반복하면 변화 추이를 더 안정적으로 볼 수 있습니다.", - ), - LearnerDashboardAchievement( - id="persona_coverage", - label="여러 페르소나 경험", - state=_achievement_state(persona_coverage >= 3, persona_coverage >= 2), - detail=f"현재 {persona_coverage}개 페르소나에서 연습 기록이 있습니다.", - ), - ] - - -def _dashboard_feedback(sessions: list[InProcSession]) -> list[LearnerDashboardFeedbackItem]: - items: list[LearnerDashboardFeedbackItem] = [] - for item in session_metrics.recent_feedback_notes(sessions, limit=5): - score = item.get("score") - rapport = item.get("rapport") - items.append( - LearnerDashboardFeedbackItem( - session_id=str(item["session_id"]), - persona_code=str(item["persona_code"]), - persona_name=str(item["persona_name"]), - session_no=int(item["session_no"]), - stage=_stage_label(item["stage"]), - turn_seq=int(item["turn_seq"]), - created_at=str(item["created_at"]), - score=float(score) if isinstance(score, (int, float)) else None, - rapport=float(rapport) if isinstance(rapport, (int, float)) else None, - note=str(item["note"]), - techniques=[str(label) for label in item.get("techniques", [])], - ) - ) - return items - - -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=_stage_label(turn.stage), - text=turn.text_masked, - created_at=_iso(turn.created_at) or "", - ) - for turn in turns - ], - review_ready=review_ready, - ) - - async def _session_archive_response( sess: InProcSession, principal: Principal, @@ -1740,23 +732,13 @@ async def learner_dashboard(principal: CurrentPrincipal) -> LearnerDashboardResp for session_id, ready in review_ready.items() if session_id not in archives } - overview = LearnerDashboardOverview( - total_sessions=len(sessions), - completed_sessions=sum(1 for sess in sessions if sess.ended), - active_sessions=sum(1 for sess in sessions if not sess.ended), - review_ready_sessions=sum(1 for ready in visible_review_ready.values() if ready), - archived_sessions=len(archives), - learner_turns=sum(1 for sess in sessions for turn in _learner_visible_turns(sess) if turn.speaker == "counselor"), - client_turns=sum(1 for sess in sessions for turn in _learner_visible_turns(sess) if turn.speaker == "client"), - last_practiced_at=session_metrics.iso_datetime( - max((session_metrics.session_activity_time(sess) for sess in sessions), default=0.0) - ) - if sessions - else None, - ) return LearnerDashboardResponse( source="database" if durable else "runtime", - overview=overview, + overview=_dashboard_overview( + sessions, + visible_review_ready=visible_review_ready, + archived_sessions=len(archives), + ), growth=_dashboard_growth(sessions), persona_progress=_dashboard_persona_progress(sessions, visible_review_ready), achievements=_dashboard_achievements(sessions, visible_review_ready), @@ -1924,156 +906,33 @@ async def get_session_review( principal, 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 = [_stage_label(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, turn.text_masked), - ) - ) - - 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 += " 현재 평가는 런타임 캐시에서 복원되었습니다." - - generated_worksheet = _case_worksheet_from_turns(turns) saved_worksheet_payload, _ = await session_persistence.load_case_worksheet( session_id, principal, ) - case_worksheet = _saved_case_worksheet_from_payload(saved_worksheet_payload) or generated_worksheet - teacher_review: SessionTeacherReviewStatus | None = None - if principal.role in {Role.TEACHER, Role.ADMIN}: - review_status, _ = await session_persistence.load_session_review_status(session_id, principal) - review_status_value = str((review_status or {}).get("status") or "pending") - if review_status_value not in {"viewed", "closed"}: - review_status_value = "pending" - teacher_review = SessionTeacherReviewStatus( - status=review_status_value, # type: ignore[arg-type] - note=str((review_status or {}).get("note") or ""), - reviewerId=str((review_status or {}).get("reviewer_id") or "") or None, - reviewedAt=str((review_status or {}).get("reviewed_at") or "") or None, - updatedAt=str((review_status or {}).get("updated_at") or "") or None, + include_teacher_review = principal.role in {Role.TEACHER, Role.ADMIN} + teacher_review_record = None + if include_teacher_review: + teacher_review_record, _ = await session_persistence.load_session_review_status( + session_id, + principal, ) - 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, - nextLine=next_line, - clientFeedback=client_feedback, - audioUrl=None, - pdfExportUrl=None, - degraded=review_degraded, - reviewReady=evaluation_ready, - teacherReview=teacher_review, + return build_session_review( + SessionReviewReadInput( + session=sess, + evaluation_record=evaluation_record, + evaluation_durable=evaluation_durable, + saved_worksheet_payload=saved_worksheet_payload, + include_teacher_review=include_teacher_review, + teacher_review_record=teacher_review_record, + ) ) - @router.post("/{session_id}/share", response_model=SessionShareResponse) async def create_session_share( session_id: str, @@ -2267,7 +1126,7 @@ async def live_coach_turn( persona_name=sess.persona.display_name, learner_text=body.learner_text, client_reply=body.client_reply, - recent_turns=sess.recent_turns(k=8, visible_to=_LEARNER_VISIBLE_AI_ROLE), + recent_turns=sess.recent_turns(k=8, visible_to=LEARNER_VISIBLE_AI_ROLE), evaluation=_latest_turn_evaluation(sess, body.turn_seq), ) suggestion = await live_coach.generate_live_coaching( diff --git a/apps/api/app/routes/voice.py b/apps/api/app/routes/voice.py index 60e2d40..9a26c8a 100644 --- a/apps/api/app/routes/voice.py +++ b/apps/api/app/routes/voice.py @@ -48,6 +48,56 @@ WS_CLOSE_UNAUTHORIZED = 1008 # Per-utterance audio cap to avoid unbounded memory growth. _MAX_AUDIO_BYTES = 10 * 1024 * 1024 +_PROVIDER_EVENT_MAX_ITEMS = 12 +_PROVIDER_EVENT_MAX_STRING = 80 +_PROVIDER_EVENT_ALLOWED_KEYS = { + "category", + "event_type", + "type", + "kind", + "label", + "source", + "provider", + "start_ms", + "end_ms", + "duration_ms", + "confidence", + "score", + "is_final", +} +_PROVIDER_EVENT_TAXONOMY = { + "barge_in": ("barge_in", "turn_taking"), + "interrupt": ("barge_in", "turn_taking"), + "interruption": ("barge_in", "turn_taking"), + "overlap": ("barge_in", "turn_taking"), + "sigh": ("sigh", "paralinguistic"), + "sighing": ("sigh", "paralinguistic"), + "sob": ("cry", "paralinguistic"), + "cry": ("cry", "paralinguistic"), + "crying": ("cry", "paralinguistic"), + "weep": ("cry", "paralinguistic"), + "laugh": ("laugh", "paralinguistic"), + "laughter": ("laugh", "paralinguistic"), + "breath": ("breath", "paralinguistic"), + "breathing": ("breath", "paralinguistic"), + "voice_activity": ("voice_activity", "speech_activity"), + "vad": ("voice_activity", "speech_activity"), + "speech_start": ("speech_start", "speech_activity"), + "speech_end": ("speech_end", "speech_activity"), + "speech_final": ("speech_final", "speech_activity"), + "silence": ("silence", "timing"), + "pause": ("silence", "timing"), + "long_pause": ("silence", "timing"), + "speech_rate": ("speech_rate", "prosody"), + "fast_speech": ("speech_rate", "prosody"), + "slow_speech": ("speech_rate", "prosody"), + "pitch": ("pitch", "prosody"), + "intonation": ("intonation", "prosody"), + "prosody": ("prosody", "prosody"), + "noise": ("background_noise", "audio_quality"), + "background_noise": ("background_noise", "audio_quality"), +} +_PROVIDER_EVENT_TYPE_FIELDS = ("event_type", "type", "kind", "label") @router.get("/health") @@ -187,6 +237,7 @@ async def voice_ws(websocket: WebSocket) -> None: audio_ended_at=audio_ended_at, silence_ms=silence_ms, barge_in=_safe_bool(ctrl.get("barge_in")), + provider_events=_safe_provider_events(ctrl.get("provider_events")), ) last_audio_end_at = audio_ended_at audio_started_at = None @@ -232,6 +283,7 @@ async def _handle_utterance( audio_ended_at: float | None = None, silence_ms: int | None = None, barge_in: bool | None = None, + provider_events: list[dict[str, object]] | None = None, ) -> None: """Transcribe one utterance, generate the client reply, then synthesize TTS.""" if not audio: @@ -259,6 +311,7 @@ async def _handle_utterance( 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) + provider_events = _merge_provider_events(provider_events, getattr(stt, "provider_events", [])) await _safe_send_json( websocket, {"type": "transcript", "text": learner_text, "final": True, "speaker": "counselor"}, @@ -277,6 +330,7 @@ async def _handle_utterance( silence_ms=silence_ms, speech_rate=speech_rate, barge_in=barge_in, + provider_events=provider_events, ) @@ -291,6 +345,7 @@ async def _run_turn_and_speak( silence_ms: int | None = None, speech_rate: float | None = None, barge_in: bool | None = None, + provider_events: list[dict[str, object]] | None = None, ) -> None: """Run one counseling turn and stream synthesized client speech.""" sess, err = await _load_voice_session(session_id, principal) @@ -351,6 +406,7 @@ async def _run_turn_and_speak( silence_ms=silence_ms, speech_rate=speech_rate, barge_in=barge_in, + provider_events=provider_events or [], evaluation=result.evaluation, ), ) @@ -646,6 +702,71 @@ def _safe_bool(value: object) -> bool | None: return bool(value) +def _provider_event_slug(value: object) -> str: + text = str(value or "").strip().lower() + if not text: + return "" + chars = [] + previous_underscore = False + for char in text: + if char.isalnum(): + chars.append(char) + previous_underscore = False + elif not previous_underscore: + chars.append("_") + previous_underscore = True + return "".join(chars).strip("_")[:_PROVIDER_EVENT_MAX_STRING] + + +def _provider_event_taxonomy(event: dict[str, object]) -> tuple[str, str]: + for field in _PROVIDER_EVENT_TYPE_FIELDS: + slug = _provider_event_slug(event.get(field)) + if not slug: + continue + canonical = _PROVIDER_EVENT_TAXONOMY.get(slug) + if canonical: + return canonical + return slug, "unknown" + return "", "" + + +def _safe_provider_events(value: object) -> list[dict[str, object]]: + if not isinstance(value, list): + return [] + events: list[dict[str, object]] = [] + for item in value[:_PROVIDER_EVENT_MAX_ITEMS]: + if not isinstance(item, dict): + continue + safe: dict[str, object] = {} + for key in _PROVIDER_EVENT_ALLOWED_KEYS: + raw = item.get(key) + if isinstance(raw, bool): + safe[key] = raw + elif isinstance(raw, (int, float)): + safe[key] = raw + elif isinstance(raw, str): + text = raw.strip() + if text: + safe[key] = text[:_PROVIDER_EVENT_MAX_STRING] + if safe: + event_type, category = _provider_event_taxonomy(safe) + if event_type: + safe["event_type"] = event_type + if category: + safe["category"] = category + events.append(safe) + return events + + +def _merge_provider_events(*values: object) -> list[dict[str, object]]: + merged: list[dict[str, object]] = [] + for value in values: + merged.extend(_safe_provider_events(value)) + if len(merged) >= _PROVIDER_EVENT_MAX_ITEMS: + return merged[:_PROVIDER_EVENT_MAX_ITEMS] + return merged + + async def _safe_send_json(websocket: WebSocket, payload: dict) -> None: if websocket.client_state != WebSocketState.CONNECTED: return diff --git a/apps/api/app/services/memory.py b/apps/api/app/services/memory.py index bf5578e..a016ef5 100644 --- a/apps/api/app/services/memory.py +++ b/apps/api/app/services/memory.py @@ -18,12 +18,26 @@ DB 미가용(Docker off) 시에도 동작하도록 입력은 plain dict/list 로 from __future__ import annotations +import re from dataclasses import dataclass, field from typing import Any, Callable, Optional from .state_machine import SessionState +_SESSION_DIGEST_EXCERPT_CHARS = 90 +_CASE_DIGEST_MAX_ENTRIES = 12 +_RAPPORT_TRAJECTORY_MAX_ENTRIES = 24 +_PINNED_FACT_MAX_VALUE_CHARS = 160 +_COUNSELING_AGREEMENT_RE = re.compile(r"(상담|회기).*(주\s*\d+\s*회|매주|약속|계속|이어)|" + r"(주\s*\d+\s*회|매주).*(상담|회기)") +_COUNSELING_AGREEMENT_WITHDRAWAL_RE = re.compile( + r"(상담|회기).*(그만|중단|취소|철회|안\s*하|하지\s*않|못\s*하|이어\s*가지\s*않|계속\s*하지\s*않)|" + r"(약속).*(취소|철회|못\s*지키|지키지\s*않|안\s*지키)|" + r"더\s*이상.*(상담|회기).*(안\s*하|하지\s*않|못\s*하)" +) + + # ════════════════════════════════════════════════════════════════════════════ # 회기 시작 — 회상 (큰그림 → 세부) # ════════════════════════════════════════════════════════════════════════════ @@ -114,6 +128,18 @@ class CompressionJob: open_threads: list[str] = field(default_factory=list) +@dataclass(frozen=True, slots=True) +class PinnedFactCandidate: + """Rule-derived client-visible fact candidate for app.pinned_fact.""" + + key: str + value: str + fact_type: str + status: str = "stable" + confidence: float = 0.7 + source_turn_id: str | None = None + + def make_carry_over( *, state: SessionState, @@ -170,6 +196,174 @@ def build_compression_messages(job: CompressionJob) -> list[dict[str, str]]: return [{"role": "system", "content": system}, {"role": "user", "content": user}] +def _compact(value: Any, *, limit: int = _SESSION_DIGEST_EXCERPT_CHARS) -> str: + text = " ".join(str(value or "").split()) + if len(text) <= limit: + return text + return text[: max(0, limit - 3)].rstrip() + "..." + + +def build_fallback_session_digest( + *, + session_no: int, + masked_turns: list[dict[str, str]], + end_state: dict, +) -> str: + """마스킹 축어록 기반 임시 회기 digest. + + LLM 압축/embedding writer가 붙기 전에도 다음 회기 recall이 빈 문자열로 남지 않도록 + client-visible 마스킹 발화와 결정론 상태 수치만 사용한다. + """ + if not masked_turns: + return f"S{session_no}: 실제 발화가 없어 요약을 생성하지 않았다." + + counselor_count = sum(1 for turn in masked_turns if turn.get("speaker") == "counselor") + client_turns = [turn for turn in masked_turns if turn.get("speaker") == "client"] + client_count = len(client_turns) + last_client = _compact(client_turns[-1].get("text") if client_turns else "") + stage = str(end_state.get("stage") or "미확인") + openness = end_state.get("effective_openness") + rapport = end_state.get("rapport_credit") + status_bits = [f"종료 단계 {stage}"] + if openness is not None: + status_bits.append(f"개방도 {openness}") + if rapport is not None: + status_bits.append(f"라포 {rapport}") + status = ", ".join(status_bits) + if last_client: + return ( + f"S{session_no}: 마스킹 축어록 기준 상담자 {counselor_count}회, " + f"내담자 {client_count}회 발화. 마지막 내담자 반응은 \"{last_client}\". {status}." + ) + return ( + f"S{session_no}: 마스킹 축어록 기준 상담자 {counselor_count}회, " + f"내담자 {client_count}회 발화. {status}." + ) + + +def _fact_value(text: Any) -> str: + return _compact(text, limit=_PINNED_FACT_MAX_VALUE_CHARS) + + +def extract_pinned_fact_candidates( + masked_turns: list[dict[str, Any]], +) -> list[PinnedFactCandidate]: + """Extract conservative pinned facts from masked client-visible text. + + The first pass intentionally avoids clinical inference. It only preserves + explicit facts already surfaced by the client AI and already masked for + learner visibility. + """ + by_key: dict[str, PinnedFactCandidate] = {} + for turn in masked_turns: + if turn.get("speaker") != "client": + continue + text = _fact_value(turn.get("text")) + if not text: + continue + source_turn_id = turn.get("turn_id") + if "[NAME]" in text: + by_key["identity:name"] = PinnedFactCandidate( + key="identity:name", + value="[NAME]", + fact_type="identity", + confidence=0.85, + source_turn_id=str(source_turn_id) if source_turn_id else None, + ) + if "[ORG]" in text: + by_key["identity:org"] = PinnedFactCandidate( + key="identity:org", + value="[ORG]", + fact_type="identity", + confidence=0.85, + source_turn_id=str(source_turn_id) if source_turn_id else None, + ) + if _COUNSELING_AGREEMENT_WITHDRAWAL_RE.search(text): + by_key["agreement:counseling"] = PinnedFactCandidate( + key="agreement:counseling", + value=text, + fact_type="agreement", + status="contradicted", + confidence=0.8, + source_turn_id=str(source_turn_id) if source_turn_id else None, + ) + elif _COUNSELING_AGREEMENT_RE.search(text): + by_key["agreement:counseling"] = PinnedFactCandidate( + key="agreement:counseling", + value=text, + fact_type="agreement", + confidence=0.75, + source_turn_id=str(source_turn_id) if source_turn_id else None, + ) + return list(by_key.values()) + + +def merge_case_digest( + *, + existing_digest: str | None, + session_no: int, + session_digest: str, + max_entries: int = _CASE_DIGEST_MAX_ENTRIES, +) -> str: + """case_profile.case_digest를 session_no 기준으로 idempotent append한다.""" + prefix = f"S{session_no}:" + lines = [ + line.strip() + for line in str(existing_digest or "").splitlines() + if line.strip() and not line.strip().startswith(prefix) + ] + next_line = session_digest.strip() + if next_line and not next_line.startswith(prefix): + next_line = f"{prefix} {next_line}" + if next_line: + lines.append(next_line) + return "\n".join(lines[-max_entries:]) + + +def rapport_trajectory_point(*, session_no: int, end_state: dict) -> dict[str, Any]: + """case_profile.rapport_trajectory에 저장할 최소 무손실 수치 포인트.""" + return { + "session_no": int(session_no), + "stage": end_state.get("stage"), + "end_rapport": end_state.get("rapport_credit"), + "end_openness": end_state.get("effective_openness"), + "resistance": end_state.get("resistance"), + } + + +def merge_rapport_trajectory( + existing: Any, + point: dict[str, Any], + *, + max_entries: int = _RAPPORT_TRAJECTORY_MAX_ENTRIES, +) -> list[dict[str, Any]]: + """session_no 기준으로 trajectory를 덮어쓰기 가능하게 append한다.""" + session_no = point.get("session_no") + merged: list[dict[str, Any]] = [] + if isinstance(existing, list): + for item in existing: + if not isinstance(item, dict): + continue + if item.get("session_no") == session_no: + continue + merged.append(dict(item)) + merged.append(dict(point)) + return merged[-max_entries:] + + +def update_alliance_level(previous: Any, end_rapport: Any) -> float: + """case_profile.alliance_level EWMA. 이전 값이 없으면 schema default 0.2 기준.""" + try: + prev = float(previous) + except (TypeError, ValueError): + prev = 0.2 + try: + rapport = float(end_rapport) + except (TypeError, ValueError): + rapport = prev + return round(max(0.0, min(1.0, prev * 0.7 + rapport * 0.3)), 4) + + __all__ = [ "RecallContext", "build_recall_context", @@ -177,4 +371,11 @@ __all__ = [ "CompressionJob", "make_carry_over", "build_compression_messages", + "build_fallback_session_digest", + "PinnedFactCandidate", + "extract_pinned_fact_candidates", + "merge_case_digest", + "rapport_trajectory_point", + "merge_rapport_trajectory", + "update_alliance_level", ] diff --git a/apps/api/app/services/voice.py b/apps/api/app/services/voice.py index cdd27ac..1a0f70e 100644 --- a/apps/api/app/services/voice.py +++ b/apps/api/app/services/voice.py @@ -19,7 +19,7 @@ PRESET_TO_OPENAI_VOICE 테이블이 흡수. 새 preset 추가는 이 테이블 from __future__ import annotations import re -from dataclasses import dataclass +from dataclasses import dataclass, field from pathlib import Path from typing import Any, AsyncIterator, Mapping, Optional @@ -137,6 +137,7 @@ class TranscriptResult: language: Optional[str] = None model: str = STT_MODEL duration: Optional[float] = None + provider_events: list[dict[str, object]] = field(default_factory=list) @dataclass(frozen=True, slots=True) diff --git a/apps/api/app/session_persistence.py b/apps/api/app/session_persistence.py index 9f64cc1..6401ec2 100644 --- a/apps/api/app/session_persistence.py +++ b/apps/api/app/session_persistence.py @@ -551,6 +551,7 @@ def _turn_from_row(row, evaluation: dict[str, Any] | None = None) -> TurnRecord: silence_ms=_row_value(row, "silence_ms"), speech_rate=_row_value(row, "speech_rate"), barge_in=_row_value(row, "barge_in"), + provider_events=_dict_items(_row_value(row, "provider_events")), evaluation=evaluation, visible_to=tuple(_row_value(row, "visible_to") or DEFAULT_TURN_VISIBLE_TO), ) @@ -1890,7 +1891,7 @@ async def load_session( """ SELECT id, seq, speaker, stage, text, text_masked, created_at, llm_provider, model, tokens_in, tokens_out, cost_usd, - audio_ref, silence_ms, speech_rate, barge_in, visible_to + audio_ref, silence_ms, speech_rate, barge_in, provider_events, visible_to FROM app.turns WHERE session_id = $1::uuid ORDER BY seq @@ -1945,12 +1946,12 @@ async def append_turn( INSERT INTO app.turns ( session_id, seq, speaker, stage, text, text_masked, actor_kind, llm_provider, model, tokens_in, tokens_out, cost_usd, - audio_ref, silence_ms, speech_rate, barge_in, visible_to + audio_ref, silence_ms, speech_rate, barge_in, provider_events, visible_to ) VALUES ( $1::uuid, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, - $13, $14, $15, $16, $17::text[] + $13, $14, $15, $16, $17::jsonb, $18::text[] ) ON CONFLICT (session_id, seq) DO NOTHING RETURNING id @@ -1971,6 +1972,7 @@ async def append_turn( turn.silence_ms, turn.speech_rate, turn.barge_in, + turn.provider_events or [], list(turn.visible_to or DEFAULT_TURN_VISIBLE_TO), ) if inserted_turn_id is None: @@ -2000,13 +2002,176 @@ async def update_state( return False +async def _insert_pinned_fact_history( + conn: Any, + *, + fact_id: Any, + case_id: str, + old_value: Any, + new_value: Any, + reason: str, + session_no: int, + turn_id: str | None, +) -> None: + await conn.execute( + """ + INSERT INTO app.pinned_fact_history ( + fact_id, case_id, old_value, new_value, reason, session_no, turn_id + ) + VALUES ($1::uuid, $2::uuid, $3, $4, $5, $6, $7::uuid) + """, + fact_id, + case_id, + old_value, + new_value, + reason, + session_no, + turn_id, + ) + + +async def _upsert_pinned_fact_candidates(conn: Any, sess: InProcSession) -> None: + turn_rows = [ + { + "speaker": turn.speaker, + "text": turn.text_masked, + "turn_id": turn.turn_id, + } + for turn in sess.turns_visible_to("client") + ] + candidates = memory.extract_pinned_fact_candidates(turn_rows) + for fact in candidates: + if fact.status == "contradicted": + row = await conn.fetchrow( + """ + WITH existing AS ( + SELECT id, value + FROM app.pinned_fact + WHERE case_id = $1::uuid + AND key = $2 + AND status <> 'locked' + FOR UPDATE + ), + updated AS ( + UPDATE app.pinned_fact + SET value = $3, + fact_type = $4, + status = 'contradicted', + source_turn = COALESCE($5::uuid, app.pinned_fact.source_turn), + confidence = GREATEST(app.pinned_fact.confidence, $6), + version = app.pinned_fact.version + 1, + updated_session_no = $7, + visible_to = $8::text[], + updated_at = now() + FROM existing + WHERE app.pinned_fact.id = existing.id + RETURNING + app.pinned_fact.id, + existing.value AS old_value, + app.pinned_fact.value AS new_value + ) + SELECT id, old_value, new_value FROM updated + """, + sess.case_id, + fact.key, + fact.value, + fact.fact_type, + fact.source_turn_id, + fact.confidence, + sess.session_no, + ["evaluator"], + ) + if not row: + continue + old_value = row["old_value"] + new_value = row["new_value"] + if old_value == new_value: + continue + await _insert_pinned_fact_history( + conn, + fact_id=row["id"], + case_id=sess.case_id, + old_value=old_value, + new_value=new_value, + reason="contradiction", + session_no=sess.session_no, + turn_id=fact.source_turn_id, + ) + continue + row = await conn.fetchrow( + """ + WITH existing AS ( + SELECT id, value, status + FROM app.pinned_fact + WHERE case_id = $1::uuid AND key = $2 + FOR UPDATE + ), + upserted AS ( + INSERT INTO app.pinned_fact ( + case_id, key, value, fact_type, status, source_turn, + confidence, updated_session_no, visible_to, updated_at + ) + VALUES ( + $1::uuid, $2, $3, $4, $5, $6::uuid, + $7, $8, $9::text[], now() + ) + ON CONFLICT (case_id, key) DO UPDATE SET + value = EXCLUDED.value, + fact_type = EXCLUDED.fact_type, + status = EXCLUDED.status, + source_turn = COALESCE(EXCLUDED.source_turn, app.pinned_fact.source_turn), + confidence = GREATEST(app.pinned_fact.confidence, EXCLUDED.confidence), + version = CASE + WHEN app.pinned_fact.value IS DISTINCT FROM EXCLUDED.value + THEN app.pinned_fact.version + 1 + ELSE app.pinned_fact.version + END, + updated_session_no = EXCLUDED.updated_session_no, + visible_to = EXCLUDED.visible_to, + updated_at = now() + WHERE app.pinned_fact.status <> 'locked' + RETURNING + app.pinned_fact.id, + (SELECT value FROM existing) AS old_value, + app.pinned_fact.value AS new_value + ) + SELECT id, old_value, new_value FROM upserted + """, + sess.case_id, + fact.key, + fact.value, + fact.fact_type, + fact.status, + fact.source_turn_id, + fact.confidence, + sess.session_no, + ["client", "evaluator"], + ) + if not row: + continue + old_value = row["old_value"] + new_value = row["new_value"] + if old_value is not None and old_value == new_value: + continue + reason = "progression" if old_value is None else "clarification" + await _insert_pinned_fact_history( + conn, + fact_id=row["id"], + case_id=sess.case_id, + old_value=old_value, + new_value=new_value, + reason=reason, + session_no=sess.session_no, + turn_id=fact.source_turn_id, + ) + async def end_session(sess: InProcSession, carry: memory.CarryOver) -> bool: try: get_pool() - digest = ( - f"회기 축어록 {len(sess.turns)}개가 저장되었습니다. 정밀 리뷰는 생성 대기 중입니다." - if sess.turns - else "실제 발화가 없어 요약을 생성하지 않았습니다." + digest = memory.build_fallback_session_digest( + session_no=sess.session_no, + masked_turns=sess.masked_turns(visible_to="client"), + end_state=carry.end_state, ) async with acquire(role="learner", user_id=sess.learner_id) as conn: await conn.execute( @@ -2039,6 +2204,51 @@ async def end_session(sess: InProcSession, carry: memory.CarryOver) -> bool: digest, list(carry.compression_job.open_threads if carry.compression_job else []), ) + case_row = await conn.fetchrow( + """ + SELECT case_digest, rapport_trajectory, alliance_level + FROM app.case_profile + WHERE case_id = $1::uuid + AND learner_id = $2::uuid + """, + sess.case_id, + sess.learner_id, + ) + if case_row is not None: + trajectory_point = memory.rapport_trajectory_point( + session_no=sess.session_no, + end_state=carry.end_state, + ) + case_digest = memory.merge_case_digest( + existing_digest=case_row["case_digest"], + session_no=sess.session_no, + session_digest=digest, + ) + rapport_trajectory = memory.merge_rapport_trajectory( + case_row["rapport_trajectory"], + trajectory_point, + ) + alliance_level = memory.update_alliance_level( + case_row["alliance_level"], + trajectory_point.get("end_rapport"), + ) + await conn.execute( + """ + UPDATE app.case_profile + SET case_digest = $3, + rapport_trajectory = $4::jsonb, + alliance_level = $5, + updated_at = now() + WHERE case_id = $1::uuid + AND learner_id = $2::uuid + """, + sess.case_id, + sess.learner_id, + case_digest, + rapport_trajectory, + alliance_level, + ) + await _upsert_pinned_fact_candidates(conn, sess) return True except Exception: require_runtime_fallback_allowed("session end") @@ -2110,7 +2320,7 @@ async def list_sessions( """ SELECT id, seq, speaker, stage, text, text_masked, created_at, llm_provider, model, tokens_in, tokens_out, cost_usd, - audio_ref, silence_ms, speech_rate, barge_in, visible_to + audio_ref, silence_ms, speech_rate, barge_in, provider_events, visible_to FROM app.turns WHERE session_id = $1::uuid ORDER BY seq diff --git a/apps/api/app/session_read_model.py b/apps/api/app/session_read_model.py new file mode 100644 index 0000000..b827f32 --- /dev/null +++ b/apps/api/app/session_read_model.py @@ -0,0 +1,1369 @@ +"""Browser-facing session read models and deterministic builders. + +Routes keep ownership of auth, RLS-backed DB reads, and persistence. This module +owns the response DTOs plus pure projection logic so a future Node.js read API +has one contract surface to mirror. +""" + +from __future__ import annotations + +import re +from collections import Counter +from dataclasses import dataclass +from datetime import datetime +from typing import Literal, Optional + +from pydantic import BaseModel, Field + +from . import turn_runtime +from .config import settings +from .services import session_metrics +from .store import InProcSession, TurnRecord + +StageLabel = Literal["라포", "탐색", "개입", "정리"] + +LEARNER_VISIBLE_AI_ROLE = "counselor" +_PHASE_KEY_BY_LABEL = { + "라포": "rapport", + "탐색": "explore", + "개입": "intervene", + "정리": "closing", +} + + +class LearnerSessionSummary(BaseModel): + session_id: str + persona_code: str + persona_name: str + session_no: int + status: Literal["active", "ended"] + stage: StageLabel + turn_count: int + learner_turn_count: int + client_turn_count: int + started_at: str + ended_at: str | None = None + review_ready: bool = False + archived: bool = False + archived_at: str | None = None + + +class LearnerSessionsResponse(BaseModel): + source: str = "runtime" + sessions: list[LearnerSessionSummary] = Field(default_factory=list) + + +class LearnerDashboardOverview(BaseModel): + total_sessions: int = 0 + completed_sessions: int = 0 + active_sessions: int = 0 + review_ready_sessions: int = 0 + archived_sessions: int = 0 + learner_turns: int = 0 + client_turns: int = 0 + last_practiced_at: str | None = None + + +class LearnerDashboardGrowthPoint(BaseModel): + session_id: str + session_no: int + persona_code: str + stage: StageLabel + started_at: str + ended_at: str | None = None + score: float | None = None + rapport: float | None = None + technique_count: int = 0 + watch_count: int = 0 + + +class LearnerDashboardGrowth(BaseModel): + first_score: float | None = None + latest_score: float | None = None + score_delta: float | None = None + avg_score: float | None = None + avg_rapport: float | None = None + trend: str = "insufficient" + evaluated_sessions: int = 0 + top_techniques: list[str] = Field(default_factory=list) + points: list[LearnerDashboardGrowthPoint] = Field(default_factory=list) + + +class LearnerDashboardPersonaProgress(BaseModel): + persona_code: str + persona_name: str + sessions: int = 0 + completed_sessions: int = 0 + active_sessions: int = 0 + review_ready_sessions: int = 0 + latest_at: str | None = None + latest_stage: StageLabel | None = None + latest_score: float | None = None + trend: str = "insufficient" + + +class LearnerDashboardAchievement(BaseModel): + id: str + label: str + state: Literal["done", "available", "locked"] = "locked" + detail: str + + +class LearnerDashboardFeedbackItem(BaseModel): + session_id: str + persona_code: str + persona_name: str + session_no: int + stage: StageLabel + turn_seq: int + created_at: str + score: float | None = None + rapport: float | None = None + note: str + techniques: list[str] = Field(default_factory=list) + + +class LearnerDashboardResponse(BaseModel): + source: str = "runtime" + overview: LearnerDashboardOverview + growth: LearnerDashboardGrowth + persona_progress: list[LearnerDashboardPersonaProgress] = Field(default_factory=list) + achievements: list[LearnerDashboardAchievement] = Field(default_factory=list) + recent_feedback: list[LearnerDashboardFeedbackItem] = Field(default_factory=list) + message: str + + +class SessionArchiveResponse(BaseModel): + session_id: str + archived: bool + archived_at: str | None = None + source: str = "runtime" + session: LearnerSessionSummary + + +class SessionDetailTurn(BaseModel): + turn_seq: int + speaker: Literal["learner", "client"] + stage: StageLabel + 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: StageLabel + 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", "paralinguistic", "prosody", "audio_quality"] + 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", "saved_by_learner"] = "empty" + generatedBy: str = "rule-based transcript extractor" + sections: list[ReviewWorksheetSection] = Field(default_factory=list) + limitations: list[str] = Field(default_factory=list) + savedAt: Optional[str] = None + + +class ReviewCaseWorksheetSaveRequest(BaseModel): + sections: list[ReviewWorksheetSection] = Field(default_factory=list) + limitations: list[str] = Field(default_factory=list) + + +class SessionTeacherReviewStatus(BaseModel): + status: Literal["pending", "viewed", "closed"] = "pending" + note: str = "" + reviewerId: str | None = None + reviewedAt: str | None = None + updatedAt: str | None = None + + +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 + teacherReview: SessionTeacherReviewStatus | None = None + + +class SessionShareResponse(BaseModel): + shareUrl: str + title: str + description: str + imageUrl: str + createdAt: str + + +class SessionShareDeleteResponse(BaseModel): + revoked: bool + + +@dataclass(frozen=True) +class SessionReviewReadInput: + session: InProcSession + evaluation_record: dict[str, object] | None = None + evaluation_durable: bool = False + saved_worksheet_payload: dict[str, object] | None = None + include_teacher_review: bool = False + teacher_review_record: dict[str, object] | None = None + now_ts: float | None = None + + +def stage_label(stage: object) -> str: + return turn_runtime.stage_label(stage) + + +def iso(ts: float | None) -> str | None: + if ts is None: + return None + return datetime.fromtimestamp(ts).isoformat(timespec="seconds") + + +def learner_visible_turns(sess: InProcSession) -> list[TurnRecord]: + return sess.turns_visible_to(LEARNER_VISIBLE_AI_ROLE) + + +def learner_summary( + sess: InProcSession, + *, + review_ready: bool = False, + archived: bool = False, + archived_at: str | None = None, +) -> 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, + archived=archived, + archived_at=archived_at, + ) + + +def dashboard_overview( + sessions: list[InProcSession], + *, + visible_review_ready: dict[str, bool], + archived_sessions: int, +) -> LearnerDashboardOverview: + return LearnerDashboardOverview( + total_sessions=len(sessions), + completed_sessions=sum(1 for sess in sessions if sess.ended), + active_sessions=sum(1 for sess in sessions if not sess.ended), + review_ready_sessions=sum(1 for ready in visible_review_ready.values() if ready), + archived_sessions=archived_sessions, + learner_turns=sum( + 1 + for sess in sessions + for turn in learner_visible_turns(sess) + if turn.speaker == "counselor" + ), + client_turns=sum( + 1 + for sess in sessions + for turn in learner_visible_turns(sess) + if turn.speaker == "client" + ), + last_practiced_at=session_metrics.iso_datetime( + max((session_metrics.session_activity_time(sess) for sess in sessions), default=0.0) + ) + if sessions + else None, + ) + + +def _dashboard_growth_point(point: session_metrics.SessionGrowthPoint) -> LearnerDashboardGrowthPoint: + return LearnerDashboardGrowthPoint( + session_id=point.session_id, + session_no=point.session_no, + persona_code=point.persona_code, + stage=stage_label(point.stage), + started_at=point.started_at, + ended_at=point.ended_at, + score=point.score, + rapport=point.rapport, + technique_count=point.technique_count, + watch_count=point.watch_count, + ) + + +def dashboard_growth(sessions: list[InProcSession]) -> LearnerDashboardGrowth: + metrics = session_metrics.build_learner_growth( + sessions, + learner_label=lambda _learner_id: "나", + limit=1, + ) + if not metrics: + return LearnerDashboardGrowth() + item = metrics[0] + points = [_dashboard_growth_point(point) for point in item.points] + return LearnerDashboardGrowth( + first_score=item.first_score, + latest_score=item.latest_score, + score_delta=item.score_delta, + avg_score=item.avg_score, + avg_rapport=item.avg_rapport, + trend=item.trend, + evaluated_sessions=sum(1 for point in points if point.score is not None), + top_techniques=item.top_techniques, + points=points, + ) + + +def dashboard_persona_progress( + sessions: list[InProcSession], + review_ready: dict[str, bool], +) -> list[LearnerDashboardPersonaProgress]: + grouped: dict[str, list[InProcSession]] = {} + for sess in sessions: + grouped.setdefault(sess.persona_code, []).append(sess) + + rows: list[LearnerDashboardPersonaProgress] = [] + for persona_code, items in grouped.items(): + ordered = sorted(items, key=session_metrics.session_activity_time) + latest = ordered[-1] + metrics = session_metrics.build_learner_growth( + ordered, + learner_label=lambda _learner_id: "나", + limit=1, + ) + growth = metrics[0] if metrics else None + rows.append( + LearnerDashboardPersonaProgress( + persona_code=persona_code, + persona_name=latest.persona.display_name, + sessions=len(ordered), + completed_sessions=sum(1 for sess in ordered if sess.ended), + active_sessions=sum(1 for sess in ordered if not sess.ended), + review_ready_sessions=sum( + 1 for sess in ordered if review_ready.get(sess.session_id, False) + ), + latest_at=session_metrics.iso_datetime( + session_metrics.session_activity_time(latest) + ), + latest_stage=stage_label(latest.state.stage), + latest_score=growth.latest_score if growth else None, + trend=growth.trend if growth else "insufficient", + ) + ) + return sorted( + rows, + key=lambda row: row.latest_at or "", + reverse=True, + ) + + +def _achievement_state(done: bool, available: bool) -> Literal["done", "available", "locked"]: + if done: + return "done" + if available: + return "available" + return "locked" + + +def dashboard_achievements( + sessions: list[InProcSession], + review_ready: dict[str, bool], +) -> list[LearnerDashboardAchievement]: + completed = sum(1 for sess in sessions if sess.ended) + active = sum(1 for sess in sessions if not sess.ended) + review_count = sum(1 for ready in review_ready.values() if ready) + by_persona = Counter(sess.persona_code for sess in sessions) + max_persona_sessions = max(by_persona.values(), default=0) + persona_coverage = len(by_persona) + return [ + LearnerDashboardAchievement( + id="first_session_complete", + label="첫 회기 완료", + state=_achievement_state(completed >= 1, active >= 1), + detail="한 회기를 종료하면 리뷰와 워크시트 흐름이 열립니다.", + ), + LearnerDashboardAchievement( + id="review_ready", + label="리뷰 확인 가능", + state=_achievement_state(review_count >= 1, completed >= 1), + detail=f"현재 리뷰 가능한 회기 {review_count}건입니다.", + ), + LearnerDashboardAchievement( + id="persona_repeat", + label="같은 내담자 반복 연습", + state=_achievement_state(max_persona_sessions >= 3, max_persona_sessions >= 1), + detail="같은 페르소나를 반복하면 변화 추이를 더 안정적으로 볼 수 있습니다.", + ), + LearnerDashboardAchievement( + id="persona_coverage", + label="여러 페르소나 경험", + state=_achievement_state(persona_coverage >= 3, persona_coverage >= 2), + detail=f"현재 {persona_coverage}개 페르소나에서 연습 기록이 있습니다.", + ), + ] + + +def dashboard_feedback(sessions: list[InProcSession]) -> list[LearnerDashboardFeedbackItem]: + items: list[LearnerDashboardFeedbackItem] = [] + for item in session_metrics.recent_feedback_notes(sessions, limit=5): + score = item.get("score") + rapport = item.get("rapport") + items.append( + LearnerDashboardFeedbackItem( + session_id=str(item["session_id"]), + persona_code=str(item["persona_code"]), + persona_name=str(item["persona_name"]), + session_no=int(item["session_no"]), + stage=stage_label(item["stage"]), + turn_seq=int(item["turn_seq"]), + created_at=str(item["created_at"]), + score=float(score) if isinstance(score, (int, float)) else None, + rapport=float(rapport) if isinstance(rapport, (int, float)) else None, + note=str(item["note"]), + techniques=[str(label) for label in item.get("techniques", [])], + ) + ) + return items + + +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=stage_label(turn.stage), + text=turn.text_masked, + created_at=iso(turn.created_at) or "", + ) + for turn in turns + ], + review_ready=review_ready, + ) + + +def _offset_label(seconds: float) -> str: + whole = max(0, int(round(seconds))) + minutes, sec = divmod(whole, 60) + return f"{minutes}:{sec:02d}" + + +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 saved_case_worksheet_from_payload(payload: dict[str, object] | None) -> ReviewCaseWorksheet | None: + if not payload: + return None + try: + worksheet = ReviewCaseWorksheet.model_validate(payload) + except Exception: + return None + return worksheet.model_copy(update={"status": "saved_by_learner"}) + + +def _worksheet_share_highlights(worksheet: ReviewCaseWorksheet, *, limit: int = 4) -> list[dict[str, str]]: + highlights: list[dict[str, str]] = [] + for section in worksheet.sections: + for item in section.items: + value = _compact_text(item.value or "") + if not value: + continue + highlights.append( + { + "section": section.title, + "label": item.label, + "value": "비공개 요약 항목", + } + ) + if len(highlights) >= limit: + return highlights + return highlights + + +def _review_point_titles(points: list[ReviewPoint], *, limit: int = 3) -> list[str]: + return [_clip_text(point.title or point.body, 72) for point in points[:limit] if (point.title or point.body)] + + +def _share_image_url() -> str: + return f"{settings.frontend_base_url.rstrip('/')}/design-elements/clinical-paper-ambient.png" + + +def session_share_payload(review: SessionReviewResponse) -> dict[str, object]: + title = f"Vignette 회기 리뷰 · {review.client.name} {review.date}" + description = _clip_text(review.summary, 156) + return { + "version": 1, + "title": title, + "description": description, + "summary": _clip_text(review.summary, 420), + "clientName": review.client.name, + "persona": review.client.persona, + "date": review.date, + "durationLabel": review.durationLabel, + "reachedPhase": review.reachedPhase, + "sessionSignal": review.sessionSignal, + "reviewReady": review.reviewReady, + "goodMoments": _review_point_titles(review.goodMoments), + "growthPoints": _review_point_titles(review.growthPoints), + "worksheetHighlights": _worksheet_share_highlights(review.caseWorksheet), + "imageUrl": _share_image_url(), + "appUrl": settings.frontend_base_url.rstrip("/"), + "privacy": "공유 카드에는 회기 원문 축어록과 학습자 식별 정보를 포함하지 않습니다.", + } + + +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 {} + + +_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]: + 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_body_markdown(text: str) -> str: + body = text.strip() + body = re.sub( + r"`?\beffective[_\s-]?openness\b`?(?!\(유효 개방도\))", + "`effective openness(유효 개방도)`", + body, + flags=re.IGNORECASE, + ) + body = re.sub( + r"(? str | None: + clean = " ".join(str(text or "").split()) + if not clean: + return None + sentences = [part.strip() for part in re.split(r"(?<=[.!?。!?])\s+", clean) if part.strip()] + for sentence in sentences: + if "가장 큰 마음" in sentence: + return sentence if len(sentence) <= limit else f"{sentence[: limit - 3].rstrip()}..." + for sentence in sentences: + if "?" in sentence: + return sentence if len(sentence) <= limit else f"{sentence[: limit - 3].rstrip()}..." + return clean if len(clean) <= limit else f"{clean[: limit - 3].rstrip()}..." + + +def _review_note_from_turn_eval( + ev: dict[str, object] | None, + learner_text: str | None = None, +) -> Optional[ReviewNote]: + if not isinstance(ev, dict): + return None + quote = _review_quote_excerpt(learner_text) + 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_parts = ( + f"- 권장: {expected}" if expected else "", + f"- 실제: {actual}" if actual else "", + ) + body = "\n".join(p for p in body_parts if p) + return ReviewNote( + author="평가 AI", + tone="warn", + title=f"의도와 다른 부분 · {dimension}".rstrip(" ·") or "의도와 다른 부분", + body=_review_note_body_markdown(body or "권장 반응과 실제 반응에 차이가 있었어요."), + quote=quote, + ) + 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=_review_note_body_markdown( + note_text + or ( + "타당화·공감·탐색이 회기 흐름에 맞았습니다.\n\n" + "`effective openness(유효 개방도)`가 낮은 내담자라면 다음 질문은 " + "더 작고 구체적인 선택지로 낮춰도 좋습니다." + ) + ), + quote=quote, + ) + if appropriateness == "warn" and note_text: + return ReviewNote( + author="평가 AI", + tone="warn", + title="점검해볼 지점", + body=_review_note_body_markdown(note_text), + quote=quote, + ) + 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}초" + + +_PROVIDER_REVIEW_EVENT_LABELS: dict[str, tuple[str, str, str]] = { + "sigh": ("paralinguistic", "음성 단서", "한숨 감지"), + "cry": ("paralinguistic", "음성 단서", "울음 감지"), + "laugh": ("paralinguistic", "음성 단서", "웃음 감지"), + "breath": ("paralinguistic", "음성 단서", "호흡 변화"), + "pitch": ("prosody", "운율", "피치 변화"), + "intonation": ("prosody", "운율", "억양 변화"), + "prosody": ("prosody", "운율", "운율 변화"), + "background_noise": ("audio_quality", "오디오 품질", "배경 소음"), +} + + +def _provider_confidence_label(event: dict[str, object]) -> str: + raw = event.get("confidence", event.get("score")) + if not isinstance(raw, (int, float)): + return "" + value = float(raw) + if 0 <= value <= 1: + return f"신뢰도 {value * 100:.0f}%" + if 1 < value <= 100: + return f"신뢰도 {value:.0f}%" + return "" + + +def _provider_duration_label(event: dict[str, object]) -> str: + raw = event.get("duration_ms") + if not isinstance(raw, (int, float)): + return "" + milliseconds = int(raw) + if milliseconds <= 0: + return "" + return _seconds_label(milliseconds) + + +def _review_provider_event(event: dict[str, object]) -> ReviewNonverbalEvent | None: + event_type = str(event.get("event_type") or "").strip() + if not event_type: + return None + if event_type == "barge_in": + return ReviewNonverbalEvent(kind="barge_in", label="끼어듦", detail="provider 감지") + if event_type == "silence": + detail = _provider_duration_label(event) or "provider 감지" + return ReviewNonverbalEvent(kind="silence", label="침묵", detail=detail) + if event_type == "speech_rate": + return ReviewNonverbalEvent(kind="pace", label="발화 속도", detail="provider 감지") + + mapped = _PROVIDER_REVIEW_EVENT_LABELS.get(event_type) + if mapped is None: + return None + kind, label, detail = mapped + extras = [item for item in (_provider_confidence_label(event), _provider_duration_label(event)) if item] + if extras: + detail = f"{detail} · {' · '.join(extras)}" + return ReviewNonverbalEvent(kind=kind, label=label, detail=detail) + + +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="오디오 메타 저장됨", + ) + ) + for event in turn.provider_events: + if not isinstance(event, dict): + continue + review_event = _review_provider_event(event) + if review_event is not None: + events.append(review_event) + return events + + +def build_session_review(read_input: SessionReviewReadInput) -> SessionReviewResponse: + sess = read_input.session + visible_turns = learner_visible_turns(sess) + hidden_turns = len(visible_turns) != len(sess.turns) + + end_ts = sess.ended_at or read_input.now_ts 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 = [stage_label(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 = read_input.evaluation_record + 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" + ) + 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, turn.text_masked), + ) + ) + + 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 read_input.evaluation_durable: + summary += " 현재 평가는 런타임 캐시에서 복원되었습니다." + + generated_worksheet = case_worksheet_from_turns(turns) + case_worksheet = ( + saved_case_worksheet_from_payload(read_input.saved_worksheet_payload) + or generated_worksheet + ) + + teacher_review: SessionTeacherReviewStatus | None = None + if read_input.include_teacher_review: + review_status = read_input.teacher_review_record or {} + review_status_value = str(review_status.get("status") or "pending") + if review_status_value not in {"viewed", "closed"}: + review_status_value = "pending" + teacher_review = SessionTeacherReviewStatus( + status=review_status_value, # type: ignore[arg-type] + note=str(review_status.get("note") or ""), + reviewerId=str(review_status.get("reviewer_id") or "") or None, + reviewedAt=str(review_status.get("reviewed_at") or "") or None, + updatedAt=str(review_status.get("updated_at") or "") or None, + ) + + return SessionReviewResponse( + session_id=sess.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, + nextLine=next_line, + clientFeedback=client_feedback, + audioUrl=None, + pdfExportUrl=None, + degraded=review_degraded, + reviewReady=evaluation_ready, + teacherReview=teacher_review, + ) diff --git a/apps/api/app/store.py b/apps/api/app/store.py index 519658a..2bca7d4 100644 --- a/apps/api/app/store.py +++ b/apps/api/app/store.py @@ -42,6 +42,7 @@ class TurnRecord: silence_ms: int | None = None speech_rate: float | None = None barge_in: bool | None = None + provider_events: list[dict[str, object]] = field(default_factory=list) # fast-loop 턴 평가(TurnEvaluation.to_hook_dict). 학습자(상담자) 발화에 부착. evaluation: Optional[dict] = None visible_to: tuple[str, ...] = DEFAULT_TURN_VISIBLE_TO diff --git a/apps/api/app/test_session_memory.py b/apps/api/app/test_session_memory.py new file mode 100644 index 0000000..4345d40 --- /dev/null +++ b/apps/api/app/test_session_memory.py @@ -0,0 +1,573 @@ +from __future__ import annotations + +import unittest +from unittest.mock import patch + +from . import session_persistence +from .routes import sessions +from .services import memory, rag, state_machine +from .services.persona import P1 +from .store import InProcSession, TurnRecord + + +class SessionMemoryPureTest(unittest.TestCase): + def test_case_digest_merge_is_idempotent_by_session_no(self) -> None: + digest = memory.merge_case_digest( + existing_digest="S1: 이전 회기 요약\nS2: 오래된 요약", + session_no=2, + session_digest="S2: 새 요약", + ) + + self.assertEqual(digest, "S1: 이전 회기 요약\nS2: 새 요약") + + def test_rapport_trajectory_merge_replaces_same_session(self) -> None: + merged = memory.merge_rapport_trajectory( + [{"session_no": 1, "end_rapport": 0.2}, {"session_no": 2, "end_rapport": 0.3}], + {"session_no": 2, "end_rapport": 0.7, "end_openness": 0.5}, + ) + + self.assertEqual(len(merged), 2) + self.assertEqual(merged[-1]["session_no"], 2) + self.assertEqual(merged[-1]["end_rapport"], 0.7) + + def test_fallback_session_digest_uses_masked_client_visible_text(self) -> None: + digest = memory.build_fallback_session_digest( + session_no=3, + masked_turns=[ + {"speaker": "counselor", "text": "그때 마음이 어땠나요?"}, + {"speaker": "client", "text": "저는 [NAME]이고 [ORG]에 다녀요."}, + ], + end_state={"stage": "탐색", "effective_openness": 0.42, "rapport_credit": 0.31}, + ) + + self.assertIn("S3:", digest) + self.assertIn("[NAME]", digest) + self.assertIn("[ORG]", digest) + self.assertNotIn("김서연", digest) + + def test_extract_pinned_fact_candidates_is_conservative_and_masked(self) -> None: + facts = memory.extract_pinned_fact_candidates( + [ + {"speaker": "counselor", "text": "이름을 말해줄 수 있나요?"}, + { + "speaker": "client", + "text": "저는 [NAME]이고 [ORG]에 다녀요. 동생과 자주 싸워요. 주 1회 상담 약속은 지키고 싶어요.", + "turn_id": "00000000-0000-0000-0000-000000000201", + }, + ] + ) + + by_key = {fact.key: fact for fact in facts} + self.assertEqual(by_key["identity:name"].value, "[NAME]") + self.assertEqual(by_key["identity:org"].value, "[ORG]") + self.assertEqual(by_key["agreement:counseling"].fact_type, "agreement") + self.assertNotIn("relationship:sibling", by_key) + self.assertNotIn("김서연", " ".join(fact.value for fact in facts)) + + def test_extract_pinned_fact_candidates_skips_inferred_or_transient_content(self) -> None: + facts = memory.extract_pinned_fact_candidates( + [ + {"speaker": "client", "text": "오늘은 그냥 기분이 좀 나빴어요."}, + {"speaker": "client", "text": "동생과 자주 싸워요."}, + {"speaker": "client", "text": "죽고 싶다는 생각이 스쳐갔어요."}, + ] + ) + + self.assertEqual(facts, []) + + def test_extract_pinned_fact_candidates_marks_explicit_agreement_withdrawal_only(self) -> None: + facts = memory.extract_pinned_fact_candidates( + [ + { + "speaker": "client", + "text": "상담 약속은 이제 못 지키겠어요. 회기는 그만하고 싶어요.", + "turn_id": "00000000-0000-0000-0000-000000000203", + }, + {"speaker": "client", "text": "동생과 자주 싸워요."}, + ] + ) + + by_key = {fact.key: fact for fact in facts} + self.assertEqual(list(by_key), ["agreement:counseling"]) + self.assertEqual(by_key["agreement:counseling"].status, "contradicted") + self.assertEqual(by_key["agreement:counseling"].fact_type, "agreement") + self.assertIn("못 지키겠어요", by_key["agreement:counseling"].value) + self.assertEqual( + by_key["agreement:counseling"].source_turn_id, + "00000000-0000-0000-0000-000000000203", + ) + + def test_episodic_turn_inputs_use_masked_client_visible_client_turns_only(self) -> None: + inputs = rag.episodic_turn_inputs_from_records( + session_id="00000000-0000-0000-0000-00000000feed", + case_id="00000000-0000-0000-0000-00000000ca5e", + turns=[ + TurnRecord( + turn_seq=1, + speaker="client", + stage="라포", + text="저는 김서연이고 한신대학교에 다녀요.", + text_masked="저는 [NAME]이고 [ORG]에 다녀요.", + turn_id="00000000-0000-0000-0000-000000000201", + ), + TurnRecord( + turn_seq=2, + speaker="counselor", + stage="라포", + text="상담자 발화", + text_masked="상담자 발화", + turn_id="00000000-0000-0000-0000-000000000202", + ), + TurnRecord( + turn_seq=3, + speaker="client", + stage="라포", + text="평가자만 볼 발화", + text_masked="평가자만 볼 발화", + turn_id="00000000-0000-0000-0000-000000000203", + visible_to=("evaluator",), + ), + TurnRecord( + turn_seq=4, + speaker="client", + stage="라포", + text="DB turn id 없음", + text_masked="DB turn id 없음", + ), + ], + ) + + self.assertEqual(len(inputs), 1) + self.assertEqual(inputs[0].turn_id, "00000000-0000-0000-0000-000000000201") + self.assertEqual(inputs[0].seq, 1) + self.assertEqual(inputs[0].text_masked, "저는 [NAME]이고 [ORG]에 다녀요.") + self.assertNotIn("김서연", inputs[0].text_masked) + self.assertNotIn("한신대학교", inputs[0].text_masked) + + +class SessionMemoryPersistenceTest(unittest.IsolatedAsyncioTestCase): + async def test_write_persona_turn_embeddings_is_masked_and_idempotent(self) -> None: + class FakeConn: + def __init__(self) -> None: + self.executed: list[tuple[str, tuple[object, ...]]] = [] + + async def execute(self, query: str, *args: object) -> str: + self.executed.append((query, args)) + return "INSERT 0 1" + + conn = FakeConn() + turn = rag.EpisodicTurnInput( + turn_id="00000000-0000-0000-0000-000000000201", + case_id="00000000-0000-0000-0000-00000000ca5e", + session_id="00000000-0000-0000-0000-00000000feed", + seq=7, + text_masked="저는 [NAME]이고 [ORG]에 다녀요.", + ) + captured_texts: list[str] = [] + + def fake_embed_query(text: str) -> rag.EmbeddedQuery: + captured_texts.append(text) + return rag.EmbeddedQuery(dense=[0.1] * rag.EMBED_DIM, sparse={"42": 0.7}) + + with patch.object(rag, "embed_query", fake_embed_query): + result = await rag.write_persona_turn_embeddings(conn, turns=[turn]) + + self.assertEqual(result.inserted, 1) + self.assertEqual(captured_texts, ["저는 [NAME]이고 [ORG]에 다녀요."]) + self.assertEqual(len(conn.executed), 1) + query, args = conn.executed[0] + self.assertIn("INSERT INTO app.turn_embedding", query) + self.assertIn("ON CONFLICT (turn_id) DO NOTHING", query) + self.assertEqual(args[0], turn.turn_id) + self.assertEqual(args[1], turn.case_id) + self.assertEqual(args[2], turn.session_id) + self.assertEqual(args[3], 7) + self.assertIn("[0.1,0.1", str(args[4])) + self.assertEqual(args[5], '{"42": 0.7}') + self.assertNotIn("김서연", " ".join(str(arg) for arg in args)) + + async def test_end_persisted_session_schedules_episodic_embedding_writer(self) -> None: + scheduled: list[object] = [] + + def fake_create_task(coro): + scheduled.append(coro) + coro.close() + return None + + sess = InProcSession( + session_id="00000000-0000-0000-0000-00000000feed", + case_id="00000000-0000-0000-0000-00000000ca5e", + learner_id="00000000-0000-0000-0000-000000000101", + persona_code=P1.code, + theory_mode="humanistic", + persona=P1, + state=state_machine.SessionState(), + session_no=1, + ) + carry = memory.CarryOver( + end_state={}, + rapport_delta=0.0, + compression_job=None, + ) + + with patch.object(session_persistence, "end_session", return_value=True), patch.object( + sessions.asyncio, + "create_task", + fake_create_task, + ): + await sessions._end_persisted_session(sess, carry) + + self.assertTrue(sess.ended) + self.assertEqual(len(scheduled), 1) + + async def test_seed_recall_loads_case_digest_and_client_visible_pinned_facts(self) -> None: + test_case = self + + class FakeConn: + def __init__(self) -> None: + self.fetch_queries: list[tuple[str, tuple[object, ...]]] = [] + + async def fetchrow(self, query: str, *args: object): + self.fetch_queries.append((query, args)) + if "FROM app.case_profile" in query: + return {"case_digest": "S1: 케이스 큰그림"} + if "FROM app.session_summary" in query: + return { + "digest": "직전 회기 요약", + "open_threads": ["가족 이야기를 이어가기"], + "end_state": {"rapport_credit": 0.5}, + } + return None + + async def fetch(self, query: str, *args: object): + self.fetch_queries.append((query, args)) + test_case.assertIn("FROM app.pinned_fact", query) + test_case.assertIn("status IN ('stable', 'evolving', 'locked')", query) + test_case.assertIn("$2 = ANY(visible_to)", query) + return [{"value": "동생과의 갈등"}, {"value": "주 1회 상담 약속"}] + + class FakeAcquire: + def __init__(self, conn: FakeConn) -> None: + self.conn = conn + + async def __aenter__(self) -> FakeConn: + return self.conn + + async def __aexit__(self, exc_type: object, exc: object, tb: object) -> None: + return None + + conn = FakeConn() + with patch.object(sessions.db, "get_pool", return_value=object()), patch.object( + sessions.db, + "acquire", + return_value=FakeAcquire(conn), + ): + recall = await sessions._build_seed_recall( + case_id="00000000-0000-0000-0000-00000000ca5e" + ) + + self.assertIn("[케이스 큰그림]", recall.recall_summary or "") + self.assertIn("S1: 케이스 큰그림", recall.recall_summary or "") + self.assertIn("직전 회기 요약", recall.recall_summary or "") + self.assertEqual(recall.pinned_facts, ["동생과의 갈등", "주 1회 상담 약속"]) + self.assertEqual(recall.carry, {"rapport_credit": 0.5}) + + async def test_end_session_updates_session_summary_and_case_profile(self) -> None: + test_case = self + + class FakeConn: + def __init__(self) -> None: + self.executed: list[tuple[str, tuple[object, ...]]] = [] + + async def execute(self, query: str, *args: object) -> str: + self.executed.append((query, args)) + return "OK" + + async def fetchrow(self, query: str, *args: object): + self.executed.append((query, args)) + if "INSERT INTO app.pinned_fact" in query: + return { + "id": f"00000000-0000-0000-0000-00000000fa{len(self.executed):02d}", + "old_value": None, + "new_value": args[2], + } + test_case.assertIn("FROM app.case_profile", query) + return { + "case_digest": "S1: 이전 회기", + "rapport_trajectory": [{"session_no": 1, "end_rapport": 0.2}], + "alliance_level": 0.2, + } + + class FakeAcquire: + def __init__(self, conn: FakeConn) -> None: + self.conn = conn + + async def __aenter__(self) -> FakeConn: + return self.conn + + async def __aexit__(self, exc_type: object, exc: object, tb: object) -> None: + return None + + learner_id = "00000000-0000-0000-0000-000000000101" + case_id = "00000000-0000-0000-0000-00000000ca5e" + state = state_machine.SessionState( + stage=state_machine.Stage.EXPLORE, + turn_seq=2, + effective_openness=0.42, + rapport_credit=0.31, + resistance=0.5, + ) + sess = InProcSession( + session_id="00000000-0000-0000-0000-00000000feed", + case_id=case_id, + learner_id=learner_id, + persona_code=P1.code, + theory_mode="humanistic", + persona=P1, + state=state, + session_no=2, + prev_rapport_credit=0.1, + turns=[ + TurnRecord( + turn_seq=1, + speaker="counselor", + stage="라포", + text="실명 질문", + text_masked="이름을 말해줄 수 있나요?", + ), + TurnRecord( + turn_seq=1, + speaker="client", + stage="라포", + text="저는 김서연이고 한신대학교에 다녀요. 동생과 자주 싸워요. 주 1회 상담 약속은 지키고 싶어요.", + text_masked="저는 [NAME]이고 [ORG]에 다녀요. 동생과 자주 싸워요. 주 1회 상담 약속은 지키고 싶어요.", + turn_id="00000000-0000-0000-0000-000000000201", + ), + ], + ) + carry = memory.make_carry_over( + state=state, + session_id=sess.session_id, + case_id=sess.case_id, + session_no=sess.session_no, + masked_turns=sess.masked_turns(), + prev_rapport_credit=sess.prev_rapport_credit, + ) + + conn = FakeConn() + with patch.object(session_persistence, "get_pool", return_value=object()), patch.object( + session_persistence, + "acquire", + return_value=FakeAcquire(conn), + ): + persisted = await session_persistence.end_session(sess, carry) + + self.assertTrue(persisted) + case_updates = [ + args for query, args in conn.executed if "UPDATE app.case_profile" in query + ] + self.assertEqual(len(case_updates), 1) + _, _, case_digest, trajectory, alliance_level = case_updates[0] + self.assertIn("S1: 이전 회기", case_digest) + self.assertIn("S2:", case_digest) + self.assertIn("[NAME]", case_digest) + self.assertNotIn("김서연", case_digest) + self.assertEqual(trajectory[-1]["session_no"], 2) + self.assertEqual(trajectory[-1]["end_rapport"], 0.31) + self.assertGreater(alliance_level, 0.2) + pinned_writes = [ + args + for query, args in conn.executed + if "INSERT INTO app.pinned_fact (" in query + ] + self.assertEqual(len(pinned_writes), 3) + by_key = {args[1]: args for args in pinned_writes} + self.assertEqual(by_key["identity:name"][2], "[NAME]") + self.assertEqual(by_key["identity:org"][2], "[ORG]") + self.assertEqual(by_key["agreement:counseling"][3], "agreement") + self.assertNotIn("relationship:sibling", by_key) + self.assertEqual(by_key["identity:name"][5], "00000000-0000-0000-0000-000000000201") + self.assertEqual(by_key["identity:name"][7], 2) + self.assertEqual(by_key["identity:name"][8], ["client", "evaluator"]) + self.assertNotIn("김서연", " ".join(str(args[2]) for args in pinned_writes)) + history_writes = [ + args for query, args in conn.executed if "INSERT INTO app.pinned_fact_history" in query + ] + self.assertEqual(len(history_writes), 3) + self.assertEqual(history_writes[0][1], case_id) + self.assertIsNone(history_writes[0][2]) + self.assertEqual(history_writes[0][3], "[NAME]") + self.assertEqual(history_writes[0][4], "progression") + self.assertEqual(history_writes[0][5], 2) + self.assertEqual( + history_writes[0][6], + "00000000-0000-0000-0000-000000000201", + ) + self.assertNotIn("김서연", " ".join(str(args[3]) for args in history_writes)) + + async def test_pinned_fact_history_skips_same_value_refresh(self) -> None: + class FakeConn: + def __init__(self) -> None: + self.executed: list[tuple[str, tuple[object, ...]]] = [] + + async def fetchrow(self, query: str, *args: object): + self.executed.append((query, args)) + return { + "id": "00000000-0000-0000-0000-00000000fa11", + "old_value": args[2], + "new_value": args[2], + } + + async def execute(self, query: str, *args: object) -> str: + self.executed.append((query, args)) + return "OK" + + conn = FakeConn() + sess = InProcSession( + session_id="00000000-0000-0000-0000-00000000feed", + case_id="00000000-0000-0000-0000-00000000ca5e", + learner_id="00000000-0000-0000-0000-000000000101", + persona_code=P1.code, + theory_mode="humanistic", + persona=P1, + state=state_machine.SessionState(), + session_no=3, + turns=[ + TurnRecord( + turn_seq=1, + speaker="client", + stage="라포", + text="저는 김서연입니다.", + text_masked="저는 [NAME]입니다.", + turn_id="00000000-0000-0000-0000-000000000201", + ) + ], + ) + + await session_persistence._upsert_pinned_fact_candidates(conn, sess) + + pinned_writes = [ + args + for query, args in conn.executed + if "INSERT INTO app.pinned_fact (" in query + ] + history_writes = [ + args for query, args in conn.executed if "INSERT INTO app.pinned_fact_history" in query + ] + self.assertEqual(len(pinned_writes), 1) + self.assertEqual(history_writes, []) + + async def test_pinned_fact_history_records_value_change_as_clarification(self) -> None: + class FakeConn: + def __init__(self) -> None: + self.executed: list[tuple[str, tuple[object, ...]]] = [] + + async def fetchrow(self, query: str, *args: object): + self.executed.append((query, args)) + return { + "id": "00000000-0000-0000-0000-00000000fa22", + "old_value": "예전에는 격주 상담 약속을 말함.", + "new_value": args[2], + } + + async def execute(self, query: str, *args: object) -> str: + self.executed.append((query, args)) + return "OK" + + conn = FakeConn() + sess = InProcSession( + session_id="00000000-0000-0000-0000-00000000feed", + case_id="00000000-0000-0000-0000-00000000ca5e", + learner_id="00000000-0000-0000-0000-000000000101", + persona_code=P1.code, + theory_mode="humanistic", + persona=P1, + state=state_machine.SessionState(), + session_no=4, + turns=[ + TurnRecord( + turn_seq=1, + speaker="client", + stage="라포", + text="앞으로는 주 1회 상담 약속을 지키고 싶어요.", + text_masked="앞으로는 주 1회 상담 약속을 지키고 싶어요.", + turn_id="00000000-0000-0000-0000-000000000202", + ) + ], + ) + + await session_persistence._upsert_pinned_fact_candidates(conn, sess) + + history_writes = [ + args for query, args in conn.executed if "INSERT INTO app.pinned_fact_history" in query + ] + self.assertEqual(len(history_writes), 1) + self.assertEqual(history_writes[0][2], "예전에는 격주 상담 약속을 말함.") + self.assertIn("주 1회 상담 약속", str(history_writes[0][3])) + self.assertEqual(history_writes[0][4], "clarification") + self.assertEqual(history_writes[0][5], 4) + + async def test_pinned_fact_history_records_explicit_agreement_withdrawal_as_contradiction(self) -> None: + test_case = self + + class FakeConn: + def __init__(self) -> None: + self.executed: list[tuple[str, tuple[object, ...]]] = [] + + async def fetchrow(self, query: str, *args: object): + self.executed.append((query, args)) + test_case.assertIn("UPDATE app.pinned_fact", query) + test_case.assertIn("status = 'contradicted'", query) + return { + "id": "00000000-0000-0000-0000-00000000fa33", + "old_value": "주 1회 상담 약속은 지키고 싶어요.", + "new_value": args[2], + } + + async def execute(self, query: str, *args: object) -> str: + self.executed.append((query, args)) + return "OK" + + conn = FakeConn() + sess = InProcSession( + session_id="00000000-0000-0000-0000-00000000feed", + case_id="00000000-0000-0000-0000-00000000ca5e", + learner_id="00000000-0000-0000-0000-000000000101", + persona_code=P1.code, + theory_mode="humanistic", + persona=P1, + state=state_machine.SessionState(), + session_no=5, + turns=[ + TurnRecord( + turn_seq=1, + speaker="client", + stage="라포", + text="상담 약속은 이제 못 지키겠어요. 회기는 그만하고 싶어요.", + text_masked="상담 약속은 이제 못 지키겠어요. 회기는 그만하고 싶어요.", + turn_id="00000000-0000-0000-0000-000000000203", + ) + ], + ) + + await session_persistence._upsert_pinned_fact_candidates(conn, sess) + + contradiction_updates = [ + args for query, args in conn.executed if "UPDATE app.pinned_fact" in query + ] + self.assertEqual(len(contradiction_updates), 1) + self.assertEqual(contradiction_updates[0][1], "agreement:counseling") + self.assertIn("그만하고 싶어요", str(contradiction_updates[0][2])) + self.assertEqual(contradiction_updates[0][7], ["evaluator"]) + history_writes = [ + args for query, args in conn.executed if "INSERT INTO app.pinned_fact_history" in query + ] + self.assertEqual(len(history_writes), 1) + self.assertEqual(history_writes[0][2], "주 1회 상담 약속은 지키고 싶어요.") + self.assertIn("못 지키겠어요", str(history_writes[0][3])) + self.assertEqual(history_writes[0][4], "contradiction") + self.assertEqual(history_writes[0][5], 5) + + +if __name__ == "__main__": + unittest.main() diff --git a/apps/api/app/test_session_turn_persistence.py b/apps/api/app/test_session_turn_persistence.py index fd008b6..f95687f 100644 --- a/apps/api/app/test_session_turn_persistence.py +++ b/apps/api/app/test_session_turn_persistence.py @@ -7,7 +7,8 @@ import unittest from types import SimpleNamespace from unittest.mock import AsyncMock, patch -from . import turn_runtime +from . import session_persistence, turn_runtime +from .contracts.engine_gateway import EngineGatewaySseLineDecoder from .deps import Principal, Role from .engine_client import EngineError from .routes import sessions @@ -17,6 +18,14 @@ from .services.voice import TTSChunk, TranscriptResult, VoicePreset from .store import InProcSession, TurnRecord, store +async def _decoded_stream_packets(stream_engine, req): + decoder = EngineGatewaySseLineDecoder() + async for raw in stream_engine.stream(req): + packet = decoder.feed_line(raw) + if packet is not None: + yield packet + + def _principal() -> Principal: return Principal( user_id="00000000-0000-0000-0000-000000000101", @@ -67,6 +76,65 @@ class SessionTurnPersistenceTest(unittest.IsolatedAsyncioTestCase): async def asyncTearDown(self) -> None: store._sessions.clear() + async def test_append_turn_writes_provider_events_to_db(self) -> None: + class FakeConn: + def __init__(self) -> None: + self.insert_query = "" + self.insert_args: tuple[object, ...] = () + + async def fetchval(self, query: str, *args: object) -> object: + if "SELECT id FROM app.sessions" in query: + return "turn-persistence-session" + if "COALESCE(MAX(seq)" in query: + return 1 + if "INSERT INTO app.turns" in query: + self.insert_query = query + self.insert_args = args + return "00000000-0000-0000-0000-000000009999" + return None + + class FakeAcquire: + def __init__(self, conn: FakeConn) -> None: + self.conn = conn + + async def __aenter__(self) -> FakeConn: + return self.conn + + async def __aexit__(self, exc_type, exc, tb) -> None: + return None + + conn = FakeConn() + turn = TurnRecord( + turn_seq=1, + speaker="counselor", + stage="rapport", + text="voice text", + text_masked="voice text", + audio_ref="voice:webm:sha256:test", + silence_ms=1234, + speech_rate=210.0, + barge_in=True, + provider_events=[{"type": "sigh", "confidence": 0.82}], + ) + + with patch.object(session_persistence, "get_pool", return_value=object()), patch.object( + session_persistence, + "acquire", + return_value=FakeAcquire(conn), + ): + ok = await session_persistence.append_turn( + session_id="turn-persistence-session", + learner_id=_principal().user_id, + turn=turn, + ) + + self.assertTrue(ok) + self.assertIn("provider_events", conn.insert_query) + self.assertIn("$17::jsonb", conn.insert_query) + self.assertEqual(conn.insert_args[16], [{"type": "sigh", "confidence": 0.82}]) + self.assertEqual(conn.insert_args[17], list(turn.visible_to)) + self.assertEqual(turn.turn_id, "00000000-0000-0000-0000-000000009999") + async def test_generate_turn_engine_failure_does_not_append_learner_turn(self) -> None: principal = _principal() sess = _session(principal) @@ -96,7 +164,7 @@ class SessionTurnPersistenceTest(unittest.IsolatedAsyncioTestCase): turn_seq=ctx.state_after.turn_seq, stage=ctx.state_after.stage.value, effective_openness=ctx.state_after.effective_openness, - client_reply="괜찮아요. 천천히 말해볼게요.", + client_reply="저는 김서연 씨고 한신대학교 상담심리학과 학생이에요.", safety_flagged=False, state_after=ctx.state_after, llm_provider="claude_cli", @@ -113,10 +181,16 @@ class SessionTurnPersistenceTest(unittest.IsolatedAsyncioTestCase): principal, ) - self.assertEqual(response.client_reply, "괜찮아요. 천천히 말해볼게요.") + self.assertEqual(response.client_reply, "저는 김서연 씨고 한신대학교 상담심리학과 학생이에요.") self.assertEqual(len(sess.turns), 2) learner_turn, client_turn = sess.turns self.assertIsNone(learner_turn.llm_provider) + self.assertEqual(client_turn.text, "저는 김서연 씨고 한신대학교 상담심리학과 학생이에요.") + self.assertNotIn("김서연", client_turn.text_masked) + self.assertNotIn("한신대학교", client_turn.text_masked) + self.assertNotIn("상담심리학과", client_turn.text_masked) + self.assertIn("[NAME]", client_turn.text_masked) + self.assertIn("[ORG]", client_turn.text_masked) self.assertEqual(client_turn.llm_provider, "claude_cli") self.assertEqual(client_turn.model, "gateway-default") self.assertEqual(client_turn.tokens_in, 17) @@ -502,6 +576,10 @@ class SessionTurnPersistenceTest(unittest.IsolatedAsyncioTestCase): '"tokens_in":5,"tokens_out":7,"cost_usd":0.034567}' ) + async def stream_packets(self, req): + async for packet in _decoded_stream_packets(self, req): + yield packet + principal = _principal() sess = _session(principal) ctx = orchestrator.prepare_turn( @@ -537,6 +615,10 @@ class SessionTurnPersistenceTest(unittest.IsolatedAsyncioTestCase): yield "event: error" yield 'data: {"detail":"engine unavailable: gateway"}' + async def stream_packets(self, req): + async for packet in _decoded_stream_packets(self, req): + yield packet + principal = _principal() sess = _session(principal) ctx = orchestrator.prepare_turn( @@ -655,7 +737,15 @@ class SessionTurnPersistenceTest(unittest.IsolatedAsyncioTestCase): with patch.object( voice_routes.voice_service, "transcribe", - AsyncMock(return_value=TranscriptResult(text="오늘은 좀 힘들었어요.", duration=2.0)), + AsyncMock( + return_value=TranscriptResult( + text="오늘은 좀 힘들었어요.", + duration=2.0, + provider_events=[ + {"kind": "voice_activity", "start_ms": 10, "raw_text": "drop"}, + ], + ) + ), ), patch.object( voice_routes.orchestrator, "run_turn_generate", @@ -674,6 +764,7 @@ class SessionTurnPersistenceTest(unittest.IsolatedAsyncioTestCase): fmt="webm", silence_ms=1234, barge_in=True, + provider_events=[{"type": "sigh", "confidence": 0.82, "text": "drop"}], ) self.assertEqual(len(sess.turns), 2) @@ -682,7 +773,25 @@ class SessionTurnPersistenceTest(unittest.IsolatedAsyncioTestCase): self.assertEqual(learner_turn.silence_ms, 1234) self.assertGreater(learner_turn.speech_rate or 0, 0) self.assertTrue(learner_turn.barge_in) + self.assertEqual( + learner_turn.provider_events, + [ + { + "type": "sigh", + "confidence": 0.82, + "event_type": "sigh", + "category": "paralinguistic", + }, + { + "kind": "voice_activity", + "start_ms": 10, + "event_type": "voice_activity", + "category": "speech_activity", + }, + ], + ) self.assertIsNone(client_turn.audio_ref) + self.assertEqual(client_turn.provider_events, []) self.assertEqual(client_turn.llm_provider, "claude_cli") self.assertTrue(any(message.get("type") == "tts_end" for message in websocket.messages)) @@ -703,6 +812,26 @@ class SessionTurnPersistenceTest(unittest.IsolatedAsyncioTestCase): silence_ms=1234, speech_rate=420.0, barge_in=True, + provider_events=[ + { + "event_type": "sigh", + "category": "paralinguistic", + "confidence": 0.82, + "provider": "stt-provider", + "type": "raw_sigh", + }, + { + "event_type": "speech_start", + "category": "speech_activity", + "start_ms": 100, + }, + { + "event_type": "background_noise", + "category": "audio_quality", + "score": 77, + "label": "busy cafe", + }, + ], ), TurnRecord( turn_seq=2, @@ -715,6 +844,7 @@ class SessionTurnPersistenceTest(unittest.IsolatedAsyncioTestCase): silence_ms=2500, speech_rate=180.0, barge_in=True, + provider_events=[{"event_type": "cry", "category": "paralinguistic"}], ), ] ) @@ -723,10 +853,20 @@ class SessionTurnPersistenceTest(unittest.IsolatedAsyncioTestCase): self.assertEqual(len(response.turns), 2) learner_turn, client_turn = response.turns - self.assertEqual([event.kind for event in learner_turn.nonverbal], ["silence", "pace", "barge_in", "audio"]) + self.assertEqual( + [event.kind for event in learner_turn.nonverbal], + ["silence", "pace", "barge_in", "audio", "paralinguistic", "audio_quality"], + ) self.assertEqual(learner_turn.nonverbal[0].label, "침묵") self.assertEqual(learner_turn.nonverbal[0].detail, "1.2초") self.assertEqual(learner_turn.nonverbal[1].detail, "분당 420자") + self.assertEqual(learner_turn.nonverbal[4].label, "음성 단서") + self.assertEqual(learner_turn.nonverbal[4].detail, "한숨 감지 · 신뢰도 82%") + exposed_details = " ".join(event.detail for event in learner_turn.nonverbal) + self.assertNotIn("stt-provider", exposed_details) + self.assertNotIn("raw_sigh", exposed_details) + self.assertEqual(learner_turn.nonverbal[5].label, "오디오 품질") + self.assertEqual(learner_turn.nonverbal[5].detail, "배경 소음 · 신뢰도 77%") self.assertEqual(client_turn.nonverbal, []) async def test_review_includes_case_formulation_worksheet_draft(self) -> None: diff --git a/apps/api/app/test_voice_ws.py b/apps/api/app/test_voice_ws.py index a7135af..eb5850e 100644 --- a/apps/api/app/test_voice_ws.py +++ b/apps/api/app/test_voice_ws.py @@ -77,6 +77,50 @@ class VoiceWebSocketContractTest(unittest.IsolatedAsyncioTestCase): {"degraded": False, "persona_catalog_source": "session"}, ) + def test_provider_events_get_internal_taxonomy_without_raw_payload(self) -> None: + events = voice_routes._safe_provider_events( + [ + { + "type": "SIGH", + "confidence": 0.81, + "text": "raw transcript must drop", + }, + { + "kind": "voice_activity", + "start_ms": 10, + "raw_text": "drop", + }, + { + "label": "vendor custom marker", + "score": 0.44, + }, + ] + ) + + self.assertEqual( + events, + [ + { + "type": "SIGH", + "confidence": 0.81, + "event_type": "sigh", + "category": "paralinguistic", + }, + { + "kind": "voice_activity", + "start_ms": 10, + "event_type": "voice_activity", + "category": "speech_activity", + }, + { + "label": "vendor custom marker", + "score": 0.44, + "event_type": "vendor_custom_marker", + "category": "unknown", + }, + ], + ) + async def test_audio_start_binary_chunks_audio_end_ping_close_contract(self) -> None: websocket = FakeWebSocket( [ @@ -90,6 +134,15 @@ class VoiceWebSocketContractTest(unittest.IsolatedAsyncioTestCase): "format": "webm", "silence_ms": "450", "barge_in": "true", + "provider_events": [ + { + "type": "sigh", + "confidence": 0.82, + "text": "raw transcript must not persist", + }, + {"kind": "noise", "label": "x" * 120}, + "invalid", + ], } ), _control({"type": "close"}), @@ -141,6 +194,23 @@ class VoiceWebSocketContractTest(unittest.IsolatedAsyncioTestCase): self.assertEqual(kwargs["audio_ended_at"], 12.0) self.assertEqual(kwargs["silence_ms"], 450) self.assertIs(kwargs["barge_in"], True) + self.assertEqual( + kwargs["provider_events"], + [ + { + "type": "sigh", + "confidence": 0.82, + "event_type": "sigh", + "category": "paralinguistic", + }, + { + "kind": "noise", + "label": "x" * 80, + "event_type": "background_noise", + "category": "audio_quality", + }, + ], + ) async def test_text_turn_strips_text_runs_turn_and_ping_close_still_work(self) -> None: websocket = FakeWebSocket( diff --git a/infra/db/init/02_schema.sql b/infra/db/init/02_schema.sql index 5d38214..af8fd6f 100644 --- a/infra/db/init/02_schema.sql +++ b/infra/db/init/02_schema.sql @@ -221,6 +221,7 @@ CREATE TABLE IF NOT EXISTS app.turns ( silence_ms INT, speech_rate REAL, barge_in BOOLEAN, + provider_events JSONB NOT NULL DEFAULT '[]'::jsonb, created_at TIMESTAMPTZ NOT NULL DEFAULT now(), UNIQUE (session_id, seq) ); diff --git a/infra/db/init/04_audit_eval_rls.sql b/infra/db/init/04_audit_eval_rls.sql index f1e6a94..e2a7f28 100644 --- a/infra/db/init/04_audit_eval_rls.sql +++ b/infra/db/init/04_audit_eval_rls.sql @@ -328,6 +328,31 @@ CREATE POLICY p_pinned_select ON app.pinned_fact FOR SELECT USING ( AND current_setting('app.current_ai_view', true) = ANY(visible_to) ) OR app.current_role_name() IN ('admin','instructor') ); +DROP POLICY IF EXISTS p_pinned_insert ON app.pinned_fact; +CREATE POLICY p_pinned_insert ON app.pinned_fact FOR INSERT WITH CHECK ( + app.current_role_name() IN ('admin','instructor') + OR EXISTS ( + SELECT 1 FROM app.case_profile cp + WHERE cp.case_id = app.pinned_fact.case_id + AND cp.learner_id = app.current_uid() + ) +); +DROP POLICY IF EXISTS p_pinned_update ON app.pinned_fact; +CREATE POLICY p_pinned_update ON app.pinned_fact FOR UPDATE USING ( + app.current_role_name() IN ('admin','instructor') + OR EXISTS ( + SELECT 1 FROM app.case_profile cp + WHERE cp.case_id = app.pinned_fact.case_id + AND cp.learner_id = app.current_uid() + ) +) WITH CHECK ( + app.current_role_name() IN ('admin','instructor') + OR EXISTS ( + SELECT 1 FROM app.case_profile cp + WHERE cp.case_id = app.pinned_fact.case_id + AND cp.learner_id = app.current_uid() + ) +); -- ── learner_profile: 학습자=본인, persistent_gaps 는 응답단 필터 ── ALTER TABLE app.learner_profile ENABLE ROW LEVEL SECURITY;