Compare commits

...

2 commits

Author SHA1 Message Date
Yun Chan
6ab40ff3f4 감정 관측 변경의 안전 회귀 증거 갱신
Some checks failed
API contract / OpenAPI type drift (push) Failing after 52s
2026-09-23 04:47:07 +09:00
Yun Chan
acb0d26338 관리자 감정 관측 기록과 조회 API 추가 2026-09-23 04:45:50 +09:00
17 changed files with 1449 additions and 22 deletions

View file

@ -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

View file

@ -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",
]

View file

@ -25,6 +25,7 @@ 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,
@ -90,6 +91,7 @@ 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()
@ -102,6 +104,10 @@ 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,
@ -134,6 +140,14 @@ 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,

View file

@ -25,6 +25,10 @@ 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,
@ -38,6 +42,7 @@ 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 (
@ -1934,6 +1939,58 @@ 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,

View file

@ -26,6 +26,23 @@ 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=(

View file

@ -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,
)

View file

@ -6,6 +6,12 @@ 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
@ -40,7 +46,13 @@ _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
@ -83,7 +95,10 @@ 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 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: def _baseline_value(affect_baseline: Mapping[str, Any], key: str) -> float | None:
@ -163,15 +178,15 @@ def transition_emotions(
held.append(dimension) held.append(dimension)
continue continue
if confidence >= threshold: if confidence >= threshold:
alpha = 0.35 alpha = _ACCEPTED_ALPHA
cap = 0.15 cap = _ACCEPTED_CAP
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 = 0.15 alpha = _TENTATIVE_ALPHA
cap = 0.075 cap = _TENTATIVE_CAP
tentative.append(dimension) tentative.append(dimension)
else: else:
updated[f"emotion_{dimension}"] = old 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( def _mask_text(
value: Any, value: Any,
*, *,
@ -384,6 +487,7 @@ 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",

View file

@ -23,6 +23,7 @@ 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,
@ -94,6 +95,8 @@ 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
@ -378,13 +381,27 @@ 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(
ctx.state_after.affect_state, state_before_transition.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,

View file

@ -14,6 +14,7 @@ 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,
@ -44,6 +45,10 @@ class SessionCreationPersistenceError(RuntimeError):
"""fail-closed 환경에서 영속 세션 생성이 실패했다.""" """fail-closed 환경에서 영속 세션 생성이 실패했다."""
class ClientAffectTracePersistenceError(RuntimeError):
"""내담자 턴·감정 trace·상태를 같은 트랜잭션으로 기록하지 못했다."""
class ActiveSessionExistsError(RuntimeError): class ActiveSessionExistsError(RuntimeError):
"""같은 learner-persona 전체에 미종료 회기가 이미 존재한다.""" """같은 learner-persona 전체에 미종료 회기가 이미 존재한다."""
@ -2892,6 +2897,100 @@ 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,

View file

@ -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)

View file

@ -464,6 +464,11 @@ 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(

View file

@ -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()

View file

@ -385,6 +385,43 @@ 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",

View file

@ -8,6 +8,7 @@ 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,
@ -30,6 +31,7 @@ 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,

View file

@ -158,9 +158,7 @@ async def record_completed_turn(
client_identity=sess.persona.display_name, client_identity=sess.persona.display_name,
synthetic_generated=True, synthetic_generated=True,
) )
await append_completed_turn( client_turn = TurnRecord(
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),
@ -171,7 +169,22 @@ 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(

View file

@ -1,9 +1,9 @@
{ {
"schema_version": "vignette.p1_crisis_technical_observations.v2", "schema_version": "vignette.p1_crisis_technical_observations.v2",
"generated_at": "2026-09-22T12:32:40.3710740Z", "generated_at": "2026-09-22T19:46:03Z",
"base_commit": "8344bc2ad22151ba1cfcea09da125f71fae0ce81", "base_commit": "acb0d26338dfc5ef73ac84150907707fdf7e0a6a",
"runtime_package": { "runtime_package": {
"provenance_commit": "8344bc2ad22151ba1cfcea09da125f71fae0ce81", "provenance_commit": "acb0d26338dfc5ef73ac84150907707fdf7e0a6a",
"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": "b9d62fc2fb0dcd7066193745cb35d78019c382e406150dc4a080845dd028e003", "sha256": "5bc0edbdfbd2f99a5e2024b5a4376cb371b1ab24bb797b639500fddf0829f7e6",
"matches_head": true "matches_head": true
}, },
{ {

View file

@ -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()
)
);