상담자 발화 판정(A)·감정(B)·표현(C) 20문항 질문 세트, 감쇠 없는 이번 턴 반응과 비대칭 기분 전이, 개방도 게이트로 생성 지시를 만들고 ccd.coping_strategy 전달 누락을 고친다. trace v2와 고정 문구 속마음 요약(migration 24, AI 경로 차단 RLS)을 같은 트랜잭션에 저장하고 피드백 정책이 켜진 경우에만 done·TurnResponse·음성 reply·리뷰로 노출한다. 회기 화면 속마음 보기 토글, 리뷰 접힘 블록, 관리자 감정 관측 v2 표시를 추가한다.
2291 lines
82 KiB
Python
2291 lines
82 KiB
Python
"""Counseling session routes.
|
|
|
|
The DB-backed source of truth is still pending, so this route uses the existing
|
|
in-process session store when DB is degraded. Unlike the previous dev fallback,
|
|
all browser calls now require a verified server-side auth session and every
|
|
session operation checks learner ownership.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import json
|
|
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 ..contracts.client_affect import ClientInnerReactionV1
|
|
from ..services import (
|
|
client_affect,
|
|
evaluator,
|
|
feedback_policy,
|
|
guardrail,
|
|
inner_reaction_exposure,
|
|
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
|
|
# 학습자·교수자용 속마음 요약(§8.2·§9). 저장 성공 && 피드백 정책 켜짐일 때만 값.
|
|
inner_reaction: Optional[ClientInnerReactionV1] = 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,
|
|
)
|
|
inner_reactions = await session_persistence.list_client_inner_reactions(
|
|
session_id,
|
|
review_principal,
|
|
)
|
|
else:
|
|
evaluation_record, evaluation_durable = None, True
|
|
inner_reactions = {}
|
|
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,
|
|
inner_reactions=inner_reactions,
|
|
)
|
|
)
|
|
|
|
|
|
@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 []),
|
|
),
|
|
inner_reaction=inner_reaction_exposure.expose_client_inner_reaction(
|
|
ctx.client_inner_reaction,
|
|
stored=True,
|
|
feedback_enabled=feedback_policy.effective_learner_feedback_enabled(
|
|
sess, principal
|
|
),
|
|
),
|
|
)
|
|
|
|
|
|
@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,
|
|
)
|
|
exposed_inner_reaction = (
|
|
inner_reaction_exposure.expose_client_inner_reaction(
|
|
ctx.client_inner_reaction,
|
|
stored=True,
|
|
feedback_enabled=(
|
|
feedback_policy.effective_learner_feedback_enabled(
|
|
sess, principal
|
|
)
|
|
),
|
|
)
|
|
)
|
|
data["inner_reaction"] = (
|
|
exposed_inner_reaction.model_dump(mode="json")
|
|
if exposed_inner_reaction is not None
|
|
else None
|
|
)
|
|
_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=client_affect.public_end_state(carry.end_state),
|
|
)
|