192 lines
6.4 KiB
Python
192 lines
6.4 KiB
Python
"""관리자 감정 관측용 영속 조회."""
|
|
|
|
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, ClientAffectTraceV1
|
|
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
|
|
|
|
|
|
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=ClientAffectTraceV1.model_validate(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,
|
|
)
|