"""Counseling session routes. The DB-backed source of truth is still pending, so this route uses the existing in-process session store when DB is degraded. Unlike the previous dev fallback, all browser calls now require a verified server-side auth session and every session operation checks learner ownership. """ from __future__ import annotations import asyncio import json import logging import secrets import time from datetime import datetime from typing import Literal, Optional from uuid import UUID from fastapi import APIRouter, HTTPException, Request, status from pydantic import BaseModel, Field, field_validator from sse_starlette.sse import EventSourceResponse from .. import db, session_evaluation_repository, session_persistence, turn_runtime from ..auth_sessions import user_has_consent, user_onboarding_complete from ..config import settings from ..deps import CurrentPrincipal, Principal, Role from ..engine_client import EngineError, engine_client from ..persona_repository import get_catalog_persona from ..runtime_policy import require_runtime_fallback_allowed from ..session_evaluation_input import enriched_masked_turns from ..session_turn_memory import build_turn_memory from ..session_evaluation_timeout import ( session_evaluation_outer_timeout_seconds, session_evaluation_timeout_seconds as _session_evaluation_timeout_seconds, session_evaluation_stale_after_seconds, session_evaluation_transport_timeout_seconds, ) from ..services import ( evaluator, feedback_policy, guardrail, live_coach, memory, notifications, orchestrator, rag, rupture_runtime, rupture_scenario_director, session_digest_worker, session_learning_producer, state_machine, ) from ..session_dashboard_projection import ( LearnerDashboardResponse, 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, dashboard_training_exposure as _dashboard_training_exposure, ) from ..session_projection import ( LEARNER_VISIBLE_AI_ROLE, iso as _iso, learner_visible_turns as _learner_visible_turns, stage_label as _stage_label, ) from ..session_read_model import ( CaseMemoryPreview, CaseProgressStats, LearnerCaseListResponse, LearnerCaseSummary, LearnerSessionsResponse, LearnerSessionSummary, ReviewCaseWorksheet, ReviewCaseWorksheetSaveRequest, ReviewWorksheetItem as ReviewWorksheetItem, ReviewWorksheetSection as ReviewWorksheetSection, SessionArchiveResponse, SessionDetailResponse, SessionReviewReadInput, SessionReviewResponse, SessionProgress, SessionShareDeleteResponse, SessionShareResponse, StageLabel, build_session_progress, build_session_review, learner_summary as _learner_summary, session_detail as _session_detail, session_share_payload as _session_share_payload, ) from ..store import InProcSession, TurnRecord, store router = APIRouter(prefix="/sessions", tags=["sessions"]) logger = logging.getLogger(__name__) _SESSION_EVALUATION_IN_FLIGHT: set[str] = set() _SESSION_EVALUATION_RECOVERY_TASK: asyncio.Task[int] | None = None _STREAM_TURN_EVALUATION_TASKS: set[asyncio.Task[None]] = set() # Text, SSE, and voice all call this module function after both durable turn UUIDs # exist. Install once here (main imports sessions before voice) so no route can miss # the same-process evaluator background boundary. rupture_runtime.install_turn_finalize_hook(turn_runtime) TheoryMode = Literal["humanistic", "cbt", "integrative"] EndStateValue = str | int | float | bool | None | dict[str, float] class SessionStartRequest(BaseModel): persona_code: str = Field(..., examples=["P1"]) theory_mode: TheoryMode = "humanistic" # continue는 선택한 사례의 압축 기억을 이어 받고, fresh는 새 case_id/S1으로 시작한다. # 구클라이언트는 기존 동작을 보존하도록 continue가 기본이다. start_mode: Literal["continue", "fresh"] = "continue" case_id: UUID | None = None # 이번 회기 목표 단계(2026-07-13 회의 P1). 회의 권장은 2개 수준이지만 # 소유자 지시(2026-07-15)로 1~4개까지 자유 선택을 허용한다. # 빈 리스트는 구계약 클라이언트 호환용 — 준비 페이지는 항상 1개 이상을 보낸다. goal_stages: list[StageLabel] = Field(default_factory=list, max_length=4) @field_validator("goal_stages") @classmethod def _dedupe_goal_stages(cls, value: list[StageLabel]) -> list[StageLabel]: seen: list[StageLabel] = [] for stage in value: if stage not in seen: seen.append(stage) return seen[:4] class SessionStartResponse(BaseModel): session_id: str case_id: str persona_id: str persona_version: int session_no: int stage: StageLabel effective_openness: float recall_summary: Optional[str] = None degraded: bool = False started_at: str = "" goal_stages: list[StageLabel] = Field(default_factory=list) # 시간 기반 회기 종료 계약(회의 P1): 프론트 타이머·10분 전 알람의 기준값. duration_limit_seconds: int = 0 warning_before_end_seconds: int = 0 learner_feedback_enabled: bool = True start_mode: Literal["continue", "fresh"] = "continue" class TurnRequest(BaseModel): text: str = Field(..., min_length=1) class LiveCoachRequest(BaseModel): learner_text: str = Field(..., min_length=1) question: Optional[str] = None client_reply: Optional[str] = None turn_seq: Optional[int] = Field(default=None, ge=1) class LiveCoachHistoryResponse(BaseModel): source: Literal["database", "runtime"] = "runtime" quota: live_coach.LiveCoachQuota = Field( default_factory=lambda: live_coach.LiveCoachQuota( remaining=session_persistence.LIVE_COACH_INITIAL_CREDITS, max=session_persistence.LIVE_COACH_MAX_CREDITS, ) ) events: list[live_coach.LiveCoachEvent] = Field(default_factory=list) credit_events: list[live_coach.LiveCoachCreditEvent] = Field(default_factory=list) class CrisisResourceResponse(BaseModel): title: str number: str message: str class TurnResponse(BaseModel): turn_seq: int stage: StageLabel effective_openness: float client_reply: Optional[str] = None safety_flagged: bool = False crisis_kind: str = "none" crisis_resource: Optional[CrisisResourceResponse] = None conversation_stopped: bool = False output_error: Optional[str] = None # P2 단계 누적 게이지·상세 수치 — 턴마다 갱신된 파생값. progress: Optional[SessionProgress] = None class SessionEndResponse(BaseModel): session_id: str session_no: int digest_pending: bool end_state: dict[str, EndStateValue] _RECALL_CACHE: dict[str, memory.RecallContext] = {} # 세션별 KB 증상 행동단서(회기 1회 산출·캐시). 빈 list 캐시 = 회기 내 재시도 안 함(안정성). _KB_CUES_CACHE: dict[str, list[str]] = {} _RAG_WARM_SEMAPHORE = asyncio.Semaphore(1) def cached_kb_cues(session_id: str) -> list[str]: """Return a defensive copy of the session-scoped, process-lifetime KB cues.""" return list(_KB_CUES_CACHE.get(session_id) or []) def invalidate_session_context_cache(session_id: str) -> None: """Invalidate all derived turn context when a session reaches its terminal state.""" _RECALL_CACHE.pop(session_id, None) _KB_CUES_CACHE.pop(session_id, None) # ──────────────────────────────────────────────────────────────────────────── # RAG 배선 헬퍼 — 내담자(CLIENT) 뷰. 임베더/KB/DB 풀 미가용 시 빈 값으로 graceful # degradation: 상담 루프를 절대 막지 않는다(라이브 루프 비차단이 계약). routes/kb.py가 # 같은 예외를 503으로 올리는 것과 의도적으로 다르다. 임베딩은 rag가 스레드풀로 offload. # ──────────────────────────────────────────────────────────────────────────── _RAG_RECALL_K = 5 _KB_CUES_K = 4 def _persona_kb_query(card) -> str: """페르소나 증상·호소 → KB 행동단서 검색 질의(임베더/tsquery 입력 전용, LLM 미주입). 질의는 프롬프트에 들어가지 않는다. 회수된 behavior_cue만 L2로 주입되고, CLIENT 정책 (expose_body=False)이 본문을 잘라 '행동단서'만 돌려준다(CCD 본문 비노출 자동 보존). """ parts: list[str] = [] presenting = getattr(card, "presenting", None) or {} if presenting.get("주호소"): parts.append(str(presenting["주호소"])) if presenting.get("표층"): parts.append(str(presenting["표층"])) dsm = getattr(card, "dsm5_dimensional", None) or {} parts.extend(str(key) for key in dsm.keys() if key != "note") return " ".join(p for p in parts if p).strip() async def _retrieve_kb_behavior_cues(card) -> list[str]: """KB 증상 행동단서 회수(CLIENT 정책). 미가용 시 빈 리스트(비차단).""" query = _persona_kb_query(card) if not query: return [] try: async with db.acquire(ai_view=rag.AIRole.CLIENT.value) as conn: result = await rag.search_kb( conn, query=query, role=rag.AIRole.CLIENT, k=_KB_CUES_K, ) return [c.behavior_cue for c in result.chunks if c.behavior_cue] except Exception: # rag.NotConfigured(임베더/KB 미가용)·RuntimeError(풀 미초기화)·DB 오류 포함. # 비치명적: 빈 단서로 진행. CancelledError는 BaseException이라 미포착. return [] async def _retrieve_live_coach_grounding( *, learner_text: str, client_reply: str | None, stage: str, theory_mode: str, ) -> list[live_coach.LiveCoachGrounding]: """라이브 코치용 평가 근거 회수. 미가용 시 빈 리스트로 진행한다.""" learner_masked = guardrail.mask_pii(learner_text).text_masked # LiveCoachRequest.client_reply는 브라우저가 보내는 제어 입력이므로 합성 출력이 # 아니다. 일반 입력용 PII 게이트를 유지해 문맥 없는 실제 이름도 차단한다. client_masked = guardrail.mask_pii(client_reply or "").text_masked query = " ".join( part for part in [stage, theory_mode, learner_masked, client_masked] if part ).strip() if not query: return [] try: async with db.acquire(ai_view=rag.AIRole.EVALUATOR.value) as conn: result = await rag.retrieve_eval_grounding( conn, query=query, k=4, kinds=( "theory", "technique", "supervisor_pattern", "microskill", "taxonomy", ), ) try: await rag.log_retrieval( conn, result=result, ai_role="evaluator", used_in_answer=True, ) except Exception: pass except Exception: return [] out: list[live_coach.LiveCoachGrounding] = [] for chunk in result.chunks: body = chunk.body or chunk.behavior_cue or chunk.context_prefix or "" if not body: continue meta = chunk.meta if isinstance(chunk.meta, dict) else {} title = str( meta.get("source_title") or meta.get("title") or chunk.source_id or "Vignette KB" ).strip() source_type = str(meta.get("source_type") or "").strip() source_version = str( meta.get("source_version") or meta.get("version") or "" ).strip() citation = str(meta.get("citation") or "").strip() out.append( live_coach.LiveCoachGrounding( source_id=chunk.source_id or f"kb:{chunk.chunk_id}", title=title or "Vignette KB", locator=chunk.heading_path, kb_kind=chunk.kb_kind, source_type=source_type or None, version=source_version or None, citation=citation or None, license_class=chunk.license_class, external_llm_ok=chunk.external_llm_ok, summary=body[:500], ) ) return out def _latest_turn_evaluation( sess: InProcSession, turn_seq: int | None ) -> Optional[dict]: """방금 상담자 발화에 붙은 fast-loop 평가를 찾는다.""" for turn in reversed(sess.turns): if turn.speaker != "counselor": continue if turn_seq is not None and turn.turn_seq != turn_seq: continue if isinstance(turn.evaluation, dict): return turn.evaluation return None return None async def _ensure_kb_cues(session_id: str, card) -> list[str]: """세션별 KB 행동단서(회기 1회 산출·캐시, 서버 재시작/재개 시 lazy 재계산).""" cached = _KB_CUES_CACHE.get(session_id) if cached is not None: return cached cues = await _retrieve_kb_behavior_cues(card) _KB_CUES_CACHE[session_id] = cues return cues async def _load_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: 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 WHERE case_id = $1::uuid ORDER BY session_no DESC, created_at DESC LIMIT 1 """, 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 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 { "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"]], } async def _hydrate_episodic_text(conn, result) -> list[str]: """retrieve_persona_memory가 돌려준 turn_id → app.turns 마스킹 본문 조인(내담자 발화).""" turn_ids = [c.meta.get("turn_id") for c in result.chunks if c.meta.get("turn_id")] if not turn_ids: return [] rows = await conn.fetch( """ SELECT id, text_masked FROM app.turns WHERE id = ANY($1::uuid[]) AND speaker = 'client' """, turn_ids, ) by_id = {str(r["id"]): r["text_masked"] for r in rows} return [by_id[t] for t in turn_ids if by_id.get(t)] def _recall_query(prev_summary: Optional[dict], card) -> str: """episodic recall 질의: 직전 open_threads 우선, 없으면 주호소.""" if prev_summary: threads = prev_summary.get("open_threads") or [] if threads: return " ".join(str(t) for t in threads) presenting = getattr(card, "presenting", None) or {} return str(presenting.get("주호소") or "").strip() async def _episodic_recall_snippets(case_id: str, query: str) -> list[str]: """case 스코프 episodic 벡터 recall → 내담자 발화 단편(마스킹본). 미가용 시 [].""" if not query: return [] try: async with db.acquire(ai_view=rag.AIRole.CLIENT.value) as conn: result = await rag.retrieve_persona_memory( conn, case_id=case_id, query=query, k=_RAG_RECALL_K, ) return await _hydrate_episodic_text(conn, result) except Exception: return [] async def _build_start_recall(*, case_id: str, card) -> memory.RecallContext: """회기 시작 회상 조립: prev_summary(case) + episodic recall을 build_recall_context로 합본. 전 구간 graceful(미가용 시 빈 회상). """ try: db.get_pool() # 풀 미초기화 시 RuntimeError → 첫 회기와 동일한 빈 회상 except RuntimeError: return memory.build_recall_context() 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) return memory.build_recall_context( case_digest=case_memory.get("case_digest"), prev_summary=prev_summary, episodic_snippets=episodic, pinned_facts=case_memory.get("pinned_facts") or [], ) async def _build_seed_recall(*, case_id: str | None) -> memory.RecallContext: if not case_id: return memory.build_recall_context() try: db.get_pool() except RuntimeError: return memory.build_recall_context() 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: cached = _RECALL_CACHE.get(sess.session_id) if cached is not None: return cached recall = await _build_seed_recall(case_id=sess.case_id) _RECALL_CACHE[sess.session_id] = recall return recall def session_time_over(sess: InProcSession) -> bool: """시간 기반 회기 종료(회의 P1): 제한 + 마무리 유예까지 지난 세션인지 판정. 제한 시간(기본 60분) 도달 자체는 프론트가 정리 유도·자동 종료로 처리하고, 서버는 유예(기본 +10분)까지 지난 뒤의 새 턴만 거부한다(마무리 인사 허용). """ if settings.session_duration_minutes <= 0: return False limit_seconds = ( settings.session_duration_minutes + settings.session_overtime_grace_minutes ) * 60 return (time.time() - sess.created_at) > limit_seconds def _ensure_turn_time_allowed(sess: InProcSession) -> None: if session_time_over(sess): raise HTTPException( status.HTTP_409_CONFLICT, detail="session_time_over", ) async def _prepare_turn_context( *, session_id: str, sess: InProcSession, learner_text: str, ) -> orchestrator.TurnContext: _ensure_turn_time_allowed(sess) recall = await ensure_recall_context(sess) kb_cues = ( _KB_CUES_CACHE.get(session_id) or [] ) # 비차단: warm 전이면 빈 단서(graceful) scenario_context = ( await rupture_scenario_director.load_stored_scenario_context( session_id=session_id, case_id=sess.case_id, ) ) ctx = orchestrator.prepare_turn( session_id=session_id, case_id=sess.case_id, card=sess.persona, state=sess.state, learner_text=learner_text, learner_identity=sess.learner_label, memory=build_turn_memory(sess, recall, kb_cues), theory_mode=sess.theory_mode, scenario_context=scenario_context, ) assert ctx.state_after is not None return ctx async def _warm_rag_caches(session_id: str, case_id: str, card) -> None: """RAG 회상·KB 행동단서를 **백그라운드**로 산출해 캐시한다(요청 경로 비차단). BGE-M3 임베더 첫 로드(~수 초)가 회기 시작/턴 응답을 막지 않도록 create_task로 띄운다. warm 완료 전 턴은 빈 회상/단서로 진행(graceful), 이후 턴부터 RAG 주입. 전 구간 비치명적. """ async with _RAG_WARM_SEMAPHORE: try: _RECALL_CACHE[session_id] = await _build_start_recall( case_id=case_id, card=card ) except Exception: pass try: _KB_CUES_CACHE[session_id] = await _retrieve_kb_behavior_cues(card) except Exception: pass def _ensure_learner(principal: Principal) -> Principal: if principal.role == Role.LEARNER: return principal if principal.can_access_role(Role.LEARNER): return principal.with_role(Role.LEARNER) raise HTTPException( status.HTTP_403_FORBIDDEN, detail="only learners can use sessions" ) async def _ensure_practice_consent(principal: Principal) -> None: if principal.consent_at is not None: return if await user_has_consent(principal.user_id): return raise HTTPException(status.HTTP_403_FORBIDDEN, detail="consent_required") async def _ensure_onboarding_complete(principal: Principal) -> None: if principal.profile_completed_at is not None: return if await user_onboarding_complete(principal.user_id): return raise HTTPException(status.HTTP_403_FORBIDDEN, detail="onboarding_required") async def _load_session_or_404( session_id: str, principal: Principal, *, allow_ended: bool = False, include_turn_evaluation: bool = False, ) -> InProcSession: sess, err = await turn_runtime.load_owned_session( session_id, principal, allow_ended=allow_ended, include_turn_evaluation=include_turn_evaluation, ) if err == turn_runtime.SessionAccessError.NOT_FOUND: raise HTTPException(status.HTTP_404_NOT_FOUND, detail="session not found") if err == turn_runtime.SessionAccessError.FORBIDDEN: raise HTTPException( status.HTTP_403_FORBIDDEN, detail="session does not belong to user" ) if err == turn_runtime.SessionAccessError.ENDED: raise HTTPException(status.HTTP_409_CONFLICT, detail="session already ended") assert sess is not None return sess def _review_supervisor_principal(principal: Principal) -> Principal | None: if principal.role in {Role.TEACHER, Role.ADMIN}: return principal if principal.super_admin: return principal.with_role(Role.ADMIN) return None async def _load_supervisor_review_session_or_404( session_id: str, principal: Principal, *, include_turn_evaluation: bool = False, ) -> InProcSession: sess = await session_persistence.load_session( session_id, principal, allow_ended=True, include_turn_evaluation=include_turn_evaluation, ) if sess is None and turn_runtime.runtime_fallback_allowed(): sess = store.get(session_id) if sess is None: raise HTTPException(status.HTTP_404_NOT_FOUND, detail="session not found") return sess async def _load_review_session_or_404( session_id: str, principal: Principal, *, include_turn_evaluation: bool = False, ) -> tuple[InProcSession, Principal]: if principal.role == Role.LEARNER: try: sess = await _load_session_or_404( session_id, principal, allow_ended=True, include_turn_evaluation=include_turn_evaluation, ) return sess, principal except HTTPException as exc: supervisor = _review_supervisor_principal(principal) if supervisor is None or exc.status_code not in { status.HTTP_403_FORBIDDEN, status.HTTP_404_NOT_FOUND, }: raise return ( await _load_supervisor_review_session_or_404( session_id, supervisor, include_turn_evaluation=include_turn_evaluation, ), supervisor, ) supervisor = _review_supervisor_principal(principal) if supervisor is None: raise HTTPException( status.HTTP_403_FORBIDDEN, detail="session review access denied" ) return ( await _load_supervisor_review_session_or_404( session_id, supervisor, include_turn_evaluation=include_turn_evaluation, ), supervisor, ) async def _end_persisted_session(sess: InProcSession, carry: memory.CarryOver) -> None: if await session_persistence.end_session(sess, carry): sess.ended = True sess.ended_at = datetime.now().timestamp() store.put(sess) if _should_schedule_session_digest_worker(carry): asyncio.create_task(_run_session_digest_worker_for_session(sess.session_id)) asyncio.create_task(_write_episodic_embeddings(sess)) asyncio.create_task(engine_client.close_session(sess.session_id)) return require_runtime_fallback_allowed("session end") store.end(sess.session_id) asyncio.create_task(engine_client.close_session(sess.session_id)) def _should_schedule_session_digest_worker(carry: memory.CarryOver) -> bool: return bool( settings.session_digest_worker_enabled and carry.compression_job is not None ) async def _run_session_digest_worker_for_session(session_id: str) -> None: """Best-effort M2 LLM digest compressor. The DB connection is held only for load/apply. Engine generation runs outside the transaction so a slow provider cannot pin the pool. """ try: db.get_pool() async with db.acquire(role="admin") as conn: loaded = await session_digest_worker.load_session_digest_job( conn, session_id ) if loaded is None: return model = settings.session_digest_worker_model.strip() or None worker = await session_digest_worker.run_session_digest_worker( loaded.job, engine_client, existing_case_digest=loaded.existing_case_digest, model=model, audit_hook=session_persistence.record_llm_call_audit, ) if worker.apply_plan is None: return async with db.acquire(role="admin") as conn: await session_digest_worker.apply_session_digest_plan( conn, worker.apply_plan, learner_id=loaded.learner_id, ) except Exception: logger.warning( "session digest worker failed for session_id=%s", session_id, exc_info=True ) return async def _write_episodic_embeddings(sess: InProcSession) -> None: """Best-effort M2 episodic writer. 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, ) if not inputs: return try: 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: pass def _public_share_url(request: Request, token: str) -> str: base = str(request.base_url).rstrip("/") return f"{base}/share/session/{token}" async def _evaluate_stream_turn( ctx: orchestrator.TurnContext, final_reply: str ) -> Optional[dict]: """stream 경로 완료 후 fast-loop 평가를 계산한다. 실패는 턴 저장을 막지 않는다.""" if not final_reply: return None try: hook = evaluator.make_eval_hook( engine_client, audit_hook=session_persistence.record_llm_call_audit, ) return await hook(ctx, final_reply) except Exception as exc: logger.warning( "turn fast-loop evaluation failed: session_id=%s", ctx.session_id, exc_info=True, ) return orchestrator.turn_evaluation_error_payload(ctx, exc) async def _evaluate_and_persist_stream_turn( *, sess: InProcSession, ctx: orchestrator.TurnContext, final_reply: str, result: orchestrator.TurnResult, learner_turn: TurnRecord, ) -> None: """응답 완료 뒤 fast-loop 평가를 저장해 다음 발화의 임계 경로에서 분리한다.""" evaluation = await _evaluate_stream_turn(ctx, final_reply) if evaluation is None: return if learner_turn.turn_id is not None: saved = await session_persistence.replace_turn_evaluation( turn_id=learner_turn.turn_id, evaluation=evaluation, ) if not saved: logger.warning( "turn fast-loop evaluation was not saved: session_id=%s turn_id=%s", ctx.session_id, learner_turn.turn_id, ) return learner_turn.evaluation = evaluation cached = store.get(ctx.session_id) if cached is not None: for turn in cached.turns: if learner_turn.turn_id and turn.turn_id == learner_turn.turn_id: turn.evaluation = evaluation break if ( learner_turn.turn_id is None and turn.speaker == "counselor" and turn.turn_seq == learner_turn.turn_seq ): turn.evaluation = evaluation break result.evaluation = evaluation await turn_runtime.maybe_recharge_live_coach_credit(sess, ctx, result) rupture_runtime.schedule_session_scan( ctx.session_id, trigger="fast_evaluation_persisted", ) def _observe_stream_turn_evaluation_task(task: asyncio.Task[None]) -> None: _STREAM_TURN_EVALUATION_TASKS.discard(task) try: task.result() except asyncio.CancelledError: logger.info("turn fast-loop evaluation background task cancelled") except Exception: logger.exception("turn fast-loop evaluation background task crashed") def _schedule_stream_turn_evaluation( *, sess: InProcSession, ctx: orchestrator.TurnContext, final_reply: str, result: orchestrator.TurnResult, learner_turn: TurnRecord, ) -> asyncio.Task[None]: task = asyncio.create_task( _evaluate_and_persist_stream_turn( sess=sess, ctx=ctx, final_reply=final_reply, result=result, learner_turn=learner_turn, ), name=f"turn-evaluation:{ctx.session_id}:{result.turn_seq}", ) _STREAM_TURN_EVALUATION_TASKS.add(task) task.add_done_callback(_observe_stream_turn_evaluation_task) return task def _stream_result_from_done( ctx: orchestrator.TurnContext, final_reply: str, data: dict[str, object], evaluation: Optional[dict], ) -> orchestrator.TurnResult: assert ctx.state_after is not None output_error = str(data.get("output_error") or "") or None return orchestrator.TurnResult( turn_seq=ctx.state_after.turn_seq, stage=_stage_label(ctx.state_after.stage), effective_openness=ctx.state_after.effective_openness, client_reply=None if output_error else final_reply or None, safety_flagged=bool(data.get("safety_flagged")), state_after=ctx.state_after, evaluation=evaluation, crisis_kind=ctx.crisis.kind.value if ctx.crisis else "none", crisis_resource=data.get("crisis_resource") if isinstance(data.get("crisis_resource"), dict) else None, conversation_stopped=bool(data.get("conversation_stopped")), llm_provider=str(data.get("llm_provider") or "") or None, model=str(data.get("model") or "") or None, tokens_in=int(data.get("tokens_in") or 0), tokens_out=int(data.get("tokens_out") or 0), cost_usd=float(data.get("cost_usd") or 0.0), output_error=output_error, ) async def _generate_and_save_session_evaluation(sess: InProcSession) -> None: if not sess.turns: return timeout_seconds = _session_evaluation_timeout_seconds() transport_timeout_seconds = session_evaluation_transport_timeout_seconds() enriched = enriched_masked_turns( sess.masked_turns(), counselor_identity=sess.learner_label, client_identity=sess.persona.display_name, ) try: result = await asyncio.wait_for( evaluator.evaluate_session( session_id=sess.session_id, stage=_stage_label(sess.state.stage), masked_turns=enriched, engine=engine_client, technique_codes=[], theory_mode=sess.theory_mode, scope="session_end", audit_hook=session_persistence.record_llm_call_audit, timeout=transport_timeout_seconds, ), # gateway의 생성 deadline과 HTTP deadline을 같게 두면 요청 초기화·응답 수신 # 비용만으로 app이 먼저 취소될 수 있다. transport grace와 durable audit 기록 # 예산 뒤에 outer grace를 두어 정상 결과를 timeout error로 바꾸지 않는다. timeout=session_evaluation_outer_timeout_seconds(), ) write = session_evaluation_repository.SessionEvaluationWrite.from_result( session_id=sess.session_id, learner_id=sess.learner_id, result=result, counselor_identity=sess.learner_label, client_identity=sess.persona.display_name, ) saved = await session_evaluation_repository.save_session_evaluation(write) if not saved: logger.error( "session evaluation save did not reach durable store: session_id=%s status=%s scope=%s", sess.session_id, write.status, write.scope, ) if saved and write.status == "ready": try: await session_learning_producer.produce_session_learning_artifacts( sess.session_id ) except Exception: # 평가 원장은 이미 커밋됐다. 후속 학습 원장 장애가 ready 평가를 # error로 덮어쓰거나 알림 생성을 막아서는 안 된다. logger.exception( "session learning artifacts failed after evaluation save: session_id=%s", sess.session_id, ) if saved: await _enqueue_session_review_ready_notification(sess.session_id) except asyncio.TimeoutError: message = f"session evaluation timeout after {timeout_seconds:g}s" logger.exception("%s: session_id=%s", message, sess.session_id) write = session_evaluation_repository.SessionEvaluationWrite.from_error( session_id=sess.session_id, learner_id=sess.learner_id, scope="session_end", stage=_stage_label(sess.state.stage), error=message, counselor_identity=sess.learner_label, client_identity=sess.persona.display_name, ) saved = await session_evaluation_repository.save_session_evaluation(write) if not saved: logger.error( "session evaluation error save did not reach durable store: session_id=%s error=%s", sess.session_id, write.error, ) if saved: await _enqueue_session_review_ready_notification(sess.session_id) except Exception as exc: logger.exception("session evaluation failed: session_id=%s", sess.session_id) write = session_evaluation_repository.SessionEvaluationWrite.from_error( session_id=sess.session_id, learner_id=sess.learner_id, scope="session_end", stage=_stage_label(sess.state.stage), error=exc, counselor_identity=sess.learner_label, client_identity=sess.persona.display_name, ) saved = await session_evaluation_repository.save_session_evaluation(write) if not saved: logger.error( "session evaluation failure record did not reach durable store: session_id=%s error=%s", sess.session_id, write.error, ) if saved: await _enqueue_session_review_ready_notification(sess.session_id) def _observe_session_evaluation_task(task: asyncio.Task[None], session_id: str) -> None: _SESSION_EVALUATION_IN_FLIGHT.discard(session_id) try: task.result() except asyncio.CancelledError: logger.warning( "session evaluation background task cancelled: session_id=%s", session_id ) except Exception: logger.exception( "session evaluation background task crashed: session_id=%s", session_id ) def _schedule_session_evaluation(sess: InProcSession) -> asyncio.Task[None] | None: if not sess.turns: return None if sess.session_id in _SESSION_EVALUATION_IN_FLIGHT: logger.info( "session evaluation already scheduled: session_id=%s", sess.session_id, ) return None _SESSION_EVALUATION_IN_FLIGHT.add(sess.session_id) task = asyncio.create_task( _generate_and_save_session_evaluation(sess), name=f"session-evaluation:{sess.session_id}", ) task.add_done_callback( lambda done, session_id=sess.session_id: _observe_session_evaluation_task( done, session_id ) ) return task async def recover_missing_session_evaluations(*, limit: int | None = None) -> int: recovery_limit = ( settings.session_evaluation_recovery_limit if limit is None else limit ) if recovery_limit <= 0: return 0 stale_after_seconds = session_evaluation_stale_after_seconds() ( candidates, durable, ) = await session_persistence.list_sessions_missing_session_evaluation( older_than_seconds=stale_after_seconds, limit=recovery_limit, ) if not durable: logger.warning("session evaluation recovery skipped: durable store unavailable") return 0 scheduled = 0 for sess in candidates: if _schedule_session_evaluation(sess) is not None: scheduled += 1 if scheduled: logger.info("session evaluation recovery scheduled %d session(s)", scheduled) return scheduled def _observe_session_evaluation_recovery_task(task: asyncio.Task[int]) -> None: global _SESSION_EVALUATION_RECOVERY_TASK if _SESSION_EVALUATION_RECOVERY_TASK is task: _SESSION_EVALUATION_RECOVERY_TASK = None try: task.result() except asyncio.CancelledError: logger.warning("session evaluation recovery task cancelled") except Exception: logger.exception("session evaluation recovery task crashed") def schedule_missing_session_evaluation_recovery() -> asyncio.Task[int] | None: global _SESSION_EVALUATION_RECOVERY_TASK if settings.session_evaluation_recovery_limit <= 0: return None if ( _SESSION_EVALUATION_RECOVERY_TASK is not None and not _SESSION_EVALUATION_RECOVERY_TASK.done() ): return _SESSION_EVALUATION_RECOVERY_TASK task = asyncio.create_task( recover_missing_session_evaluations(), name="session-evaluation-recovery", ) _SESSION_EVALUATION_RECOVERY_TASK = task task.add_done_callback(_observe_session_evaluation_recovery_task) return task def cancel_missing_session_evaluation_recovery() -> None: task = _SESSION_EVALUATION_RECOVERY_TASK if task is not None and not task.done(): task.cancel() async def _enqueue_session_review_ready_notification(session_id: str) -> None: try: await notifications.enqueue_session_review_ready(session_id=session_id) except Exception as exc: logger.warning("session review notification enqueue failed: %s", exc) async def _load_learner_sessions( principal: Principal, *, include_turn_evaluation: bool = False, ) -> tuple[list[InProcSession], bool]: if include_turn_evaluation: sessions, durable = await session_persistence.list_recent_sessions( principal, include_turn_evaluation=True, ) else: sessions, durable = await session_persistence.list_recent_sessions(principal) if not durable: require_runtime_fallback_allowed("session list") sessions = [ sess for sess in store.list() if sess.learner_id == principal.user_id ] sessions.sort(key=lambda sess: sess.created_at, reverse=True) return sessions, durable async def _review_ready(sess: InProcSession, principal: Principal) -> bool: if not feedback_policy.can_expose_principal_learner_feedback(sess, principal): return False turns = _learner_visible_turns(sess) if not sess.ended or not turns: return False if len(turns) != len(sess.turns): return False evaluation_record, _ = await session_evaluation_repository.load_session_evaluation( sess.session_id, principal, ) return bool(evaluation_record and evaluation_record.get("status") == "ready") async def _review_ready_map( sessions: list[InProcSession], principal: Principal, ) -> dict[str, bool]: results = await asyncio.gather( *[_review_ready(sess, principal) for sess in sessions] ) return {sess.session_id: ready for sess, ready in zip(sessions, results)} async def _archive_map( sessions: list[InProcSession], principal: Principal, ) -> dict[str, dict[str, object]]: records, _ = await session_persistence.list_session_archives( [sess.session_id for sess in sessions], principal, ) return records async def _session_archive_response( sess: InProcSession, principal: Principal, *, archived: bool, archived_at: str | None, source: str, ) -> SessionArchiveResponse: return SessionArchiveResponse( session_id=sess.session_id, archived=archived, archived_at=archived_at, source=source, session=_learner_summary( sess, review_ready=await _review_ready(sess, principal), learner_feedback_enabled=( feedback_policy.effective_learner_feedback_enabled(sess, principal) ), archived=archived, archived_at=archived_at, ), ) @router.get("", response_model=LearnerSessionsResponse) async def list_learner_sessions(principal: CurrentPrincipal) -> LearnerSessionsResponse: """Return the current learner's real practice sessions.""" principal = _ensure_learner(principal) sessions, durable = await _load_learner_sessions(principal) archives = await _archive_map(sessions, principal) summaries: list[LearnerSessionSummary] = [] for sess in sessions[:20]: archive_record = archives.get(sess.session_id) archived_at = archive_record.get("archived_at") if archive_record else None summaries.append( _learner_summary( sess, review_ready=await _review_ready(sess, principal), learner_feedback_enabled=( feedback_policy.effective_learner_feedback_enabled( sess, principal, ) ), archived=archive_record is not None, archived_at=_iso(float(archived_at)) if isinstance(archived_at, (int, float)) else None, ) ) return LearnerSessionsResponse( source="database" if durable else "runtime", sessions=summaries, ) def _preview_text(value: object, *, limit: int) -> str | None: """Keep the learner foldout bounded even when a legacy digest is verbose.""" text = str(value or "").strip() if not text: return None if len(text) <= limit: return text return f"{text[: max(1, limit - 1)].rstrip()}…" def _preview_items(value: object, *, limit: int, item_limit: int) -> list[str]: if not isinstance(value, (list, tuple)): return [] items: list[str] = [] for raw in value: clipped = _preview_text(raw, limit=item_limit) if clipped: items.append(clipped) if len(items) >= limit: break return items def _db_datetime_iso(value: object) -> str | None: return value.isoformat() if isinstance(value, datetime) else None @router.get("/cases", response_model=LearnerCaseListResponse) async def list_learner_cases( persona_code: str, principal: CurrentPrincipal, ) -> LearnerCaseListResponse: """Return complete case-local progress for one NPC, not a capped history slice.""" principal = _ensure_learner(principal) try: catalog_persona = await get_catalog_persona(persona_code) except Exception as exc: raise HTTPException( status.HTTP_503_SERVICE_UNAVAILABLE, detail="persona catalog database unavailable", ) from exc if catalog_persona is None: raise HTTPException( status.HTTP_404_NOT_FOUND, detail=f"unknown persona {persona_code}" ) try: rows = await session_persistence.list_case_summaries( learner_id=principal.user_id, persona_id=catalog_persona.persona_id, ) except session_persistence.CaseProgressUnavailableError as exc: raise HTTPException( status.HTTP_503_SERVICE_UNAVAILABLE, detail="case_progress_unavailable", ) from exc card = catalog_persona.card return LearnerCaseListResponse( cases=[ LearnerCaseSummary( case_id=str(row["case_id"]), persona_code=card.code, persona_name=card.display_name, last_session_no=int(row["last_session_no"] or 0), progress=CaseProgressStats( total_sessions=int(row["total_sessions"] or 0), completed_sessions=int(row["completed_sessions"] or 0), total_turns=int(row["total_turns"] or 0), total_duration_seconds=int(row["total_duration_seconds"] or 0), active_session_id=( str(row["active_session_id"]) if row.get("active_session_id") is not None else None ), active_session_no=( int(row["active_session_no"]) if row.get("active_session_no") is not None else None ), active_started_at=_db_datetime_iso(row.get("active_started_at")), last_activity_at=_db_datetime_iso(row.get("last_activity_at")), ), ) for row in rows ] ) @router.get("/cases/{case_id}/memory", response_model=CaseMemoryPreview) async def get_learner_case_memory_preview( case_id: UUID, principal: CurrentPrincipal, ) -> CaseMemoryPreview: """Load only the learner-safe compact memory when its foldout is opened.""" principal = _ensure_learner(principal) case_key = str(case_id) try: async with db.acquire(role="learner", user_id=principal.user_id) as conn: case_row = await conn.fetchrow( """ SELECT case_digest FROM app.case_profile WHERE case_id = $1::uuid AND learner_id = $2::uuid """, case_key, principal.user_id, ) if case_row is None: raise HTTPException(status.HTTP_404_NOT_FOUND, detail="case_not_found") summary_row = await conn.fetchrow( """ SELECT ss.digest, ss.open_threads FROM app.session_summary AS ss JOIN app.sessions AS s ON s.id = ss.session_id WHERE ss.case_id = $1::uuid AND s.learner_id = $2::uuid ORDER BY ss.session_no DESC, ss.created_at DESC LIMIT 1 """, case_key, principal.user_id, ) fact_rows = await conn.fetch( """ SELECT LEFT(pf.value, 160) AS value FROM app.pinned_fact AS pf JOIN app.case_profile AS cp ON cp.case_id = pf.case_id WHERE pf.case_id = $1::uuid AND cp.learner_id = $2::uuid AND pf.status IN ('stable', 'evolving', 'locked') AND 'client' = ANY(pf.visible_to) ORDER BY pf.updated_at DESC LIMIT 8 """, case_key, principal.user_id, ) except HTTPException: raise except Exception as exc: logger.exception("case memory preview read failed", extra={"case_id": case_key}) raise HTTPException( status.HTTP_503_SERVICE_UNAVAILABLE, detail="case_memory_unavailable", ) from exc case_digest = _preview_text(case_row["case_digest"], limit=600) latest_session_digest = _preview_text( summary_row["digest"] if summary_row is not None else None, limit=600, ) open_threads = _preview_items( summary_row["open_threads"] if summary_row is not None else [], limit=6, item_limit=160, ) pinned_facts = _preview_items( [row["value"] for row in fact_rows], limit=8, item_limit=160, ) return CaseMemoryPreview( case_id=case_key, memory_available=bool( case_digest or latest_session_digest or open_threads or pinned_facts ), case_digest=case_digest, latest_session_digest=latest_session_digest, open_threads=open_threads, pinned_facts=pinned_facts, ) @router.get("/dashboard", response_model=LearnerDashboardResponse) async def learner_dashboard(principal: CurrentPrincipal) -> LearnerDashboardResponse: """Return the current learner's real practice dashboard aggregates.""" principal = _ensure_learner(principal) sessions, durable = await _load_learner_sessions( principal, include_turn_evaluation=principal.learner_feedback_enabled, ) review_ready = await _review_ready_map(sessions, principal) archives = await _archive_map(sessions, principal) learner_feedback_enabled = { sess.session_id: feedback_policy.effective_learner_feedback_enabled( sess, principal, ) for sess in sessions } feedback_sessions = [ sess for sess in sessions if learner_feedback_enabled.get(sess.session_id, True) ] visible_review_ready = { session_id: ready for session_id, ready in review_ready.items() if session_id not in archives } return LearnerDashboardResponse( source="database" if durable else "runtime", overview=_dashboard_overview( sessions, visible_review_ready=visible_review_ready, archived_sessions=len(archives), ), growth=_dashboard_growth(feedback_sessions), persona_progress=_dashboard_persona_progress( sessions, visible_review_ready, learner_feedback_enabled, ), training_exposure=_dashboard_training_exposure(sessions), achievements=_dashboard_achievements(sessions, visible_review_ready), recent_feedback=_dashboard_feedback(feedback_sessions), message=( "실제 연습 기록을 기준으로 개인 학습 흐름을 표시합니다." if sessions else "아직 표시할 실제 연습 기록이 없습니다." ), ) @router.get("/{session_id}", response_model=SessionDetailResponse) async def get_session_detail( session_id: str, principal: CurrentPrincipal, ) -> SessionDetailResponse: """Return a learner-owned session with transcript for resume/history.""" principal = _ensure_learner(principal) sess = await _load_session_or_404( session_id, principal, allow_ended=True, ) return _session_detail( sess, review_ready=await _review_ready(sess, principal), learner_feedback_enabled=( feedback_policy.effective_learner_feedback_enabled(sess, principal) ), ) @router.post("/{session_id}/archive", response_model=SessionArchiveResponse) async def archive_session( session_id: str, principal: CurrentPrincipal, ) -> SessionArchiveResponse: """Archive an ended learner-owned session without deleting transcript or review evidence.""" principal = _ensure_learner(principal) sess = await _load_session_or_404( session_id, principal, allow_ended=True, ) if not sess.ended: raise HTTPException( status.HTTP_409_CONFLICT, detail="active sessions cannot be archived" ) record, durable = await session_persistence.set_session_archived( session_id=sess.session_id, learner_id=principal.user_id, archived=True, ) archived_at = record.get("archived_at") if record else None return await _session_archive_response( sess, principal, archived=True, archived_at=_iso(float(archived_at)) if isinstance(archived_at, (int, float)) else None, source="database" if durable else "runtime", ) @router.post("/{session_id}/restore", response_model=SessionArchiveResponse) async def restore_archived_session( session_id: str, principal: CurrentPrincipal, ) -> SessionArchiveResponse: """Restore an archived learner-owned session to the normal history/review queues.""" principal = _ensure_learner(principal) sess = await _load_session_or_404( session_id, principal, allow_ended=True, ) _, durable = await session_persistence.set_session_archived( session_id=sess.session_id, learner_id=principal.user_id, archived=False, ) return await _session_archive_response( sess, principal, archived=False, archived_at=None, source="database" if durable else "runtime", ) @router.post( "", response_model=SessionStartResponse, status_code=status.HTTP_201_CREATED ) async def start_session( body: SessionStartRequest, principal: CurrentPrincipal, ) -> SessionStartResponse: """Start a learner-owned practice session.""" principal = _ensure_learner(principal) await _ensure_onboarding_complete(principal) await _ensure_practice_consent(principal) try: catalog_persona = await get_catalog_persona(body.persona_code) except Exception as exc: raise HTTPException( status.HTTP_503_SERVICE_UNAVAILABLE, detail="persona catalog database unavailable", ) from exc if catalog_persona is None: raise HTTPException( status.HTTP_404_NOT_FOUND, detail=f"unknown persona {body.persona_code}" ) card = catalog_persona.card if body.start_mode == "fresh" and body.case_id is not None: raise HTTPException( status.HTTP_422_UNPROCESSABLE_ENTITY, detail="fresh_start_must_not_select_case", ) # The durable transaction chooses or creates the case only after active-case # validation. Start with an empty recall here so a fresh request cannot see # any legacy case before its own empty case exists. recall = memory.build_recall_context() st = state_machine.init_state( params=card.openness_params(), carry=recall.carry, ) carry_rapport = st.rapport_credit goal_stages = [str(stage) for stage in body.goal_stages] learner_feedback_enabled = principal.learner_feedback_enabled async def build_locked_start_state( stable_case_id: str, _session_no: int, ) -> state_machine.SessionState: nonlocal recall recall = ( memory.build_recall_context() if body.start_mode == "fresh" else await _build_seed_recall(case_id=stable_case_id) ) return state_machine.init_state( params=card.openness_params(), carry=recall.carry, ) try: sess = await session_persistence.create_session( learner_id=principal.user_id, card=card, theory_mode=body.theory_mode, state=st, session_no=1, carry_rapport=carry_rapport, persona_id=catalog_persona.persona_id, persona_version=catalog_persona.version, case_id=str(body.case_id) if body.case_id is not None else None, start_mode=body.start_mode, goal_stages=goal_stages, learner_feedback_enabled=learner_feedback_enabled, locked_state_factory=build_locked_start_state, ) except session_persistence.ActiveSessionExistsError as exc: raise HTTPException( status.HTTP_409_CONFLICT, detail={ "code": "active_session_exists", "session_id": exc.session_id, }, ) from exc except session_persistence.CaseNotFoundError as exc: raise HTTPException( status.HTTP_404_NOT_FOUND, detail="case_not_found", ) from exc except session_persistence.SessionCreationPersistenceError as exc: raise HTTPException( status.HTTP_503_SERVICE_UNAVAILABLE, detail="session_persistence_unavailable", ) from exc degraded = catalog_persona.degraded or sess is None if sess is None: require_runtime_fallback_allowed("session creation") # Runtime fallback cannot prove durable case memory. Keep it empty rather # than leaking a guessed legacy recall, while preserving a selected # in-process case ID when one is available. recall = memory.build_recall_context() st = state_machine.init_state( params=card.openness_params(), carry=recall.carry, ) carry_rapport = st.rapport_credit active_session = store.find_active( learner_id=principal.user_id, persona_id=catalog_persona.persona_id, persona_code=card.code, ) if active_session is not None: raise HTTPException( status.HTTP_409_CONFLICT, detail={ "code": "active_session_exists", "session_id": active_session.session_id, }, ) runtime_case_id: str | None = None runtime_session_no = 1 if body.start_mode == "continue": related = [ candidate for candidate in store.list() if candidate.learner_id == principal.user_id and ( candidate.persona_id == catalog_persona.persona_id or candidate.persona_code == card.code ) ] if body.case_id is not None: runtime_case_id = str(body.case_id) related = [ candidate for candidate in related if candidate.case_id == runtime_case_id ] if not related: raise HTTPException( status.HTTP_404_NOT_FOUND, detail="case_not_found", ) elif related: newest = max(related, key=lambda candidate: candidate.created_at) runtime_case_id = newest.case_id related = [ candidate for candidate in related if candidate.case_id == runtime_case_id ] if related: runtime_session_no = max( candidate.session_no for candidate in related ) + 1 sess = store.create( learner_id=principal.user_id, persona=card, theory_mode=body.theory_mode, state=st, persona_id=catalog_persona.persona_id, persona_version=catalog_persona.version, case_id=runtime_case_id, session_no=runtime_session_no, carry_rapport=carry_rapport, goal_stages=goal_stages, learner_feedback_enabled=learner_feedback_enabled, ) else: st = sess.state store.put(sess) # DB 재조회 전의 첫 턴과 runtime fallback에서도 인증된 학습자 표시명을 # 역할 기반 비식별화에 사용할 수 있게 in-process 세션에만 보존한다. sess.learner_label = principal.display_name or principal.user_id store.put(sess) # 즉시 빈/carry 회상으로 응답을 막지 않는다. RAG 회상·KB 단서(임베더 로드 수 초)는 # 백그라운드 warm으로 캐시 — 회기 시작/턴 응답이 임베더 로드에 블로킹되지 않게(성능 회귀 방지). _RECALL_CACHE[sess.session_id] = recall asyncio.create_task(_warm_rag_caches(sess.session_id, sess.case_id, card)) return SessionStartResponse( session_id=sess.session_id, case_id=sess.case_id, persona_id=catalog_persona.persona_id, persona_version=catalog_persona.version, session_no=sess.session_no, stage=_stage_label(st.stage), effective_openness=round(st.effective_openness, 4), recall_summary=recall.recall_summary, degraded=degraded, started_at=_iso(sess.created_at) or "", goal_stages=body.goal_stages, duration_limit_seconds=settings.session_duration_minutes * 60, warning_before_end_seconds=settings.session_warning_minutes * 60, learner_feedback_enabled=sess.learner_feedback_enabled, start_mode=body.start_mode, ) @router.get("/{session_id}/review", response_model=SessionReviewResponse) async def get_session_review( session_id: str, principal: CurrentPrincipal, ) -> SessionReviewResponse: """Return a role-safe review built only from the stored session transcript.""" sess, review_principal = await _load_review_session_or_404( session_id, principal, include_turn_evaluation=True, ) expose_learner_feedback = ( feedback_policy.can_expose_principal_learner_feedback( sess, review_principal, ) ) if expose_learner_feedback: ( evaluation_record, evaluation_durable, ) = await session_evaluation_repository.load_session_evaluation( session_id, review_principal, ) else: evaluation_record, evaluation_durable = None, True saved_worksheet_payload, _ = await session_persistence.load_case_worksheet( session_id, review_principal, ) include_teacher_review = review_principal.role in {Role.TEACHER, Role.ADMIN} learner_feedback_enabled = feedback_policy.effective_learner_feedback_enabled( sess, review_principal, ) teacher_review_record = None if include_teacher_review: teacher_review_record, _ = await session_persistence.load_session_review_status( session_id, review_principal, ) 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, learner_feedback_enabled=learner_feedback_enabled, expose_learner_feedback=expose_learner_feedback, ) ) @router.post("/{session_id}/share", response_model=SessionShareResponse) async def create_session_share( session_id: str, request: Request, principal: CurrentPrincipal, ) -> SessionShareResponse: """Create a public unfurl URL for a learner-owned ended session review.""" principal = _ensure_learner(principal) sess = await _load_session_or_404( session_id, principal, allow_ended=True, include_turn_evaluation=True, ) if not sess.ended: raise HTTPException( status.HTTP_409_CONFLICT, detail="session must be ended before sharing" ) if not feedback_policy.can_expose_principal_learner_feedback(sess, principal): raise HTTPException( status_code=status.HTTP_403_FORBIDDEN, detail=feedback_policy.LEARNER_FEEDBACK_DISABLED_DETAIL, ) review = await get_session_review(session_id, principal) token = secrets.token_urlsafe(32) payload = _session_share_payload(review) saved = await session_persistence.save_session_share( session_id=session_id, learner_id=principal.user_id, token_hash=session_persistence.share_token_hash(token), payload=payload, ) if saved is None: raise HTTPException( status.HTTP_503_SERVICE_UNAVAILABLE, detail="session share persistence unavailable", ) created_at = saved.get("created_at") created_label = ( _iso( created_at if isinstance(created_at, (int, float)) else datetime.now().timestamp() ) or "" ) return SessionShareResponse( shareUrl=_public_share_url(request, token), title=str(payload["title"]), description=str(payload["description"]), imageUrl=str(payload["imageUrl"]), createdAt=created_label, ) @router.delete("/{session_id}/share", response_model=SessionShareDeleteResponse) async def revoke_session_share( session_id: str, principal: CurrentPrincipal, ) -> SessionShareDeleteResponse: """Revoke the public share URL for a learner-owned session.""" principal = _ensure_learner(principal) await _load_session_or_404(session_id, principal, allow_ended=True) revoked = await session_persistence.revoke_session_share( session_id=session_id, learner_id=principal.user_id, ) return SessionShareDeleteResponse(revoked=revoked) @router.put("/{session_id}/review/worksheet", response_model=ReviewCaseWorksheet) async def save_session_review_worksheet( session_id: str, body: ReviewCaseWorksheetSaveRequest, principal: CurrentPrincipal, ) -> ReviewCaseWorksheet: """Persist the learner's edited case formulation worksheet for this session.""" principal = _ensure_learner(principal) await _load_session_or_404( session_id, principal, allow_ended=True, include_turn_evaluation=False, ) worksheet = ReviewCaseWorksheet( status="saved_by_learner", generatedBy="learner-edited worksheet", sections=body.sections, limitations=body.limitations, ) ok = await session_persistence.save_case_worksheet( session_id=session_id, learner_id=principal.user_id, payload=worksheet.model_dump(mode="json"), ) if not ok: raise HTTPException( status.HTTP_503_SERVICE_UNAVAILABLE, detail="case worksheet persistence unavailable", ) return worksheet @router.post("/{session_id}/turn", response_model=TurnResponse) async def submit_turn( session_id: str, body: TurnRequest, principal: CurrentPrincipal, ) -> TurnResponse: """Submit one trainee utterance and return the generated client reply.""" principal = _ensure_learner(principal) sess = await _load_session_or_404(session_id, principal) ctx = await _prepare_turn_context( session_id=session_id, learner_text=body.text, sess=sess, ) try: result = await orchestrator.run_turn_generate( ctx, engine_client, eval_hook=evaluator.make_eval_hook( engine_client, audit_hook=session_persistence.record_llm_call_audit, ), audit_hook=session_persistence.record_llm_call_audit, ) except EngineError as exc: raise HTTPException( status.HTTP_503_SERVICE_UNAVAILABLE, detail=f"engine unavailable: {exc}", ) from exc await turn_runtime.finalize_completed_turn( sess, ctx, result, context_prefix="session", ) return TurnResponse( turn_seq=result.turn_seq, stage=_stage_label(result.state_after.stage), effective_openness=round(result.effective_openness, 4), client_reply=result.client_reply, safety_flagged=result.safety_flagged, crisis_kind=result.crisis_kind, crisis_resource=result.crisis_resource, conversation_stopped=result.conversation_stopped, output_error=result.output_error, progress=build_session_progress( result.state_after, prev_rapport_credit=sess.prev_rapport_credit, goal_stages=list(sess.goal_stages or []), ), ) @router.get("/{session_id}/live-coach", response_model=LiveCoachHistoryResponse) async def list_live_coach_history( session_id: str, principal: CurrentPrincipal, ) -> LiveCoachHistoryResponse: """현재 회기에서 학습자에게 실제로 전달된 라이브 코칭 이력을 반환한다.""" principal = _ensure_learner(principal) sess = await _load_session_or_404(session_id, principal) if not feedback_policy.can_expose_principal_learner_feedback(sess, principal): raise HTTPException( status_code=status.HTTP_403_FORBIDDEN, detail=feedback_policy.LEARNER_FEEDBACK_DISABLED_DETAIL, ) events, durable = await session_persistence.list_live_coach_events( session_id, principal ) quota, quota_durable = await session_persistence.get_live_coach_quota( session_id, principal ) ( credit_events, credit_durable, ) = await session_persistence.list_live_coach_credit_events( session_id, principal, ) return LiveCoachHistoryResponse( source="database" if durable and quota_durable and credit_durable else "runtime", quota=live_coach.LiveCoachQuota(**quota), events=[live_coach.LiveCoachEvent(**event) for event in events], credit_events=[ live_coach.LiveCoachCreditEvent(**event) for event in credit_events ], ) @router.post("/{session_id}/live-coach", response_model=live_coach.LiveCoachSuggestion) async def live_coach_turn( session_id: str, body: LiveCoachRequest, principal: CurrentPrincipal, ) -> live_coach.LiveCoachSuggestion: """방금 완료된 턴에 대한 비차단 라이브 코칭을 반환한다.""" principal = _ensure_learner(principal) sess = await _load_session_or_404(session_id, principal) if not feedback_policy.can_expose_principal_learner_feedback(sess, principal): raise HTTPException( status_code=status.HTTP_403_FORBIDDEN, detail=feedback_policy.LEARNER_FEEDBACK_DISABLED_DETAIL, ) quota, _ = await session_persistence.get_live_coach_quota(session_id, principal) if int(quota.get("remaining", 0)) <= 0: raise HTTPException( status_code=status.HTTP_409_CONFLICT, detail="live coach credit exhausted", ) turn_seq = body.turn_seq or max(1, int(getattr(sess.state, "turn_seq", 1) or 1)) stage = _stage_label(sess.state.stage) grounding = await _retrieve_live_coach_grounding( learner_text=body.learner_text, client_reply=body.client_reply, stage=stage, theory_mode=sess.theory_mode, ) prior_events, _ = await session_persistence.list_live_coach_events( session_id, principal ) prior_coach = [ { "title": str((event.get("suggestion") or {}).get("title") or ""), "focus": str((event.get("suggestion") or {}).get("focus") or ""), } for event in prior_events[-2:] ] item = live_coach.LiveCoachInput( session_id=sess.session_id, turn_seq=turn_seq, stage=stage, effective_openness=sess.state.effective_openness, theory_mode=sess.theory_mode, persona_code=sess.persona_code, persona_name=sess.persona.display_name, learner_text=body.learner_text, question=body.question, client_reply=body.client_reply, recent_turns=sess.recent_turns(k=8, visible_to=LEARNER_VISIBLE_AI_ROLE), evaluation=_latest_turn_evaluation(sess, body.turn_seq), goal_stages=list(sess.goal_stages or []), prior_coach=prior_coach, ) suggestion = await live_coach.generate_live_coaching( item, engine=engine_client, grounding=grounding, audit_hook=session_persistence.record_llm_call_audit, ) try: _, coach_event_durable = await session_persistence.save_live_coach_event( session_id=sess.session_id, learner_id=sess.learner_id, turn_seq=turn_seq, stage=stage, learner_text=body.learner_text, client_reply=body.client_reply, suggestion=suggestion, ) except session_persistence.LiveCoachCreditExhausted as exc: raise HTTPException( status_code=status.HTTP_409_CONFLICT, detail="live coach credit exhausted", ) from exc quota_after, quota_durable = await session_persistence.get_live_coach_quota( session_id, principal ) ( credit_events, credit_durable, ) = await session_persistence.list_live_coach_credit_events(session_id, principal) turn_credit_events = [ live_coach.LiveCoachCreditEvent(**event) for event in credit_events if int(event.get("turn_seq") or 0) == int(turn_seq) ] return suggestion.model_copy( update={ "persistence_source": "database" if coach_event_durable and quota_durable and credit_durable else "runtime", "quota": live_coach.LiveCoachQuota(**quota_after), "credit_events": turn_credit_events[-2:], } ) @router.post("/{session_id}/stream") async def stream_turn( session_id: str, body: TurnRequest, principal: CurrentPrincipal, ): """Stream a generated client reply for one trainee utterance.""" principal = _ensure_learner(principal) sess = await _load_session_or_404(session_id, principal) ctx = await _prepare_turn_context( session_id=session_id, learner_text=body.text, sess=sess, ) async def event_generator(): finalized_turn = False last_beat = asyncio.get_running_loop().time() final_reply = "" try: async for ev in orchestrator.run_turn_stream( ctx, engine_client, audit_hook=session_persistence.record_llm_call_audit, ): if ev.event == "token": text = str(ev.data.get("text", "")) final_reply += text yield {"event": "token", "data": text} elif ev.event == "done": data = { **ev.data, "stage": _stage_label(ctx.state_after.stage), "progress": build_session_progress( ctx.state_after, prev_rapport_credit=sess.prev_rapport_credit, goal_stages=list(sess.goal_stages or []), ).model_dump(), } if not finalized_turn: result = _stream_result_from_done( ctx, final_reply, data, None ) learner_turn = await turn_runtime.finalize_completed_turn( sess, ctx, result, context_prefix="session", recharge_live_coach=False, ) _schedule_stream_turn_evaluation( sess=sess, ctx=ctx, final_reply=final_reply, result=result, learner_turn=learner_turn, ) finalized_turn = True yield { "event": "done", "data": json.dumps(data, ensure_ascii=False), } elif ( ev.event == "safety" and bool(ev.data.get("conversation_stopped")) and ctx.crisis is not None and ctx.crisis.escalate and not finalized_turn ): safety_data = { "session_id": ctx.session_id, "stage": _stage_label(ctx.state_after.stage), "effective_openness": round( ctx.state_after.effective_openness, 4 ), "turn_seq": ctx.state_after.turn_seq, "safety_flagged": True, "crisis_kind": ctx.crisis.kind.value, "crisis_resource": ev.data.get("crisis_resource"), "conversation_stopped": True, } result = _stream_result_from_done( ctx, final_reply, safety_data, None ) await turn_runtime.finalize_completed_turn( sess, ctx, result, context_prefix="session", ) finalized_turn = True yield { "event": ev.event, "data": json.dumps(ev.data, ensure_ascii=False), } else: yield { "event": ev.event, "data": json.dumps(ev.data, ensure_ascii=False), } now = asyncio.get_running_loop().time() if now - last_beat >= settings.sse_heartbeat_seconds: yield {"event": "ping", "data": "{}"} last_beat = now except Exception as exc: yield { "event": "error", "data": json.dumps({"detail": str(exc)}, ensure_ascii=False), } return return EventSourceResponse(event_generator()) @router.post("/{session_id}/end", response_model=SessionEndResponse) async def end_session( session_id: str, principal: CurrentPrincipal, ) -> SessionEndResponse: """End a learner-owned session and prepare carry-over state.""" principal = _ensure_learner(principal) sess = await _load_session_or_404(session_id, principal, allow_ended=True) was_ended = bool(sess.ended) recall = _RECALL_CACHE.get(session_id) or memory.RecallContext() carry = memory.make_carry_over( state=sess.state, session_id=session_id, case_id=sess.case_id, session_no=sess.session_no, masked_turns=sess.masked_turns(visible_to="client"), prev_rapport_credit=sess.prev_rapport_credit, open_threads=recall.open_threads, ) await _end_persisted_session(sess, carry) invalidate_session_context_cache(session_id) rupture_runtime.schedule_session_scan(session_id, trigger="session_ended") has_counselor_turns = any(t.speaker in ("counselor", "learner") for t in sess.turns) if not was_ended and has_counselor_turns: _schedule_session_evaluation(sess) return SessionEndResponse( session_id=session_id, session_no=sess.session_no, digest_pending=carry.compression_job is not None, end_state=carry.end_state, )