vignette/apps/api/app/test_session_turn_persistence.py
2026-08-28 19:58:41 +09:00

2775 lines
108 KiB
Python

"""Regression tests for session turn persistence ordering."""
from __future__ import annotations
import asyncio
import json
import unittest
from types import SimpleNamespace
from unittest.mock import AsyncMock, patch
from fastapi import HTTPException
from .config import settings
from . import session_persistence, turn_runtime
from .contracts.engine_gateway import EngineGatewaySseLineDecoder
from .deps import Principal, Role
from .engine_client import EngineError, GenerateResponse
from .paths import repo_path
from .routes import eval as eval_routes
from .routes import sessions
from .routes import voice as voice_routes
from .services import (
guardrail,
live_coach,
memory,
orchestrator,
persona as persona_service,
rag,
state_machine,
)
from .services.voice import TTSChunk, TranscriptResult, VoicePreset
from .store import InProcSession, TurnRecord, store
async def _decoded_stream_packets(stream_engine, req):
decoder = EngineGatewaySseLineDecoder()
async for raw in stream_engine.stream(req):
packet = decoder.feed_line(raw)
if packet is not None:
yield packet
class _AsyncConnContext:
def __init__(self, conn):
self.conn = conn
async def __aenter__(self):
return self.conn
async def __aexit__(self, exc_type, exc, tb):
return False
def _principal() -> Principal:
return Principal(
user_id="00000000-0000-0000-0000-000000000101",
role=Role.LEARNER,
cohort_ids=[],
email="turn-test@hs.ac.kr",
display_name="Turn Test",
consent_at=1.0,
profile_completed_at=1.0,
)
def _session(principal: Principal) -> InProcSession:
card = persona_service.P1
sess = InProcSession(
session_id="turn-persistence-session",
case_id="turn-persistence-case",
learner_id=principal.user_id,
persona_code=card.code,
theory_mode="humanistic",
persona=card,
state=state_machine.SessionState(
resistance=card.base_resistance(),
ideation_stage=card.ideation_baseline(),
),
)
store.put(sess)
return sess
def _crisis_case_signal(case_id: str) -> str:
case_set = json.loads(
repo_path(
"data",
"clinical",
"p1-crisis-review-cases.json",
).read_text(encoding="utf-8")
)
return next(
case["synthetic_scenario"]["signal"]
for case in case_set["cases"]
if case["case_id"] == case_id
)
async def _consume_event_source(response: object) -> bytes:
body = bytearray()
iterator = getattr(response, "body_iterator")
async for chunk in iterator:
if isinstance(chunk, str):
body.extend(chunk.encode("utf-8"))
elif isinstance(chunk, (bytes, bytearray)):
body.extend(chunk)
else:
body.extend(str(chunk).encode("utf-8"))
return bytes(body)
class SessionTurnPersistenceTest(unittest.IsolatedAsyncioTestCase):
async def asyncSetUp(self) -> None:
store._sessions.clear()
sessions._RECALL_CACHE.clear()
sessions._KB_CUES_CACHE.clear()
session_persistence._LIVE_COACH_EVENT_CACHE.clear()
async def asyncTearDown(self) -> None:
store._sessions.clear()
sessions._RECALL_CACHE.clear()
sessions._KB_CUES_CACHE.clear()
session_persistence._LIVE_COACH_EVENT_CACHE.clear()
async def test_end_session_does_not_reschedule_evaluation_for_already_ended_session(
self,
) -> None:
principal = _principal()
sess = _session(principal)
sess.ended = True
sess.ended_at = 1_000.0
with (
patch.object(
sessions, "_load_session_or_404", AsyncMock(return_value=sess)
),
patch.object(
sessions, "_end_persisted_session", AsyncMock(return_value=None)
),
patch.object(
sessions, "_schedule_session_evaluation"
) as schedule_session_evaluation,
):
response = await sessions.end_session(sess.session_id, principal)
self.assertEqual(response.session_id, sess.session_id)
schedule_session_evaluation.assert_not_called()
async def test_failed_end_evaluation_keeps_old_review_and_next_session_writable(
self,
) -> None:
principal = _principal()
card = persona_service.P1
case_id = "00000000-0000-0000-0000-00000000ca5e"
old_session = _session(principal)
old_session.case_id = case_id
old_session.session_no = 1
old_session.turns.extend(
[
TurnRecord(
turn_seq=1,
speaker="learner",
stage=old_session.state.stage.value,
text="지금 가장 버거운 마음이 어떤 건가요?",
text_masked="지금 가장 버거운 마음이 어떤 건가요?",
),
TurnRecord(
turn_seq=2,
speaker="client",
stage=old_session.state.stage.value,
text="아무것도 하고 싶지 않아요.",
text_masked="아무것도 하고 싶지 않아요.",
),
]
)
durable_sessions = {old_session.session_id: old_session}
end_order: list[str] = []
async def persist_end(sess: InProcSession, _carry: memory.CarryOver) -> None:
end_order.append("persisted")
sess.ended = True
sess.ended_at = sess.created_at + 120
durable_sessions[sess.session_id] = sess
def schedule_evaluation(sess: InProcSession) -> None:
self.assertTrue(durable_sessions[sess.session_id].ended)
end_order.append("evaluation_scheduled")
with (
patch.object(
sessions,
"_load_session_or_404",
AsyncMock(return_value=old_session),
),
patch.object(sessions, "_end_persisted_session", persist_end),
patch.object(
sessions.rupture_runtime,
"schedule_session_scan",
),
patch.object(
sessions,
"_schedule_session_evaluation",
schedule_evaluation,
),
):
ended = await sessions.end_session(old_session.session_id, principal)
self.assertEqual(ended.session_id, old_session.session_id)
self.assertEqual(ended.session_no, 1)
self.assertEqual(end_order, ["persisted", "evaluation_scheduled"])
self.assertTrue(durable_sessions[old_session.session_id].ended)
self.assertIsNotNone(durable_sessions[old_session.session_id].ended_at)
durable_evaluations: dict[str, dict[str, object]] = {}
async def save_evaluation(
write: session_persistence.SessionEvaluationWrite,
) -> bool:
durable_evaluations[write.session_id] = write.cache_record()
return True
with (
patch.object(
sessions.evaluator,
"evaluate_session",
AsyncMock(side_effect=asyncio.TimeoutError),
),
patch.object(
sessions.session_persistence,
"save_session_evaluation",
save_evaluation,
),
patch.object(
sessions,
"_enqueue_session_review_ready_notification",
AsyncMock(return_value=None),
),
):
await sessions._generate_and_save_session_evaluation(old_session)
old_evaluation = durable_evaluations[old_session.session_id]
self.assertEqual(old_evaluation["status"], "error")
self.assertIn("session evaluation timeout", str(old_evaluation["error"]))
self.assertTrue(durable_sessions[old_session.session_id].ended)
catalog_persona = SimpleNamespace(
card=card,
persona_id="00000000-0000-0000-0000-0000000000a1",
version=3,
degraded=False,
)
case_context = sessions.session_persistence.CaseContext(
case_id=case_id,
last_session_no=1,
)
next_session_id = "00000000-0000-0000-0000-000000000702"
async def create_next_session(**kwargs: object) -> InProcSession:
self.assertEqual(kwargs["case_id"], case_id)
self.assertEqual(kwargs["session_no"], 2)
created = InProcSession(
session_id=next_session_id,
case_id=case_id,
learner_id=principal.user_id,
persona_code=card.code,
theory_mode=str(kwargs["theory_mode"]),
persona=card,
state=kwargs["state"],
session_no=int(kwargs["session_no"]),
prev_rapport_credit=float(kwargs["carry_rapport"]),
)
durable_sessions[created.session_id] = created
return created
def close_background(coro: object) -> None:
close = getattr(coro, "close", None)
if callable(close):
close()
with (
patch.object(
sessions,
"get_catalog_persona",
AsyncMock(return_value=catalog_persona),
),
patch.object(
sessions.session_persistence,
"get_case_context",
AsyncMock(return_value=case_context),
),
patch.object(
sessions,
"_build_seed_recall",
AsyncMock(return_value=memory.RecallContext()),
),
patch.object(
sessions.session_persistence,
"create_session",
create_next_session,
),
patch.object(sessions.asyncio, "create_task", close_background),
):
started = await sessions.start_session(
sessions.SessionStartRequest(persona_code=card.code),
principal,
)
self.assertNotEqual(started.session_id, old_session.session_id)
self.assertEqual(started.session_id, next_session_id)
self.assertEqual(started.case_id, case_id)
self.assertEqual(started.session_no, 2)
next_session = durable_sessions[next_session_id]
persisted_turns: list[TurnRecord] = []
async def persist_turn(**kwargs: object) -> bool:
turn = kwargs["turn"]
assert isinstance(turn, TurnRecord)
persisted_turns.append(turn)
return True
async def successful_turn(
ctx: orchestrator.TurnContext,
_engine: object,
**_kwargs: object,
) -> orchestrator.TurnResult:
assert ctx.state_after is not None
return orchestrator.TurnResult(
turn_seq=ctx.state_after.turn_seq,
stage=ctx.state_after.stage.value,
effective_openness=ctx.state_after.effective_openness,
client_reply="조금 더 이야기해볼게요.",
safety_flagged=False,
state_after=ctx.state_after,
)
with (
patch.object(
sessions,
"_load_session_or_404",
AsyncMock(return_value=next_session),
),
patch.object(
sessions.orchestrator,
"run_turn_generate",
successful_turn,
),
patch.object(
sessions.session_persistence,
"append_turn",
persist_turn,
),
patch.object(
sessions.session_persistence,
"update_state",
AsyncMock(return_value=True),
),
patch.object(
sessions.rupture_runtime,
"schedule_session_scan",
),
):
first_turn = await sessions.submit_turn(
next_session_id,
sessions.TurnRequest(text="지난 이야기부터 이어가도 괜찮을까요?"),
principal,
)
self.assertEqual(first_turn.client_reply, "조금 더 이야기해볼게요.")
self.assertEqual([turn.speaker for turn in persisted_turns], ["counselor", "client"])
self.assertEqual([turn.speaker for turn in next_session.turns], ["counselor", "client"])
with (
patch.object(
sessions,
"_load_session_or_404",
AsyncMock(return_value=old_session),
),
patch.object(
sessions.session_persistence,
"load_session_evaluation",
AsyncMock(return_value=(old_evaluation, True)),
),
patch.object(
sessions.session_persistence,
"load_case_worksheet",
AsyncMock(return_value=(None, True)),
),
):
old_review = await sessions.get_session_review(
old_session.session_id,
principal,
)
self.assertFalse(old_review.reviewReady)
self.assertTrue(old_review.degraded)
self.assertEqual(old_review.supervisorState, "평가 실패")
self.assertIn("AI 평가 재시도가 필요합니다", old_review.summary)
teacher = Principal(
user_id="00000000-0000-0000-0000-000000000202",
role=Role.TEACHER,
cohort_ids=[],
email="teacher@hs.ac.kr",
display_name="Teacher",
)
retry_save = AsyncMock(return_value=True)
with (
patch.object(
eval_routes,
"_load_session_or_404",
AsyncMock(return_value=old_session),
),
patch.object(
eval_routes.evaluator,
"evaluate_session",
AsyncMock(side_effect=EngineError("controlled evaluator 500")),
),
patch.object(
eval_routes.session_persistence,
"save_session_evaluation",
retry_save,
),
):
with self.assertRaises(HTTPException) as retry_error:
await eval_routes.reevaluate_session(
old_session.session_id,
eval_routes.ReevaluateRequest(scope="session_end"),
teacher,
)
self.assertEqual(retry_error.exception.status_code, 503)
retry_save.assert_awaited_once()
retry_write = retry_save.await_args.args[0]
self.assertEqual(retry_write.status, "error")
async def test_append_turn_writes_provider_events_to_db(self) -> None:
class FakeConn:
def __init__(self) -> None:
self.insert_query = ""
self.insert_args: tuple[object, ...] = ()
async def fetchval(self, query: str, *args: object) -> object:
if "SELECT id FROM app.sessions" in query:
return "turn-persistence-session"
if "COALESCE(MAX(seq)" in query:
return 1
if "INSERT INTO app.turns" in query:
self.insert_query = query
self.insert_args = args
return "00000000-0000-0000-0000-000000009999"
return None
class FakeAcquire:
def __init__(self, conn: FakeConn) -> None:
self.conn = conn
async def __aenter__(self) -> FakeConn:
return self.conn
async def __aexit__(self, exc_type, exc, tb) -> None:
return None
conn = FakeConn()
turn = TurnRecord(
turn_seq=1,
speaker="counselor",
stage="rapport",
text="voice text",
text_masked="voice text",
audio_ref="voice:webm:sha256:test",
silence_ms=1234,
speech_rate=210.0,
barge_in=True,
provider_events=[{"type": "sigh", "confidence": 0.82}],
)
with (
patch.object(session_persistence, "get_pool", return_value=object()),
patch.object(
session_persistence,
"acquire",
return_value=FakeAcquire(conn),
),
):
ok = await session_persistence.append_turn(
session_id="turn-persistence-session",
learner_id=_principal().user_id,
turn=turn,
)
self.assertTrue(ok)
self.assertIn("provider_events", conn.insert_query)
self.assertIn("$17::jsonb", conn.insert_query)
self.assertEqual(conn.insert_args[16], [{"type": "sigh", "confidence": 0.82}])
self.assertEqual(conn.insert_args[17], list(turn.visible_to))
self.assertEqual(turn.turn_id, "00000000-0000-0000-0000-000000009999")
async def test_create_session_binds_goal_stages_as_jsonb_array(self) -> None:
class FakeTransaction:
async def __aenter__(self) -> None:
return None
async def __aexit__(self, exc_type, exc, tb) -> None:
return None
class FakeConn:
def __init__(self) -> None:
self.session_insert_args: tuple[object, ...] = ()
def transaction(self) -> FakeTransaction:
return FakeTransaction()
async def fetchrow(self, query: str, *args: object) -> dict[str, object]:
if "UPDATE app.case_profile" in query:
return {"last_session_no": 2}
if "INSERT INTO app.sessions" in query:
self.session_insert_args = args
return {
"id": "00000000-0000-0000-0000-000000000301",
"runtime_case_id": args[0],
"case_id": args[1],
"learner_id": args[2],
"persona_code": args[5],
"session_no": args[8],
"theory_mode": args[9],
"started_at": session_persistence.datetime.fromtimestamp(
1_000.0, tz=session_persistence.timezone.utc
),
"ended_at": None,
"prev_rapport_credit": args[10],
}
raise AssertionError(f"unexpected query: {query}")
async def execute(self, query: str, *args: object) -> str:
self.assert_jsonb_object(query, args)
return "INSERT 0 1"
@staticmethod
def assert_jsonb_object(query: str, args: tuple[object, ...]) -> None:
if "INSERT INTO app.session_state" in query:
assert isinstance(args[8], dict)
conn = FakeConn()
goals = ["라포", "탐색"]
with (
patch.object(session_persistence, "get_pool", return_value=object()),
patch.object(
session_persistence,
"acquire",
return_value=_AsyncConnContext(conn),
),
):
created = await session_persistence.create_session(
learner_id=_principal().user_id,
card=persona_service.P1,
theory_mode="humanistic",
state=state_machine.SessionState(),
session_no=2,
carry_rapport=0.25,
persona_id="00000000-0000-0000-0000-000000000201",
persona_version=1,
case_id="00000000-0000-0000-0000-000000000202",
goal_stages=goals,
learner_feedback_enabled=False,
)
self.assertIsNotNone(created)
self.assertIsInstance(conn.session_insert_args[11], list)
self.assertEqual(conn.session_insert_args[11], goals)
self.assertNotIsInstance(conn.session_insert_args[11], str)
self.assertFalse(conn.session_insert_args[12])
assert created is not None
self.assertFalse(created.learner_feedback_enabled)
async def test_generate_turn_engine_failure_does_not_append_learner_turn(
self,
) -> None:
principal = _principal()
sess = _session(principal)
with patch.object(
sessions.orchestrator,
"run_turn_generate",
AsyncMock(side_effect=EngineError("engine unavailable: test")),
):
with self.assertRaises(sessions.HTTPException) as caught:
await sessions.submit_turn(
sess.session_id,
sessions.TurnRequest(text="실패한 발화"),
principal,
)
self.assertEqual(caught.exception.status_code, 503)
self.assertEqual(sess.turns, [])
async def test_live_coach_turn_marks_runtime_persistence_source(self) -> None:
principal = _principal()
sess = _session(principal)
sess.turns.extend(
[
TurnRecord(
turn_seq=1,
speaker="counselor",
stage="explore",
text="요즘 가장 크게 남는 마음은 무엇인가요?",
text_masked="요즘 가장 크게 남는 마음은 무엇인가요?",
),
TurnRecord(
turn_seq=2,
speaker="client",
stage="explore",
text="잘 모르겠지만 계속 답답해요.",
text_masked="잘 모르겠지만 계속 답답해요.",
),
]
)
suggestion = live_coach.LiveCoachSuggestion(
status="ready",
tone="pos",
focus="emotion",
title="감정 반영이 선명합니다",
message="방금 응답의 정서를 먼저 붙잡았습니다.",
)
with (
patch.object(
sessions,
"_retrieve_live_coach_grounding",
AsyncMock(return_value=[]),
),
patch.object(
sessions.live_coach,
"generate_live_coaching",
AsyncMock(return_value=suggestion),
),
patch.object(
sessions.session_persistence,
"get_live_coach_quota",
AsyncMock(
side_effect=[
({"remaining": 3, "max": 3}, True),
({"remaining": 2, "max": 3}, True),
]
),
),
patch.object(
sessions.session_persistence,
"save_live_coach_event",
AsyncMock(return_value=({"event_id": "cached-live-coach"}, False)),
),
patch.object(
sessions.session_persistence,
"list_live_coach_credit_events",
AsyncMock(return_value=([], True)),
),
):
response = await sessions.live_coach_turn(
sess.session_id,
sessions.LiveCoachRequest(
learner_text="요즘 가장 크게 남는 마음은 무엇인가요?",
client_reply="잘 모르겠지만 계속 답답해요.",
turn_seq=1,
),
principal,
)
self.assertEqual(response.persistence_source, "runtime")
self.assertEqual(response.quota.remaining, 2)
async def test_generate_turn_persists_client_engine_telemetry(self) -> None:
principal = _principal()
sess = _session(principal)
sessions._RECALL_CACHE[sess.session_id] = memory.RecallContext(
recall_summary="직전 회기에서 김서연은 가족 이야기를 열어두었다.",
pinned_facts=["박민수와 주 1회 상담 약속"],
)
async def successful_turn(ctx, engine, **kwargs):
assert ctx.state_after is not None
self.assertIn("[NAME]", ctx.memory.recall_summary or "")
self.assertNotIn("김서연", ctx.memory.recall_summary or "")
self.assertEqual(ctx.memory.pinned_facts, ["[NAME]와 주 1회 상담 약속"])
return orchestrator.TurnResult(
turn_seq=ctx.state_after.turn_seq,
stage=ctx.state_after.stage.value,
effective_openness=ctx.state_after.effective_openness,
client_reply=(
"저는 김서연 씨고 한신대학교 상담심리학과 학생이에요. "
"여기서 뭘 해야 하는지 잘 모르겠는데요."
),
safety_flagged=False,
state_after=ctx.state_after,
llm_provider="claude_cli",
model="gateway-default",
tokens_in=17,
tokens_out=23,
cost_usd=0.012345,
)
with patch.object(sessions.orchestrator, "run_turn_generate", successful_turn):
response = await sessions.submit_turn(
sess.session_id,
sessions.TurnRequest(text="요즘 많이 힘들었겠어요."),
principal,
)
self.assertEqual(
response.client_reply,
"저는 김서연 씨고 한신대학교 상담심리학과 학생이에요. "
"여기서 뭘 해야 하는지 잘 모르겠는데요.",
)
self.assertEqual(len(sess.turns), 2)
learner_turn, client_turn = sess.turns
self.assertIsNone(learner_turn.llm_provider)
self.assertEqual(
client_turn.text,
"저는 김서연 씨고 한신대학교 상담심리학과 학생이에요. "
"여기서 뭘 해야 하는지 잘 모르겠는데요.",
)
self.assertNotIn("김서연", client_turn.text_masked)
self.assertNotIn("한신대학교", client_turn.text_masked)
self.assertNotIn("상담심리학과", client_turn.text_masked)
self.assertIn("[NAME]", client_turn.text_masked)
self.assertIn("[ORG]", client_turn.text_masked)
self.assertIn("잘 모르겠는데요", client_turn.text_masked)
self.assertNotIn("[NAME]는데요", client_turn.text_masked)
self.assertEqual(client_turn.llm_provider, "claude_cli")
self.assertEqual(client_turn.model, "gateway-default")
self.assertEqual(client_turn.tokens_in, 17)
self.assertEqual(client_turn.tokens_out, 23)
self.assertEqual(client_turn.cost_usd, 0.012345)
async def test_generate_real_crisis_stops_before_engine_and_returns_109_resource(
self,
) -> None:
principal = _principal()
sess = _session(principal)
async def should_not_generate(*args, **kwargs):
raise AssertionError("engine must not be called for learner_real crisis")
with patch.object(sessions.engine_client, "generate", should_not_generate):
response = await sessions.submit_turn(
sess.session_id,
sessions.TurnRequest(text=_crisis_case_signal("P1-CRISIS-004")),
principal,
)
self.assertTrue(response.safety_flagged)
self.assertTrue(response.conversation_stopped)
self.assertEqual(response.crisis_kind, "learner_real")
self.assertIsNotNone(response.crisis_resource)
self.assertEqual(response.crisis_resource.number, "109")
self.assertIsNone(response.client_reply)
self.assertEqual(len(sess.turns), 1)
self.assertEqual(sess.turns[0].speaker, "counselor")
async def test_record_safety_event_writes_teacher_alert_payload(self) -> None:
principal = _principal()
sess = _session(principal)
ctx = orchestrator.prepare_turn(
session_id=sess.session_id,
case_id=sess.case_id,
card=sess.persona,
state=sess.state,
learner_text="저 지금 자살하고 싶어요. 도와주세요.",
theory_mode=sess.theory_mode,
)
result = orchestrator.TurnResult(
turn_seq=ctx.state_after.turn_seq if ctx.state_after else 1,
stage=ctx.state_after.stage.value if ctx.state_after else "라포",
effective_openness=ctx.state_after.effective_openness
if ctx.state_after
else 0.0,
client_reply=None,
safety_flagged=True,
state_after=ctx.state_after or sess.state,
crisis_kind="learner_real",
crisis_resource={"title": "자살예방상담전화 109", "number": "109"},
conversation_stopped=True,
)
calls: list[tuple[str, tuple[object, ...]]] = []
class FakeConn:
async def execute(self, query: str, *args: object) -> str:
calls.append((query, args))
return "INSERT 0 1"
class FakeAcquire:
async def __aenter__(self) -> FakeConn:
return FakeConn()
async def __aexit__(
self, exc_type: object, exc: object, tb: object
) -> None:
return None
with patch.object(
turn_runtime.db, "acquire", return_value=FakeAcquire()
) as acquire:
await turn_runtime.record_safety_event(sess, ctx, result)
acquire.assert_called_once_with(ai_context=True)
self.assertEqual(len(calls), 1)
query, args = calls[0]
self.assertIn("INSERT INTO app.safety_events", query)
self.assertEqual(args[0], sess.session_id)
self.assertEqual(args[1], "learner_real")
self.assertGreaterEqual(args[2], 4)
detail = args[3]
self.assertTrue(detail["conversation_stopped"])
self.assertEqual(detail["crisis_resource"]["number"], "109")
self.assertEqual(detail["alert_status"], "teacher_dashboard")
def test_safety_alert_detail_parser_handles_legacy_json_string(self) -> None:
detail = session_persistence._json_object_payload(
'"{\\"crisis_resource\\":{\\"number\\":\\"109\\"},\\"alert_status\\":\\"teacher_dashboard\\"}"'
)
self.assertEqual(detail["crisis_resource"]["number"], "109")
self.assertEqual(detail["alert_status"], "teacher_dashboard")
async def test_record_safety_event_fails_closed_when_insert_fails_outside_dev(
self,
) -> None:
principal = _principal()
sess = _session(principal)
ctx = orchestrator.prepare_turn(
session_id=sess.session_id,
case_id=sess.case_id,
card=sess.persona,
state=sess.state,
learner_text="저 지금 자살하고 싶어요. 도와주세요.",
theory_mode=sess.theory_mode,
)
result = orchestrator.TurnResult(
turn_seq=ctx.state_after.turn_seq if ctx.state_after else 1,
stage=ctx.state_after.stage.value if ctx.state_after else "라포",
effective_openness=ctx.state_after.effective_openness
if ctx.state_after
else 0.0,
client_reply=None,
safety_flagged=True,
state_after=ctx.state_after or sess.state,
crisis_kind="learner_real",
crisis_resource={"title": "자살예방상담전화 109", "number": "109"},
conversation_stopped=True,
)
class FakeConn:
async def execute(self, query: str, *args: object) -> str:
raise RuntimeError("safety_events unavailable")
class FakeAcquire:
async def __aenter__(self) -> FakeConn:
return FakeConn()
async def __aexit__(
self, exc_type: object, exc: object, tb: object
) -> None:
return None
previous_environment = settings.environment
settings.environment = "staging"
try:
with patch.object(turn_runtime.db, "acquire", return_value=FakeAcquire()):
with self.assertRaises(HTTPException) as raised:
await turn_runtime.record_safety_event(sess, ctx, result)
finally:
settings.environment = previous_environment
self.assertEqual(raised.exception.status_code, 503)
self.assertIn("safety event persistence unavailable", raised.exception.detail)
async def test_stream_turn_persists_client_engine_telemetry(self) -> None:
principal = _principal()
sess = _session(principal)
async def successful_stream(ctx, engine, **kwargs):
assert ctx.state_after is not None
yield orchestrator.StreamEvent("token", {"text": "괜찮아요."})
yield orchestrator.StreamEvent(
"done",
{
"session_id": ctx.session_id,
"stage": ctx.state_after.stage.value,
"effective_openness": ctx.state_after.effective_openness,
"turn_seq": ctx.state_after.turn_seq,
"safety_flagged": False,
"llm_provider": "claude_cli",
"model": "gateway-default",
"tokens_in": 31,
"tokens_out": 37,
"cost_usd": 0.023456,
},
)
with patch.object(sessions.orchestrator, "run_turn_stream", successful_stream):
response = await sessions.stream_turn(
sess.session_id,
sessions.TurnRequest(text="스트림 성공 발화"),
principal,
)
body = await _consume_event_source(response)
self.assertIn(b"done", body)
self.assertEqual(len(sess.turns), 2)
client_turn = sess.turns[1]
self.assertEqual(client_turn.llm_provider, "claude_cli")
self.assertEqual(client_turn.model, "gateway-default")
self.assertEqual(client_turn.tokens_in, 31)
self.assertEqual(client_turn.tokens_out, 37)
self.assertEqual(client_turn.cost_usd, 0.023456)
async def test_stream_real_crisis_stops_before_engine_and_persists_learner_only(
self,
) -> None:
principal = _principal()
sess = _session(principal)
def should_not_stream(*args, **kwargs):
raise AssertionError(
"stream engine must not be called for learner_real crisis"
)
with (
patch.object(sessions.engine_client, "stream", should_not_stream),
patch.object(
turn_runtime, "record_completed_turn", new_callable=AsyncMock
) as completed_turn,
patch.object(
turn_runtime, "record_safety_event", new_callable=AsyncMock
) as safety_event,
):
response = await sessions.stream_turn(
sess.session_id,
sessions.TurnRequest(text=_crisis_case_signal("P1-CRISIS-004")),
principal,
)
body = await _consume_event_source(response)
rendered = body.decode("utf-8")
self.assertIn("'event': 'safety'", rendered)
self.assertIn("'event': 'done'", rendered)
self.assertIn("109", rendered)
self.assertIn("conversation_stopped", rendered)
completed_turn.assert_awaited_once()
safety_event.assert_awaited_once()
saved_result = completed_turn.await_args.args[2]
self.assertTrue(saved_result.safety_flagged)
self.assertTrue(saved_result.conversation_stopped)
self.assertEqual(saved_result.crisis_resource["number"], "109")
async def test_stream_turn_persists_fast_loop_evaluation_on_learner_turn(
self,
) -> None:
principal = _principal()
sess = _session(principal)
async def successful_stream(ctx, engine, **kwargs):
assert ctx.state_after is not None
yield orchestrator.StreamEvent("token", {"text": "조금 말해볼게요."})
yield orchestrator.StreamEvent(
"done",
{
"session_id": ctx.session_id,
"stage": ctx.state_after.stage.value,
"effective_openness": ctx.state_after.effective_openness,
"turn_seq": ctx.state_after.turn_seq,
"safety_flagged": False,
"llm_provider": "claude_cli",
"model": "gateway-default",
"tokens_in": 9,
"tokens_out": 10,
"cost_usd": 0.001,
},
)
async def fake_eval_hook(ctx, client_reply):
return {
"loop": "fast",
"turn_seq": ctx.state_after.turn_seq,
"stage": ctx.state_after.stage.value,
"appropriateness": "pos",
"appropriateness_note": f"응답 반영: {client_reply}",
}
with (
patch.object(sessions.orchestrator, "run_turn_stream", successful_stream),
patch.object(
sessions.evaluator,
"make_eval_hook",
return_value=fake_eval_hook,
),
):
response = await sessions.stream_turn(
sess.session_id,
sessions.TurnRequest(text="스트림 평가 발화"),
principal,
)
await _consume_event_source(response)
if sessions._STREAM_TURN_EVALUATION_TASKS:
await asyncio.gather(*tuple(sessions._STREAM_TURN_EVALUATION_TASKS))
self.assertEqual(len(sess.turns), 2)
learner_turn, client_turn = sess.turns
self.assertEqual(learner_turn.speaker, "counselor")
self.assertIsNotNone(learner_turn.evaluation)
self.assertEqual(learner_turn.evaluation["appropriateness"], "pos")
self.assertIn(
"조금 말해볼게요", learner_turn.evaluation["appropriateness_note"]
)
self.assertIsNone(client_turn.evaluation)
async def test_stream_turn_done_does_not_wait_for_fast_loop_evaluation(
self,
) -> None:
principal = _principal()
sess = _session(principal)
evaluation_started = asyncio.Event()
release_evaluation = asyncio.Event()
async def successful_stream(ctx, engine, **kwargs):
assert ctx.state_after is not None
yield orchestrator.StreamEvent("token", {"text": "지금 답할게요."})
yield orchestrator.StreamEvent(
"done",
{
"session_id": ctx.session_id,
"turn_seq": ctx.state_after.turn_seq,
"safety_flagged": False,
"llm_provider": "claude_cli",
"model": "gateway-default",
},
)
async def slow_eval_hook(ctx, client_reply):
evaluation_started.set()
await release_evaluation.wait()
return {
"loop": "fast",
"turn_seq": ctx.state_after.turn_seq,
"stage": ctx.state_after.stage.value,
"appropriateness": "pos",
}
with (
patch.object(sessions.orchestrator, "run_turn_stream", successful_stream),
patch.object(
sessions.evaluator,
"make_eval_hook",
return_value=slow_eval_hook,
),
):
response = await sessions.stream_turn(
sess.session_id,
sessions.TurnRequest(text="평가를 기다리지 않는 발화"),
principal,
)
consume_task = asyncio.create_task(_consume_event_source(response))
await asyncio.wait_for(evaluation_started.wait(), timeout=1)
body = await asyncio.wait_for(consume_task, timeout=0.2)
self.assertIn("'event': 'done'", body.decode("utf-8"))
self.assertEqual(len(sess.turns), 2)
self.assertIsNone(sess.turns[0].evaluation)
release_evaluation.set()
if sessions._STREAM_TURN_EVALUATION_TASKS:
await asyncio.gather(*tuple(sessions._STREAM_TURN_EVALUATION_TASKS))
self.assertEqual(sess.turns[0].evaluation["appropriateness"], "pos")
async def test_stream_turn_surfaces_fast_loop_evaluation_failure_on_review(
self,
) -> None:
principal = _principal()
sess = _session(principal)
async def successful_stream(ctx, engine, **kwargs):
assert ctx.state_after is not None
yield orchestrator.StreamEvent("token", {"text": "조금 말해볼게요."})
yield orchestrator.StreamEvent(
"done",
{
"session_id": ctx.session_id,
"stage": ctx.state_after.stage.value,
"effective_openness": ctx.state_after.effective_openness,
"turn_seq": ctx.state_after.turn_seq,
"safety_flagged": False,
"llm_provider": "claude_cli",
"model": "gateway-default",
"tokens_in": 9,
"tokens_out": 10,
"cost_usd": 0.001,
},
)
async def failing_eval_hook(ctx, client_reply):
raise RuntimeError("김서연 평가 timeout 010-1234-5678")
with (
patch.object(sessions.orchestrator, "run_turn_stream", successful_stream),
patch.object(
sessions.evaluator,
"make_eval_hook",
return_value=failing_eval_hook,
),
):
response = await sessions.stream_turn(
sess.session_id,
sessions.TurnRequest(text="스트림 평가 실패 발화"),
principal,
)
await _consume_event_source(response)
if sessions._STREAM_TURN_EVALUATION_TASKS:
await asyncio.gather(*tuple(sessions._STREAM_TURN_EVALUATION_TASKS))
self.assertEqual(len(sess.turns), 2)
learner_turn = sess.turns[0]
self.assertIsNotNone(learner_turn.evaluation)
assert learner_turn.evaluation is not None
self.assertEqual(learner_turn.evaluation["loop"], "fast")
self.assertEqual(learner_turn.evaluation["appropriateness"], "neutral")
self.assertEqual(learner_turn.evaluation["error"], "RuntimeError")
self.assertNotIn("김서연", learner_turn.evaluation["error"])
self.assertNotIn("010-1234-5678", learner_turn.evaluation["error"])
review = sessions.build_session_review(
sessions.SessionReviewReadInput(
session=sess,
evaluation_record=None,
evaluation_durable=False,
)
)
review_turn = review.turns[0]
self.assertEqual(review_turn.id, "t1")
self.assertEqual(review_turn.turn_id, learner_turn.turn_id)
self.assertIsNotNone(review_turn.note)
assert review_turn.note is not None
self.assertEqual(review_turn.note.title, "턴 직후 평가 실패")
self.assertIn(
"fast-loop(턴 직후) 평가를 완료하지 못했습니다", review_turn.note.body
)
self.assertIn("AI 평가 재시도가 필요합니다", review_turn.note.body)
self.assertNotIn("RuntimeError", review_turn.note.body)
self.assertNotIn("김서연", review_turn.note.body)
self.assertNotIn("010-1234-5678", review_turn.note.body)
async def test_live_coach_degrades_to_rule_based_suggestion_when_engine_fails(
self,
) -> None:
principal = _principal()
sess = _session(principal)
sess.turns.append(
TurnRecord(
turn_seq=1,
speaker="counselor",
stage=sess.state.stage.value,
text="그냥 학교는 가야 하는 거 아닐까요?",
text_masked="그냥 학교는 가야 하는 거 아닐까요?",
evaluation={
"appropriateness": "warn",
"appropriateness_note": "조언이 빠름",
},
)
)
with (
patch.object(
sessions,
"_retrieve_live_coach_grounding",
AsyncMock(return_value=[]),
),
patch.object(
sessions.engine_client,
"generate",
AsyncMock(side_effect=EngineError("offline")),
),
):
response = await sessions.live_coach_turn(
sess.session_id,
sessions.LiveCoachRequest(
learner_text="그냥 학교는 가야 하는 거 아닐까요?",
client_reply="몰라요. 그런 말 들으려고 온 건 아닌데요.",
turn_seq=1,
),
principal,
)
self.assertEqual(response.status, "degraded")
self.assertEqual(response.tone, "warn")
self.assertEqual(response.focus, "rapport")
self.assertIn("조언", response.title + response.message)
self.assertTrue(response.next_utterance)
self.assertTrue(response.sources)
self.assertEqual(
response.sources[0].source_id, "official_counseling_guideline_seed"
)
self.assertNotIn(
"workbook_0615_case_conceptualization",
{source.source_id for source in response.sources},
)
history = await sessions.list_live_coach_history(sess.session_id, principal)
self.assertEqual(history.source, "runtime")
self.assertEqual(len(history.events), 1)
self.assertEqual(history.events[0].turn_seq, 1)
self.assertEqual(history.events[0].stage, "라포")
self.assertEqual(history.events[0].suggestion.status, "degraded")
self.assertEqual(history.events[0].suggestion.title, response.title)
self.assertIn("학교", history.events[0].learner_text_excerpt or "")
session_persistence._LIVE_COACH_EVENT_CACHE[sess.session_id][0]["stage"] = (
"unknown-stage"
)
legacy_history = await sessions.list_live_coach_history(
sess.session_id, principal
)
self.assertIsNone(legacy_history.events[0].stage)
async def test_live_coach_degrades_when_llm_audit_is_unavailable(self) -> None:
principal = _principal()
sess = _session(principal)
sess.turns.append(
TurnRecord(
turn_seq=1,
speaker="counselor",
stage=sess.state.stage.value,
text="그 마음이 컸겠네요.",
text_masked="그 마음이 컸겠네요.",
evaluation={
"appropriateness": "pos",
"appropriateness_note": "감정 반영",
},
)
)
engine_response = GenerateResponse(
text="",
model="fake-live-coach-model",
provider="fake-provider",
tokens_in=11,
tokens_out=7,
cost_usd=0.01,
structured={
"tone": "pos",
"focus": "emotion",
"title": "엔진 생성 코칭",
"message": "감정 반영을 이어가세요.",
"next_utterance": "그 마음이 가장 컸던 순간이 언제였나요?",
},
)
with (
patch.object(
sessions,
"_retrieve_live_coach_grounding",
AsyncMock(return_value=[]),
),
patch.object(
sessions.engine_client,
"generate",
AsyncMock(return_value=engine_response),
) as generate,
patch.object(
session_persistence,
"record_llm_call_audit",
AsyncMock(return_value=False),
) as audit,
):
response = await sessions.live_coach_turn(
sess.session_id,
sessions.LiveCoachRequest(
learner_text="그 마음이 컸겠네요.",
client_reply="네, 아무도 몰라주는 것 같았어요.",
turn_seq=1,
),
principal,
)
generate.assert_awaited_once()
audit.assert_awaited_once()
self.assertEqual(response.status, "degraded")
self.assertNotEqual(response.title, "엔진 생성 코칭")
self.assertIn("응답 검증 기록", response.rationale or "")
history = await sessions.list_live_coach_history(sess.session_id, principal)
self.assertEqual(len(history.events), 1)
self.assertEqual(history.events[0].suggestion.status, "degraded")
async def test_live_coach_rag_grounding_preserves_source_pack_metadata(
self,
) -> None:
retrieval = rag.RetrievalResult(
chunks=[
rag.RetrievedChunk(
chunk_id=17,
score=0.91,
kb_kind="supervisor_pattern",
heading_path="risk/protective-factors",
context_prefix="공식 자살위험 평가 지침 / 보호요인 / 2026-06-15",
body="최근성, 보호요인, 안전계획을 확인한다.",
meta={
"source_title": "공식 자살위험 평가 지침",
"source_type": "official_guideline_summary",
"source_version": "2026-06-15",
"citation": "허가된 공식 지침 요약",
},
source_id="official_suicide_risk_guidelines",
license_class="B",
external_llm_ok=True,
)
],
policy_name="evaluator:k=supervisor_pattern:s<=2",
top1_score=0.91,
latency_ms=12,
query_text="정리 humanistic 자살 위험 최근성 보호요인",
)
retrieve_grounding = AsyncMock(return_value=retrieval)
with (
patch.object(
sessions.db, "acquire", return_value=_AsyncConnContext(object())
),
patch.object(
sessions.rag,
"retrieve_eval_grounding",
retrieve_grounding,
),
patch.object(sessions.rag, "log_retrieval", AsyncMock()),
):
grounding = await sessions._retrieve_live_coach_grounding(
learner_text="자살 위험 최근성과 보호요인을 확인했어요.",
client_reply=(
"박민수에게 연락할지 잘 모르겠는데요. 그래도 오늘은 "
"친구에게 연락할 수 있어요."
),
stage="정리",
theory_mode="humanistic",
)
self.assertEqual(len(grounding), 1)
source = grounding[0]
self.assertEqual(source.source_id, "official_suicide_risk_guidelines")
self.assertEqual(source.title, "공식 자살위험 평가 지침")
self.assertEqual(source.source_type, "official_guideline_summary")
self.assertEqual(source.version, "2026-06-15")
self.assertEqual(source.citation, "허가된 공식 지침 요약")
self.assertEqual(source.license_class, "B")
self.assertTrue(source.external_llm_ok)
query = retrieve_grounding.await_args.kwargs["query"]
self.assertNotIn("박민수", query)
self.assertIn("[NAME]에게", query)
self.assertIn("잘 모르겠는데요", query)
self.assertNotIn("[NAME]는데요", query)
async def test_live_coach_persistence_failure_is_not_swallowed(self) -> None:
principal = _principal()
sess = _session(principal)
suggestion = live_coach.LiveCoachSuggestion(
status="ready",
tone="pos",
focus="emotion",
title="감정 반영이 선명합니다",
message="학습자가 내담자의 감정을 먼저 되짚었습니다.",
)
save_error = sessions.HTTPException(
status_code=sessions.status.HTTP_503_SERVICE_UNAVAILABLE,
detail="live coach event save persistence unavailable; runtime fallback is disabled in prod",
)
with (
patch.object(
sessions,
"_retrieve_live_coach_grounding",
AsyncMock(return_value=[]),
),
patch.object(
sessions.live_coach,
"generate_live_coaching",
AsyncMock(return_value=suggestion),
),
patch.object(
session_persistence,
"save_live_coach_event",
AsyncMock(side_effect=save_error),
) as save_event,
):
with self.assertRaises(sessions.HTTPException) as caught:
await sessions.live_coach_turn(
sess.session_id,
sessions.LiveCoachRequest(
learner_text="그 마음이 컸겠네요.",
client_reply="네, 아무도 몰라주는 것 같았어요.",
turn_seq=1,
),
principal,
)
self.assertEqual(caught.exception.status_code, 503)
self.assertIn("live coach event save", str(caught.exception.detail))
save_event.assert_awaited_once()
async def test_live_coach_event_normalizes_legacy_stage_values(self) -> None:
suggestion = live_coach.LiveCoachSuggestion(
status="degraded",
tone="neutral",
focus="exploration",
title="코칭",
message="다음 발화를 준비하세요.",
)
coded = live_coach.LiveCoachEvent(
event_id="event-1",
session_id="session-1",
turn_seq=1,
stage="rapport",
created_at="2026-06-28T00:00:00Z",
suggestion=suggestion,
)
invalid = live_coach.LiveCoachEvent(
event_id="event-2",
session_id="session-1",
turn_seq=2,
stage="unknown-stage",
created_at="2026-06-28T00:00:00Z",
suggestion=suggestion,
)
self.assertEqual(coded.stage, "라포")
self.assertIsNone(invalid.stage)
async def test_live_coach_uses_official_risk_reference_pack_for_crisis_signal(
self,
) -> None:
item = live_coach.LiveCoachInput(
session_id="risk-coach-session",
turn_seq=3,
stage="exploration",
effective_openness=0.45,
theory_mode="humanistic",
persona_code="P1",
persona_name="서연",
learner_text="죽고 싶다는 생각이 들 때도 있나요?",
client_reply="가끔 그런 생각이 들어요.",
recent_turns=[],
evaluation={
"appropriateness": "warn",
"appropriateness_note": "위험사정 필요",
},
)
engine = SimpleNamespace(generate=AsyncMock(side_effect=EngineError("offline")))
suggestion = await live_coach.generate_live_coaching(item, engine=engine)
self.assertEqual(suggestion.focus, "risk")
source_ids = [source.source_id for source in suggestion.sources]
self.assertIn("official_suicide_risk_guidelines", source_ids)
official = next(
source
for source in suggestion.sources
if source.source_id == "official_suicide_risk_guidelines"
)
self.assertEqual(official.source_type, "official_guideline")
self.assertTrue(official.citation)
engine.generate.assert_called_once()
async def test_exact_persona_crisis_cases_fail_closed_to_risk_coaching(self) -> None:
expected_locators = {
"P1-CRISIS-001": "사고·계획·의도·행동 구분",
"P1-CRISIS-002": "사고·계획·의도·행동 구분",
"P1-CRISIS-005": "대한민국 109·119 긴급 연결",
"P1-CRISIS-006": "자해 이후 심리사회적 평가",
}
for case_id, expected_locator in expected_locators.items():
with self.subTest(case_id=case_id):
item = live_coach.LiveCoachInput(
session_id=f"{case_id.lower()}-coach",
turn_seq=4,
stage="exploration",
effective_openness=0.45,
theory_mode="humanistic",
persona_code="P1",
persona_name="서연",
learner_text="그 마음을 조금 더 들려줄 수 있을까요?",
client_reply=_crisis_case_signal(case_id),
recent_turns=[],
)
engine = SimpleNamespace(
generate=AsyncMock(side_effect=AssertionError("crisis coaching must bypass engine"))
)
suggestion = await live_coach.generate_live_coaching(item, engine=engine)
self.assertEqual(suggestion.status, "ready")
self.assertEqual(suggestion.focus, "risk")
self.assertEqual(suggestion.tone, "warn")
self.assertIsNotNone(suggestion.safety_note)
self.assertIn("가상내담자 위기 신호", suggestion.safety_note or "")
self.assertNotEqual(suggestion.safety_note, guardrail.CRISIS_RESOURCE_MESSAGE)
self.assertIn(
expected_locator,
[source.locator for source in suggestion.sources],
)
engine.generate.assert_not_called()
if case_id == "P1-CRISIS-005":
self.assertIn("109", suggestion.message)
self.assertIn("119", suggestion.message)
self.assertIn("응급실", suggestion.message)
self.assertIn("혼자 두지", suggestion.message)
if case_id == "P1-CRISIS-002":
for required in (
"최근성",
"빈도",
"강도",
"계획 유무",
"의도",
"접근 가능성",
"지지자",
"살아갈 이유",
"도움 요청 가능성",
):
self.assertIn(required, suggestion.message)
if case_id == "P1-CRISIS-006":
self.assertIn("비판단", suggestion.message)
self.assertIn("심리사회적", suggestion.message)
self.assertIn("구체적인 수단", suggestion.message)
self.assertIn("전문기관", suggestion.message)
self.assertIn("추후 확인", suggestion.message)
async def test_live_coach_does_not_short_circuit_nonclinical_keyword_contexts(self) -> None:
client_replies = (
"이번 발표에 목숨 걸 만큼 열심히 했어요.",
"자살 예방 교육에서 109를 배웠어요.",
"게임에서 캐릭터 목숨이 하나 남았어요.",
"저는 지금 자살하고 싶지 않아요.",
"요즘은 죽고 싶지 않아요.",
"저는 지금 죽고 싶은 건 아니에요.",
"“저는 지금 죽고 싶어요” 문장을 번역해 주세요.",
"저는 지금 죽고 싶은 생각이 없어요.",
"저는 지금 죽고 싶은 마음이 전혀 없습니다.",
"제가 지금 죽고 싶다고 말한 학생을 돕고 있어요.",
"제가 지금 죽고 싶다고 한 친구를 돕고 있어요.",
"저는 지금 자살 생각이 없어요.",
"저는 지금 자해 충동이 없어요.",
"저는 예전에 죽고 싶었지만 지금은 괜찮아요.",
"제가 지금 여쭤보는 건 자살 생각이 있는지예요.",
"제가 지금 묻는 질문은 혹시 죽고 싶은 생각이 있는지예요.",
"누구나 지금 자살하고 싶다고 말하는 것은 아니에요.",
"저는 지금 자살할 것 같지 않아요.",
"저는 지금 자해할 것 같지 않아요.",
"저는 지금 살기 싫지 않아요.",
"저는 지금 없어지고 싶지 않아요.",
"저는 지금 죽어야 할 것 같지 않아요.",
"저는 지금 살아갈 이유가 없는 건 아니에요.",
"저는 지금 자해 충동이 있지는 않아요.",
)
for index, client_reply in enumerate(client_replies, start=1):
with self.subTest(client_reply=client_reply):
item = live_coach.LiveCoachInput(
session_id=f"nonclinical-keyword-{index}",
turn_seq=2,
stage="exploration",
effective_openness=0.5,
theory_mode="humanistic",
persona_code="P1",
persona_name="서연",
learner_text="그 경험을 조금 더 들려줄 수 있을까요?",
client_reply=client_reply,
recent_turns=[],
)
engine = SimpleNamespace(generate=AsyncMock(side_effect=EngineError("offline")))
suggestion = await live_coach.generate_live_coaching(item, engine=engine)
engine.generate.assert_awaited_once()
self.assertNotEqual(suggestion.focus, "risk")
async def test_live_coach_short_circuits_strong_crisis_variants(self) -> None:
client_replies = (
"저는 지금 자살하려고 해요.",
"지금 당장 죽으려고 해요.",
"제 목숨을 끊으려고 해요.",
"자살 계획이 있어요.",
"자해 시도를 하려고 해요.",
"자살 충동이 있어요.",
"자살할지도 모르겠어요.",
"스스로를 해칠 계획이 있어요.",
"제 안전을 유지할 수 없어요.",
)
for index, client_reply in enumerate(client_replies, start=1):
with self.subTest(client_reply=client_reply):
item = live_coach.LiveCoachInput(
session_id=f"strong-crisis-variant-{index}",
turn_seq=2,
stage="exploration",
effective_openness=0.5,
theory_mode="humanistic",
persona_code="P1",
persona_name="서연",
learner_text="그 경험을 조금 더 들려줄 수 있을까요?",
client_reply=client_reply,
recent_turns=[],
)
engine = SimpleNamespace(
generate=AsyncMock(
side_effect=AssertionError("strong crisis signal must bypass engine")
)
)
suggestion = await live_coach.generate_live_coaching(item, engine=engine)
engine.generate.assert_not_called()
self.assertEqual(suggestion.focus, "risk")
self.assertEqual(suggestion.status, "ready")
async def test_live_coach_prompt_uses_client_role_token_not_persona_name(
self,
) -> None:
item = live_coach.LiveCoachInput(
session_id="masked-coach-session",
turn_seq=1,
stage="rapport",
effective_openness=0.25,
theory_mode="humanistic",
persona_code="P1",
persona_name="서연",
learner_text="오늘 어떤 마음으로 오셨어요?",
client_reply="조금 긴장돼요.",
prior_coach=[{"title": "서연에게 감정을 반영하세요", "focus": "emotion"}],
)
prompt = "\n".join(
message.content for message in live_coach._messages(item, grounding=[])
)
self.assertIn("[내담자] [CLIENT] (P1)", prompt)
self.assertIn("[CLIENT]에게 감정을 반영하세요", prompt)
self.assertNotIn("서연", prompt)
async def test_live_coach_structured_output_masks_client_name_before_storage(
self,
) -> None:
principal = _principal()
sess = _session(principal)
engine_response = GenerateResponse(
text="",
model="fake-live-coach-model",
provider="fake-provider",
tokens_in=12,
tokens_out=9,
cost_usd=0.01,
structured={
"tone": "pos",
"focus": "emotion",
"title": "서연의 감정을 반영하세요",
"message": "서연이 말한 긴장을 한 번 더 따라가세요.",
"next_utterance": "서연님, 그 긴장이 언제 가장 커지나요?",
"rationale": "서연의 표현을 그대로 짚으면 초점이 선명해집니다.",
"safety_note": "서연의 안전 신호도 확인하세요.",
},
)
with (
patch.object(
sessions,
"_retrieve_live_coach_grounding",
AsyncMock(return_value=[]),
),
patch.object(
sessions.engine_client,
"generate",
AsyncMock(return_value=engine_response),
),
patch.object(
session_persistence,
"record_llm_call_audit",
AsyncMock(return_value=True),
),
):
response = await sessions.live_coach_turn(
sess.session_id,
sessions.LiveCoachRequest(
learner_text="그 마음이 컸겠네요.",
client_reply="네, 조금 긴장돼요.",
turn_seq=1,
),
principal,
)
rendered = " ".join(
value
for value in (
response.title,
response.message,
response.next_utterance,
response.rationale,
response.safety_note,
)
if value
)
self.assertNotIn("서연", rendered)
self.assertIn("[CLIENT]", rendered)
history = await sessions.list_live_coach_history(sess.session_id, principal)
persisted = history.events[0].suggestion.model_dump_json()
self.assertNotIn("서연", persisted)
self.assertIn("[CLIENT]", persisted)
async def test_start_session_uses_stable_case_context_and_seed_recall(self) -> None:
principal = _principal()
card = persona_service.P1
case_context = sessions.session_persistence.CaseContext(
case_id="00000000-0000-0000-0000-00000000ca5e",
last_session_no=1,
)
catalog_persona = SimpleNamespace(
card=card,
persona_id="00000000-0000-0000-0000-0000000000a1",
version=3,
degraded=False,
)
recall = memory.RecallContext(
recall_summary="지난 회기에서 가족 이야기를 열어두었다.",
carry={
"rapport_credit": 0.6,
"resistance": card.base_resistance(),
"ideation_stage": card.ideation_baseline(),
},
)
async def fake_create_session(**kwargs):
self.assertEqual(kwargs["case_id"], case_context.case_id)
self.assertEqual(kwargs["session_no"], 2)
self.assertGreater(kwargs["state"].rapport_credit, 0)
return InProcSession(
session_id="stable-case-session",
case_id=kwargs["case_id"],
learner_id=principal.user_id,
persona_code=card.code,
theory_mode=kwargs["theory_mode"],
persona=card,
state=kwargs["state"],
session_no=kwargs["session_no"],
prev_rapport_credit=kwargs["carry_rapport"],
)
def close_background(coro):
coro.close()
return None
with (
patch.object(
sessions, "get_catalog_persona", AsyncMock(return_value=catalog_persona)
),
patch.object(
sessions.session_persistence,
"get_case_context",
AsyncMock(return_value=case_context),
),
patch.object(
sessions,
"_build_seed_recall",
AsyncMock(return_value=recall),
),
patch.object(
sessions.session_persistence,
"create_session",
fake_create_session,
),
patch.object(sessions.asyncio, "create_task", close_background),
):
response = await sessions.start_session(
sessions.SessionStartRequest(persona_code=card.code),
principal,
)
self.assertEqual(response.case_id, case_context.case_id)
self.assertEqual(response.session_no, 2)
self.assertEqual(response.recall_summary, recall.recall_summary)
self.assertIs(sessions._RECALL_CACHE[response.session_id], recall)
async def test_next_session_turn_injects_seed_recall_into_engine_messages(
self,
) -> None:
principal = _principal()
card = persona_service.P1
case_context = sessions.session_persistence.CaseContext(
case_id="00000000-0000-0000-0000-00000000ca5e",
last_session_no=1,
)
catalog_persona = SimpleNamespace(
card=card,
persona_id="00000000-0000-0000-0000-0000000000a1",
version=3,
degraded=False,
)
test_case = self
class FakeConn:
async def fetchrow(self, query: str, *args: object):
if "FROM app.case_profile" in query:
return {"case_digest": "S1: 김서연은 가족 이야기를 열어두었다."}
if "FROM app.session_summary" in query:
return {
"digest": "직전 회기에서 김서연은 침묵 이후 학교 이야기를 꺼냈다.",
"open_threads": ["다음 회기에서 상담 지속 의사를 확인하기"],
"end_state": {
"rapport_credit": 0.55,
"resistance": card.base_resistance(),
"ideation_stage": card.ideation_baseline(),
},
}
return None
async def fetch(self, query: str, *args: object):
test_case.assertIn("FROM app.pinned_fact", query)
test_case.assertIn("$2 = ANY(visible_to)", query)
return [{"value": "박민수와 주 1회 상담 약속"}]
class FakeAcquire:
def __init__(self, conn: FakeConn) -> None:
self.conn = conn
async def __aenter__(self) -> FakeConn:
return self.conn
async def __aexit__(
self, exc_type: object, exc: object, tb: object
) -> None:
return None
async def fake_create_session(**kwargs):
return InProcSession(
session_id="db-seed-recall-session",
case_id=kwargs["case_id"],
learner_id=principal.user_id,
persona_code=card.code,
theory_mode=kwargs["theory_mode"],
persona=card,
state=kwargs["state"],
session_no=kwargs["session_no"],
prev_rapport_credit=kwargs["carry_rapport"],
)
def close_background(coro):
coro.close()
return None
with (
patch.object(
sessions, "get_catalog_persona", AsyncMock(return_value=catalog_persona)
),
patch.object(
sessions.session_persistence,
"get_case_context",
AsyncMock(return_value=case_context),
),
patch.object(sessions.db, "get_pool", return_value=object()),
patch.object(
sessions.db,
"acquire",
return_value=FakeAcquire(FakeConn()),
),
patch.object(
sessions.session_persistence,
"create_session",
fake_create_session,
),
patch.object(sessions.asyncio, "create_task", close_background),
):
response = await sessions.start_session(
sessions.SessionStartRequest(persona_code=card.code),
principal,
)
started = store.get(response.session_id)
self.assertIsNotNone(started)
captured_messages: list[str] = []
async def successful_turn(ctx, engine, **kwargs):
assert ctx.state_after is not None
captured_messages.extend(message.content for message in ctx.messages)
return orchestrator.TurnResult(
turn_seq=ctx.state_after.turn_seq,
stage=ctx.state_after.stage.value,
effective_openness=ctx.state_after.effective_openness,
client_reply="조금 더 이야기해볼게요.",
safety_flagged=False,
state_after=ctx.state_after,
)
with (
patch.object(
sessions, "_load_session_or_404", AsyncMock(return_value=started)
),
patch.object(
sessions.orchestrator,
"run_turn_generate",
successful_turn,
),
):
await sessions.submit_turn(
response.session_id,
sessions.TurnRequest(text="지난번 이야기를 이어가도 괜찮을까요?"),
principal,
)
message_blob = "\n".join(captured_messages)
self.assertIn("[L2 회상", message_blob)
self.assertIn("[케이스 큰그림]", message_blob)
self.assertIn("[직전 회기 요약]", message_blob)
self.assertIn("다음 회기에서 상담 지속 의사를 확인하기", message_blob)
self.assertIn("[L4 고정 사실", message_blob)
self.assertIn("[NAME]와 주 1회 상담 약속", message_blob)
self.assertNotIn("김서연", message_blob)
self.assertNotIn("박민수", message_blob)
async def test_start_session_requires_learner_consent_before_catalog_lookup(
self,
) -> None:
principal = _principal()
principal.consent_at = None
with patch.object(
sessions,
"get_catalog_persona",
AsyncMock(
side_effect=AssertionError(
"consent gate must run before catalog lookup"
)
),
) as get_persona:
with self.assertRaises(sessions.HTTPException) as caught:
await sessions.start_session(
sessions.SessionStartRequest(persona_code=persona_service.P1.code),
principal,
)
self.assertEqual(caught.exception.status_code, 403)
self.assertEqual(caught.exception.detail, "consent_required")
get_persona.assert_not_awaited()
async def test_start_session_requires_onboarding_before_consent_and_catalog_lookup(
self,
) -> None:
principal = _principal()
principal.profile_completed_at = None
with patch.object(
sessions,
"get_catalog_persona",
AsyncMock(
side_effect=AssertionError(
"onboarding gate must run before catalog lookup"
)
),
) as get_persona:
with self.assertRaises(sessions.HTTPException) as caught:
await sessions.start_session(
sessions.SessionStartRequest(persona_code=persona_service.P1.code),
principal,
)
self.assertEqual(caught.exception.status_code, 403)
self.assertEqual(caught.exception.detail, "onboarding_required")
get_persona.assert_not_awaited()
async def test_run_turn_stream_parses_gateway_done_telemetry(self) -> None:
class FakeStreamEngine:
engine_mode = "claude_cli"
default_model = None
async def stream(self, req):
yield "event: token"
yield '{"ignored":"not data"}'
yield 'data: {"text":"부분 응답"}'
yield "event: done"
yield (
'data: {"provider":"claude_cli","model":"gateway-default",'
'"tokens_in":5,"tokens_out":7,"cost_usd":0.034567}'
)
async def stream_packets(self, req):
async for packet in _decoded_stream_packets(self, req):
yield packet
principal = _principal()
sess = _session(principal)
ctx = orchestrator.prepare_turn(
session_id=sess.session_id,
case_id=sess.case_id,
card=sess.persona,
state=sess.state,
learner_text="게이트웨이 스트림 테스트",
memory=orchestrator.TurnMemory(recent_turns=[]),
)
events = [
event
async for event in orchestrator.run_turn_stream(ctx, FakeStreamEngine()) # type: ignore[arg-type]
]
self.assertEqual([event.event for event in events], ["token", "done"])
self.assertEqual(events[0].data["text"], "부분 응답")
self.assertEqual(events[1].data["llm_provider"], "claude_cli")
self.assertEqual(events[1].data["model"], "gateway-default")
self.assertEqual(events[1].data["tokens_in"], 5)
self.assertEqual(events[1].data["tokens_out"], 7)
self.assertEqual(events[1].data["cost_usd"], 0.034567)
async def test_run_turn_stream_rejects_clean_eof_without_gateway_done(self) -> None:
class FakeStreamEngine:
engine_mode = "claude_cli"
default_model = None
async def stream(self, req):
yield "event: token"
yield 'data: {"text":"완료 전 부분 응답"}'
async def stream_packets(self, req):
async for packet in _decoded_stream_packets(self, req):
yield packet
principal = _principal()
sess = _session(principal)
ctx = orchestrator.prepare_turn(
session_id=sess.session_id,
case_id=sess.case_id,
card=sess.persona,
state=sess.state,
learner_text="완료 이벤트 없는 스트림 테스트",
memory=orchestrator.TurnMemory(recent_turns=[]),
)
events = [
event
async for event in orchestrator.run_turn_stream(ctx, FakeStreamEngine()) # type: ignore[arg-type]
]
self.assertEqual([event.event for event in events], ["error"])
self.assertEqual(events[0].data["detail"], "client_stream_incomplete")
async def test_run_turn_stream_treats_gateway_error_event_as_error(self) -> None:
class FakeStreamEngine:
engine_mode = "claude_cli"
default_model = None
async def stream(self, req):
yield "event: token"
yield 'data: {"text":"부분 응답"}'
yield "event: error"
yield 'data: {"detail":"engine unavailable: gateway"}'
async def stream_packets(self, req):
async for packet in _decoded_stream_packets(self, req):
yield packet
principal = _principal()
sess = _session(principal)
ctx = orchestrator.prepare_turn(
session_id=sess.session_id,
case_id=sess.case_id,
card=sess.persona,
state=sess.state,
learner_text="게이트웨이 오류 테스트",
memory=orchestrator.TurnMemory(recent_turns=[]),
)
events = [
event
async for event in orchestrator.run_turn_stream(ctx, FakeStreamEngine()) # type: ignore[arg-type]
]
self.assertEqual([event.event for event in events], ["error"])
self.assertIn("engine unavailable", events[0].data["detail"])
async def test_stream_turn_engine_error_event_does_not_append_partial_turns(
self,
) -> None:
principal = _principal()
sess = _session(principal)
async def failing_stream(*args, **kwargs):
yield orchestrator.StreamEvent("token", {"text": "부분 응답"})
yield orchestrator.StreamEvent(
"error", {"detail": "engine unavailable: stream"}
)
with patch.object(sessions.orchestrator, "run_turn_stream", failing_stream):
response = await sessions.stream_turn(
sess.session_id,
sessions.TurnRequest(text="스트림 실패 발화"),
principal,
)
body = await _consume_event_source(response)
self.assertIn(b"engine unavailable: stream", body)
self.assertEqual(sess.turns, [])
async def test_voice_turn_engine_failure_does_not_append_learner_turn(self) -> None:
class FakeWebSocket:
def __init__(self) -> None:
self.messages: list[dict[str, object]] = []
self.client_state = voice_routes.WebSocketState.CONNECTED
async def send_text(self, data: str) -> None:
self.messages.append(json.loads(data))
principal = _principal()
sess = _session(principal)
websocket = FakeWebSocket()
with patch.object(
voice_routes.orchestrator,
"run_turn_generate",
AsyncMock(side_effect=EngineError("voice engine unavailable")),
):
await voice_routes._run_turn_and_speak(
websocket, # type: ignore[arg-type]
voice_routes.VoiceSessionContext(
session_id=sess.session_id,
principal=principal,
voice_preset=VoicePreset(preset="neutral", openai_voice="sage"),
),
voice_routes.VoiceTurnInput(learner_text="음성 실패 발화"),
)
self.assertTrue(
any(
message.get("type") == "error"
and "engine unavailable" in str(message.get("detail"))
for message in websocket.messages
),
websocket.messages,
)
self.assertEqual(sess.turns, [])
async def test_voice_audio_turn_persists_paralinguistic_metadata(self) -> None:
class FakeWebSocket:
def __init__(self) -> None:
self.messages: list[dict[str, object]] = []
self.binary: list[bytes] = []
self.client_state = voice_routes.WebSocketState.CONNECTED
async def send_text(self, data: str) -> None:
self.messages.append(json.loads(data))
async def send_bytes(self, data: bytes) -> None:
self.binary.append(data)
async def successful_turn(ctx, engine, **kwargs):
assert ctx.state_after is not None
return orchestrator.TurnResult(
turn_seq=ctx.state_after.turn_seq,
stage=ctx.state_after.stage.value,
effective_openness=ctx.state_after.effective_openness,
client_reply="천천히 말해줘서 고마워요.",
safety_flagged=False,
state_after=ctx.state_after,
llm_provider="claude_cli",
model="gateway-default",
tokens_in=11,
tokens_out=13,
cost_usd=0.0012,
)
async def fake_synthesize_stream(text, voice_preset):
yield TTSChunk(audio=b"tts-audio")
principal = _principal()
sess = _session(principal)
websocket = FakeWebSocket()
audio = b"\x00\x80" * 1600
with (
patch.object(
voice_routes.voice_service,
"transcribe",
AsyncMock(
return_value=TranscriptResult(
text="오늘은 좀 힘들었어요.",
duration=2.0,
provider_events=[
{
"kind": "voice_activity",
"start_ms": 10,
"raw_text": "drop",
},
],
)
),
),
patch.object(
voice_routes.orchestrator,
"run_turn_generate",
successful_turn,
),
patch.object(
voice_routes.voice_service,
"synthesize_stream",
fake_synthesize_stream,
),
):
await voice_routes._handle_utterance(
websocket, # type: ignore[arg-type]
voice_routes.VoiceSessionContext(
session_id=sess.session_id,
principal=principal,
voice_preset=VoicePreset(preset="neutral", openai_voice="sage"),
),
voice_routes.VoiceAudioInput(
audio=audio,
fmt="webm",
prosody=voice_routes.VoiceProsody(
silence_ms=1234,
barge_in=True,
provider_events=[
{"type": "sigh", "confidence": 0.82, "text": "drop"}
],
),
),
)
self.assertEqual(len(sess.turns), 2)
learner_turn, client_turn = sess.turns
self.assertTrue(str(learner_turn.audio_ref).startswith("voice:webm:sha256:"))
self.assertEqual(learner_turn.silence_ms, 1234)
self.assertGreater(learner_turn.speech_rate or 0, 0)
self.assertTrue(learner_turn.barge_in)
self.assertEqual(
learner_turn.provider_events,
[
{
"type": "sigh",
"confidence": 0.82,
"event_type": "sigh",
"category": "paralinguistic",
},
{
"kind": "voice_activity",
"start_ms": 10,
"event_type": "voice_activity",
"category": "speech_activity",
},
],
)
self.assertIsNone(client_turn.audio_ref)
self.assertEqual(client_turn.provider_events, [])
self.assertEqual(client_turn.llm_provider, "claude_cli")
self.assertTrue(
any(message.get("type") == "tts_end" for message in websocket.messages)
)
async def test_review_exposes_voice_nonverbal_events_on_learner_turn(self) -> None:
principal = _principal()
sess = _session(principal)
created_at = sess.created_at
sess.turns.extend(
[
TurnRecord(
turn_seq=1,
speaker="counselor",
stage=sess.state.stage.value,
text="learner voice turn",
text_masked="learner voice turn",
created_at=created_at + 1,
audio_ref="voice:webm:sha256:test",
silence_ms=1234,
speech_rate=420.0,
barge_in=True,
provider_events=[
{
"event_type": "sigh",
"category": "paralinguistic",
"confidence": 0.82,
"provider": "stt-provider",
"type": "raw_sigh",
},
{
"event_type": "speech_start",
"category": "speech_activity",
"start_ms": 100,
},
{
"event_type": "silence",
"category": "timing",
"duration_ms": 1234,
},
{
"event_type": "background_noise",
"category": "audio_quality",
"score": 77,
"label": "busy cafe",
},
],
),
TurnRecord(
turn_seq=2,
speaker="client",
stage=sess.state.stage.value,
text="client reply",
text_masked="client reply",
created_at=created_at + 2,
audio_ref="voice:webm:sha256:client",
silence_ms=2500,
speech_rate=180.0,
barge_in=True,
provider_events=[
{"event_type": "cry", "category": "paralinguistic"}
],
),
]
)
response = await sessions.get_session_review(sess.session_id, principal)
self.assertEqual(len(response.turns), 2)
learner_turn, client_turn = response.turns
self.assertEqual(
[event.kind for event in learner_turn.nonverbal],
["silence", "pace", "barge_in", "audio", "paralinguistic", "audio_quality"],
)
self.assertEqual(learner_turn.nonverbal[0].label, "침묵")
self.assertEqual(learner_turn.nonverbal[0].detail, "1.2초")
self.assertEqual(
[event.kind for event in learner_turn.nonverbal].count("silence"), 1
)
self.assertEqual(learner_turn.nonverbal[1].detail, "분당 420자")
self.assertEqual(learner_turn.nonverbal[4].label, "음성 단서")
self.assertEqual(learner_turn.nonverbal[4].detail, "한숨 감지 · 신뢰도 82%")
exposed_details = " ".join(event.detail for event in learner_turn.nonverbal)
self.assertNotIn("stt-provider", exposed_details)
self.assertNotIn("raw_sigh", exposed_details)
self.assertEqual(learner_turn.nonverbal[5].label, "오디오 품질")
self.assertEqual(learner_turn.nonverbal[5].detail, "배경 소음 · 신뢰도 77%")
self.assertEqual(client_turn.nonverbal, [])
async def test_review_includes_case_formulation_worksheet_draft(self) -> None:
principal = _principal()
sess = _session(principal)
created_at = sess.created_at
sess.turns.extend(
[
TurnRecord(
turn_seq=1,
speaker="counselor",
stage=sess.state.stage.value,
text="오늘은 어떤 목표로 이야기해보고 싶으세요?",
text_masked="오늘은 어떤 목표로 이야기해보고 싶으세요?",
created_at=created_at + 1,
),
TurnRecord(
turn_seq=2,
speaker="client",
stage=sess.state.stage.value,
text="요즘 너무 불안하고 친구 관계 스트레스 때문에 잠을 잘 못 자요.",
text_masked="요즘 너무 불안하고 친구 관계 스트레스 때문에 잠을 잘 못 자요.",
created_at=created_at + 2,
),
]
)
response = await sessions.get_session_review(sess.session_id, principal)
worksheet = response.caseWorksheet
self.assertEqual(worksheet.status, "draft_from_transcript")
self.assertGreaterEqual(len(worksheet.sections), 5)
exploration = worksheet.sections[0]
self.assertEqual(exploration.key, "exploration_11")
complaint = next(
item for item in exploration.items if item.key == "presenting_complaint"
)
self.assertEqual(complaint.confidence, "medium")
self.assertEqual(complaint.evidence[0].turnId, "t2")
self.assertIn("불안", complaint.value or "")
risk = next(item for item in exploration.items if item.key == "risk")
self.assertEqual(risk.confidence, "none")
self.assertEqual(risk.evidence, [])
self.assertIn("명시 근거", risk.emptyReason or "")
async def test_review_prefers_saved_case_formulation_worksheet(self) -> None:
principal = _principal()
sess = _session(principal)
sess.turns.append(
TurnRecord(
turn_seq=1,
speaker="client",
stage=sess.state.stage.value,
text="자동 초안 대신 저장본을 확인합니다.",
text_masked="자동 초안 대신 저장본을 확인합니다.",
created_at=sess.created_at + 1,
)
)
saved_payload = {
"status": "saved_by_learner",
"generatedBy": "learner-edited worksheet",
"sections": [
{
"key": "exploration_11",
"title": "탐색 11항목",
"items": [
{
"key": "presenting_complaint",
"label": "주호소",
"value": "학습자가 저장한 주호소",
"evidence": [],
"confidence": "medium",
"emptyReason": None,
}
],
}
],
"limitations": ["학습자 저장본"],
"savedAt": "2026-06-27T10:00:00+00:00",
}
with patch.object(
sessions.session_persistence,
"load_case_worksheet",
AsyncMock(return_value=(saved_payload, True)),
):
response = await sessions.get_session_review(sess.session_id, principal)
worksheet = response.caseWorksheet
self.assertEqual(worksheet.status, "saved_by_learner")
self.assertEqual(worksheet.generatedBy, "learner-edited worksheet")
self.assertEqual(worksheet.sections[0].items[0].value, "학습자가 저장한 주호소")
self.assertEqual(worksheet.limitations, ["학습자 저장본"])
self.assertEqual(worksheet.savedAt, "2026-06-27T10:00:00+00:00")
async def test_teacher_review_includes_manual_worksheet_decision(self) -> None:
learner = _principal()
teacher_principal = Principal(
user_id="00000000-0000-0000-0000-000000000902",
role=Role.TEACHER,
)
sess = _session(learner)
sess.ended = True
sess.ended_at = sess.created_at + 600
sess.turns.append(
TurnRecord(
turn_seq=1,
speaker="client",
stage=sess.state.stage.value,
text="저장본 검수를 확인합니다.",
text_masked="저장본 검수를 확인합니다.",
created_at=sess.created_at + 1,
)
)
saved_payload = {
"status": "saved_by_learner",
"generatedBy": "learner-edited worksheet",
"sections": [],
"limitations": ["학습자 저장본"],
"savedAt": "2026-06-27T10:00:00+00:00",
}
review_status = {
"session_id": sess.session_id,
"reviewer_id": teacher_principal.user_id,
"status": "viewed",
"note": "회기 전체 검토 메모",
"worksheet_status": "changes_requested",
"worksheet_note": "주호소 근거를 더 명확히 쓰도록 지도",
"worksheet_reviewed_at": "2026-06-27T10:05:00Z",
"reviewed_at": "",
"updated_at": "2026-06-27T10:05:00Z",
}
with (
patch.object(
sessions.session_persistence,
"load_session",
AsyncMock(return_value=sess),
),
patch.object(
sessions.session_persistence,
"load_case_worksheet",
AsyncMock(return_value=(saved_payload, True)),
),
patch.object(
sessions.session_persistence,
"load_session_evaluation",
AsyncMock(return_value=(None, False)),
),
patch.object(
sessions.session_persistence,
"load_session_review_status",
AsyncMock(return_value=(review_status, True)),
),
):
response = await sessions.get_session_review(
sess.session_id, teacher_principal
)
self.assertIsNotNone(response.teacherReview)
assert response.teacherReview is not None
self.assertEqual(response.caseWorksheet.status, "saved_by_learner")
self.assertEqual(response.teacherReview.worksheetStatus, "changes_requested")
self.assertIn("주호소", response.teacherReview.worksheetNote)
self.assertEqual(
response.teacherReview.worksheetReviewedAt, "2026-06-27T10:05:00Z"
)
async def test_review_remasks_legacy_session_evaluation_payload(self) -> None:
class FakeKoRecognizer:
def analyze(self, text: str):
spans = []
for entity_type, value in (
("NAME", "보라별"),
("ORG", "미래학교상담연구랩"),
):
start = text.find(value)
if start >= 0:
spans.append(
guardrail.PiiEntitySpan(
entity_type, start, start + len(value)
)
)
return spans
guardrail.set_ko_pii_recognizer(FakeKoRecognizer())
self.addCleanup(guardrail.set_ko_pii_recognizer, None)
principal = _principal()
sess = _session(principal)
sess.ended = True
sess.ended_at = sess.created_at + 180
sess.turns.extend(
[
TurnRecord(
turn_seq=1,
speaker="counselor",
stage=sess.state.stage.value,
text="감정을 먼저 확인해보겠습니다.",
text_masked="감정을 먼저 확인해보겠습니다.",
created_at=sess.created_at + 1,
),
TurnRecord(
turn_seq=2,
speaker="client",
stage=sess.state.stage.value,
text="연락처 이야기는 부담스러워요.",
text_masked="연락처 이야기는 부담스러워요.",
created_at=sess.created_at + 2,
),
]
)
raw_record = {
"status": "ready",
"source": "engine",
"scope": "session_end",
"stage": "정리",
"payload": {
"loop": "deep",
"session_id": sess.session_id,
"stage": "정리",
"scope": "session_end",
"turns_evaluated": 2,
"distribution": {},
"strengths": ["보라별의 감정을 반영했다."],
"improvements": ["미래학교상담연구랩과 010-1234-5678 재확인을 줄인다."],
"intent_deviations": [
{
"dimension": "pacing",
"expected": "보라별의 감정 확인",
"actual": "010-1234-5678 연락처 재질문",
"severity": "moderate",
}
],
"supervisor_rationale": "보라별의 호소를 요약했다.",
"supervisor_critique": "미래학교상담연구랩 언급이 반복됐다.",
"alternative_utterances": ["보라별님, 지금 감정부터 천천히 볼까요?"],
},
"error": None,
}
with (
patch.object(
sessions.session_persistence,
"load_session",
AsyncMock(return_value=sess),
),
patch.object(
sessions.session_persistence,
"load_case_worksheet",
AsyncMock(return_value=(None, False)),
),
patch.object(
sessions.session_persistence,
"load_session_evaluation",
AsyncMock(return_value=(raw_record, True)),
),
):
response = await sessions.get_session_review(sess.session_id, principal)
blob = response.model_dump_json()
for raw in ("보라별", "미래학교상담연구랩", "010-1234-5678"):
self.assertNotIn(raw, blob)
for masked in ("[NAME]", "[ORG]", "[PHONE]"):
self.assertIn(masked, blob)
self.assertTrue(response.reviewReady)
self.assertEqual(response.supervisorState, "평가 완료")
async def test_review_surfaces_session_evaluation_error_record(self) -> None:
principal = _principal()
sess = _session(principal)
sess.ended = True
sess.ended_at = sess.created_at + 180
sess.turns.extend(
[
TurnRecord(
turn_seq=1,
speaker="learner",
stage=sess.state.stage.value,
text="많이 지쳐 보였어요. 지금 제일 버거운 마음이 뭔가요?",
text_masked="많이 지쳐 보였어요. 지금 제일 버거운 마음이 뭔가요?",
created_at=sess.created_at + 1,
),
TurnRecord(
turn_seq=2,
speaker="client",
stage=sess.state.stage.value,
text="그냥 아무것도 하고 싶지 않아요.",
text_masked="그냥 아무것도 하고 싶지 않아요.",
created_at=sess.created_at + 2,
),
]
)
error_record = session_persistence.SessionEvaluationWrite.from_error(
session_id=sess.session_id,
learner_id=principal.user_id,
scope="session_end",
stage="정리",
error="session evaluation timeout after 45s",
).cache_record()
with (
patch.object(
sessions.session_persistence,
"load_session",
AsyncMock(return_value=sess),
),
patch.object(
sessions.session_persistence,
"load_case_worksheet",
AsyncMock(return_value=(None, False)),
),
patch.object(
sessions.session_persistence,
"load_session_evaluation",
AsyncMock(return_value=(error_record, True)),
),
):
response = await sessions.get_session_review(sess.session_id, principal)
self.assertFalse(response.reviewReady)
self.assertTrue(response.degraded)
self.assertEqual(response.supervisorState, "평가 실패")
self.assertIn(
"deep-loop 평가 AI 산출물을 표시하지 못했습니다", response.summary
)
self.assertIn("AI 평가 재시도가 필요합니다", response.summary)
self.assertNotIn("session evaluation timeout after 45s", response.summary)
self.assertEqual(response.rubric, [])
self.assertEqual(response.goodMoments, [])
self.assertEqual(response.growthPoints, [])
def test_review_marks_stale_missing_session_evaluation_as_failed(self) -> None:
principal = _principal()
sess = _session(principal)
sess.ended = True
sess.ended_at = sess.created_at + 180
sess.turns.extend(
[
TurnRecord(
turn_seq=1,
speaker="counselor",
stage=sess.state.stage.value,
text="많이 지쳐 보였어요. 지금 제일 버거운 마음이 뭔가요?",
text_masked="많이 지쳐 보였어요. 지금 제일 버거운 마음이 뭔가요?",
created_at=sess.created_at + 1,
),
TurnRecord(
turn_seq=2,
speaker="client",
stage=sess.state.stage.value,
text="그냥 아무것도 하고 싶지 않아요.",
text_masked="그냥 아무것도 하고 싶지 않아요.",
created_at=sess.created_at + 2,
),
]
)
previous_timeout = settings.session_evaluation_timeout
settings.session_evaluation_timeout = 10.0
try:
response = sessions.build_session_review(
sessions.SessionReviewReadInput(
session=sess,
evaluation_record=None,
evaluation_durable=True,
now_ts=(sess.ended_at or sess.created_at) + 41.0,
)
)
finally:
settings.session_evaluation_timeout = previous_timeout
self.assertFalse(response.reviewReady)
self.assertTrue(response.degraded)
self.assertEqual(response.supervisorState, "평가 실패")
self.assertIn("AI 평가 재시도가 필요합니다", response.summary)
self.assertEqual(response.rubric, [])
self.assertEqual(response.goodMoments, [])
self.assertEqual(response.growthPoints, [])
async def test_learner_can_save_case_formulation_worksheet(self) -> None:
principal = _principal()
sess = _session(principal)
request = sessions.ReviewCaseWorksheetSaveRequest(
sections=[
sessions.ReviewWorksheetSection(
key="exploration_11",
title="탐색 11항목",
items=[
sessions.ReviewWorksheetItem(
key="presenting_complaint",
label="주호소",
value="수정한 주호소",
confidence="medium",
)
],
)
],
limitations=["임상 루브릭 전"],
)
with (
patch.object(
sessions,
"_load_session_or_404",
AsyncMock(return_value=sess),
) as load_session,
patch.object(
sessions.session_persistence,
"save_case_worksheet",
AsyncMock(return_value=True),
) as save_worksheet,
):
response = await sessions.save_session_review_worksheet(
sess.session_id,
request,
principal,
)
load_session.assert_awaited_once_with(
sess.session_id,
principal,
allow_ended=True,
include_turn_evaluation=False,
)
save_worksheet.assert_awaited_once()
save_kwargs = save_worksheet.await_args.kwargs
self.assertEqual(save_kwargs["session_id"], sess.session_id)
self.assertEqual(save_kwargs["learner_id"], principal.user_id)
self.assertEqual(save_kwargs["payload"]["status"], "saved_by_learner")
self.assertEqual(
save_kwargs["payload"]["sections"][0]["items"][0]["value"], "수정한 주호소"
)
self.assertEqual(response.status, "saved_by_learner")
self.assertEqual(response.sections[0].items[0].value, "수정한 주호소")
if __name__ == "__main__":
unittest.main()