Compare commits

..

No commits in common. "6ab40ff3f4f0198d1b7b6b6f0d5eeeb060209249" and "d22cd9883dc0a02b662b7865fd6314fbde407b7e" have entirely different histories.

17 changed files with 22 additions and 1449 deletions

View file

@ -1,58 +0,0 @@
"""관리자 감정 관측 조회 API 계약."""
from __future__ import annotations
from datetime import datetime
from pydantic import BaseModel, ConfigDict, Field
from .client_affect import ClientAffectTraceV1
class AdminAffectRuntimeResponse(BaseModel):
model_config = ConfigDict(extra="forbid", protected_namespaces=())
enabled: bool
provider: str
model: str
configured: bool
class AdminAffectSessionSummary(BaseModel):
model_config = ConfigDict(extra="forbid", protected_namespaces=())
session_id: str
persona_code: str
started_at: datetime
ended: bool
trace_count: int = Field(ge=0)
class AdminAffectSessionListResponse(BaseModel):
model_config = ConfigDict(extra="forbid", protected_namespaces=())
runtime: AdminAffectRuntimeResponse
sessions: list[AdminAffectSessionSummary]
total: int = Field(ge=0)
limit: int = Field(ge=1, le=100)
offset: int = Field(ge=0)
class AdminAffectTraceRecord(BaseModel):
model_config = ConfigDict(extra="forbid", protected_namespaces=())
turn_id: str
seq: int = Field(ge=1)
created_at: datetime
trace: ClientAffectTraceV1
class AdminAffectSessionDetailResponse(BaseModel):
model_config = ConfigDict(extra="forbid", protected_namespaces=())
session_id: str
persona_code: str
current_emotions: dict[str, float | None]
traces: list[AdminAffectTraceRecord]
total_traces: int = Field(ge=0)
has_more: bool

View file

@ -1,100 +0,0 @@
"""관리자 감정 관측용 Jev 감정 전이 trace 계약."""
from __future__ import annotations
import math
from typing import Literal
from pydantic import BaseModel, ConfigDict, Field, model_validator
CLIENT_AFFECT_DIMENSIONS = (
"anxiety",
"sadness",
"anger",
"shame",
"guilt",
"loneliness",
"relief",
"hope",
"trust",
)
ClientAffectDecision = Literal["accepted", "tentative", "held"]
class ClientAffectPolicyV1(BaseModel):
model_config = ConfigDict(extra="forbid", frozen=True, protected_namespaces=())
version: Literal["jev-affect-v1"]
min_confidence: float = Field(ge=0.0, le=1.0)
accepted_alpha: float = Field(ge=0.0, le=1.0)
accepted_cap: float = Field(ge=0.0, le=1.0)
tentative_alpha: float = Field(ge=0.0, le=1.0)
tentative_cap: float = Field(ge=0.0, le=1.0)
tentative_confidence_floor: float = Field(ge=0.0, le=1.0)
adjacent_probability_threshold: float = Field(ge=0.0, le=1.0)
class ClientAffectContextV1(BaseModel):
model_config = ConfigDict(extra="forbid", frozen=True, protected_namespaces=())
stage: str
resistance: float = Field(ge=0.0, le=1.0)
effective_openness: float = Field(ge=0.0, le=1.0)
rapport_credit: float = Field(ge=0.0)
class ClientAffectDimensionTraceV1(BaseModel):
model_config = ConfigDict(extra="forbid", frozen=True, protected_namespaces=())
key: str
before: float = Field(ge=0.0, le=1.0)
target: float | None = Field(default=None, ge=0.0, le=1.0)
after: float = Field(ge=0.0, le=1.0)
confidence: float | None = Field(default=None, ge=0.0, le=1.0)
probabilities: tuple[float, float, float, float, float] | None = None
decision: ClientAffectDecision
@model_validator(mode="after")
def require_probability_distribution(self) -> "ClientAffectDimensionTraceV1":
if self.probabilities is None:
return self
if any(
not math.isfinite(value) or value < 0.0 or value > 1.0
for value in self.probabilities
):
raise ValueError("probabilities must be finite values within 0..1")
return self
class ClientAffectTraceV1(BaseModel):
model_config = ConfigDict(extra="forbid", frozen=True, protected_namespaces=())
schema_version: Literal[1]
provider: str
model: str
latency_ms: int = Field(ge=0)
input_tokens: int = Field(ge=0)
output_tokens: int = Field(ge=0)
cost_usd: float | None = Field(default=None, ge=0.0)
turn_seq: int = Field(ge=1)
policy: ClientAffectPolicyV1
context: ClientAffectContextV1
dimensions: tuple[ClientAffectDimensionTraceV1, ...] = Field(min_length=9, max_length=9)
@model_validator(mode="after")
def require_fixed_dimension_order(self) -> "ClientAffectTraceV1":
if tuple(dimension.key for dimension in self.dimensions) != CLIENT_AFFECT_DIMENSIONS:
raise ValueError("dimensions must use the fixed client affect order")
return self
__all__ = [
"CLIENT_AFFECT_DIMENSIONS",
"ClientAffectContextV1",
"ClientAffectDecision",
"ClientAffectDimensionTraceV1",
"ClientAffectPolicyV1",
"ClientAffectTraceV1",
]

View file

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

View file

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

View file

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

View file

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

View file

@ -6,12 +6,6 @@ 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
@ -46,13 +40,7 @@ _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
@ -95,10 +83,7 @@ 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))
>= _ADJACENT_PROBABILITY_THRESHOLD
)
return max(normalized[index] + normalized[index + 1] for index in range(4)) >= 0.80
def _baseline_value(affect_baseline: Mapping[str, Any], key: str) -> float | None:
@ -178,15 +163,15 @@ def transition_emotions(
held.append(dimension)
continue
if confidence >= threshold:
alpha = _ACCEPTED_ALPHA
cap = _ACCEPTED_CAP
alpha = 0.35
cap = 0.15
elif (
confidence >= _TENTATIVE_CONFIDENCE_FLOOR
and _tentative_distribution_is_concentrated(estimate.probabilities)
):
# confidence는 정답 확률이 아니라 분포 집중도 요약이다.
alpha = _TENTATIVE_ALPHA
cap = _TENTATIVE_CAP
alpha = 0.15
cap = 0.075
tentative.append(dimension)
else:
updated[f"emotion_{dimension}"] = old
@ -204,94 +189,6 @@ def transition_emotions(
)
def _trace_probabilities(value: Any) -> tuple[float, float, float, float, float] | None:
if not isinstance(value, tuple) or len(value) != 5:
return None
normalized = tuple(_unit_number(item) for item in value)
if any(item is None for item in normalized):
return None
return (
normalized[0],
normalized[1],
normalized[2],
normalized[3],
normalized[4],
)
def build_client_affect_trace(
*,
affect_state_before: Mapping[str, Any],
affect_baseline: Mapping[str, Any],
affect_state_after: Mapping[str, Any],
appraisal: AppraisalResult,
transition: AffectTransition,
turn_seq: int,
stage: str,
resistance: float,
effective_openness: float,
rapport_credit: float,
min_confidence: float,
) -> ClientAffectTraceV1:
"""전이와 같은 입력으로 관리자 전용 trace를 고정 순서로 만든다."""
before = resolve_emotions(affect_state_before, affect_baseline)
after = resolve_emotions(affect_state_after, affect_baseline)
tentative = set(transition.tentative_dimensions)
accepted = set(transition.accepted_dimensions)
dimensions: list[ClientAffectDimensionTraceV1] = []
for key in EMOTION_DIMENSIONS:
estimate = appraisal.emotions.get(key)
target = _unit_number(estimate.score) if estimate is not None else None
confidence = _unit_number(estimate.confidence) if estimate is not None else None
probabilities = (
_trace_probabilities(estimate.probabilities) if estimate is not None else None
)
if key in tentative:
decision = "tentative"
elif key in accepted:
decision = "accepted"
else:
decision = "held"
dimensions.append(
ClientAffectDimensionTraceV1(
key=key,
before=before[key],
target=target,
after=after[key],
confidence=confidence,
probabilities=probabilities,
decision=decision,
)
)
return ClientAffectTraceV1(
schema_version=1,
provider=appraisal.provider,
model=appraisal.model,
latency_ms=appraisal.latency_ms,
input_tokens=appraisal.input_tokens,
output_tokens=appraisal.output_tokens,
cost_usd=appraisal.cost_usd,
turn_seq=turn_seq,
policy=ClientAffectPolicyV1(
version=_AFFECT_POLICY_VERSION,
min_confidence=min_confidence,
accepted_alpha=_ACCEPTED_ALPHA,
accepted_cap=_ACCEPTED_CAP,
tentative_alpha=_TENTATIVE_ALPHA,
tentative_cap=_TENTATIVE_CAP,
tentative_confidence_floor=_TENTATIVE_CONFIDENCE_FLOOR,
adjacent_probability_threshold=_ADJACENT_PROBABILITY_THRESHOLD,
),
context=ClientAffectContextV1(
stage=stage,
resistance=resistance,
effective_openness=effective_openness,
rapport_credit=rapport_credit,
),
dimensions=tuple(dimensions),
)
def _mask_text(
value: Any,
*,
@ -487,7 +384,6 @@ 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",

View file

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

View file

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

View file

@ -1,327 +0,0 @@
"""관리자 감정 관측 조회의 권한·페이지·영속 경계 검증."""
from __future__ import annotations
import unittest
from datetime import datetime, timezone
from unittest.mock import AsyncMock, patch
from uuid import UUID
from fastapi import FastAPI, HTTPException
from fastapi.testclient import TestClient
from .contracts.admin_affect import (
AdminAffectRuntimeResponse,
AdminAffectSessionListResponse,
)
from .deps import Principal, Role, get_current_principal
from .routes import admin as admin_routes
from .services import admin_affect
SESSION_ID = UUID("00000000-0000-0000-0000-000000000101")
TURN_ID = UUID("00000000-0000-0000-0000-000000000201")
ADMIN_ID = "00000000-0000-0000-0000-000000000901"
OBSERVED_AT = datetime(2026, 9, 23, 1, 2, 3, tzinfo=timezone.utc)
class _Acquire:
def __init__(self, conn: object) -> None:
self.conn = conn
async def __aenter__(self) -> object:
return self.conn
async def __aexit__(self, exc_type: object, exc: object, tb: object) -> None:
return None
def _trace(*, seq: int) -> dict[str, object]:
dimensions = []
for key in (
"anxiety",
"sadness",
"anger",
"shame",
"guilt",
"loneliness",
"relief",
"hope",
"trust",
):
dimensions.append(
{
"key": key,
"before": 0.4,
"target": 0.6,
"after": 0.47,
"confidence": 0.8,
"probabilities": [0.05, 0.1, 0.2, 0.35, 0.3],
"decision": "accepted",
}
)
return {
"schema_version": 1,
"provider": "openrouter",
"model": "~typesafe/jev-latest",
"latency_ms": 91,
"input_tokens": 12,
"output_tokens": 18,
"cost_usd": None,
"turn_seq": seq,
"policy": {
"version": "jev-affect-v1",
"min_confidence": 0.65,
"accepted_alpha": 0.35,
"accepted_cap": 0.15,
"tentative_alpha": 0.15,
"tentative_cap": 0.075,
"tentative_confidence_floor": 0.35,
"adjacent_probability_threshold": 0.8,
},
"context": {
"stage": "탐색",
"resistance": 0.4,
"effective_openness": 0.6,
"rapport_credit": 0.2,
},
"dimensions": dimensions,
}
def _admin_principal() -> Principal:
return Principal(user_id=ADMIN_ID, role=Role.ADMIN)
class AdminAffectStoreTest(unittest.IsolatedAsyncioTestCase):
async def test_list_uses_admin_scope_and_excludes_learner_identity(self) -> None:
case = self
class Conn:
async def fetchrow(self, query: str, *args: object) -> dict[str, object]:
case.assertIn("FROM app.sessions", query)
case.assertNotIn("learner", query.lower())
return {"total": 1}
async def fetch(self, query: str, *args: object) -> list[dict[str, object]]:
case.assertIn("app.client_affect_trace", query)
case.assertNotIn("learner_id", query)
case.assertNotIn("app.app_user", query)
case.assertNotIn(" text", query.lower())
case.assertEqual(args, (30, 0))
return [
{
"session_id": SESSION_ID,
"persona_code": "P4",
"started_at": OBSERVED_AT,
"ended": False,
"trace_count": 2,
}
]
runtime = AdminAffectRuntimeResponse(
enabled=True,
provider="openrouter",
model="~typesafe/jev-latest",
configured=True,
)
with (
patch.object(admin_affect, "acquire", return_value=_Acquire(Conn())) as acquire,
patch.object(admin_affect, "runtime_snapshot", return_value=runtime),
):
response = await admin_affect.list_sessions(
user_id=ADMIN_ID,
limit=30,
offset=0,
)
acquire.assert_called_once_with(role="admin", user_id=ADMIN_ID)
self.assertEqual(response.total, 1)
self.assertEqual(response.sessions[0].session_id, str(SESSION_ID))
self.assertEqual(response.sessions[0].trace_count, 2)
async def test_list_preserves_total_beyond_the_last_page(self) -> None:
case = self
class Conn:
async def fetchrow(self, query: str, *args: object) -> dict[str, object]:
return {"total": 4}
async def fetch(self, query: str, *args: object) -> list[dict[str, object]]:
case.assertEqual(args, (30, 30))
return []
with patch.object(admin_affect, "acquire", return_value=_Acquire(Conn())):
response = await admin_affect.list_sessions(
user_id=ADMIN_ID,
limit=30,
offset=30,
)
self.assertEqual(response.total, 4)
self.assertEqual(response.sessions, [])
async def test_detail_returns_actual_snapshot_and_ascending_page(self) -> None:
case = self
class Conn:
async def fetchrow(self, query: str, *args: object) -> dict[str, object]:
if "FROM app.sessions AS s" in query:
case.assertNotIn("learner", query.lower())
case.assertNotIn("text", query.lower())
return {
"session_id": SESSION_ID,
"persona_code": "P4",
"affect_state": {
"emotion_anxiety": 0.6,
"emotion_hope": 0.2,
"emotion_anger": 1.5,
},
}
case.assertIn("count(*)", query)
return {"total_traces": 4}
async def fetch(self, query: str, *args: object) -> list[dict[str, object]]:
case.assertIn("turn.seq < $2", query)
case.assertIn("ORDER BY turn.seq DESC", query)
case.assertNotIn("turn.text", query)
case.assertEqual(args, (SESSION_ID, 9, 3))
return [
{"turn_id": TURN_ID, "seq": 8, "created_at": OBSERVED_AT, "trace": _trace(seq=8)},
{"turn_id": TURN_ID, "seq": 7, "created_at": OBSERVED_AT, "trace": _trace(seq=7)},
{"turn_id": TURN_ID, "seq": 6, "created_at": OBSERVED_AT, "trace": _trace(seq=6)},
]
with patch.object(admin_affect, "acquire", return_value=_Acquire(Conn())) as acquire:
response = await admin_affect.get_session_detail(
user_id=ADMIN_ID,
session_id=SESSION_ID,
limit=2,
before_seq=9,
)
acquire.assert_called_once_with(role="admin", user_id=ADMIN_ID)
self.assertEqual([record.seq for record in response.traces], [7, 8])
self.assertTrue(response.has_more)
self.assertEqual(response.total_traces, 4)
self.assertEqual(response.current_emotions["anxiety"], 0.6)
self.assertEqual(response.current_emotions["hope"], 0.2)
self.assertIsNone(response.current_emotions["sadness"])
self.assertIsNone(response.current_emotions["anger"])
async def test_detail_keeps_trace_free_legacy_session_observable(self) -> None:
class Conn:
async def fetchrow(self, query: str, *args: object) -> dict[str, object]:
if "FROM app.sessions AS s" in query:
return {
"session_id": SESSION_ID,
"persona_code": "P4",
"affect_state": {},
}
return {"total_traces": 0}
async def fetch(self, query: str, *args: object) -> list[dict[str, object]]:
return []
with patch.object(admin_affect, "acquire", return_value=_Acquire(Conn())):
response = await admin_affect.get_session_detail(
user_id=ADMIN_ID,
session_id=SESSION_ID,
limit=100,
before_seq=None,
)
self.assertEqual(response.traces, [])
self.assertEqual(response.total_traces, 0)
self.assertFalse(response.has_more)
self.assertTrue(all(value is None for value in response.current_emotions.values()))
async def test_detail_maps_missing_session_and_database_failure_separately(self) -> None:
class MissingConn:
async def fetchrow(self, query: str, *args: object) -> None:
return None
with patch.object(admin_affect, "acquire", return_value=_Acquire(MissingConn())):
with self.assertRaises(admin_affect.AdminAffectSessionNotFoundError):
await admin_affect.get_session_detail(
user_id=ADMIN_ID,
session_id=SESSION_ID,
limit=100,
before_seq=None,
)
with patch.object(admin_affect, "acquire", side_effect=RuntimeError("database down")):
with self.assertRaises(admin_affect.AdminAffectPersistenceError):
await admin_affect.list_sessions(user_id=ADMIN_ID, limit=30, offset=0)
class AdminAffectHttpBoundaryTest(unittest.TestCase):
def _client(self, principal: Principal) -> TestClient:
app = FastAPI()
app.include_router(admin_routes.router)
app.dependency_overrides[get_current_principal] = lambda: principal
return TestClient(app)
def test_non_admin_is_rejected_before_service_access(self) -> None:
principal = Principal(user_id="learner-1", role=Role.LEARNER)
with patch.object(admin_affect, "list_sessions", AsyncMock()) as list_sessions:
response = self._client(principal).get("/admin/affect/sessions")
self.assertEqual(response.status_code, 403)
list_sessions.assert_not_awaited()
def test_invalid_uuid_and_pagination_are_validation_errors(self) -> None:
client = self._client(_admin_principal())
self.assertEqual(client.get("/admin/affect/sessions?limit=0").status_code, 422)
self.assertEqual(
client.get("/admin/affect/sessions/not-a-uuid").status_code,
422,
)
self.assertEqual(
client.get(f"/admin/affect/sessions/{SESSION_ID}?before_seq=0").status_code,
422,
)
def test_route_maps_persistence_error_to_503(self) -> None:
with patch.object(
admin_affect,
"list_sessions",
AsyncMock(side_effect=admin_affect.AdminAffectPersistenceError("database down")),
):
response = self._client(_admin_principal()).get("/admin/affect/sessions")
self.assertEqual(response.status_code, 503)
def test_route_maps_missing_session_to_404(self) -> None:
with patch.object(
admin_affect,
"get_session_detail",
AsyncMock(side_effect=admin_affect.AdminAffectSessionNotFoundError("missing")),
):
response = self._client(_admin_principal()).get(
f"/admin/affect/sessions/{SESSION_ID}"
)
self.assertEqual(response.status_code, 404)
def test_route_serializes_list_response(self) -> None:
result = AdminAffectSessionListResponse(
runtime=AdminAffectRuntimeResponse(
enabled=True,
provider="openrouter",
model="~typesafe/jev-latest",
configured=True,
),
sessions=[],
total=0,
limit=30,
offset=0,
)
with patch.object(admin_affect, "list_sessions", AsyncMock(return_value=result)) as list_sessions:
response = self._client(_admin_principal()).get("/admin/affect/sessions")
self.assertEqual(response.status_code, 200)
self.assertEqual(response.json()["runtime"]["provider"], "openrouter")
self.assertEqual(response.json()["sessions"], [])
list_sessions.assert_awaited_once_with(user_id=ADMIN_ID, limit=30, offset=0)

View file

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

View file

@ -1,346 +0,0 @@
"""관리자 전용 Jev 감정 trace와 원자 영속화 회귀."""
from __future__ import annotations
import math
import unittest
from dataclasses import replace
from unittest.mock import patch
from unittest.mock import AsyncMock
from pydantic import ValidationError
from . import session_persistence, turn_runtime
from .contracts.client_affect import ClientAffectDimensionTraceV1, ClientAffectTraceV1
from .services import client_affect, orchestrator, persona, state_machine
from .services.jev_client import AppraisalResult, EMOTION_DIMENSIONS, EmotionEstimate
from .store import InProcSession, TurnRecord
def _appraisal(
*,
confidence: float = 0.5,
probabilities: tuple[float, ...] | None = (0.0, 0.5, 0.5, 0.0, 0.0),
) -> AppraisalResult:
return AppraisalResult(
emotions={
dimension: EmotionEstimate(
score=0.375,
confidence=confidence,
probabilities=probabilities,
)
for dimension in EMOTION_DIMENSIONS
},
model="jev-test",
latency_ms=11,
input_tokens=13,
output_tokens=17,
provider="typesafe",
cost_usd=None,
)
def _trace() -> ClientAffectTraceV1:
appraisal = _appraisal()
return _trace_from_appraisal(appraisal)
def _trace_from_appraisal(appraisal: AppraisalResult) -> ClientAffectTraceV1:
before = {f"emotion_{dimension}": 0.5 for dimension in EMOTION_DIMENSIONS}
transition = client_affect.transition_emotions(
before,
{},
appraisal,
min_confidence=0.65,
)
return client_affect.build_client_affect_trace(
affect_state_before=before,
affect_baseline={},
affect_state_after=transition.affect_state,
appraisal=appraisal,
transition=transition,
turn_seq=1,
stage="라포",
resistance=0.65,
effective_openness=0.15,
rapport_credit=1.25,
min_confidence=0.65,
)
class ClientAffectTraceContractTest(unittest.TestCase):
def test_trace_preserves_transition_values_and_tentative_is_not_accepted(self) -> None:
trace = _trace()
self.assertEqual(trace.schema_version, 1)
self.assertEqual(trace.context.rapport_credit, 1.25)
self.assertEqual(
tuple(dimension.key for dimension in trace.dimensions),
EMOTION_DIMENSIONS,
)
self.assertTrue(
all(dimension.decision == "tentative" for dimension in trace.dimensions)
)
self.assertEqual(trace.dimensions[0].before, 0.5)
self.assertEqual(trace.dimensions[0].target, 0.375)
self.assertEqual(trace.dimensions[0].after, 0.48125)
self.assertEqual(trace.dimensions[0].probabilities, (0.0, 0.5, 0.5, 0.0, 0.0))
def test_dimension_contract_rejects_nonfinite_probability(self) -> None:
for invalid in (math.nan, math.inf):
with self.subTest(invalid=invalid), self.assertRaises(ValidationError):
ClientAffectDimensionTraceV1(
key="anxiety",
before=0.0,
target=None,
after=0.0,
confidence=None,
probabilities=(0.0, invalid, 0.0, 0.0, 1.0),
decision="held",
)
def test_accepted_and_held_trace_values_preserve_nullable_inputs(self) -> None:
accepted = _trace_from_appraisal(
_appraisal(
confidence=0.9,
probabilities=(0.0, 0.0, 0.4, 0.6, 0.0),
)
)
held_appraisal = AppraisalResult(
emotions={
dimension: EmotionEstimate(
score=math.nan,
confidence=None,
probabilities=None,
)
for dimension in EMOTION_DIMENSIONS
},
model="jev-test",
latency_ms=11,
input_tokens=13,
output_tokens=17,
provider="typesafe",
cost_usd=None,
)
held = _trace_from_appraisal(held_appraisal)
self.assertTrue(all(item.decision == "accepted" for item in accepted.dimensions))
self.assertEqual(
accepted.dimensions[0].probabilities,
(0.0, 0.0, 0.4, 0.6, 0.0),
)
self.assertEqual(accepted.dimensions[0].target, 0.375)
self.assertEqual(accepted.dimensions[0].confidence, 0.9)
self.assertTrue(all(item.decision == "held" for item in held.dimensions))
self.assertTrue(
all(
item.target is None
and item.confidence is None
and item.probabilities is None
for item in held.dimensions
)
)
class _Transaction:
def __init__(self) -> None:
self.error: type[BaseException] | None = None
async def __aenter__(self) -> None:
return None
async def __aexit__(self, exc_type, exc, tb) -> bool:
self.error = exc_type
return False
class _Connection:
def __init__(self, *, fail_trace_insert: bool = False) -> None:
self.fail_trace_insert = fail_trace_insert
self.transaction_context = _Transaction()
self.executed: list[str] = []
def transaction(self) -> _Transaction:
return self.transaction_context
async def fetchval(self, query: str, *args: object) -> object:
if "FROM app.sessions" in query:
return "00000000-0000-0000-0000-000000000111"
if "COALESCE(MAX(seq)" in query:
return 2
if "INSERT INTO app.turns" in query:
return "00000000-0000-0000-0000-000000000222"
raise AssertionError(f"unexpected query: {query}")
async def execute(self, query: str, *args: object) -> str:
self.executed.append(query)
if self.fail_trace_insert and "INSERT INTO app.client_affect_trace" in query:
raise RuntimeError("trace insert failed")
return "INSERT 0 1"
class _Acquire:
def __init__(self, conn: _Connection) -> None:
self.conn = conn
async def __aenter__(self) -> _Connection:
return self.conn
async def __aexit__(self, exc_type, exc, tb) -> bool:
return False
class ClientAffectTracePersistenceTest(unittest.IsolatedAsyncioTestCase):
async def test_atomic_write_assigns_turn_id_only_after_trace_and_state_write(self) -> None:
conn = _Connection()
turn = TurnRecord(
turn_seq=1,
speaker="client",
stage="라포",
text="조금 더 이야기해볼게요.",
text_masked="조금 더 이야기해볼게요.",
)
state = state_machine.SessionState(turn_seq=1)
with (
patch.object(session_persistence, "get_pool", return_value=object()),
patch.object(session_persistence, "acquire", return_value=_Acquire(conn)),
):
stored = await session_persistence.append_client_turn_with_affect_trace(
session_id="00000000-0000-0000-0000-000000000111",
learner_id="00000000-0000-0000-0000-000000000101",
turn=turn,
state=state,
trace=_trace(),
)
self.assertTrue(stored)
self.assertEqual(turn.turn_id, "00000000-0000-0000-0000-000000000222")
self.assertIsNone(conn.transaction_context.error)
self.assertIn("INSERT INTO app.client_affect_trace", conn.executed[0])
self.assertIn("INSERT INTO app.session_state", conn.executed[1])
async def test_atomic_write_keeps_turn_identifier_unpublished_when_trace_insert_fails(self) -> None:
conn = _Connection(fail_trace_insert=True)
turn = TurnRecord(
turn_seq=1,
speaker="client",
stage="라포",
text="조금 더 이야기해볼게요.",
text_masked="조금 더 이야기해볼게요.",
)
with (
patch.object(session_persistence, "get_pool", return_value=object()),
patch.object(session_persistence, "acquire", return_value=_Acquire(conn)),
):
with self.assertRaises(session_persistence.ClientAffectTracePersistenceError):
await session_persistence.append_client_turn_with_affect_trace(
session_id="00000000-0000-0000-0000-000000000111",
learner_id="00000000-0000-0000-0000-000000000101",
turn=turn,
state=state_machine.SessionState(turn_seq=1),
trace=_trace(),
)
self.assertIsNone(turn.turn_id)
self.assertIs(conn.transaction_context.error, RuntimeError)
class ClientAffectTraceRuntimeTest(unittest.IsolatedAsyncioTestCase):
def _session_and_result(
self,
) -> tuple[InProcSession, orchestrator.TurnContext, orchestrator.TurnResult]:
state_before = state_machine.SessionState()
state_after = replace(state_before, turn_seq=1)
sess = InProcSession(
session_id="trace-runtime-session",
case_id="trace-runtime-case",
learner_id="00000000-0000-0000-0000-000000000101",
persona_code=persona.P1.code,
theory_mode="humanistic",
persona=persona.P1,
state=state_before,
)
ctx = orchestrator.TurnContext(
session_id=sess.session_id,
case_id=sess.case_id,
persona=sess.persona,
state_before=state_before,
learner_text_raw="그 마음을 조금 더 들려주실 수 있을까요?",
learner_text_masked="그 마음을 조금 더 들려주실 수 있을까요?",
state_after=state_after,
client_affect_trace=_trace(),
)
result = orchestrator.TurnResult(
turn_seq=1,
stage=state_after.stage.value,
effective_openness=state_after.effective_openness,
client_reply="조금 더 이야기해볼게요.",
safety_flagged=False,
state_after=state_after,
)
return sess, ctx, result
async def test_trace_path_updates_runtime_mirrors_only_after_atomic_success(self) -> None:
sess, ctx, result = self._session_and_result()
append_counselor = AsyncMock()
append_atomic = AsyncMock(return_value=True)
update_state = AsyncMock()
with (
patch.object(turn_runtime, "append_completed_turn", append_counselor),
patch.object(
session_persistence,
"append_client_turn_with_affect_trace",
append_atomic,
),
patch.object(turn_runtime, "update_session_state", update_state),
):
await turn_runtime.record_completed_turn(
sess,
ctx,
result,
context_prefix="trace test",
)
append_atomic.assert_awaited_once()
append_counselor.assert_awaited_once()
update_state.assert_not_awaited()
self.assertIs(sess.state, result.state_after)
self.assertEqual([turn.speaker for turn in sess.turns], ["client"])
async def test_trace_path_keeps_runtime_mirrors_unchanged_when_atomic_write_fails(self) -> None:
sess, ctx, result = self._session_and_result()
append_counselor = AsyncMock()
update_state = AsyncMock()
with (
patch.object(turn_runtime, "append_completed_turn", append_counselor),
patch.object(
session_persistence,
"append_client_turn_with_affect_trace",
AsyncMock(
side_effect=session_persistence.ClientAffectTracePersistenceError(
"atomic write failed"
)
),
),
patch.object(turn_runtime, "update_session_state", update_state),
):
with self.assertRaises(session_persistence.ClientAffectTracePersistenceError):
await turn_runtime.record_completed_turn(
sess,
ctx,
result,
context_prefix="trace test",
)
append_counselor.assert_awaited_once()
update_state.assert_not_awaited()
self.assertIs(sess.state, ctx.state_before)
self.assertEqual(sess.turns, [])
if __name__ == "__main__":
unittest.main()

View file

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

View file

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

View file

@ -158,33 +158,20 @@ 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,
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,
),
context=f"{context_prefix} turn append",
)
await update_session_state(

View file

@ -1,9 +1,9 @@
{
"schema_version": "vignette.p1_crisis_technical_observations.v2",
"generated_at": "2026-09-22T19:46:03Z",
"base_commit": "acb0d26338dfc5ef73ac84150907707fdf7e0a6a",
"generated_at": "2026-09-22T12:32:40.3710740Z",
"base_commit": "8344bc2ad22151ba1cfcea09da125f71fae0ce81",
"runtime_package": {
"provenance_commit": "acb0d26338dfc5ef73ac84150907707fdf7e0a6a",
"provenance_commit": "8344bc2ad22151ba1cfcea09da125f71fae0ce81",
"matches_head": true,
"files": [
{
@ -18,7 +18,7 @@
},
{
"path": "apps/api/app/services/orchestrator.py",
"sha256": "5bc0edbdfbd2f99a5e2024b5a4376cb371b1ab24bb797b639500fddf0829f7e6",
"sha256": "b9d62fc2fb0dcd7066193745cb35d78019c382e406150dc4a080845dd028e003",
"matches_head": true
},
{

View file

@ -1,39 +0,0 @@
-- =============================================================================
-- Vignette · migration 23 — 관리자 감정 관측용 Jev 전이 trace
-- =============================================================================
-- 기존 turns/session_state 행은 수정하지 않는다. trace는 성공한 client 턴에만
-- 연결되며, 앱 런타임은 client turn INSERT와 state UPSERT를 같은 트랜잭션에 둔다.
CREATE TABLE IF NOT EXISTS app.client_affect_trace (
turn_id UUID PRIMARY KEY REFERENCES app.turns(id) ON DELETE CASCADE,
session_id UUID NOT NULL REFERENCES app.sessions(id) ON DELETE CASCADE,
trace JSONB NOT NULL,
created_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE INDEX IF NOT EXISTS idx_client_affect_trace_session_created
ON app.client_affect_trace (session_id, created_at DESC);
ALTER TABLE app.client_affect_trace ENABLE ROW LEVEL SECURITY;
DROP POLICY IF EXISTS p_client_affect_trace_select_admin ON app.client_affect_trace;
DROP POLICY IF EXISTS p_client_affect_trace_insert_learner ON app.client_affect_trace;
CREATE POLICY p_client_affect_trace_select_admin
ON app.client_affect_trace FOR SELECT
USING (app.current_role_name() = 'admin');
CREATE POLICY p_client_affect_trace_insert_learner
ON app.client_affect_trace FOR INSERT
WITH CHECK (
app.current_role_name() = 'learner'
AND EXISTS (
SELECT 1
FROM app.turns AS turn_row
JOIN app.sessions AS session_row ON session_row.id = turn_row.session_id
WHERE turn_row.id = app.client_affect_trace.turn_id
AND turn_row.session_id = app.client_affect_trace.session_id
AND turn_row.speaker = 'client'
AND session_row.learner_id = app.current_uid()
)
);