"""관리자 감정 관측용 영속 조회.""" from __future__ import annotations import math from typing import Any from uuid import UUID from pydantic import ValidationError from ..config import settings from ..contracts.admin_affect import ( AdminAffectRuntimeResponse, AdminAffectSessionDetailResponse, AdminAffectSessionListResponse, AdminAffectSessionSummary, AdminAffectTraceRecord, ) from ..contracts.client_affect import ( CLIENT_AFFECT_DIMENSIONS, ClientAffectTrace, ClientAffectTraceV1, ClientAffectTraceV2, ) from ..db import acquire from .jev_client import jev_client class AdminAffectSessionNotFoundError(LookupError): """요청한 회기가 존재하지 않는다.""" class AdminAffectPersistenceError(RuntimeError): """관리자 감정 관측용 영속 조회를 완료할 수 없다.""" def runtime_snapshot() -> AdminAffectRuntimeResponse: """현재 Jev 클라이언트 설정만 노출한다. 연결 검증 결과는 포함하지 않는다.""" return AdminAffectRuntimeResponse( enabled=settings.client_affect_provider == "jev", provider=jev_client.provider, model=jev_client.model, configured=jev_client.configured, ) def _current_emotions(affect_state: Any) -> dict[str, float | None]: state = affect_state if isinstance(affect_state, dict) else {} emotions: dict[str, float | None] = {} for dimension in CLIENT_AFFECT_DIMENSIONS: value = state.get(f"emotion_{dimension}") if isinstance(value, bool) or not isinstance(value, (int, float)): emotions[dimension] = None continue number = float(value) emotions[dimension] = number if math.isfinite(number) and 0.0 <= number <= 1.0 else None return emotions def _parse_trace(raw: Any) -> ClientAffectTrace: """schema_version으로 v1/v2를 구분해 검증한다. 알 수 없는 버전은 예외로 503 처리된다.""" if not isinstance(raw, dict): raise ValueError("client affect trace must be an object") version = raw.get("schema_version") if version == 1: return ClientAffectTraceV1.model_validate(raw) if version == 2: return ClientAffectTraceV2.model_validate(raw) raise ValueError("unsupported client affect trace schema_version") async def list_sessions( *, user_id: str, limit: int, offset: int, ) -> AdminAffectSessionListResponse: """최근 회기 순으로 민감 식별자 없이 감정 trace 수를 조회한다.""" try: async with acquire(role="admin", user_id=user_id) as conn: total_row = await conn.fetchrow( """ SELECT count(*)::int AS total FROM app.sessions """ ) rows = await conn.fetch( """ SELECT s.id AS session_id, COALESCE(s.persona_code, '') AS persona_code, s.started_at, s.ended_at IS NOT NULL AS ended, COALESCE(trace_count.trace_count, 0)::int AS trace_count FROM app.sessions AS s LEFT JOIN ( SELECT session_id, count(*)::int AS trace_count FROM app.client_affect_trace GROUP BY session_id ) AS trace_count ON trace_count.session_id = s.id ORDER BY s.started_at DESC, s.id DESC LIMIT $1 OFFSET $2 """, limit, offset, ) except Exception as exc: raise AdminAffectPersistenceError("admin affect sessions are unavailable") from exc sessions = [ AdminAffectSessionSummary( session_id=str(row["session_id"]), persona_code=str(row["persona_code"]), started_at=row["started_at"], ended=bool(row["ended"]), trace_count=int(row["trace_count"]), ) for row in rows ] return AdminAffectSessionListResponse( runtime=runtime_snapshot(), sessions=sessions, total=int(total_row["total"] if total_row is not None else 0), limit=limit, offset=offset, ) async def get_session_detail( *, user_id: str, session_id: UUID, limit: int, before_seq: int | None, ) -> AdminAffectSessionDetailResponse: """저장된 trace와 현재 snapshot만 조회하며 과거 상태를 재구성하지 않는다.""" try: async with acquire(role="admin", user_id=user_id) as conn: session = await conn.fetchrow( """ SELECT s.id AS session_id, COALESCE(s.persona_code, '') AS persona_code, state.affect_state FROM app.sessions AS s LEFT JOIN app.session_state AS state ON state.session_id = s.id WHERE s.id = $1::uuid """, session_id, ) if session is None: raise AdminAffectSessionNotFoundError("admin affect session not found") total_row = await conn.fetchrow( """ SELECT count(*)::int AS total_traces FROM app.client_affect_trace WHERE session_id = $1::uuid """, session_id, ) trace_rows = await conn.fetch( """ SELECT turn.id AS turn_id, turn.seq, trace.created_at, trace.trace FROM app.client_affect_trace AS trace JOIN app.turns AS turn ON turn.id = trace.turn_id WHERE trace.session_id = $1::uuid AND ($2::int IS NULL OR turn.seq < $2) ORDER BY turn.seq DESC, turn.id DESC LIMIT $3 """, session_id, before_seq, limit + 1, ) except AdminAffectSessionNotFoundError: raise except Exception as exc: raise AdminAffectPersistenceError("admin affect session is unavailable") from exc has_more = len(trace_rows) > limit selected_rows = list(trace_rows[:limit]) try: traces = [ AdminAffectTraceRecord( turn_id=str(row["turn_id"]), seq=int(row["seq"]), created_at=row["created_at"], trace=_parse_trace(row["trace"]), ) for row in reversed(selected_rows) ] except (KeyError, TypeError, ValidationError, ValueError) as exc: raise AdminAffectPersistenceError("admin affect trace is invalid") from exc return AdminAffectSessionDetailResponse( session_id=str(session["session_id"]), persona_code=str(session["persona_code"]), current_emotions=_current_emotions(session["affect_state"]), traces=traces, total_traces=int(total_row["total_traces"] if total_row is not None else 0), has_more=has_more, )