From acb0d26338dfc5ef73ac84150907707fdf7e0a6a Mon Sep 17 00:00:00 2001 From: Yun Chan Date: Wed, 23 Sep 2026 04:45:50 +0900 Subject: [PATCH 1/2] =?UTF-8?q?=EA=B4=80=EB=A6=AC=EC=9E=90=20=EA=B0=90?= =?UTF-8?q?=EC=A0=95=20=EA=B4=80=EC=B8=A1=20=EA=B8=B0=EB=A1=9D=EA=B3=BC=20?= =?UTF-8?q?=EC=A1=B0=ED=9A=8C=20API=20=EC=B6=94=EA=B0=80?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- apps/api/app/contracts/admin_affect.py | 58 ++++ apps/api/app/contracts/client_affect.py | 100 +++++++ apps/api/app/main.py | 14 + apps/api/app/routes/admin.py | 57 ++++ apps/api/app/runtime_schema.py | 17 ++ apps/api/app/services/admin_affect.py | 192 ++++++++++++ apps/api/app/services/client_affect.py | 114 ++++++- apps/api/app/services/orchestrator.py | 19 +- apps/api/app/session_persistence.py | 99 +++++++ apps/api/app/test_admin_affect.py | 327 ++++++++++++++++++++ apps/api/app/test_client_affect.py | 5 + apps/api/app/test_client_affect_trace.py | 346 ++++++++++++++++++++++ apps/api/app/test_orchestrator_masking.py | 37 +++ apps/api/app/test_runtime_schema_ssot.py | 2 + apps/api/app/turn_runtime.py | 37 ++- infra/db/init/23_client_affect_trace.sql | 39 +++ 16 files changed, 1445 insertions(+), 18 deletions(-) create mode 100644 apps/api/app/contracts/admin_affect.py create mode 100644 apps/api/app/contracts/client_affect.py create mode 100644 apps/api/app/services/admin_affect.py create mode 100644 apps/api/app/test_admin_affect.py create mode 100644 apps/api/app/test_client_affect_trace.py create mode 100644 infra/db/init/23_client_affect_trace.sql diff --git a/apps/api/app/contracts/admin_affect.py b/apps/api/app/contracts/admin_affect.py new file mode 100644 index 0000000..27e707c --- /dev/null +++ b/apps/api/app/contracts/admin_affect.py @@ -0,0 +1,58 @@ +"""관리자 감정 관측 조회 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 diff --git a/apps/api/app/contracts/client_affect.py b/apps/api/app/contracts/client_affect.py new file mode 100644 index 0000000..7f75432 --- /dev/null +++ b/apps/api/app/contracts/client_affect.py @@ -0,0 +1,100 @@ +"""관리자 감정 관측용 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", +] diff --git a/apps/api/app/main.py b/apps/api/app/main.py index 562e4e5..e22e6ea 100644 --- a/apps/api/app/main.py +++ b/apps/api/app/main.py @@ -25,6 +25,7 @@ from .session_persistence import ensure_review_tables from .services.jev_client import jev_client from .runtime_schema import ( CALIBRATION_TRANSFER_SCHEMA_CONTRACT, + CLIENT_AFFECT_TRACE_SCHEMA_CONTRACT, CONTINUOUS_IMPROVEMENT_SCHEMA_CONTRACT, DELIBERATE_PRACTICE_SCHEMA_CONTRACT, MEASUREMENT_SCHEMA_CONTRACT, @@ -90,6 +91,7 @@ async def lifespan(app: FastAPI): supervision_research_schema_ready = False continuous_improvement_schema_ready = False multimodal_alliance_schema_ready = False + client_affect_trace_schema_ready = False app.state.upload_database_proof = None try: await init_pool() @@ -102,6 +104,10 @@ async def lifespan(app: FastAPI): await ensure_notification_tables() await ensure_protocol_tables() 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( conn, MEASUREMENT_SCHEMA_CONTRACT, @@ -134,6 +140,14 @@ async def lifespan(app: FastAPI): conn, 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( MEASUREMENT_SCHEMA_CONTRACT, ready=measurement_schema_ready, diff --git a/apps/api/app/routes/admin.py b/apps/api/app/routes/admin.py index 4b636d0..3daf028 100644 --- a/apps/api/app/routes/admin.py +++ b/apps/api/app/routes/admin.py @@ -25,6 +25,10 @@ from ..auth_sessions import ( upsert_managed_user, ) from ..config import settings +from ..contracts.admin_affect import ( + AdminAffectSessionDetailResponse, + AdminAffectSessionListResponse, +) from ..contracts.engine_gateway import ( ENGINE_PROVIDER_DEFAULTS, ENGINE_PROVIDERS, @@ -38,6 +42,7 @@ from ..deps import Principal, require_admin_access from ..engine_client import engine_client from ..runtime_policy import require_runtime_fallback_allowed 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_oauth as provider_oauth_service from ..services.llm_pricing import ( @@ -1934,6 +1939,58 @@ async def admin_voice_runtime( 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) async def admin_usage( principal: AdminPrincipal, diff --git a/apps/api/app/runtime_schema.py b/apps/api/app/runtime_schema.py index a0e014f..11cb527 100644 --- a/apps/api/app/runtime_schema.py +++ b/apps/api/app/runtime_schema.py @@ -26,6 +26,23 @@ class RuntimeSchemaContract: 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( component="review/evaluation", relations=( diff --git a/apps/api/app/services/admin_affect.py b/apps/api/app/services/admin_affect.py new file mode 100644 index 0000000..a747d0f --- /dev/null +++ b/apps/api/app/services/admin_affect.py @@ -0,0 +1,192 @@ +"""관리자 감정 관측용 영속 조회.""" + +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, + ) diff --git a/apps/api/app/services/client_affect.py b/apps/api/app/services/client_affect.py index d02fc84..0e458df 100644 --- a/apps/api/app/services/client_affect.py +++ b/apps/api/app/services/client_affect.py @@ -6,6 +6,12 @@ import math from dataclasses import dataclass from typing import Any, Iterable, Mapping +from ..contracts.client_affect import ( + ClientAffectContextV1, + ClientAffectDimensionTraceV1, + ClientAffectPolicyV1, + ClientAffectTraceV1, +) from . import guardrail from .jev_client import AppraisalResult, EMOTION_DIMENSIONS @@ -40,7 +46,13 @@ _NEGATIVE_EMOTIONS = frozenset( {"anxiety", "sadness", "anger", "shame", "guilt", "loneliness"} ) _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 +_ADJACENT_PROBABILITY_THRESHOLD = 0.8 _PROBABILITY_SUM_TOLERANCE = 0.025000001 @@ -83,7 +95,10 @@ def _tentative_distribution_is_concentrated(probabilities: Any) -> bool: if not math.isclose(total, 1.0, abs_tol=_PROBABILITY_SUM_TOLERANCE): return False normalized = tuple(value / total for value in values if value is not None) - return max(normalized[index] + normalized[index + 1] for index in range(4)) >= 0.80 + return ( + 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: @@ -163,15 +178,15 @@ def transition_emotions( held.append(dimension) continue if confidence >= threshold: - alpha = 0.35 - cap = 0.15 + alpha = _ACCEPTED_ALPHA + cap = _ACCEPTED_CAP elif ( confidence >= _TENTATIVE_CONFIDENCE_FLOOR and _tentative_distribution_is_concentrated(estimate.probabilities) ): # confidence는 정답 확률이 아니라 분포 집중도 요약이다. - alpha = 0.15 - cap = 0.075 + alpha = _TENTATIVE_ALPHA + cap = _TENTATIVE_CAP tentative.append(dimension) else: updated[f"emotion_{dimension}"] = old @@ -189,6 +204,94 @@ 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( value: Any, *, @@ -384,6 +487,7 @@ def public_end_state(end_state: Mapping[str, Any]) -> dict[str, Any]: __all__ = [ "AffectTransition", "baseline_emotions", + "build_client_affect_trace", "build_appraisal_state", "public_end_state", "resolve_emotions", diff --git a/apps/api/app/services/orchestrator.py b/apps/api/app/services/orchestrator.py index 7c3842a..a91ac18 100644 --- a/apps/api/app/services/orchestrator.py +++ b/apps/api/app/services/orchestrator.py @@ -23,6 +23,7 @@ import time from dataclasses import dataclass, field, replace from typing import Any, AsyncIterator, Awaitable, Callable, Optional +from ..contracts.client_affect import ClientAffectTraceV1 from ..config import settings from ..engine_client import ( EngineClient, @@ -94,6 +95,8 @@ class TurnContext: scenario_directive: Optional[rupture_scenario_director.ScenarioDirective] = None # 외부 감정 평가의 안전한 provenance. 원문·점수·확률은 넣지 않는다. client_affect_metadata: Optional[dict[str, Any]] = None + # 관리자 관측 전용 Jev 전이 trace. 공개 결과나 provider event에는 넣지 않는다. + client_affect_trace: ClientAffectTraceV1 | None = None def to_state_context(self) -> PersonaStateContext: st = self.state_after or self.state_before @@ -378,13 +381,27 @@ async def _apply_client_affect( client_identity=ctx.client_identity, ) appraisal = await jev_client.appraise(state) + state_before_transition = ctx.state_after transition = client_affect.transition_emotions( - ctx.state_after.affect_state, + state_before_transition.affect_state, ctx.persona.affect_baseline, appraisal, min_confidence=settings.jev_min_confidence, ) 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 = { "provider": appraisal.provider, "model": appraisal.model, diff --git a/apps/api/app/session_persistence.py b/apps/api/app/session_persistence.py index 7011f17..47ce53b 100644 --- a/apps/api/app/session_persistence.py +++ b/apps/api/app/session_persistence.py @@ -14,6 +14,7 @@ from typing import Any, Awaitable, Callable, Iterable, Literal from .db import acquire, get_pool from .deps import Principal from .config import settings +from .contracts.client_affect import ClientAffectTraceV1 from .persona_repository import ( SEED_VERSION, card_from_row, @@ -44,6 +45,10 @@ class SessionCreationPersistenceError(RuntimeError): """fail-closed 환경에서 영속 세션 생성이 실패했다.""" +class ClientAffectTracePersistenceError(RuntimeError): + """내담자 턴·감정 trace·상태를 같은 트랜잭션으로 기록하지 못했다.""" + + class ActiveSessionExistsError(RuntimeError): """같은 learner-persona 전체에 미종료 회기가 이미 존재한다.""" @@ -2892,6 +2897,100 @@ async def append_turn( 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( *, session_id: str, diff --git a/apps/api/app/test_admin_affect.py b/apps/api/app/test_admin_affect.py new file mode 100644 index 0000000..825880f --- /dev/null +++ b/apps/api/app/test_admin_affect.py @@ -0,0 +1,327 @@ +"""관리자 감정 관측 조회의 권한·페이지·영속 경계 검증.""" + +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) diff --git a/apps/api/app/test_client_affect.py b/apps/api/app/test_client_affect.py index b2c951a..c001430 100644 --- a/apps/api/app/test_client_affect.py +++ b/apps/api/app/test_client_affect.py @@ -464,6 +464,11 @@ class ClientAffectRuntimeTest(unittest.IsolatedAsyncioTestCase): "load_stored_scenario_context", 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"), ): response = await sessions.stream_turn( diff --git a/apps/api/app/test_client_affect_trace.py b/apps/api/app/test_client_affect_trace.py new file mode 100644 index 0000000..a15ecf7 --- /dev/null +++ b/apps/api/app/test_client_affect_trace.py @@ -0,0 +1,346 @@ +"""관리자 전용 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() diff --git a/apps/api/app/test_orchestrator_masking.py b/apps/api/app/test_orchestrator_masking.py index 198860e..84e2758 100644 --- a/apps/api/app/test_orchestrator_masking.py +++ b/apps/api/app/test_orchestrator_masking.py @@ -385,6 +385,43 @@ class OrchestratorMaskingGateTest(unittest.IsolatedAsyncioTestCase): for key in ("messages", "prompt", "text"): 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: ctx = orchestrator.prepare_turn( session_id="masking-session", diff --git a/apps/api/app/test_runtime_schema_ssot.py b/apps/api/app/test_runtime_schema_ssot.py index ffb6060..2e0ef6d 100644 --- a/apps/api/app/test_runtime_schema_ssot.py +++ b/apps/api/app/test_runtime_schema_ssot.py @@ -8,6 +8,7 @@ from pathlib import Path from .runtime_schema import ( CALIBRATION_TRANSFER_SCHEMA_CONTRACT, + CLIENT_AFFECT_TRACE_SCHEMA_CONTRACT, CONTINUOUS_IMPROVEMENT_SCHEMA_CONTRACT, DELIBERATE_PRACTICE_SCHEMA_CONTRACT, MEASUREMENT_SCHEMA_CONTRACT, @@ -30,6 +31,7 @@ INFRA_SQL = "\n".join( class RuntimeSchemaSsotTest(unittest.TestCase): def test_runtime_contract_objects_are_owned_by_infra_sql(self) -> None: for contract in ( + CLIENT_AFFECT_TRACE_SCHEMA_CONTRACT, REVIEW_SCHEMA_CONTRACT, NOTIFICATION_SCHEMA_CONTRACT, MEASUREMENT_SCHEMA_CONTRACT, diff --git a/apps/api/app/turn_runtime.py b/apps/api/app/turn_runtime.py index 91a777f..cc9a08a 100644 --- a/apps/api/app/turn_runtime.py +++ b/apps/api/app/turn_runtime.py @@ -158,20 +158,33 @@ async def record_completed_turn( client_identity=sess.persona.display_name, synthetic_generated=True, ) + client_turn = TurnRecord( + turn_seq=result.turn_seq, + speaker="client", + stage=stage_label(result.state_after.stage), + text=result.client_reply, + text_masked=client_mask.text_masked, + llm_provider=result.llm_provider, + model=result.model, + tokens_in=result.tokens_in, + tokens_out=result.tokens_out, + 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, - TurnRecord( - turn_seq=result.turn_seq, - speaker="client", - stage=stage_label(result.state_after.stage), - text=result.client_reply, - text_masked=client_mask.text_masked, - llm_provider=result.llm_provider, - model=result.model, - tokens_in=result.tokens_in, - tokens_out=result.tokens_out, - cost_usd=result.cost_usd, - ), + client_turn, context=f"{context_prefix} turn append", ) await update_session_state( diff --git a/infra/db/init/23_client_affect_trace.sql b/infra/db/init/23_client_affect_trace.sql new file mode 100644 index 0000000..bf9dfa1 --- /dev/null +++ b/infra/db/init/23_client_affect_trace.sql @@ -0,0 +1,39 @@ +-- ============================================================================= +-- 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() + ) + ); From 6ab40ff3f4f0198d1b7b6b6f0d5eeeb060209249 Mon Sep 17 00:00:00 2001 From: Yun Chan Date: Wed, 23 Sep 2026 04:47:07 +0900 Subject: [PATCH 2/2] =?UTF-8?q?=EA=B0=90=EC=A0=95=20=EA=B4=80=EC=B8=A1=20?= =?UTF-8?q?=EB=B3=80=EA=B2=BD=EC=9D=98=20=EC=95=88=EC=A0=84=20=ED=9A=8C?= =?UTF-8?q?=EA=B7=80=20=EC=A6=9D=EA=B1=B0=20=EA=B0=B1=EC=8B=A0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- ...c-001-crisis-case-runtime-observations-2026-09-22.json | 8 ++++---- 1 file changed, 4 insertions(+), 4 deletions(-) diff --git a/docs/ops/evidence/c-001-crisis-case-runtime-observations-2026-09-22.json b/docs/ops/evidence/c-001-crisis-case-runtime-observations-2026-09-22.json index ddbde1b..32b63d6 100644 --- a/docs/ops/evidence/c-001-crisis-case-runtime-observations-2026-09-22.json +++ b/docs/ops/evidence/c-001-crisis-case-runtime-observations-2026-09-22.json @@ -1,9 +1,9 @@ { "schema_version": "vignette.p1_crisis_technical_observations.v2", - "generated_at": "2026-09-22T12:32:40.3710740Z", - "base_commit": "8344bc2ad22151ba1cfcea09da125f71fae0ce81", + "generated_at": "2026-09-22T19:46:03Z", + "base_commit": "acb0d26338dfc5ef73ac84150907707fdf7e0a6a", "runtime_package": { - "provenance_commit": "8344bc2ad22151ba1cfcea09da125f71fae0ce81", + "provenance_commit": "acb0d26338dfc5ef73ac84150907707fdf7e0a6a", "matches_head": true, "files": [ { @@ -18,7 +18,7 @@ }, { "path": "apps/api/app/services/orchestrator.py", - "sha256": "b9d62fc2fb0dcd7066193745cb35d78019c382e406150dc4a080845dd028e003", + "sha256": "5bc0edbdfbd2f99a5e2024b5a4376cb371b1ab24bb797b639500fddf0829f7e6", "matches_head": true }, {