Compare commits
No commits in common. "6ab40ff3f4f0198d1b7b6b6f0d5eeeb060209249" and "d22cd9883dc0a02b662b7865fd6314fbde407b7e" have entirely different histories.
6ab40ff3f4
...
d22cd9883d
17 changed files with 22 additions and 1449 deletions
|
|
@ -1,58 +0,0 @@
|
||||||
"""관리자 감정 관측 조회 API 계약."""
|
|
||||||
|
|
||||||
from __future__ import annotations
|
|
||||||
|
|
||||||
from datetime import datetime
|
|
||||||
|
|
||||||
from pydantic import BaseModel, ConfigDict, Field
|
|
||||||
|
|
||||||
from .client_affect import ClientAffectTraceV1
|
|
||||||
|
|
||||||
|
|
||||||
class AdminAffectRuntimeResponse(BaseModel):
|
|
||||||
model_config = ConfigDict(extra="forbid", protected_namespaces=())
|
|
||||||
|
|
||||||
enabled: bool
|
|
||||||
provider: str
|
|
||||||
model: str
|
|
||||||
configured: bool
|
|
||||||
|
|
||||||
|
|
||||||
class AdminAffectSessionSummary(BaseModel):
|
|
||||||
model_config = ConfigDict(extra="forbid", protected_namespaces=())
|
|
||||||
|
|
||||||
session_id: str
|
|
||||||
persona_code: str
|
|
||||||
started_at: datetime
|
|
||||||
ended: bool
|
|
||||||
trace_count: int = Field(ge=0)
|
|
||||||
|
|
||||||
|
|
||||||
class AdminAffectSessionListResponse(BaseModel):
|
|
||||||
model_config = ConfigDict(extra="forbid", protected_namespaces=())
|
|
||||||
|
|
||||||
runtime: AdminAffectRuntimeResponse
|
|
||||||
sessions: list[AdminAffectSessionSummary]
|
|
||||||
total: int = Field(ge=0)
|
|
||||||
limit: int = Field(ge=1, le=100)
|
|
||||||
offset: int = Field(ge=0)
|
|
||||||
|
|
||||||
|
|
||||||
class AdminAffectTraceRecord(BaseModel):
|
|
||||||
model_config = ConfigDict(extra="forbid", protected_namespaces=())
|
|
||||||
|
|
||||||
turn_id: str
|
|
||||||
seq: int = Field(ge=1)
|
|
||||||
created_at: datetime
|
|
||||||
trace: ClientAffectTraceV1
|
|
||||||
|
|
||||||
|
|
||||||
class AdminAffectSessionDetailResponse(BaseModel):
|
|
||||||
model_config = ConfigDict(extra="forbid", protected_namespaces=())
|
|
||||||
|
|
||||||
session_id: str
|
|
||||||
persona_code: str
|
|
||||||
current_emotions: dict[str, float | None]
|
|
||||||
traces: list[AdminAffectTraceRecord]
|
|
||||||
total_traces: int = Field(ge=0)
|
|
||||||
has_more: bool
|
|
||||||
|
|
@ -1,100 +0,0 @@
|
||||||
"""관리자 감정 관측용 Jev 감정 전이 trace 계약."""
|
|
||||||
|
|
||||||
from __future__ import annotations
|
|
||||||
|
|
||||||
import math
|
|
||||||
from typing import Literal
|
|
||||||
|
|
||||||
from pydantic import BaseModel, ConfigDict, Field, model_validator
|
|
||||||
|
|
||||||
|
|
||||||
CLIENT_AFFECT_DIMENSIONS = (
|
|
||||||
"anxiety",
|
|
||||||
"sadness",
|
|
||||||
"anger",
|
|
||||||
"shame",
|
|
||||||
"guilt",
|
|
||||||
"loneliness",
|
|
||||||
"relief",
|
|
||||||
"hope",
|
|
||||||
"trust",
|
|
||||||
)
|
|
||||||
|
|
||||||
ClientAffectDecision = Literal["accepted", "tentative", "held"]
|
|
||||||
|
|
||||||
|
|
||||||
class ClientAffectPolicyV1(BaseModel):
|
|
||||||
model_config = ConfigDict(extra="forbid", frozen=True, protected_namespaces=())
|
|
||||||
|
|
||||||
version: Literal["jev-affect-v1"]
|
|
||||||
min_confidence: float = Field(ge=0.0, le=1.0)
|
|
||||||
accepted_alpha: float = Field(ge=0.0, le=1.0)
|
|
||||||
accepted_cap: float = Field(ge=0.0, le=1.0)
|
|
||||||
tentative_alpha: float = Field(ge=0.0, le=1.0)
|
|
||||||
tentative_cap: float = Field(ge=0.0, le=1.0)
|
|
||||||
tentative_confidence_floor: float = Field(ge=0.0, le=1.0)
|
|
||||||
adjacent_probability_threshold: float = Field(ge=0.0, le=1.0)
|
|
||||||
|
|
||||||
|
|
||||||
class ClientAffectContextV1(BaseModel):
|
|
||||||
model_config = ConfigDict(extra="forbid", frozen=True, protected_namespaces=())
|
|
||||||
|
|
||||||
stage: str
|
|
||||||
resistance: float = Field(ge=0.0, le=1.0)
|
|
||||||
effective_openness: float = Field(ge=0.0, le=1.0)
|
|
||||||
rapport_credit: float = Field(ge=0.0)
|
|
||||||
|
|
||||||
|
|
||||||
class ClientAffectDimensionTraceV1(BaseModel):
|
|
||||||
model_config = ConfigDict(extra="forbid", frozen=True, protected_namespaces=())
|
|
||||||
|
|
||||||
key: str
|
|
||||||
before: float = Field(ge=0.0, le=1.0)
|
|
||||||
target: float | None = Field(default=None, ge=0.0, le=1.0)
|
|
||||||
after: float = Field(ge=0.0, le=1.0)
|
|
||||||
confidence: float | None = Field(default=None, ge=0.0, le=1.0)
|
|
||||||
probabilities: tuple[float, float, float, float, float] | None = None
|
|
||||||
decision: ClientAffectDecision
|
|
||||||
|
|
||||||
@model_validator(mode="after")
|
|
||||||
def require_probability_distribution(self) -> "ClientAffectDimensionTraceV1":
|
|
||||||
if self.probabilities is None:
|
|
||||||
return self
|
|
||||||
if any(
|
|
||||||
not math.isfinite(value) or value < 0.0 or value > 1.0
|
|
||||||
for value in self.probabilities
|
|
||||||
):
|
|
||||||
raise ValueError("probabilities must be finite values within 0..1")
|
|
||||||
return self
|
|
||||||
|
|
||||||
|
|
||||||
class ClientAffectTraceV1(BaseModel):
|
|
||||||
model_config = ConfigDict(extra="forbid", frozen=True, protected_namespaces=())
|
|
||||||
|
|
||||||
schema_version: Literal[1]
|
|
||||||
provider: str
|
|
||||||
model: str
|
|
||||||
latency_ms: int = Field(ge=0)
|
|
||||||
input_tokens: int = Field(ge=0)
|
|
||||||
output_tokens: int = Field(ge=0)
|
|
||||||
cost_usd: float | None = Field(default=None, ge=0.0)
|
|
||||||
turn_seq: int = Field(ge=1)
|
|
||||||
policy: ClientAffectPolicyV1
|
|
||||||
context: ClientAffectContextV1
|
|
||||||
dimensions: tuple[ClientAffectDimensionTraceV1, ...] = Field(min_length=9, max_length=9)
|
|
||||||
|
|
||||||
@model_validator(mode="after")
|
|
||||||
def require_fixed_dimension_order(self) -> "ClientAffectTraceV1":
|
|
||||||
if tuple(dimension.key for dimension in self.dimensions) != CLIENT_AFFECT_DIMENSIONS:
|
|
||||||
raise ValueError("dimensions must use the fixed client affect order")
|
|
||||||
return self
|
|
||||||
|
|
||||||
|
|
||||||
__all__ = [
|
|
||||||
"CLIENT_AFFECT_DIMENSIONS",
|
|
||||||
"ClientAffectContextV1",
|
|
||||||
"ClientAffectDecision",
|
|
||||||
"ClientAffectDimensionTraceV1",
|
|
||||||
"ClientAffectPolicyV1",
|
|
||||||
"ClientAffectTraceV1",
|
|
||||||
]
|
|
||||||
|
|
@ -25,7 +25,6 @@ from .session_persistence import ensure_review_tables
|
||||||
from .services.jev_client import jev_client
|
from .services.jev_client import jev_client
|
||||||
from .runtime_schema import (
|
from .runtime_schema import (
|
||||||
CALIBRATION_TRANSFER_SCHEMA_CONTRACT,
|
CALIBRATION_TRANSFER_SCHEMA_CONTRACT,
|
||||||
CLIENT_AFFECT_TRACE_SCHEMA_CONTRACT,
|
|
||||||
CONTINUOUS_IMPROVEMENT_SCHEMA_CONTRACT,
|
CONTINUOUS_IMPROVEMENT_SCHEMA_CONTRACT,
|
||||||
DELIBERATE_PRACTICE_SCHEMA_CONTRACT,
|
DELIBERATE_PRACTICE_SCHEMA_CONTRACT,
|
||||||
MEASUREMENT_SCHEMA_CONTRACT,
|
MEASUREMENT_SCHEMA_CONTRACT,
|
||||||
|
|
@ -91,7 +90,6 @@ async def lifespan(app: FastAPI):
|
||||||
supervision_research_schema_ready = False
|
supervision_research_schema_ready = False
|
||||||
continuous_improvement_schema_ready = False
|
continuous_improvement_schema_ready = False
|
||||||
multimodal_alliance_schema_ready = False
|
multimodal_alliance_schema_ready = False
|
||||||
client_affect_trace_schema_ready = False
|
|
||||||
app.state.upload_database_proof = None
|
app.state.upload_database_proof = None
|
||||||
try:
|
try:
|
||||||
await init_pool()
|
await init_pool()
|
||||||
|
|
@ -104,10 +102,6 @@ async def lifespan(app: FastAPI):
|
||||||
await ensure_notification_tables()
|
await ensure_notification_tables()
|
||||||
await ensure_protocol_tables()
|
await ensure_protocol_tables()
|
||||||
async with acquire(role="admin") as conn:
|
async with acquire(role="admin") as conn:
|
||||||
client_affect_trace_schema_ready = await schema_contract_ready(
|
|
||||||
conn,
|
|
||||||
CLIENT_AFFECT_TRACE_SCHEMA_CONTRACT,
|
|
||||||
)
|
|
||||||
measurement_schema_ready = await schema_contract_ready(
|
measurement_schema_ready = await schema_contract_ready(
|
||||||
conn,
|
conn,
|
||||||
MEASUREMENT_SCHEMA_CONTRACT,
|
MEASUREMENT_SCHEMA_CONTRACT,
|
||||||
|
|
@ -140,14 +134,6 @@ async def lifespan(app: FastAPI):
|
||||||
conn,
|
conn,
|
||||||
MULTIMODAL_ALLIANCE_SCHEMA_CONTRACT,
|
MULTIMODAL_ALLIANCE_SCHEMA_CONTRACT,
|
||||||
)
|
)
|
||||||
if runtime_schema_bootstrap_required(
|
|
||||||
CLIENT_AFFECT_TRACE_SCHEMA_CONTRACT,
|
|
||||||
ready=client_affect_trace_schema_ready,
|
|
||||||
):
|
|
||||||
logger.warning(
|
|
||||||
"감정 관측 trace 스키마가 불완전해 Jev 전이 trace를 저장할 수 없음; "
|
|
||||||
"infra/db/init/23_client_affect_trace.sql 적용 필요"
|
|
||||||
)
|
|
||||||
if runtime_schema_bootstrap_required(
|
if runtime_schema_bootstrap_required(
|
||||||
MEASUREMENT_SCHEMA_CONTRACT,
|
MEASUREMENT_SCHEMA_CONTRACT,
|
||||||
ready=measurement_schema_ready,
|
ready=measurement_schema_ready,
|
||||||
|
|
|
||||||
|
|
@ -25,10 +25,6 @@ from ..auth_sessions import (
|
||||||
upsert_managed_user,
|
upsert_managed_user,
|
||||||
)
|
)
|
||||||
from ..config import settings
|
from ..config import settings
|
||||||
from ..contracts.admin_affect import (
|
|
||||||
AdminAffectSessionDetailResponse,
|
|
||||||
AdminAffectSessionListResponse,
|
|
||||||
)
|
|
||||||
from ..contracts.engine_gateway import (
|
from ..contracts.engine_gateway import (
|
||||||
ENGINE_PROVIDER_DEFAULTS,
|
ENGINE_PROVIDER_DEFAULTS,
|
||||||
ENGINE_PROVIDERS,
|
ENGINE_PROVIDERS,
|
||||||
|
|
@ -42,7 +38,6 @@ from ..deps import Principal, require_admin_access
|
||||||
from ..engine_client import engine_client
|
from ..engine_client import engine_client
|
||||||
from ..runtime_policy import require_runtime_fallback_allowed
|
from ..runtime_policy import require_runtime_fallback_allowed
|
||||||
from ..services import evaluator, notifications, rag
|
from ..services import evaluator, notifications, rag
|
||||||
from ..services import admin_affect as admin_affect_service
|
|
||||||
from ..services import provider_credentials as provider_credentials_service
|
from ..services import provider_credentials as provider_credentials_service
|
||||||
from ..services import provider_oauth as provider_oauth_service
|
from ..services import provider_oauth as provider_oauth_service
|
||||||
from ..services.llm_pricing import (
|
from ..services.llm_pricing import (
|
||||||
|
|
@ -1939,58 +1934,6 @@ async def admin_voice_runtime(
|
||||||
return voice_runtime_metrics.snapshot()
|
return voice_runtime_metrics.snapshot()
|
||||||
|
|
||||||
|
|
||||||
@router.get("/affect/sessions", response_model=AdminAffectSessionListResponse)
|
|
||||||
async def list_admin_affect_sessions(
|
|
||||||
principal: AdminPrincipal,
|
|
||||||
limit: Annotated[int, Query(ge=1, le=100)] = 30,
|
|
||||||
offset: Annotated[int, Query(ge=0)] = 0,
|
|
||||||
) -> AdminAffectSessionListResponse:
|
|
||||||
"""관리자 전용 감정 관측 회기 목록을 반환한다."""
|
|
||||||
|
|
||||||
try:
|
|
||||||
return await admin_affect_service.list_sessions(
|
|
||||||
user_id=principal.user_id,
|
|
||||||
limit=limit,
|
|
||||||
offset=offset,
|
|
||||||
)
|
|
||||||
except admin_affect_service.AdminAffectPersistenceError as exc:
|
|
||||||
raise HTTPException(
|
|
||||||
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
|
|
||||||
detail="admin affect sessions are unavailable",
|
|
||||||
) from exc
|
|
||||||
|
|
||||||
|
|
||||||
@router.get(
|
|
||||||
"/affect/sessions/{session_id}",
|
|
||||||
response_model=AdminAffectSessionDetailResponse,
|
|
||||||
)
|
|
||||||
async def get_admin_affect_session_detail(
|
|
||||||
session_id: UUID,
|
|
||||||
principal: AdminPrincipal,
|
|
||||||
limit: Annotated[int, Query(ge=1, le=200)] = 100,
|
|
||||||
before_seq: Annotated[int | None, Query(gt=0)] = None,
|
|
||||||
) -> AdminAffectSessionDetailResponse:
|
|
||||||
"""저장된 Jev trace와 현재 감정 snapshot을 반환한다."""
|
|
||||||
|
|
||||||
try:
|
|
||||||
return await admin_affect_service.get_session_detail(
|
|
||||||
user_id=principal.user_id,
|
|
||||||
session_id=session_id,
|
|
||||||
limit=limit,
|
|
||||||
before_seq=before_seq,
|
|
||||||
)
|
|
||||||
except admin_affect_service.AdminAffectSessionNotFoundError as exc:
|
|
||||||
raise HTTPException(
|
|
||||||
status_code=status.HTTP_404_NOT_FOUND,
|
|
||||||
detail="admin affect session not found",
|
|
||||||
) from exc
|
|
||||||
except admin_affect_service.AdminAffectPersistenceError as exc:
|
|
||||||
raise HTTPException(
|
|
||||||
status_code=status.HTTP_503_SERVICE_UNAVAILABLE,
|
|
||||||
detail="admin affect session is unavailable",
|
|
||||||
) from exc
|
|
||||||
|
|
||||||
|
|
||||||
@router.get("/usage", response_model=AdminUsageResponse)
|
@router.get("/usage", response_model=AdminUsageResponse)
|
||||||
async def admin_usage(
|
async def admin_usage(
|
||||||
principal: AdminPrincipal,
|
principal: AdminPrincipal,
|
||||||
|
|
|
||||||
|
|
@ -26,23 +26,6 @@ class RuntimeSchemaContract:
|
||||||
indexes: tuple[str, ...] = ()
|
indexes: tuple[str, ...] = ()
|
||||||
|
|
||||||
|
|
||||||
CLIENT_AFFECT_TRACE_SCHEMA_CONTRACT = RuntimeSchemaContract(
|
|
||||||
component="client affect trace",
|
|
||||||
relations=("app.client_affect_trace",),
|
|
||||||
columns=(
|
|
||||||
"app.client_affect_trace.turn_id",
|
|
||||||
"app.client_affect_trace.session_id",
|
|
||||||
"app.client_affect_trace.trace",
|
|
||||||
"app.client_affect_trace.created_at",
|
|
||||||
),
|
|
||||||
policies=(
|
|
||||||
"app.client_affect_trace.p_client_affect_trace_select_admin",
|
|
||||||
"app.client_affect_trace.p_client_affect_trace_insert_learner",
|
|
||||||
),
|
|
||||||
indexes=("app.client_affect_trace.idx_client_affect_trace_session_created",),
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
REVIEW_SCHEMA_CONTRACT = RuntimeSchemaContract(
|
REVIEW_SCHEMA_CONTRACT = RuntimeSchemaContract(
|
||||||
component="review/evaluation",
|
component="review/evaluation",
|
||||||
relations=(
|
relations=(
|
||||||
|
|
|
||||||
|
|
@ -1,192 +0,0 @@
|
||||||
"""관리자 감정 관측용 영속 조회."""
|
|
||||||
|
|
||||||
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,
|
|
||||||
)
|
|
||||||
|
|
@ -6,12 +6,6 @@ import math
|
||||||
from dataclasses import dataclass
|
from dataclasses import dataclass
|
||||||
from typing import Any, Iterable, Mapping
|
from typing import Any, Iterable, Mapping
|
||||||
|
|
||||||
from ..contracts.client_affect import (
|
|
||||||
ClientAffectContextV1,
|
|
||||||
ClientAffectDimensionTraceV1,
|
|
||||||
ClientAffectPolicyV1,
|
|
||||||
ClientAffectTraceV1,
|
|
||||||
)
|
|
||||||
from . import guardrail
|
from . import guardrail
|
||||||
from .jev_client import AppraisalResult, EMOTION_DIMENSIONS
|
from .jev_client import AppraisalResult, EMOTION_DIMENSIONS
|
||||||
|
|
||||||
|
|
@ -46,13 +40,7 @@ _NEGATIVE_EMOTIONS = frozenset(
|
||||||
{"anxiety", "sadness", "anger", "shame", "guilt", "loneliness"}
|
{"anxiety", "sadness", "anger", "shame", "guilt", "loneliness"}
|
||||||
)
|
)
|
||||||
_POSITIVE_EMOTIONS = frozenset({"relief", "hope", "trust"})
|
_POSITIVE_EMOTIONS = frozenset({"relief", "hope", "trust"})
|
||||||
_AFFECT_POLICY_VERSION = "jev-affect-v1"
|
|
||||||
_ACCEPTED_ALPHA = 0.35
|
|
||||||
_ACCEPTED_CAP = 0.15
|
|
||||||
_TENTATIVE_ALPHA = 0.15
|
|
||||||
_TENTATIVE_CAP = 0.075
|
|
||||||
_TENTATIVE_CONFIDENCE_FLOOR = 0.35
|
_TENTATIVE_CONFIDENCE_FLOOR = 0.35
|
||||||
_ADJACENT_PROBABILITY_THRESHOLD = 0.8
|
|
||||||
_PROBABILITY_SUM_TOLERANCE = 0.025000001
|
_PROBABILITY_SUM_TOLERANCE = 0.025000001
|
||||||
|
|
||||||
|
|
||||||
|
|
@ -95,10 +83,7 @@ def _tentative_distribution_is_concentrated(probabilities: Any) -> bool:
|
||||||
if not math.isclose(total, 1.0, abs_tol=_PROBABILITY_SUM_TOLERANCE):
|
if not math.isclose(total, 1.0, abs_tol=_PROBABILITY_SUM_TOLERANCE):
|
||||||
return False
|
return False
|
||||||
normalized = tuple(value / total for value in values if value is not None)
|
normalized = tuple(value / total for value in values if value is not None)
|
||||||
return (
|
return max(normalized[index] + normalized[index + 1] for index in range(4)) >= 0.80
|
||||||
max(normalized[index] + normalized[index + 1] for index in range(4))
|
|
||||||
>= _ADJACENT_PROBABILITY_THRESHOLD
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def _baseline_value(affect_baseline: Mapping[str, Any], key: str) -> float | None:
|
def _baseline_value(affect_baseline: Mapping[str, Any], key: str) -> float | None:
|
||||||
|
|
@ -178,15 +163,15 @@ def transition_emotions(
|
||||||
held.append(dimension)
|
held.append(dimension)
|
||||||
continue
|
continue
|
||||||
if confidence >= threshold:
|
if confidence >= threshold:
|
||||||
alpha = _ACCEPTED_ALPHA
|
alpha = 0.35
|
||||||
cap = _ACCEPTED_CAP
|
cap = 0.15
|
||||||
elif (
|
elif (
|
||||||
confidence >= _TENTATIVE_CONFIDENCE_FLOOR
|
confidence >= _TENTATIVE_CONFIDENCE_FLOOR
|
||||||
and _tentative_distribution_is_concentrated(estimate.probabilities)
|
and _tentative_distribution_is_concentrated(estimate.probabilities)
|
||||||
):
|
):
|
||||||
# confidence는 정답 확률이 아니라 분포 집중도 요약이다.
|
# confidence는 정답 확률이 아니라 분포 집중도 요약이다.
|
||||||
alpha = _TENTATIVE_ALPHA
|
alpha = 0.15
|
||||||
cap = _TENTATIVE_CAP
|
cap = 0.075
|
||||||
tentative.append(dimension)
|
tentative.append(dimension)
|
||||||
else:
|
else:
|
||||||
updated[f"emotion_{dimension}"] = old
|
updated[f"emotion_{dimension}"] = old
|
||||||
|
|
@ -204,94 +189,6 @@ def transition_emotions(
|
||||||
)
|
)
|
||||||
|
|
||||||
|
|
||||||
def _trace_probabilities(value: Any) -> tuple[float, float, float, float, float] | None:
|
|
||||||
if not isinstance(value, tuple) or len(value) != 5:
|
|
||||||
return None
|
|
||||||
normalized = tuple(_unit_number(item) for item in value)
|
|
||||||
if any(item is None for item in normalized):
|
|
||||||
return None
|
|
||||||
return (
|
|
||||||
normalized[0],
|
|
||||||
normalized[1],
|
|
||||||
normalized[2],
|
|
||||||
normalized[3],
|
|
||||||
normalized[4],
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def build_client_affect_trace(
|
|
||||||
*,
|
|
||||||
affect_state_before: Mapping[str, Any],
|
|
||||||
affect_baseline: Mapping[str, Any],
|
|
||||||
affect_state_after: Mapping[str, Any],
|
|
||||||
appraisal: AppraisalResult,
|
|
||||||
transition: AffectTransition,
|
|
||||||
turn_seq: int,
|
|
||||||
stage: str,
|
|
||||||
resistance: float,
|
|
||||||
effective_openness: float,
|
|
||||||
rapport_credit: float,
|
|
||||||
min_confidence: float,
|
|
||||||
) -> ClientAffectTraceV1:
|
|
||||||
"""전이와 같은 입력으로 관리자 전용 trace를 고정 순서로 만든다."""
|
|
||||||
before = resolve_emotions(affect_state_before, affect_baseline)
|
|
||||||
after = resolve_emotions(affect_state_after, affect_baseline)
|
|
||||||
tentative = set(transition.tentative_dimensions)
|
|
||||||
accepted = set(transition.accepted_dimensions)
|
|
||||||
dimensions: list[ClientAffectDimensionTraceV1] = []
|
|
||||||
for key in EMOTION_DIMENSIONS:
|
|
||||||
estimate = appraisal.emotions.get(key)
|
|
||||||
target = _unit_number(estimate.score) if estimate is not None else None
|
|
||||||
confidence = _unit_number(estimate.confidence) if estimate is not None else None
|
|
||||||
probabilities = (
|
|
||||||
_trace_probabilities(estimate.probabilities) if estimate is not None else None
|
|
||||||
)
|
|
||||||
if key in tentative:
|
|
||||||
decision = "tentative"
|
|
||||||
elif key in accepted:
|
|
||||||
decision = "accepted"
|
|
||||||
else:
|
|
||||||
decision = "held"
|
|
||||||
dimensions.append(
|
|
||||||
ClientAffectDimensionTraceV1(
|
|
||||||
key=key,
|
|
||||||
before=before[key],
|
|
||||||
target=target,
|
|
||||||
after=after[key],
|
|
||||||
confidence=confidence,
|
|
||||||
probabilities=probabilities,
|
|
||||||
decision=decision,
|
|
||||||
)
|
|
||||||
)
|
|
||||||
return ClientAffectTraceV1(
|
|
||||||
schema_version=1,
|
|
||||||
provider=appraisal.provider,
|
|
||||||
model=appraisal.model,
|
|
||||||
latency_ms=appraisal.latency_ms,
|
|
||||||
input_tokens=appraisal.input_tokens,
|
|
||||||
output_tokens=appraisal.output_tokens,
|
|
||||||
cost_usd=appraisal.cost_usd,
|
|
||||||
turn_seq=turn_seq,
|
|
||||||
policy=ClientAffectPolicyV1(
|
|
||||||
version=_AFFECT_POLICY_VERSION,
|
|
||||||
min_confidence=min_confidence,
|
|
||||||
accepted_alpha=_ACCEPTED_ALPHA,
|
|
||||||
accepted_cap=_ACCEPTED_CAP,
|
|
||||||
tentative_alpha=_TENTATIVE_ALPHA,
|
|
||||||
tentative_cap=_TENTATIVE_CAP,
|
|
||||||
tentative_confidence_floor=_TENTATIVE_CONFIDENCE_FLOOR,
|
|
||||||
adjacent_probability_threshold=_ADJACENT_PROBABILITY_THRESHOLD,
|
|
||||||
),
|
|
||||||
context=ClientAffectContextV1(
|
|
||||||
stage=stage,
|
|
||||||
resistance=resistance,
|
|
||||||
effective_openness=effective_openness,
|
|
||||||
rapport_credit=rapport_credit,
|
|
||||||
),
|
|
||||||
dimensions=tuple(dimensions),
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def _mask_text(
|
def _mask_text(
|
||||||
value: Any,
|
value: Any,
|
||||||
*,
|
*,
|
||||||
|
|
@ -487,7 +384,6 @@ def public_end_state(end_state: Mapping[str, Any]) -> dict[str, Any]:
|
||||||
__all__ = [
|
__all__ = [
|
||||||
"AffectTransition",
|
"AffectTransition",
|
||||||
"baseline_emotions",
|
"baseline_emotions",
|
||||||
"build_client_affect_trace",
|
|
||||||
"build_appraisal_state",
|
"build_appraisal_state",
|
||||||
"public_end_state",
|
"public_end_state",
|
||||||
"resolve_emotions",
|
"resolve_emotions",
|
||||||
|
|
|
||||||
|
|
@ -23,7 +23,6 @@ import time
|
||||||
from dataclasses import dataclass, field, replace
|
from dataclasses import dataclass, field, replace
|
||||||
from typing import Any, AsyncIterator, Awaitable, Callable, Optional
|
from typing import Any, AsyncIterator, Awaitable, Callable, Optional
|
||||||
|
|
||||||
from ..contracts.client_affect import ClientAffectTraceV1
|
|
||||||
from ..config import settings
|
from ..config import settings
|
||||||
from ..engine_client import (
|
from ..engine_client import (
|
||||||
EngineClient,
|
EngineClient,
|
||||||
|
|
@ -95,8 +94,6 @@ class TurnContext:
|
||||||
scenario_directive: Optional[rupture_scenario_director.ScenarioDirective] = None
|
scenario_directive: Optional[rupture_scenario_director.ScenarioDirective] = None
|
||||||
# 외부 감정 평가의 안전한 provenance. 원문·점수·확률은 넣지 않는다.
|
# 외부 감정 평가의 안전한 provenance. 원문·점수·확률은 넣지 않는다.
|
||||||
client_affect_metadata: Optional[dict[str, Any]] = None
|
client_affect_metadata: Optional[dict[str, Any]] = None
|
||||||
# 관리자 관측 전용 Jev 전이 trace. 공개 결과나 provider event에는 넣지 않는다.
|
|
||||||
client_affect_trace: ClientAffectTraceV1 | None = None
|
|
||||||
|
|
||||||
def to_state_context(self) -> PersonaStateContext:
|
def to_state_context(self) -> PersonaStateContext:
|
||||||
st = self.state_after or self.state_before
|
st = self.state_after or self.state_before
|
||||||
|
|
@ -381,27 +378,13 @@ async def _apply_client_affect(
|
||||||
client_identity=ctx.client_identity,
|
client_identity=ctx.client_identity,
|
||||||
)
|
)
|
||||||
appraisal = await jev_client.appraise(state)
|
appraisal = await jev_client.appraise(state)
|
||||||
state_before_transition = ctx.state_after
|
|
||||||
transition = client_affect.transition_emotions(
|
transition = client_affect.transition_emotions(
|
||||||
state_before_transition.affect_state,
|
ctx.state_after.affect_state,
|
||||||
ctx.persona.affect_baseline,
|
ctx.persona.affect_baseline,
|
||||||
appraisal,
|
appraisal,
|
||||||
min_confidence=settings.jev_min_confidence,
|
min_confidence=settings.jev_min_confidence,
|
||||||
)
|
)
|
||||||
ctx.state_after = replace(ctx.state_after, affect_state=transition.affect_state)
|
ctx.state_after = replace(ctx.state_after, affect_state=transition.affect_state)
|
||||||
ctx.client_affect_trace = client_affect.build_client_affect_trace(
|
|
||||||
affect_state_before=state_before_transition.affect_state,
|
|
||||||
affect_baseline=ctx.persona.affect_baseline,
|
|
||||||
affect_state_after=ctx.state_after.affect_state,
|
|
||||||
appraisal=appraisal,
|
|
||||||
transition=transition,
|
|
||||||
turn_seq=ctx.state_after.turn_seq,
|
|
||||||
stage=ctx.state_after.stage.value,
|
|
||||||
resistance=ctx.state_after.resistance,
|
|
||||||
effective_openness=ctx.state_after.effective_openness,
|
|
||||||
rapport_credit=ctx.state_after.rapport_credit,
|
|
||||||
min_confidence=settings.jev_min_confidence,
|
|
||||||
)
|
|
||||||
ctx.client_affect_metadata = {
|
ctx.client_affect_metadata = {
|
||||||
"provider": appraisal.provider,
|
"provider": appraisal.provider,
|
||||||
"model": appraisal.model,
|
"model": appraisal.model,
|
||||||
|
|
|
||||||
|
|
@ -14,7 +14,6 @@ from typing import Any, Awaitable, Callable, Iterable, Literal
|
||||||
from .db import acquire, get_pool
|
from .db import acquire, get_pool
|
||||||
from .deps import Principal
|
from .deps import Principal
|
||||||
from .config import settings
|
from .config import settings
|
||||||
from .contracts.client_affect import ClientAffectTraceV1
|
|
||||||
from .persona_repository import (
|
from .persona_repository import (
|
||||||
SEED_VERSION,
|
SEED_VERSION,
|
||||||
card_from_row,
|
card_from_row,
|
||||||
|
|
@ -45,10 +44,6 @@ class SessionCreationPersistenceError(RuntimeError):
|
||||||
"""fail-closed 환경에서 영속 세션 생성이 실패했다."""
|
"""fail-closed 환경에서 영속 세션 생성이 실패했다."""
|
||||||
|
|
||||||
|
|
||||||
class ClientAffectTracePersistenceError(RuntimeError):
|
|
||||||
"""내담자 턴·감정 trace·상태를 같은 트랜잭션으로 기록하지 못했다."""
|
|
||||||
|
|
||||||
|
|
||||||
class ActiveSessionExistsError(RuntimeError):
|
class ActiveSessionExistsError(RuntimeError):
|
||||||
"""같은 learner-persona 전체에 미종료 회기가 이미 존재한다."""
|
"""같은 learner-persona 전체에 미종료 회기가 이미 존재한다."""
|
||||||
|
|
||||||
|
|
@ -2897,100 +2892,6 @@ async def append_turn(
|
||||||
return False
|
return False
|
||||||
|
|
||||||
|
|
||||||
async def append_client_turn_with_affect_trace(
|
|
||||||
*,
|
|
||||||
session_id: str,
|
|
||||||
learner_id: str,
|
|
||||||
turn: TurnRecord,
|
|
||||||
state: state_machine.SessionState,
|
|
||||||
trace: ClientAffectTraceV1,
|
|
||||||
) -> bool:
|
|
||||||
"""성공 내담자 턴·Jev trace·전이 상태를 하나의 DB 트랜잭션에 기록한다."""
|
|
||||||
inserted_turn_id: str | None = None
|
|
||||||
try:
|
|
||||||
get_pool()
|
|
||||||
async with acquire(role="learner", user_id=learner_id) as conn:
|
|
||||||
async with conn.transaction():
|
|
||||||
locked = await conn.fetchval(
|
|
||||||
"""
|
|
||||||
SELECT id
|
|
||||||
FROM app.sessions
|
|
||||||
WHERE id = $1::uuid
|
|
||||||
AND learner_id = $2::uuid
|
|
||||||
FOR UPDATE
|
|
||||||
""",
|
|
||||||
session_id,
|
|
||||||
learner_id,
|
|
||||||
)
|
|
||||||
if locked is None:
|
|
||||||
raise ClientAffectTracePersistenceError("owned session not found")
|
|
||||||
seq = int(
|
|
||||||
await conn.fetchval(
|
|
||||||
"SELECT COALESCE(MAX(seq), 0) + 1 FROM app.turns WHERE session_id = $1::uuid",
|
|
||||||
session_id,
|
|
||||||
)
|
|
||||||
or 1
|
|
||||||
)
|
|
||||||
inserted = await conn.fetchval(
|
|
||||||
"""
|
|
||||||
INSERT INTO app.turns (
|
|
||||||
session_id, seq, speaker, stage, text, text_masked, actor_kind,
|
|
||||||
llm_provider, model, tokens_in, tokens_out, cost_usd,
|
|
||||||
audio_ref, silence_ms, speech_rate, barge_in, provider_events, visible_to
|
|
||||||
)
|
|
||||||
VALUES (
|
|
||||||
$1::uuid, $2, $3, $4, $5, $6, $7,
|
|
||||||
$8, $9, $10, $11, $12,
|
|
||||||
$13, $14, $15, $16, $17::jsonb, $18::text[]
|
|
||||||
)
|
|
||||||
ON CONFLICT (session_id, seq) DO NOTHING
|
|
||||||
RETURNING id
|
|
||||||
""",
|
|
||||||
session_id,
|
|
||||||
seq,
|
|
||||||
turn.speaker,
|
|
||||||
turn.stage,
|
|
||||||
turn.text_masked,
|
|
||||||
turn.text_masked,
|
|
||||||
"client_ai",
|
|
||||||
turn.llm_provider,
|
|
||||||
turn.model,
|
|
||||||
turn.tokens_in,
|
|
||||||
turn.tokens_out,
|
|
||||||
turn.cost_usd,
|
|
||||||
turn.audio_ref,
|
|
||||||
turn.silence_ms,
|
|
||||||
turn.speech_rate,
|
|
||||||
turn.barge_in,
|
|
||||||
turn.provider_events or [],
|
|
||||||
list(turn.visible_to or DEFAULT_TURN_VISIBLE_TO),
|
|
||||||
)
|
|
||||||
if inserted is None:
|
|
||||||
raise ClientAffectTracePersistenceError("client turn insert conflict")
|
|
||||||
inserted_turn_id = str(inserted)
|
|
||||||
await conn.execute(
|
|
||||||
"""
|
|
||||||
INSERT INTO app.client_affect_trace (turn_id, session_id, trace)
|
|
||||||
VALUES ($1::uuid, $2::uuid, $3::jsonb)
|
|
||||||
""",
|
|
||||||
inserted_turn_id,
|
|
||||||
session_id,
|
|
||||||
trace.model_dump(mode="json"),
|
|
||||||
)
|
|
||||||
await _upsert_state(conn, session_id, state)
|
|
||||||
except ClientAffectTracePersistenceError:
|
|
||||||
raise
|
|
||||||
except Exception as exc:
|
|
||||||
logger.exception(
|
|
||||||
"client affect trace persistence failed: session_id=%s", session_id
|
|
||||||
)
|
|
||||||
raise ClientAffectTracePersistenceError(
|
|
||||||
"client affect trace persistence failed"
|
|
||||||
) from exc
|
|
||||||
turn.turn_id = inserted_turn_id
|
|
||||||
return True
|
|
||||||
|
|
||||||
|
|
||||||
async def update_state(
|
async def update_state(
|
||||||
*,
|
*,
|
||||||
session_id: str,
|
session_id: str,
|
||||||
|
|
|
||||||
|
|
@ -1,327 +0,0 @@
|
||||||
"""관리자 감정 관측 조회의 권한·페이지·영속 경계 검증."""
|
|
||||||
|
|
||||||
from __future__ import annotations
|
|
||||||
|
|
||||||
import unittest
|
|
||||||
from datetime import datetime, timezone
|
|
||||||
from unittest.mock import AsyncMock, patch
|
|
||||||
from uuid import UUID
|
|
||||||
|
|
||||||
from fastapi import FastAPI, HTTPException
|
|
||||||
from fastapi.testclient import TestClient
|
|
||||||
|
|
||||||
from .contracts.admin_affect import (
|
|
||||||
AdminAffectRuntimeResponse,
|
|
||||||
AdminAffectSessionListResponse,
|
|
||||||
)
|
|
||||||
from .deps import Principal, Role, get_current_principal
|
|
||||||
from .routes import admin as admin_routes
|
|
||||||
from .services import admin_affect
|
|
||||||
|
|
||||||
|
|
||||||
SESSION_ID = UUID("00000000-0000-0000-0000-000000000101")
|
|
||||||
TURN_ID = UUID("00000000-0000-0000-0000-000000000201")
|
|
||||||
ADMIN_ID = "00000000-0000-0000-0000-000000000901"
|
|
||||||
OBSERVED_AT = datetime(2026, 9, 23, 1, 2, 3, tzinfo=timezone.utc)
|
|
||||||
|
|
||||||
|
|
||||||
class _Acquire:
|
|
||||||
def __init__(self, conn: object) -> None:
|
|
||||||
self.conn = conn
|
|
||||||
|
|
||||||
async def __aenter__(self) -> object:
|
|
||||||
return self.conn
|
|
||||||
|
|
||||||
async def __aexit__(self, exc_type: object, exc: object, tb: object) -> None:
|
|
||||||
return None
|
|
||||||
|
|
||||||
|
|
||||||
def _trace(*, seq: int) -> dict[str, object]:
|
|
||||||
dimensions = []
|
|
||||||
for key in (
|
|
||||||
"anxiety",
|
|
||||||
"sadness",
|
|
||||||
"anger",
|
|
||||||
"shame",
|
|
||||||
"guilt",
|
|
||||||
"loneliness",
|
|
||||||
"relief",
|
|
||||||
"hope",
|
|
||||||
"trust",
|
|
||||||
):
|
|
||||||
dimensions.append(
|
|
||||||
{
|
|
||||||
"key": key,
|
|
||||||
"before": 0.4,
|
|
||||||
"target": 0.6,
|
|
||||||
"after": 0.47,
|
|
||||||
"confidence": 0.8,
|
|
||||||
"probabilities": [0.05, 0.1, 0.2, 0.35, 0.3],
|
|
||||||
"decision": "accepted",
|
|
||||||
}
|
|
||||||
)
|
|
||||||
return {
|
|
||||||
"schema_version": 1,
|
|
||||||
"provider": "openrouter",
|
|
||||||
"model": "~typesafe/jev-latest",
|
|
||||||
"latency_ms": 91,
|
|
||||||
"input_tokens": 12,
|
|
||||||
"output_tokens": 18,
|
|
||||||
"cost_usd": None,
|
|
||||||
"turn_seq": seq,
|
|
||||||
"policy": {
|
|
||||||
"version": "jev-affect-v1",
|
|
||||||
"min_confidence": 0.65,
|
|
||||||
"accepted_alpha": 0.35,
|
|
||||||
"accepted_cap": 0.15,
|
|
||||||
"tentative_alpha": 0.15,
|
|
||||||
"tentative_cap": 0.075,
|
|
||||||
"tentative_confidence_floor": 0.35,
|
|
||||||
"adjacent_probability_threshold": 0.8,
|
|
||||||
},
|
|
||||||
"context": {
|
|
||||||
"stage": "탐색",
|
|
||||||
"resistance": 0.4,
|
|
||||||
"effective_openness": 0.6,
|
|
||||||
"rapport_credit": 0.2,
|
|
||||||
},
|
|
||||||
"dimensions": dimensions,
|
|
||||||
}
|
|
||||||
|
|
||||||
|
|
||||||
def _admin_principal() -> Principal:
|
|
||||||
return Principal(user_id=ADMIN_ID, role=Role.ADMIN)
|
|
||||||
|
|
||||||
|
|
||||||
class AdminAffectStoreTest(unittest.IsolatedAsyncioTestCase):
|
|
||||||
async def test_list_uses_admin_scope_and_excludes_learner_identity(self) -> None:
|
|
||||||
case = self
|
|
||||||
|
|
||||||
class Conn:
|
|
||||||
async def fetchrow(self, query: str, *args: object) -> dict[str, object]:
|
|
||||||
case.assertIn("FROM app.sessions", query)
|
|
||||||
case.assertNotIn("learner", query.lower())
|
|
||||||
return {"total": 1}
|
|
||||||
|
|
||||||
async def fetch(self, query: str, *args: object) -> list[dict[str, object]]:
|
|
||||||
case.assertIn("app.client_affect_trace", query)
|
|
||||||
case.assertNotIn("learner_id", query)
|
|
||||||
case.assertNotIn("app.app_user", query)
|
|
||||||
case.assertNotIn(" text", query.lower())
|
|
||||||
case.assertEqual(args, (30, 0))
|
|
||||||
return [
|
|
||||||
{
|
|
||||||
"session_id": SESSION_ID,
|
|
||||||
"persona_code": "P4",
|
|
||||||
"started_at": OBSERVED_AT,
|
|
||||||
"ended": False,
|
|
||||||
"trace_count": 2,
|
|
||||||
}
|
|
||||||
]
|
|
||||||
|
|
||||||
runtime = AdminAffectRuntimeResponse(
|
|
||||||
enabled=True,
|
|
||||||
provider="openrouter",
|
|
||||||
model="~typesafe/jev-latest",
|
|
||||||
configured=True,
|
|
||||||
)
|
|
||||||
with (
|
|
||||||
patch.object(admin_affect, "acquire", return_value=_Acquire(Conn())) as acquire,
|
|
||||||
patch.object(admin_affect, "runtime_snapshot", return_value=runtime),
|
|
||||||
):
|
|
||||||
response = await admin_affect.list_sessions(
|
|
||||||
user_id=ADMIN_ID,
|
|
||||||
limit=30,
|
|
||||||
offset=0,
|
|
||||||
)
|
|
||||||
|
|
||||||
acquire.assert_called_once_with(role="admin", user_id=ADMIN_ID)
|
|
||||||
self.assertEqual(response.total, 1)
|
|
||||||
self.assertEqual(response.sessions[0].session_id, str(SESSION_ID))
|
|
||||||
self.assertEqual(response.sessions[0].trace_count, 2)
|
|
||||||
|
|
||||||
async def test_list_preserves_total_beyond_the_last_page(self) -> None:
|
|
||||||
case = self
|
|
||||||
|
|
||||||
class Conn:
|
|
||||||
async def fetchrow(self, query: str, *args: object) -> dict[str, object]:
|
|
||||||
return {"total": 4}
|
|
||||||
|
|
||||||
async def fetch(self, query: str, *args: object) -> list[dict[str, object]]:
|
|
||||||
case.assertEqual(args, (30, 30))
|
|
||||||
return []
|
|
||||||
|
|
||||||
with patch.object(admin_affect, "acquire", return_value=_Acquire(Conn())):
|
|
||||||
response = await admin_affect.list_sessions(
|
|
||||||
user_id=ADMIN_ID,
|
|
||||||
limit=30,
|
|
||||||
offset=30,
|
|
||||||
)
|
|
||||||
|
|
||||||
self.assertEqual(response.total, 4)
|
|
||||||
self.assertEqual(response.sessions, [])
|
|
||||||
|
|
||||||
async def test_detail_returns_actual_snapshot_and_ascending_page(self) -> None:
|
|
||||||
case = self
|
|
||||||
|
|
||||||
class Conn:
|
|
||||||
async def fetchrow(self, query: str, *args: object) -> dict[str, object]:
|
|
||||||
if "FROM app.sessions AS s" in query:
|
|
||||||
case.assertNotIn("learner", query.lower())
|
|
||||||
case.assertNotIn("text", query.lower())
|
|
||||||
return {
|
|
||||||
"session_id": SESSION_ID,
|
|
||||||
"persona_code": "P4",
|
|
||||||
"affect_state": {
|
|
||||||
"emotion_anxiety": 0.6,
|
|
||||||
"emotion_hope": 0.2,
|
|
||||||
"emotion_anger": 1.5,
|
|
||||||
},
|
|
||||||
}
|
|
||||||
case.assertIn("count(*)", query)
|
|
||||||
return {"total_traces": 4}
|
|
||||||
|
|
||||||
async def fetch(self, query: str, *args: object) -> list[dict[str, object]]:
|
|
||||||
case.assertIn("turn.seq < $2", query)
|
|
||||||
case.assertIn("ORDER BY turn.seq DESC", query)
|
|
||||||
case.assertNotIn("turn.text", query)
|
|
||||||
case.assertEqual(args, (SESSION_ID, 9, 3))
|
|
||||||
return [
|
|
||||||
{"turn_id": TURN_ID, "seq": 8, "created_at": OBSERVED_AT, "trace": _trace(seq=8)},
|
|
||||||
{"turn_id": TURN_ID, "seq": 7, "created_at": OBSERVED_AT, "trace": _trace(seq=7)},
|
|
||||||
{"turn_id": TURN_ID, "seq": 6, "created_at": OBSERVED_AT, "trace": _trace(seq=6)},
|
|
||||||
]
|
|
||||||
|
|
||||||
with patch.object(admin_affect, "acquire", return_value=_Acquire(Conn())) as acquire:
|
|
||||||
response = await admin_affect.get_session_detail(
|
|
||||||
user_id=ADMIN_ID,
|
|
||||||
session_id=SESSION_ID,
|
|
||||||
limit=2,
|
|
||||||
before_seq=9,
|
|
||||||
)
|
|
||||||
|
|
||||||
acquire.assert_called_once_with(role="admin", user_id=ADMIN_ID)
|
|
||||||
self.assertEqual([record.seq for record in response.traces], [7, 8])
|
|
||||||
self.assertTrue(response.has_more)
|
|
||||||
self.assertEqual(response.total_traces, 4)
|
|
||||||
self.assertEqual(response.current_emotions["anxiety"], 0.6)
|
|
||||||
self.assertEqual(response.current_emotions["hope"], 0.2)
|
|
||||||
self.assertIsNone(response.current_emotions["sadness"])
|
|
||||||
self.assertIsNone(response.current_emotions["anger"])
|
|
||||||
|
|
||||||
async def test_detail_keeps_trace_free_legacy_session_observable(self) -> None:
|
|
||||||
class Conn:
|
|
||||||
async def fetchrow(self, query: str, *args: object) -> dict[str, object]:
|
|
||||||
if "FROM app.sessions AS s" in query:
|
|
||||||
return {
|
|
||||||
"session_id": SESSION_ID,
|
|
||||||
"persona_code": "P4",
|
|
||||||
"affect_state": {},
|
|
||||||
}
|
|
||||||
return {"total_traces": 0}
|
|
||||||
|
|
||||||
async def fetch(self, query: str, *args: object) -> list[dict[str, object]]:
|
|
||||||
return []
|
|
||||||
|
|
||||||
with patch.object(admin_affect, "acquire", return_value=_Acquire(Conn())):
|
|
||||||
response = await admin_affect.get_session_detail(
|
|
||||||
user_id=ADMIN_ID,
|
|
||||||
session_id=SESSION_ID,
|
|
||||||
limit=100,
|
|
||||||
before_seq=None,
|
|
||||||
)
|
|
||||||
|
|
||||||
self.assertEqual(response.traces, [])
|
|
||||||
self.assertEqual(response.total_traces, 0)
|
|
||||||
self.assertFalse(response.has_more)
|
|
||||||
self.assertTrue(all(value is None for value in response.current_emotions.values()))
|
|
||||||
|
|
||||||
async def test_detail_maps_missing_session_and_database_failure_separately(self) -> None:
|
|
||||||
class MissingConn:
|
|
||||||
async def fetchrow(self, query: str, *args: object) -> None:
|
|
||||||
return None
|
|
||||||
|
|
||||||
with patch.object(admin_affect, "acquire", return_value=_Acquire(MissingConn())):
|
|
||||||
with self.assertRaises(admin_affect.AdminAffectSessionNotFoundError):
|
|
||||||
await admin_affect.get_session_detail(
|
|
||||||
user_id=ADMIN_ID,
|
|
||||||
session_id=SESSION_ID,
|
|
||||||
limit=100,
|
|
||||||
before_seq=None,
|
|
||||||
)
|
|
||||||
|
|
||||||
with patch.object(admin_affect, "acquire", side_effect=RuntimeError("database down")):
|
|
||||||
with self.assertRaises(admin_affect.AdminAffectPersistenceError):
|
|
||||||
await admin_affect.list_sessions(user_id=ADMIN_ID, limit=30, offset=0)
|
|
||||||
|
|
||||||
|
|
||||||
class AdminAffectHttpBoundaryTest(unittest.TestCase):
|
|
||||||
def _client(self, principal: Principal) -> TestClient:
|
|
||||||
app = FastAPI()
|
|
||||||
app.include_router(admin_routes.router)
|
|
||||||
app.dependency_overrides[get_current_principal] = lambda: principal
|
|
||||||
return TestClient(app)
|
|
||||||
|
|
||||||
def test_non_admin_is_rejected_before_service_access(self) -> None:
|
|
||||||
principal = Principal(user_id="learner-1", role=Role.LEARNER)
|
|
||||||
with patch.object(admin_affect, "list_sessions", AsyncMock()) as list_sessions:
|
|
||||||
response = self._client(principal).get("/admin/affect/sessions")
|
|
||||||
|
|
||||||
self.assertEqual(response.status_code, 403)
|
|
||||||
list_sessions.assert_not_awaited()
|
|
||||||
|
|
||||||
def test_invalid_uuid_and_pagination_are_validation_errors(self) -> None:
|
|
||||||
client = self._client(_admin_principal())
|
|
||||||
|
|
||||||
self.assertEqual(client.get("/admin/affect/sessions?limit=0").status_code, 422)
|
|
||||||
self.assertEqual(
|
|
||||||
client.get("/admin/affect/sessions/not-a-uuid").status_code,
|
|
||||||
422,
|
|
||||||
)
|
|
||||||
self.assertEqual(
|
|
||||||
client.get(f"/admin/affect/sessions/{SESSION_ID}?before_seq=0").status_code,
|
|
||||||
422,
|
|
||||||
)
|
|
||||||
|
|
||||||
def test_route_maps_persistence_error_to_503(self) -> None:
|
|
||||||
with patch.object(
|
|
||||||
admin_affect,
|
|
||||||
"list_sessions",
|
|
||||||
AsyncMock(side_effect=admin_affect.AdminAffectPersistenceError("database down")),
|
|
||||||
):
|
|
||||||
response = self._client(_admin_principal()).get("/admin/affect/sessions")
|
|
||||||
|
|
||||||
self.assertEqual(response.status_code, 503)
|
|
||||||
|
|
||||||
def test_route_maps_missing_session_to_404(self) -> None:
|
|
||||||
with patch.object(
|
|
||||||
admin_affect,
|
|
||||||
"get_session_detail",
|
|
||||||
AsyncMock(side_effect=admin_affect.AdminAffectSessionNotFoundError("missing")),
|
|
||||||
):
|
|
||||||
response = self._client(_admin_principal()).get(
|
|
||||||
f"/admin/affect/sessions/{SESSION_ID}"
|
|
||||||
)
|
|
||||||
|
|
||||||
self.assertEqual(response.status_code, 404)
|
|
||||||
|
|
||||||
def test_route_serializes_list_response(self) -> None:
|
|
||||||
result = AdminAffectSessionListResponse(
|
|
||||||
runtime=AdminAffectRuntimeResponse(
|
|
||||||
enabled=True,
|
|
||||||
provider="openrouter",
|
|
||||||
model="~typesafe/jev-latest",
|
|
||||||
configured=True,
|
|
||||||
),
|
|
||||||
sessions=[],
|
|
||||||
total=0,
|
|
||||||
limit=30,
|
|
||||||
offset=0,
|
|
||||||
)
|
|
||||||
with patch.object(admin_affect, "list_sessions", AsyncMock(return_value=result)) as list_sessions:
|
|
||||||
response = self._client(_admin_principal()).get("/admin/affect/sessions")
|
|
||||||
|
|
||||||
self.assertEqual(response.status_code, 200)
|
|
||||||
self.assertEqual(response.json()["runtime"]["provider"], "openrouter")
|
|
||||||
self.assertEqual(response.json()["sessions"], [])
|
|
||||||
list_sessions.assert_awaited_once_with(user_id=ADMIN_ID, limit=30, offset=0)
|
|
||||||
|
|
@ -464,11 +464,6 @@ class ClientAffectRuntimeTest(unittest.IsolatedAsyncioTestCase):
|
||||||
"load_stored_scenario_context",
|
"load_stored_scenario_context",
|
||||||
AsyncMock(return_value=None),
|
AsyncMock(return_value=None),
|
||||||
),
|
),
|
||||||
patch.object(
|
|
||||||
sessions.turn_runtime.session_persistence,
|
|
||||||
"append_client_turn_with_affect_trace",
|
|
||||||
AsyncMock(return_value=True),
|
|
||||||
),
|
|
||||||
patch.object(sessions, "_schedule_stream_turn_evaluation"),
|
patch.object(sessions, "_schedule_stream_turn_evaluation"),
|
||||||
):
|
):
|
||||||
response = await sessions.stream_turn(
|
response = await sessions.stream_turn(
|
||||||
|
|
|
||||||
|
|
@ -1,346 +0,0 @@
|
||||||
"""관리자 전용 Jev 감정 trace와 원자 영속화 회귀."""
|
|
||||||
|
|
||||||
from __future__ import annotations
|
|
||||||
|
|
||||||
import math
|
|
||||||
import unittest
|
|
||||||
from dataclasses import replace
|
|
||||||
from unittest.mock import patch
|
|
||||||
from unittest.mock import AsyncMock
|
|
||||||
|
|
||||||
from pydantic import ValidationError
|
|
||||||
|
|
||||||
from . import session_persistence, turn_runtime
|
|
||||||
from .contracts.client_affect import ClientAffectDimensionTraceV1, ClientAffectTraceV1
|
|
||||||
from .services import client_affect, orchestrator, persona, state_machine
|
|
||||||
from .services.jev_client import AppraisalResult, EMOTION_DIMENSIONS, EmotionEstimate
|
|
||||||
from .store import InProcSession, TurnRecord
|
|
||||||
|
|
||||||
|
|
||||||
def _appraisal(
|
|
||||||
*,
|
|
||||||
confidence: float = 0.5,
|
|
||||||
probabilities: tuple[float, ...] | None = (0.0, 0.5, 0.5, 0.0, 0.0),
|
|
||||||
) -> AppraisalResult:
|
|
||||||
return AppraisalResult(
|
|
||||||
emotions={
|
|
||||||
dimension: EmotionEstimate(
|
|
||||||
score=0.375,
|
|
||||||
confidence=confidence,
|
|
||||||
probabilities=probabilities,
|
|
||||||
)
|
|
||||||
for dimension in EMOTION_DIMENSIONS
|
|
||||||
},
|
|
||||||
model="jev-test",
|
|
||||||
latency_ms=11,
|
|
||||||
input_tokens=13,
|
|
||||||
output_tokens=17,
|
|
||||||
provider="typesafe",
|
|
||||||
cost_usd=None,
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
def _trace() -> ClientAffectTraceV1:
|
|
||||||
appraisal = _appraisal()
|
|
||||||
return _trace_from_appraisal(appraisal)
|
|
||||||
|
|
||||||
|
|
||||||
def _trace_from_appraisal(appraisal: AppraisalResult) -> ClientAffectTraceV1:
|
|
||||||
before = {f"emotion_{dimension}": 0.5 for dimension in EMOTION_DIMENSIONS}
|
|
||||||
transition = client_affect.transition_emotions(
|
|
||||||
before,
|
|
||||||
{},
|
|
||||||
appraisal,
|
|
||||||
min_confidence=0.65,
|
|
||||||
)
|
|
||||||
return client_affect.build_client_affect_trace(
|
|
||||||
affect_state_before=before,
|
|
||||||
affect_baseline={},
|
|
||||||
affect_state_after=transition.affect_state,
|
|
||||||
appraisal=appraisal,
|
|
||||||
transition=transition,
|
|
||||||
turn_seq=1,
|
|
||||||
stage="라포",
|
|
||||||
resistance=0.65,
|
|
||||||
effective_openness=0.15,
|
|
||||||
rapport_credit=1.25,
|
|
||||||
min_confidence=0.65,
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
class ClientAffectTraceContractTest(unittest.TestCase):
|
|
||||||
def test_trace_preserves_transition_values_and_tentative_is_not_accepted(self) -> None:
|
|
||||||
trace = _trace()
|
|
||||||
|
|
||||||
self.assertEqual(trace.schema_version, 1)
|
|
||||||
self.assertEqual(trace.context.rapport_credit, 1.25)
|
|
||||||
self.assertEqual(
|
|
||||||
tuple(dimension.key for dimension in trace.dimensions),
|
|
||||||
EMOTION_DIMENSIONS,
|
|
||||||
)
|
|
||||||
self.assertTrue(
|
|
||||||
all(dimension.decision == "tentative" for dimension in trace.dimensions)
|
|
||||||
)
|
|
||||||
self.assertEqual(trace.dimensions[0].before, 0.5)
|
|
||||||
self.assertEqual(trace.dimensions[0].target, 0.375)
|
|
||||||
self.assertEqual(trace.dimensions[0].after, 0.48125)
|
|
||||||
self.assertEqual(trace.dimensions[0].probabilities, (0.0, 0.5, 0.5, 0.0, 0.0))
|
|
||||||
|
|
||||||
def test_dimension_contract_rejects_nonfinite_probability(self) -> None:
|
|
||||||
for invalid in (math.nan, math.inf):
|
|
||||||
with self.subTest(invalid=invalid), self.assertRaises(ValidationError):
|
|
||||||
ClientAffectDimensionTraceV1(
|
|
||||||
key="anxiety",
|
|
||||||
before=0.0,
|
|
||||||
target=None,
|
|
||||||
after=0.0,
|
|
||||||
confidence=None,
|
|
||||||
probabilities=(0.0, invalid, 0.0, 0.0, 1.0),
|
|
||||||
decision="held",
|
|
||||||
)
|
|
||||||
|
|
||||||
def test_accepted_and_held_trace_values_preserve_nullable_inputs(self) -> None:
|
|
||||||
accepted = _trace_from_appraisal(
|
|
||||||
_appraisal(
|
|
||||||
confidence=0.9,
|
|
||||||
probabilities=(0.0, 0.0, 0.4, 0.6, 0.0),
|
|
||||||
)
|
|
||||||
)
|
|
||||||
held_appraisal = AppraisalResult(
|
|
||||||
emotions={
|
|
||||||
dimension: EmotionEstimate(
|
|
||||||
score=math.nan,
|
|
||||||
confidence=None,
|
|
||||||
probabilities=None,
|
|
||||||
)
|
|
||||||
for dimension in EMOTION_DIMENSIONS
|
|
||||||
},
|
|
||||||
model="jev-test",
|
|
||||||
latency_ms=11,
|
|
||||||
input_tokens=13,
|
|
||||||
output_tokens=17,
|
|
||||||
provider="typesafe",
|
|
||||||
cost_usd=None,
|
|
||||||
)
|
|
||||||
held = _trace_from_appraisal(held_appraisal)
|
|
||||||
|
|
||||||
self.assertTrue(all(item.decision == "accepted" for item in accepted.dimensions))
|
|
||||||
self.assertEqual(
|
|
||||||
accepted.dimensions[0].probabilities,
|
|
||||||
(0.0, 0.0, 0.4, 0.6, 0.0),
|
|
||||||
)
|
|
||||||
self.assertEqual(accepted.dimensions[0].target, 0.375)
|
|
||||||
self.assertEqual(accepted.dimensions[0].confidence, 0.9)
|
|
||||||
self.assertTrue(all(item.decision == "held" for item in held.dimensions))
|
|
||||||
self.assertTrue(
|
|
||||||
all(
|
|
||||||
item.target is None
|
|
||||||
and item.confidence is None
|
|
||||||
and item.probabilities is None
|
|
||||||
for item in held.dimensions
|
|
||||||
)
|
|
||||||
)
|
|
||||||
|
|
||||||
|
|
||||||
class _Transaction:
|
|
||||||
def __init__(self) -> None:
|
|
||||||
self.error: type[BaseException] | None = None
|
|
||||||
|
|
||||||
async def __aenter__(self) -> None:
|
|
||||||
return None
|
|
||||||
|
|
||||||
async def __aexit__(self, exc_type, exc, tb) -> bool:
|
|
||||||
self.error = exc_type
|
|
||||||
return False
|
|
||||||
|
|
||||||
|
|
||||||
class _Connection:
|
|
||||||
def __init__(self, *, fail_trace_insert: bool = False) -> None:
|
|
||||||
self.fail_trace_insert = fail_trace_insert
|
|
||||||
self.transaction_context = _Transaction()
|
|
||||||
self.executed: list[str] = []
|
|
||||||
|
|
||||||
def transaction(self) -> _Transaction:
|
|
||||||
return self.transaction_context
|
|
||||||
|
|
||||||
async def fetchval(self, query: str, *args: object) -> object:
|
|
||||||
if "FROM app.sessions" in query:
|
|
||||||
return "00000000-0000-0000-0000-000000000111"
|
|
||||||
if "COALESCE(MAX(seq)" in query:
|
|
||||||
return 2
|
|
||||||
if "INSERT INTO app.turns" in query:
|
|
||||||
return "00000000-0000-0000-0000-000000000222"
|
|
||||||
raise AssertionError(f"unexpected query: {query}")
|
|
||||||
|
|
||||||
async def execute(self, query: str, *args: object) -> str:
|
|
||||||
self.executed.append(query)
|
|
||||||
if self.fail_trace_insert and "INSERT INTO app.client_affect_trace" in query:
|
|
||||||
raise RuntimeError("trace insert failed")
|
|
||||||
return "INSERT 0 1"
|
|
||||||
|
|
||||||
|
|
||||||
class _Acquire:
|
|
||||||
def __init__(self, conn: _Connection) -> None:
|
|
||||||
self.conn = conn
|
|
||||||
|
|
||||||
async def __aenter__(self) -> _Connection:
|
|
||||||
return self.conn
|
|
||||||
|
|
||||||
async def __aexit__(self, exc_type, exc, tb) -> bool:
|
|
||||||
return False
|
|
||||||
|
|
||||||
|
|
||||||
class ClientAffectTracePersistenceTest(unittest.IsolatedAsyncioTestCase):
|
|
||||||
async def test_atomic_write_assigns_turn_id_only_after_trace_and_state_write(self) -> None:
|
|
||||||
conn = _Connection()
|
|
||||||
turn = TurnRecord(
|
|
||||||
turn_seq=1,
|
|
||||||
speaker="client",
|
|
||||||
stage="라포",
|
|
||||||
text="조금 더 이야기해볼게요.",
|
|
||||||
text_masked="조금 더 이야기해볼게요.",
|
|
||||||
)
|
|
||||||
state = state_machine.SessionState(turn_seq=1)
|
|
||||||
|
|
||||||
with (
|
|
||||||
patch.object(session_persistence, "get_pool", return_value=object()),
|
|
||||||
patch.object(session_persistence, "acquire", return_value=_Acquire(conn)),
|
|
||||||
):
|
|
||||||
stored = await session_persistence.append_client_turn_with_affect_trace(
|
|
||||||
session_id="00000000-0000-0000-0000-000000000111",
|
|
||||||
learner_id="00000000-0000-0000-0000-000000000101",
|
|
||||||
turn=turn,
|
|
||||||
state=state,
|
|
||||||
trace=_trace(),
|
|
||||||
)
|
|
||||||
|
|
||||||
self.assertTrue(stored)
|
|
||||||
self.assertEqual(turn.turn_id, "00000000-0000-0000-0000-000000000222")
|
|
||||||
self.assertIsNone(conn.transaction_context.error)
|
|
||||||
self.assertIn("INSERT INTO app.client_affect_trace", conn.executed[0])
|
|
||||||
self.assertIn("INSERT INTO app.session_state", conn.executed[1])
|
|
||||||
|
|
||||||
async def test_atomic_write_keeps_turn_identifier_unpublished_when_trace_insert_fails(self) -> None:
|
|
||||||
conn = _Connection(fail_trace_insert=True)
|
|
||||||
turn = TurnRecord(
|
|
||||||
turn_seq=1,
|
|
||||||
speaker="client",
|
|
||||||
stage="라포",
|
|
||||||
text="조금 더 이야기해볼게요.",
|
|
||||||
text_masked="조금 더 이야기해볼게요.",
|
|
||||||
)
|
|
||||||
|
|
||||||
with (
|
|
||||||
patch.object(session_persistence, "get_pool", return_value=object()),
|
|
||||||
patch.object(session_persistence, "acquire", return_value=_Acquire(conn)),
|
|
||||||
):
|
|
||||||
with self.assertRaises(session_persistence.ClientAffectTracePersistenceError):
|
|
||||||
await session_persistence.append_client_turn_with_affect_trace(
|
|
||||||
session_id="00000000-0000-0000-0000-000000000111",
|
|
||||||
learner_id="00000000-0000-0000-0000-000000000101",
|
|
||||||
turn=turn,
|
|
||||||
state=state_machine.SessionState(turn_seq=1),
|
|
||||||
trace=_trace(),
|
|
||||||
)
|
|
||||||
|
|
||||||
self.assertIsNone(turn.turn_id)
|
|
||||||
self.assertIs(conn.transaction_context.error, RuntimeError)
|
|
||||||
|
|
||||||
|
|
||||||
class ClientAffectTraceRuntimeTest(unittest.IsolatedAsyncioTestCase):
|
|
||||||
def _session_and_result(
|
|
||||||
self,
|
|
||||||
) -> tuple[InProcSession, orchestrator.TurnContext, orchestrator.TurnResult]:
|
|
||||||
state_before = state_machine.SessionState()
|
|
||||||
state_after = replace(state_before, turn_seq=1)
|
|
||||||
sess = InProcSession(
|
|
||||||
session_id="trace-runtime-session",
|
|
||||||
case_id="trace-runtime-case",
|
|
||||||
learner_id="00000000-0000-0000-0000-000000000101",
|
|
||||||
persona_code=persona.P1.code,
|
|
||||||
theory_mode="humanistic",
|
|
||||||
persona=persona.P1,
|
|
||||||
state=state_before,
|
|
||||||
)
|
|
||||||
ctx = orchestrator.TurnContext(
|
|
||||||
session_id=sess.session_id,
|
|
||||||
case_id=sess.case_id,
|
|
||||||
persona=sess.persona,
|
|
||||||
state_before=state_before,
|
|
||||||
learner_text_raw="그 마음을 조금 더 들려주실 수 있을까요?",
|
|
||||||
learner_text_masked="그 마음을 조금 더 들려주실 수 있을까요?",
|
|
||||||
state_after=state_after,
|
|
||||||
client_affect_trace=_trace(),
|
|
||||||
)
|
|
||||||
result = orchestrator.TurnResult(
|
|
||||||
turn_seq=1,
|
|
||||||
stage=state_after.stage.value,
|
|
||||||
effective_openness=state_after.effective_openness,
|
|
||||||
client_reply="조금 더 이야기해볼게요.",
|
|
||||||
safety_flagged=False,
|
|
||||||
state_after=state_after,
|
|
||||||
)
|
|
||||||
return sess, ctx, result
|
|
||||||
|
|
||||||
async def test_trace_path_updates_runtime_mirrors_only_after_atomic_success(self) -> None:
|
|
||||||
sess, ctx, result = self._session_and_result()
|
|
||||||
append_counselor = AsyncMock()
|
|
||||||
append_atomic = AsyncMock(return_value=True)
|
|
||||||
update_state = AsyncMock()
|
|
||||||
|
|
||||||
with (
|
|
||||||
patch.object(turn_runtime, "append_completed_turn", append_counselor),
|
|
||||||
patch.object(
|
|
||||||
session_persistence,
|
|
||||||
"append_client_turn_with_affect_trace",
|
|
||||||
append_atomic,
|
|
||||||
),
|
|
||||||
patch.object(turn_runtime, "update_session_state", update_state),
|
|
||||||
):
|
|
||||||
await turn_runtime.record_completed_turn(
|
|
||||||
sess,
|
|
||||||
ctx,
|
|
||||||
result,
|
|
||||||
context_prefix="trace test",
|
|
||||||
)
|
|
||||||
|
|
||||||
append_atomic.assert_awaited_once()
|
|
||||||
append_counselor.assert_awaited_once()
|
|
||||||
update_state.assert_not_awaited()
|
|
||||||
self.assertIs(sess.state, result.state_after)
|
|
||||||
self.assertEqual([turn.speaker for turn in sess.turns], ["client"])
|
|
||||||
|
|
||||||
async def test_trace_path_keeps_runtime_mirrors_unchanged_when_atomic_write_fails(self) -> None:
|
|
||||||
sess, ctx, result = self._session_and_result()
|
|
||||||
append_counselor = AsyncMock()
|
|
||||||
update_state = AsyncMock()
|
|
||||||
|
|
||||||
with (
|
|
||||||
patch.object(turn_runtime, "append_completed_turn", append_counselor),
|
|
||||||
patch.object(
|
|
||||||
session_persistence,
|
|
||||||
"append_client_turn_with_affect_trace",
|
|
||||||
AsyncMock(
|
|
||||||
side_effect=session_persistence.ClientAffectTracePersistenceError(
|
|
||||||
"atomic write failed"
|
|
||||||
)
|
|
||||||
),
|
|
||||||
),
|
|
||||||
patch.object(turn_runtime, "update_session_state", update_state),
|
|
||||||
):
|
|
||||||
with self.assertRaises(session_persistence.ClientAffectTracePersistenceError):
|
|
||||||
await turn_runtime.record_completed_turn(
|
|
||||||
sess,
|
|
||||||
ctx,
|
|
||||||
result,
|
|
||||||
context_prefix="trace test",
|
|
||||||
)
|
|
||||||
|
|
||||||
append_counselor.assert_awaited_once()
|
|
||||||
update_state.assert_not_awaited()
|
|
||||||
self.assertIs(sess.state, ctx.state_before)
|
|
||||||
self.assertEqual(sess.turns, [])
|
|
||||||
|
|
||||||
|
|
||||||
if __name__ == "__main__":
|
|
||||||
unittest.main()
|
|
||||||
|
|
@ -385,43 +385,6 @@ class OrchestratorMaskingGateTest(unittest.IsolatedAsyncioTestCase):
|
||||||
for key in ("messages", "prompt", "text"):
|
for key in ("messages", "prompt", "text"):
|
||||||
self.assertNotIn(key, audit_payloads[0])
|
self.assertNotIn(key, audit_payloads[0])
|
||||||
|
|
||||||
async def test_private_affect_trace_is_absent_from_public_result_and_done_event(self) -> None:
|
|
||||||
ctx = _prepare_context()
|
|
||||||
engine = CaptureStreamEngine()
|
|
||||||
|
|
||||||
async def apply_private_trace(
|
|
||||||
context: orchestrator.TurnContext,
|
|
||||||
*,
|
|
||||||
audit_hook=None,
|
|
||||||
) -> None:
|
|
||||||
setattr(context, "client_affect_trace", {
|
|
||||||
"probabilities": [0.0, 0.0, 0.4, 0.6, 0.0],
|
|
||||||
})
|
|
||||||
|
|
||||||
with patch.object(orchestrator, "_apply_client_affect", apply_private_trace):
|
|
||||||
events = [
|
|
||||||
event
|
|
||||||
async for event in orchestrator.run_turn_stream(
|
|
||||||
ctx,
|
|
||||||
engine, # type: ignore[arg-type]
|
|
||||||
)
|
|
||||||
]
|
|
||||||
|
|
||||||
done = events[-1]
|
|
||||||
self.assertEqual(done.event, "done")
|
|
||||||
self.assertNotIn("client_affect_trace", done.data)
|
|
||||||
self.assertNotIn("probabilities", _json_blob(done.data))
|
|
||||||
|
|
||||||
generated = orchestrator.TurnResult(
|
|
||||||
turn_seq=ctx.state_after.turn_seq,
|
|
||||||
stage=ctx.state_after.stage.value,
|
|
||||||
effective_openness=ctx.state_after.effective_openness,
|
|
||||||
client_reply="Masked stream reply.",
|
|
||||||
safety_flagged=False,
|
|
||||||
state_after=ctx.state_after,
|
|
||||||
)
|
|
||||||
self.assertFalse(hasattr(generated, "client_affect_trace"))
|
|
||||||
|
|
||||||
async def test_run_turn_generate_sends_only_masked_korean_pii(self) -> None:
|
async def test_run_turn_generate_sends_only_masked_korean_pii(self) -> None:
|
||||||
ctx = orchestrator.prepare_turn(
|
ctx = orchestrator.prepare_turn(
|
||||||
session_id="masking-session",
|
session_id="masking-session",
|
||||||
|
|
|
||||||
|
|
@ -8,7 +8,6 @@ from pathlib import Path
|
||||||
|
|
||||||
from .runtime_schema import (
|
from .runtime_schema import (
|
||||||
CALIBRATION_TRANSFER_SCHEMA_CONTRACT,
|
CALIBRATION_TRANSFER_SCHEMA_CONTRACT,
|
||||||
CLIENT_AFFECT_TRACE_SCHEMA_CONTRACT,
|
|
||||||
CONTINUOUS_IMPROVEMENT_SCHEMA_CONTRACT,
|
CONTINUOUS_IMPROVEMENT_SCHEMA_CONTRACT,
|
||||||
DELIBERATE_PRACTICE_SCHEMA_CONTRACT,
|
DELIBERATE_PRACTICE_SCHEMA_CONTRACT,
|
||||||
MEASUREMENT_SCHEMA_CONTRACT,
|
MEASUREMENT_SCHEMA_CONTRACT,
|
||||||
|
|
@ -31,7 +30,6 @@ INFRA_SQL = "\n".join(
|
||||||
class RuntimeSchemaSsotTest(unittest.TestCase):
|
class RuntimeSchemaSsotTest(unittest.TestCase):
|
||||||
def test_runtime_contract_objects_are_owned_by_infra_sql(self) -> None:
|
def test_runtime_contract_objects_are_owned_by_infra_sql(self) -> None:
|
||||||
for contract in (
|
for contract in (
|
||||||
CLIENT_AFFECT_TRACE_SCHEMA_CONTRACT,
|
|
||||||
REVIEW_SCHEMA_CONTRACT,
|
REVIEW_SCHEMA_CONTRACT,
|
||||||
NOTIFICATION_SCHEMA_CONTRACT,
|
NOTIFICATION_SCHEMA_CONTRACT,
|
||||||
MEASUREMENT_SCHEMA_CONTRACT,
|
MEASUREMENT_SCHEMA_CONTRACT,
|
||||||
|
|
|
||||||
|
|
@ -158,7 +158,9 @@ async def record_completed_turn(
|
||||||
client_identity=sess.persona.display_name,
|
client_identity=sess.persona.display_name,
|
||||||
synthetic_generated=True,
|
synthetic_generated=True,
|
||||||
)
|
)
|
||||||
client_turn = TurnRecord(
|
await append_completed_turn(
|
||||||
|
sess,
|
||||||
|
TurnRecord(
|
||||||
turn_seq=result.turn_seq,
|
turn_seq=result.turn_seq,
|
||||||
speaker="client",
|
speaker="client",
|
||||||
stage=stage_label(result.state_after.stage),
|
stage=stage_label(result.state_after.stage),
|
||||||
|
|
@ -169,22 +171,7 @@ async def record_completed_turn(
|
||||||
tokens_in=result.tokens_in,
|
tokens_in=result.tokens_in,
|
||||||
tokens_out=result.tokens_out,
|
tokens_out=result.tokens_out,
|
||||||
cost_usd=result.cost_usd,
|
cost_usd=result.cost_usd,
|
||||||
)
|
),
|
||||||
if ctx.client_affect_trace is not None:
|
|
||||||
await session_persistence.append_client_turn_with_affect_trace(
|
|
||||||
session_id=sess.session_id,
|
|
||||||
learner_id=sess.learner_id,
|
|
||||||
turn=client_turn,
|
|
||||||
state=result.state_after,
|
|
||||||
trace=ctx.client_affect_trace,
|
|
||||||
)
|
|
||||||
sess.turns.append(client_turn)
|
|
||||||
sess.state = result.state_after
|
|
||||||
store.put(sess)
|
|
||||||
return learner_turn
|
|
||||||
await append_completed_turn(
|
|
||||||
sess,
|
|
||||||
client_turn,
|
|
||||||
context=f"{context_prefix} turn append",
|
context=f"{context_prefix} turn append",
|
||||||
)
|
)
|
||||||
await update_session_state(
|
await update_session_state(
|
||||||
|
|
|
||||||
|
|
@ -1,9 +1,9 @@
|
||||||
{
|
{
|
||||||
"schema_version": "vignette.p1_crisis_technical_observations.v2",
|
"schema_version": "vignette.p1_crisis_technical_observations.v2",
|
||||||
"generated_at": "2026-09-22T19:46:03Z",
|
"generated_at": "2026-09-22T12:32:40.3710740Z",
|
||||||
"base_commit": "acb0d26338dfc5ef73ac84150907707fdf7e0a6a",
|
"base_commit": "8344bc2ad22151ba1cfcea09da125f71fae0ce81",
|
||||||
"runtime_package": {
|
"runtime_package": {
|
||||||
"provenance_commit": "acb0d26338dfc5ef73ac84150907707fdf7e0a6a",
|
"provenance_commit": "8344bc2ad22151ba1cfcea09da125f71fae0ce81",
|
||||||
"matches_head": true,
|
"matches_head": true,
|
||||||
"files": [
|
"files": [
|
||||||
{
|
{
|
||||||
|
|
@ -18,7 +18,7 @@
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
"path": "apps/api/app/services/orchestrator.py",
|
"path": "apps/api/app/services/orchestrator.py",
|
||||||
"sha256": "5bc0edbdfbd2f99a5e2024b5a4376cb371b1ab24bb797b639500fddf0829f7e6",
|
"sha256": "b9d62fc2fb0dcd7066193745cb35d78019c382e406150dc4a080845dd028e003",
|
||||||
"matches_head": true
|
"matches_head": true
|
||||||
},
|
},
|
||||||
{
|
{
|
||||||
|
|
|
||||||
|
|
@ -1,39 +0,0 @@
|
||||||
-- =============================================================================
|
|
||||||
-- Vignette · migration 23 — 관리자 감정 관측용 Jev 전이 trace
|
|
||||||
-- =============================================================================
|
|
||||||
-- 기존 turns/session_state 행은 수정하지 않는다. trace는 성공한 client 턴에만
|
|
||||||
-- 연결되며, 앱 런타임은 client turn INSERT와 state UPSERT를 같은 트랜잭션에 둔다.
|
|
||||||
|
|
||||||
CREATE TABLE IF NOT EXISTS app.client_affect_trace (
|
|
||||||
turn_id UUID PRIMARY KEY REFERENCES app.turns(id) ON DELETE CASCADE,
|
|
||||||
session_id UUID NOT NULL REFERENCES app.sessions(id) ON DELETE CASCADE,
|
|
||||||
trace JSONB NOT NULL,
|
|
||||||
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
|
|
||||||
);
|
|
||||||
|
|
||||||
CREATE INDEX IF NOT EXISTS idx_client_affect_trace_session_created
|
|
||||||
ON app.client_affect_trace (session_id, created_at DESC);
|
|
||||||
|
|
||||||
ALTER TABLE app.client_affect_trace ENABLE ROW LEVEL SECURITY;
|
|
||||||
|
|
||||||
DROP POLICY IF EXISTS p_client_affect_trace_select_admin ON app.client_affect_trace;
|
|
||||||
DROP POLICY IF EXISTS p_client_affect_trace_insert_learner ON app.client_affect_trace;
|
|
||||||
|
|
||||||
CREATE POLICY p_client_affect_trace_select_admin
|
|
||||||
ON app.client_affect_trace FOR SELECT
|
|
||||||
USING (app.current_role_name() = 'admin');
|
|
||||||
|
|
||||||
CREATE POLICY p_client_affect_trace_insert_learner
|
|
||||||
ON app.client_affect_trace FOR INSERT
|
|
||||||
WITH CHECK (
|
|
||||||
app.current_role_name() = 'learner'
|
|
||||||
AND EXISTS (
|
|
||||||
SELECT 1
|
|
||||||
FROM app.turns AS turn_row
|
|
||||||
JOIN app.sessions AS session_row ON session_row.id = turn_row.session_id
|
|
||||||
WHERE turn_row.id = app.client_affect_trace.turn_id
|
|
||||||
AND turn_row.session_id = app.client_affect_trace.session_id
|
|
||||||
AND turn_row.speaker = 'client'
|
|
||||||
AND session_row.learner_id = app.current_uid()
|
|
||||||
)
|
|
||||||
);
|
|
||||||
Loading…
Add table
Add a link
Reference in a new issue