상담자 발화 판정(A)·감정(B)·표현(C) 20문항 질문 세트, 감쇠 없는 이번 턴 반응과 비대칭 기분 전이, 개방도 게이트로 생성 지시를 만들고 ccd.coping_strategy 전달 누락을 고친다. trace v2와 고정 문구 속마음 요약(migration 24, AI 경로 차단 RLS)을 같은 트랜잭션에 저장하고 피드백 정책이 켜진 경우에만 done·TurnResponse·음성 reply·리뷰로 노출한다. 회기 화면 속마음 보기 토글, 리뷰 접힘 블록, 관리자 감정 관측 v2 표시를 추가한다.
806 lines
31 KiB
Python
806 lines
31 KiB
Python
"""턴 오케스트레이터 — 상담 1턴 파이프라인 1~8단계 조립.
|
|
|
|
MASTERPLAN §2.2 / MEMORY_DESIGN §2-B 턴 사이클:
|
|
1. [입력 가드레일] PII 마스킹 + 위기분류(실제위기 vs 연기) (guardrail)
|
|
2. [상태머신] effective_openness 결정론 계산 + 단계전이 (state_machine)
|
|
3. [페르소나 컨텍스트] L0~L6 messages 조립 (persona)
|
|
4. [내담자 AI] engine_client.stream/generate (CCD 비노출) (engine_client)
|
|
5. [출력 가드레일] 자살수단 차단, ideation 상한 (guardrail)
|
|
6. [평가 훅] 주입형 — 평가 함수는 *인자로 받는다*(Features 소유) (hook)
|
|
7. [상태 갱신] working state 반영(체크포인트는 호출부가 DB/store UPSERT)
|
|
8. [로깅 훅] 주입형 — turns insert/임베딩은 호출부가 주입
|
|
|
|
설계 원칙:
|
|
- 평가(evaluator) 함수와 로깅 함수는 *주입*받는다(이 모듈은 evaluator.py 를 import 하지 않음).
|
|
- DB 는 인터페이스로 추상화하되 asyncpg conn 도 받을 수 있게 했다(현재는 hook 으로만 사용).
|
|
- generate(동기, 폴백/테스트) + stream(SSE 토큰) 두 경로 모두 제공.
|
|
- 엔진 장애는 EngineError 로 전파 → 라우트가 503/SSE error 프레임으로 변환.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import time
|
|
from dataclasses import dataclass, field, replace
|
|
from typing import Any, AsyncIterator, Awaitable, Callable, Optional
|
|
|
|
from ..contracts.client_affect import ClientAffectTraceV1, ClientAffectTraceV2, ClientInnerReactionV1
|
|
from ..config import settings
|
|
from ..engine_client import (
|
|
EngineClient,
|
|
EngineError,
|
|
EngineMessage,
|
|
GenerateRequest,
|
|
GenerateResponse,
|
|
StreamRequest,
|
|
)
|
|
from ..contracts.engine_gateway import (
|
|
ENGINE_GATEWAY_DEFAULT_MODEL_SENTINEL,
|
|
ENGINE_GATEWAY_SSE_DONE,
|
|
ENGINE_GATEWAY_SSE_ERROR,
|
|
EngineGatewaySseDecodeError,
|
|
StreamDoneEvent,
|
|
StreamErrorEvent,
|
|
StreamTokenEvent,
|
|
)
|
|
from . import client_affect, guardrail, persona, rupture_scenario_director, state_machine
|
|
from .jev_client import AppraisalResult, JevError, jev_client
|
|
from .llm_audit import LlmAuditHook, generate_with_audit, record_llm_audit
|
|
from .persona import PersonaCard, PersonaStateContext, TurnMemory
|
|
from .state_machine import SessionState
|
|
|
|
# 평가 훅 타입: U_t(수련생 마스킹 발화) + 내담자응답 + 상태 → 평가 결과(dict)
|
|
# Features evaluator 가 이 시그니처에 맞춰 함수를 주입한다(여기선 호출만).
|
|
EvalHook = Callable[["TurnContext", str], Awaitable[Optional[dict]]]
|
|
|
|
|
|
def turn_evaluation_error_payload(ctx: "TurnContext", error: BaseException | str) -> dict[str, Any]:
|
|
"""Represent a non-fatal fast-loop evaluator failure without hiding it."""
|
|
st = ctx.state_after or ctx.state_before
|
|
if isinstance(error, BaseException):
|
|
# Exception messages can contain raw learner/provider text. Keep review-facing
|
|
# evidence to the exception type; detailed trace stays in server logs.
|
|
detail = type(error).__name__
|
|
else:
|
|
detail = str(error).strip() or "unknown turn evaluation error"
|
|
masked = guardrail.mask_pii(detail).text_masked.strip() or "unknown turn evaluation error"
|
|
return {
|
|
"loop": "fast",
|
|
"turn_seq": st.turn_seq,
|
|
"stage": st.stage.value,
|
|
"appropriateness": "neutral",
|
|
"error": masked,
|
|
}
|
|
|
|
|
|
@dataclass(slots=True)
|
|
class TurnContext:
|
|
"""한 턴 파이프라인을 관통하는 컨텍스트(가드레일/상태/페르소나 산출 집약)."""
|
|
|
|
session_id: str
|
|
case_id: Optional[str]
|
|
persona: PersonaCard
|
|
state_before: SessionState
|
|
learner_text_raw: str
|
|
learner_text_masked: str = ""
|
|
counselor_identity: Optional[str] = None
|
|
client_identity: Optional[str] = None
|
|
crisis: Optional[guardrail.CrisisResult] = None
|
|
state_after: Optional[SessionState] = None
|
|
messages: list[EngineMessage] = field(default_factory=list)
|
|
# 회상/메모리 주입(memory.RecallContext 에서 옴)
|
|
memory: TurnMemory = field(default_factory=TurnMemory)
|
|
# 회기 이론모드(학습자 선택: humanistic|cbt|integrative). 평가 이론부합·생성 프레이밍에 사용.
|
|
theory_mode: Optional[str] = None
|
|
# Scenario Director 내부 선택. ID/유형/provenance는 엔진 request metadata에만 존재한다.
|
|
scenario_directive: Optional[rupture_scenario_director.ScenarioDirective] = None
|
|
# 외부 감정 평가의 안전한 provenance. 원문·점수·확률은 넣지 않는다.
|
|
client_affect_metadata: Optional[dict[str, Any]] = None
|
|
# 관리자 관측 전용 Jev 전이 trace(v1|v2). 공개 결과나 provider event에는 넣지 않는다.
|
|
client_affect_trace: ClientAffectTraceV1 | ClientAffectTraceV2 | None = None
|
|
# v2 생성 지시(§7). L3 '정서 연기 지시' 줄을 대체한다. legacy·v1 경로에선 None.
|
|
client_affect_directive: Optional[str] = None
|
|
# 학습자·교수자용 속마음 요약(§8.2). 이 패킷에서는 저장·노출하지 않고 조립만 한다.
|
|
client_inner_reaction: ClientInnerReactionV1 | None = None
|
|
|
|
def to_state_context(self) -> PersonaStateContext:
|
|
st = self.state_after or self.state_before
|
|
return PersonaStateContext(
|
|
stage=st.stage.value,
|
|
effective_openness=st.effective_openness,
|
|
resistance=st.resistance,
|
|
rapport_credit=st.rapport_credit,
|
|
ideation_stage=st.ideation_stage,
|
|
affect_state=st.affect_state,
|
|
affect_directive=self.client_affect_directive,
|
|
)
|
|
|
|
|
|
@dataclass(slots=True)
|
|
class TurnResult:
|
|
"""동기(generate) 턴 결과."""
|
|
|
|
turn_seq: int
|
|
stage: str
|
|
effective_openness: float
|
|
client_reply: Optional[str]
|
|
safety_flagged: bool
|
|
state_after: SessionState
|
|
evaluation: Optional[dict] = None
|
|
crisis_kind: str = "none"
|
|
crisis_resource: Optional[dict[str, str]] = None
|
|
conversation_stopped: bool = False
|
|
llm_provider: Optional[str] = None
|
|
model: Optional[str] = None
|
|
tokens_in: int = 0
|
|
tokens_out: int = 0
|
|
cost_usd: float = 0.0
|
|
output_error: Optional[str] = None
|
|
|
|
|
|
# ════════════════════════════════════════════════════════════════════════════
|
|
# 1~3단계 — 입력 가드레일 + 상태머신 + 페르소나 컨텍스트 (엔진 호출 전 결정론)
|
|
# ════════════════════════════════════════════════════════════════════════════
|
|
def prepare_turn(
|
|
*,
|
|
session_id: str,
|
|
case_id: Optional[str],
|
|
card: PersonaCard,
|
|
state: SessionState,
|
|
learner_text: str,
|
|
learner_identity: Optional[str] = None,
|
|
memory: Optional[TurnMemory] = None,
|
|
theory_mode: Optional[str] = None,
|
|
eval_rapport_signal: Optional[float] = None,
|
|
scenario_context: Optional[
|
|
rupture_scenario_director.StoredScenarioContext
|
|
] = None,
|
|
) -> TurnContext:
|
|
"""엔진 호출 전 결정론 전처리(1~3단계). 순수 — IO/LLM 없음.
|
|
|
|
eval_rapport_signal 이 주어지면(평가 AI fast-loop 신호) 그걸 쓰고, 없으면
|
|
state_machine 의 경량 휴리스틱으로 라포 신호를 추정한다.
|
|
"""
|
|
turn_memory = memory or TurnMemory()
|
|
client_identity = card.display_name
|
|
ctx = TurnContext(
|
|
session_id=session_id,
|
|
case_id=case_id,
|
|
persona=card,
|
|
state_before=state,
|
|
learner_text_raw=learner_text,
|
|
counselor_identity=learner_identity,
|
|
client_identity=client_identity,
|
|
memory=TurnMemory(
|
|
recall_summary=_mask_optional_text(
|
|
turn_memory.recall_summary,
|
|
counselor_identity=learner_identity,
|
|
client_identity=client_identity,
|
|
),
|
|
pinned_facts=_mask_text_list(
|
|
turn_memory.pinned_facts,
|
|
counselor_identity=learner_identity,
|
|
client_identity=client_identity,
|
|
),
|
|
recent_turns=_mask_recent_turns(
|
|
turn_memory.recent_turns,
|
|
counselor_identity=learner_identity,
|
|
client_identity=client_identity,
|
|
),
|
|
kb_behavior_cues=list(turn_memory.kb_behavior_cues or []),
|
|
),
|
|
theory_mode=theory_mode,
|
|
)
|
|
|
|
# 1) 입력 가드레일 — PII 마스킹 + 위기분류
|
|
mask = guardrail.mask_role_identities(
|
|
learner_text,
|
|
counselor_identity=learner_identity,
|
|
client_identity=client_identity,
|
|
)
|
|
ctx.learner_text_masked = mask.text_masked
|
|
ctx.crisis = guardrail.classify_crisis(learner_text, speaker_is_persona_context=True)
|
|
|
|
# 2) 상태머신 — 라포 신호 → 결정론 전이
|
|
signal = (
|
|
eval_rapport_signal
|
|
if eval_rapport_signal is not None
|
|
else state_machine.estimate_rapport_signal(ctx.learner_text_masked)
|
|
)
|
|
# 위기분류가 관측한 risk_level(>0)을 상태머신에 ideation_observed 로 전달 →
|
|
# ideation_stage 보수적 상향(절대 하향 안 함, 안전 R5). C2 위기 관측 반영.
|
|
crisis_ideation = (
|
|
ctx.crisis.risk_level
|
|
if ctx.crisis is not None and ctx.crisis.risk_level > 0
|
|
else None
|
|
)
|
|
ctx.state_after = state_machine.evolve(
|
|
state,
|
|
rapport_signal=signal,
|
|
unlock_rate=card.unlock_rate(),
|
|
decay_floor=card.decay_floor(),
|
|
ideation_observed=crisis_ideation,
|
|
)
|
|
|
|
ctx.scenario_directive = rupture_scenario_director.select_scenario_directive(
|
|
case_id=ctx.case_id,
|
|
session_id=ctx.session_id,
|
|
turn_seq=ctx.state_after.turn_seq,
|
|
safety_escalated=bool(ctx.crisis is not None and ctx.crisis.escalate),
|
|
scenario_context=scenario_context,
|
|
)
|
|
hidden_behavior_cue = rupture_scenario_director.render_hidden_behavior_prompt(
|
|
ctx.scenario_directive
|
|
)
|
|
|
|
# 3) 페르소나 컨텍스트 — L0~L6 messages 조립 (CCD 는 행동으로만, L0 가 강제)
|
|
ctx.messages = persona.build_turn_messages(
|
|
card,
|
|
ctx.to_state_context(),
|
|
ctx.learner_text_masked,
|
|
memory=ctx.memory,
|
|
theory_mode=ctx.theory_mode,
|
|
hidden_behavior_cue=hidden_behavior_cue,
|
|
)
|
|
return ctx
|
|
|
|
|
|
def _mask_optional_text(
|
|
text: Optional[str],
|
|
*,
|
|
counselor_identity: Optional[str] = None,
|
|
client_identity: Optional[str] = None,
|
|
) -> Optional[str]:
|
|
if text is None:
|
|
return None
|
|
return guardrail.mask_role_identities(
|
|
text,
|
|
counselor_identity=counselor_identity,
|
|
client_identity=client_identity,
|
|
).text_masked
|
|
|
|
|
|
def _mask_text_list(
|
|
values: Optional[list[str]],
|
|
*,
|
|
counselor_identity: Optional[str] = None,
|
|
client_identity: Optional[str] = None,
|
|
) -> list[str]:
|
|
return [
|
|
guardrail.mask_role_identities(
|
|
value,
|
|
counselor_identity=counselor_identity,
|
|
client_identity=client_identity,
|
|
).text_masked
|
|
for value in (values or [])
|
|
]
|
|
|
|
|
|
def _mask_recent_turns(
|
|
turns: Optional[list[dict[str, str]]],
|
|
*,
|
|
counselor_identity: Optional[str] = None,
|
|
client_identity: Optional[str] = None,
|
|
) -> list[dict[str, str]]:
|
|
masked: list[dict[str, str]] = []
|
|
for turn in turns or []:
|
|
item = dict(turn)
|
|
item["text"] = guardrail.mask_role_identities(
|
|
str(item.get("text", "")),
|
|
counselor_identity=counselor_identity,
|
|
client_identity=client_identity,
|
|
).text_masked
|
|
masked.append(item)
|
|
return masked
|
|
|
|
|
|
def _latest_client_reply(turns: list[dict[str, str]]) -> Optional[str]:
|
|
for turn in reversed(turns):
|
|
if turn.get("speaker") == "client":
|
|
text = str(turn.get("text", "")).strip()
|
|
if text:
|
|
return text
|
|
return None
|
|
|
|
|
|
def _client_request_metadata(ctx: TurnContext) -> dict[str, Any]:
|
|
assert ctx.state_after is not None
|
|
metadata: dict[str, Any] = {"stage": ctx.state_after.stage.value}
|
|
if ctx.scenario_directive is not None:
|
|
metadata["scenario_director"] = ctx.scenario_directive.request_metadata()
|
|
if ctx.client_affect_metadata is not None:
|
|
metadata["client_affect"] = dict(ctx.client_affect_metadata)
|
|
return metadata
|
|
|
|
|
|
def _rebuild_persona_messages(ctx: TurnContext) -> None:
|
|
"""Jev 전이 뒤 같은 시나리오·회상 계약으로 L3를 다시 조립한다."""
|
|
hidden_behavior_cue = rupture_scenario_director.render_hidden_behavior_prompt(
|
|
ctx.scenario_directive
|
|
)
|
|
ctx.messages = persona.build_turn_messages(
|
|
ctx.persona,
|
|
ctx.to_state_context(),
|
|
ctx.learner_text_masked,
|
|
memory=ctx.memory,
|
|
theory_mode=ctx.theory_mode,
|
|
hidden_behavior_cue=hidden_behavior_cue,
|
|
)
|
|
|
|
|
|
def _client_profile_inputs(card: PersonaCard) -> dict[str, Any]:
|
|
"""v2 Jev state의 client_profile 원자료(마스킹 전)를 카드에서 고른다.
|
|
|
|
v1은 ccd["coping"]을 읽었지만 카드는 coping_strategy를 쓰므로 대처 방식이
|
|
한 번도 전달되지 않았다(§1 근거 5). client_profile.coping_strategy로 고친다.
|
|
"""
|
|
ccd = card.ccd or {}
|
|
triggers = card.triggers or {}
|
|
return {
|
|
"presenting": card.presenting,
|
|
"history": card.history,
|
|
"core_belief": ccd.get("core_belief"),
|
|
"automatic_thought": ccd.get("automatic_thought"),
|
|
"coping_strategy": ccd.get("coping_strategy"),
|
|
"big5": card.big5,
|
|
"sore_spots": list(triggers.get("sore_spots") or []),
|
|
"forbidden": list(triggers.get("forbidden") or []),
|
|
"speech_style": card.speech_style,
|
|
}
|
|
|
|
|
|
async def _record_client_affect_audit(
|
|
ctx: TurnContext,
|
|
appraisal: AppraisalResult,
|
|
audit_hook: Optional[LlmAuditHook],
|
|
) -> None:
|
|
"""생성 모델과 구분한 Jev 호출 provenance를 기존 감사 계약에 남긴다."""
|
|
await record_llm_audit(
|
|
audit_hook,
|
|
session_id=ctx.session_id,
|
|
provider=appraisal.provider,
|
|
model=appraisal.model,
|
|
tokens_in=appraisal.input_tokens,
|
|
tokens_out=appraisal.output_tokens,
|
|
cost_usd=appraisal.cost_usd,
|
|
inference_geo=None,
|
|
latency_ms=appraisal.latency_ms,
|
|
)
|
|
|
|
|
|
async def _apply_client_affect(
|
|
ctx: TurnContext,
|
|
*,
|
|
audit_hook: Optional[LlmAuditHook],
|
|
) -> None:
|
|
"""활성 Jev 평가(v2)를 1회 적용하고 생성 요청 직전 L3를 갱신한다.
|
|
|
|
① 이번 턴 반응 ② 기분 비대칭 전이 ③ 표현 계획(개방도 게이트)을 조합해
|
|
trace v2·속마음 요약·생성 지시 v2를 만든다(§6~§8.1). 속마음은 이 패킷에서
|
|
아직 저장·노출하지 않고 ctx에만 보존한다.
|
|
"""
|
|
if settings.client_affect_provider != "jev":
|
|
return
|
|
assert ctx.state_after is not None
|
|
state = client_affect.build_appraisal_state(
|
|
client_profile=_client_profile_inputs(ctx.persona),
|
|
affect_state=ctx.state_after.affect_state,
|
|
affect_baseline=ctx.persona.affect_baseline,
|
|
stage=ctx.state_after.stage.value,
|
|
resistance=ctx.state_after.resistance,
|
|
effective_openness=ctx.state_after.effective_openness,
|
|
counselor_utterance=ctx.learner_text_masked,
|
|
recall_summary=ctx.memory.recall_summary,
|
|
pinned_facts=ctx.memory.pinned_facts,
|
|
recent_turns=ctx.memory.recent_turns,
|
|
counselor_identity=ctx.counselor_identity,
|
|
client_identity=ctx.client_identity,
|
|
)
|
|
appraisal = await jev_client.appraise(state)
|
|
state_before_transition = ctx.state_after
|
|
transition = client_affect.transition_mood(
|
|
state_before_transition.affect_state,
|
|
ctx.persona.affect_baseline,
|
|
appraisal,
|
|
min_confidence=settings.jev_min_confidence,
|
|
)
|
|
ctx.state_after = replace(ctx.state_after, affect_state=transition.affect_state)
|
|
expression = client_affect.build_expression_plan(
|
|
appraisal, effective_openness=ctx.state_after.effective_openness
|
|
)
|
|
ctx.client_affect_trace = client_affect.build_client_affect_trace_v2(
|
|
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,
|
|
expression=expression,
|
|
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_inner_reaction = client_affect.build_inner_reaction(
|
|
appraisal, expression, turn_seq=ctx.state_after.turn_seq
|
|
)
|
|
ctx.client_affect_directive = client_affect.render_affect_directive_v2(
|
|
appraisal,
|
|
expression,
|
|
affect_state_after=ctx.state_after.affect_state,
|
|
affect_baseline=ctx.persona.affect_baseline,
|
|
)
|
|
ctx.client_affect_metadata = {
|
|
"provider": appraisal.provider,
|
|
"model": appraisal.model,
|
|
"latency_ms": appraisal.latency_ms,
|
|
"tokens_in": appraisal.input_tokens,
|
|
"tokens_out": appraisal.output_tokens,
|
|
"accepted_dimensions": list(transition.accepted_dimensions),
|
|
"held_dimensions": list(transition.held_dimensions),
|
|
"tentative_dimensions": list(transition.tentative_dimensions),
|
|
}
|
|
_rebuild_persona_messages(ctx)
|
|
await _record_client_affect_audit(ctx, appraisal, audit_hook)
|
|
|
|
|
|
def _safe_engine_error_detail(error: BaseException | str, *, fallback: str) -> str:
|
|
detail = str(error).strip() or fallback
|
|
if rupture_scenario_director.contains_internal_scenario_leakage(detail):
|
|
return fallback
|
|
return detail
|
|
|
|
|
|
# ════════════════════════════════════════════════════════════════════════════
|
|
# 4~8단계 — 동기 생성 경로 (폴백/테스트)
|
|
# ════════════════════════════════════════════════════════════════════════════
|
|
async def run_turn_generate(
|
|
ctx: TurnContext,
|
|
engine: EngineClient,
|
|
*,
|
|
eval_hook: Optional[EvalHook] = None,
|
|
audit_hook: Optional[LlmAuditHook] = None,
|
|
) -> TurnResult:
|
|
"""동기 턴 실행(4~8). 내담자 응답을 한 번에 받아 가드레일·평가 순차 적용.
|
|
|
|
eval_hook 은 Features 가 주입(없으면 생략). 엔진 장애는 EngineError 전파.
|
|
"""
|
|
assert ctx.state_after is not None
|
|
|
|
if ctx.crisis is not None and ctx.crisis.escalate:
|
|
return _crisis_gate_result(ctx)
|
|
|
|
try:
|
|
await _apply_client_affect(ctx, audit_hook=audit_hook)
|
|
except JevError as exc:
|
|
raise EngineError(f"client_affect_{exc.code}") from exc
|
|
st = ctx.state_after
|
|
assert st is not None
|
|
|
|
# 4) 내담자 AI 생성
|
|
req = GenerateRequest(
|
|
ai_role="client",
|
|
messages=ctx.messages,
|
|
session_id=ctx.session_id,
|
|
metadata=_client_request_metadata(ctx),
|
|
)
|
|
previous_client_reply = _latest_client_reply(ctx.memory.recent_turns)
|
|
resp: GenerateResponse | None = None
|
|
reply = ""
|
|
safety_flagged = ctx.crisis is not None and ctx.crisis.escalate
|
|
for attempt in range(2):
|
|
try:
|
|
resp = await generate_with_audit(engine, req, audit_hook)
|
|
except EngineError as exc:
|
|
detail = _safe_engine_error_detail(
|
|
exc,
|
|
fallback="engine generation failed",
|
|
)
|
|
if detail != str(exc):
|
|
raise EngineError(detail) from exc
|
|
raise
|
|
|
|
# 5) 출력 가드레일 — 수단 차단 + persona 품질 재생성
|
|
if rupture_scenario_director.contains_internal_scenario_leakage(resp.text):
|
|
if attempt == 0:
|
|
continue
|
|
return TurnResult(
|
|
turn_seq=st.turn_seq,
|
|
stage=st.stage.value,
|
|
effective_openness=st.effective_openness,
|
|
client_reply=None,
|
|
safety_flagged=True,
|
|
state_after=st,
|
|
evaluation=None,
|
|
crisis_kind=ctx.crisis.kind.value if ctx.crisis else "none",
|
|
llm_provider=resp.provider,
|
|
model=resp.model,
|
|
tokens_in=resp.tokens_in,
|
|
tokens_out=resp.tokens_out,
|
|
cost_usd=resp.cost_usd,
|
|
output_error="client_reply_quality_retryable",
|
|
)
|
|
guard = guardrail.sanitize_client_reply(
|
|
resp.text,
|
|
ideation_stage=st.ideation_stage,
|
|
turn_seq=st.turn_seq,
|
|
previous_client_reply=previous_client_reply,
|
|
)
|
|
if guard.needs_regeneration:
|
|
if attempt == 0:
|
|
continue
|
|
return TurnResult(
|
|
turn_seq=st.turn_seq,
|
|
stage=st.stage.value,
|
|
effective_openness=st.effective_openness,
|
|
client_reply=None,
|
|
safety_flagged=True,
|
|
state_after=st,
|
|
evaluation=None,
|
|
crisis_kind=ctx.crisis.kind.value if ctx.crisis else "none",
|
|
llm_provider=resp.provider,
|
|
model=resp.model,
|
|
tokens_in=resp.tokens_in,
|
|
tokens_out=resp.tokens_out,
|
|
cost_usd=resp.cost_usd,
|
|
output_error="client_reply_quality_retryable",
|
|
)
|
|
safety_flagged = safety_flagged or guard.blocked
|
|
reply = guard.text
|
|
break
|
|
|
|
assert resp is not None
|
|
|
|
# 6) 평가 훅(주입형) — 평가 AI 4차원 태깅 (Features 소유)
|
|
evaluation: Optional[dict] = None
|
|
if eval_hook is not None:
|
|
try:
|
|
evaluation = await eval_hook(ctx, reply)
|
|
except Exception as exc:
|
|
# 평가 실패는 상담 루프를 막지 않되, 리뷰/대시보드에서 조용히 사라지지 않게 남긴다.
|
|
evaluation = turn_evaluation_error_payload(ctx, exc)
|
|
|
|
return TurnResult(
|
|
turn_seq=st.turn_seq,
|
|
stage=st.stage.value,
|
|
effective_openness=st.effective_openness,
|
|
client_reply=reply,
|
|
safety_flagged=safety_flagged,
|
|
state_after=st,
|
|
evaluation=evaluation,
|
|
crisis_kind=ctx.crisis.kind.value if ctx.crisis else "none",
|
|
llm_provider=resp.provider,
|
|
model=resp.model,
|
|
tokens_in=resp.tokens_in,
|
|
tokens_out=resp.tokens_out,
|
|
cost_usd=resp.cost_usd,
|
|
)
|
|
|
|
|
|
def _crisis_gate_result(ctx: TurnContext) -> TurnResult:
|
|
"""실제 위기 신호는 LLM 호출 전에 중단하고 109 리소스를 반환한다."""
|
|
assert ctx.state_after is not None
|
|
crisis = ctx.crisis
|
|
return TurnResult(
|
|
turn_seq=ctx.state_after.turn_seq,
|
|
stage=ctx.state_after.stage.value,
|
|
effective_openness=ctx.state_after.effective_openness,
|
|
client_reply=None,
|
|
safety_flagged=True,
|
|
state_after=ctx.state_after,
|
|
evaluation=None,
|
|
crisis_kind=crisis.kind.value if crisis else "learner_real",
|
|
crisis_resource=guardrail.crisis_resource(),
|
|
conversation_stopped=True,
|
|
)
|
|
|
|
|
|
# ════════════════════════════════════════════════════════════════════════════
|
|
# 4~8단계 — SSE 스트림 경로 (기본 UX)
|
|
# ════════════════════════════════════════════════════════════════════════════
|
|
@dataclass(slots=True)
|
|
class StreamEvent:
|
|
"""SSE 재방출용 이벤트. 라우트가 sse_starlette 형식으로 변환."""
|
|
|
|
event: str # 'token' | 'done' | 'safety' | 'error'
|
|
data: dict[str, Any]
|
|
|
|
|
|
async def run_turn_stream(
|
|
ctx: TurnContext,
|
|
engine: EngineClient,
|
|
*,
|
|
audit_hook: Optional[LlmAuditHook] = None,
|
|
) -> AsyncIterator[StreamEvent]:
|
|
"""스트리밍 턴 실행(4~8). 게이트웨이 SSE 를 받아 token/done/safety/error 로 재방출.
|
|
|
|
출력 가드레일은 *누적 텍스트* 기준으로 수단정보를 감지(스트림 중 발견 시 safety 이벤트 +
|
|
재생성 신호). 토큰 단위 완벽 차단은 후속(현재는 누적 스캔).
|
|
"""
|
|
assert ctx.state_after is not None
|
|
st = ctx.state_after
|
|
|
|
accumulated = ""
|
|
flagged = False
|
|
output_error: str | None = None
|
|
stream_meta: dict[str, Any] = {}
|
|
gateway_done = False
|
|
previous_client_reply = _latest_client_reply(ctx.memory.recent_turns)
|
|
if ctx.crisis is not None and ctx.crisis.escalate:
|
|
flagged = True
|
|
resource = guardrail.crisis_resource()
|
|
yield StreamEvent(
|
|
"safety",
|
|
{
|
|
"reason": "learner_real_crisis",
|
|
"level": ctx.crisis.risk_level,
|
|
"crisis_resource": resource,
|
|
"conversation_stopped": True,
|
|
},
|
|
)
|
|
yield StreamEvent(
|
|
"done",
|
|
{
|
|
"session_id": ctx.session_id,
|
|
"stage": st.stage.value,
|
|
"effective_openness": round(st.effective_openness, 4),
|
|
"turn_seq": st.turn_seq,
|
|
"safety_flagged": True,
|
|
"crisis_kind": ctx.crisis.kind.value,
|
|
"crisis_resource": resource,
|
|
"conversation_stopped": True,
|
|
},
|
|
)
|
|
return
|
|
|
|
try:
|
|
await _apply_client_affect(ctx, audit_hook=audit_hook)
|
|
except JevError as exc:
|
|
yield StreamEvent("error", {"detail": f"client_affect_{exc.code}"})
|
|
return
|
|
st = ctx.state_after
|
|
assert st is not None
|
|
req = StreamRequest(
|
|
ai_role="client",
|
|
messages=ctx.messages,
|
|
session_id=ctx.session_id,
|
|
metadata=_client_request_metadata(ctx),
|
|
)
|
|
|
|
try:
|
|
started = time.perf_counter()
|
|
async for packet in engine.stream_packets(req):
|
|
if packet.event == ENGINE_GATEWAY_SSE_ERROR:
|
|
payload = packet.payload
|
|
detail = payload.detail if isinstance(payload, StreamErrorEvent) else "engine stream error"
|
|
detail = _safe_engine_error_detail(
|
|
detail,
|
|
fallback="engine stream error",
|
|
)
|
|
yield StreamEvent("error", {"detail": detail})
|
|
return
|
|
if packet.event == ENGINE_GATEWAY_SSE_DONE:
|
|
payload = packet.payload
|
|
if isinstance(payload, StreamDoneEvent):
|
|
stream_meta = payload.model_dump()
|
|
gateway_done = True
|
|
break
|
|
|
|
payload = packet.payload
|
|
if not isinstance(payload, StreamTokenEvent):
|
|
continue
|
|
text_piece = payload.text
|
|
accumulated += text_piece
|
|
|
|
if not gateway_done:
|
|
# 토큰 일부 또는 빈 본문 뒤 연결이 끊겨도 성공 done을 합성하지 않는다.
|
|
# 라우트는 이 error 이벤트를 전달하고 durable turn을 저장하지 않는다.
|
|
yield StreamEvent("error", {"detail": "client_stream_incomplete"})
|
|
return
|
|
|
|
scenario_leakage = rupture_scenario_director.contains_internal_scenario_leakage(
|
|
accumulated
|
|
)
|
|
guard = guardrail.sanitize_client_reply(
|
|
accumulated,
|
|
ideation_stage=st.ideation_stage,
|
|
turn_seq=st.turn_seq,
|
|
previous_client_reply=previous_client_reply,
|
|
)
|
|
if scenario_leakage or guard.needs_regeneration:
|
|
flagged = True
|
|
output_error = "client_reply_quality_retryable"
|
|
yield StreamEvent(
|
|
"safety",
|
|
{
|
|
"reason": output_error,
|
|
"reasons": guard.reasons if not scenario_leakage else ["quality"],
|
|
},
|
|
)
|
|
accumulated = ""
|
|
else:
|
|
flagged = flagged or guard.blocked
|
|
if guard.text:
|
|
yield StreamEvent("token", {"text": guard.text})
|
|
|
|
latency_ms = int((time.perf_counter() - started) * 1000)
|
|
await record_llm_audit(
|
|
audit_hook,
|
|
session_id=ctx.session_id,
|
|
provider=str(stream_meta.get("provider") or engine.engine_mode),
|
|
model=str(
|
|
stream_meta.get("model")
|
|
or engine.default_model
|
|
or ENGINE_GATEWAY_DEFAULT_MODEL_SENTINEL
|
|
),
|
|
tokens_in=_safe_int(stream_meta.get("tokens_in")),
|
|
tokens_out=_safe_int(stream_meta.get("tokens_out")),
|
|
cost_usd=_safe_float(stream_meta.get("cost_usd")),
|
|
inference_geo=_optional_str(stream_meta.get("inference_geo")),
|
|
latency_ms=latency_ms,
|
|
)
|
|
|
|
yield StreamEvent(
|
|
"done",
|
|
{
|
|
"session_id": ctx.session_id,
|
|
"stage": st.stage.value,
|
|
"effective_openness": round(st.effective_openness, 4),
|
|
"turn_seq": st.turn_seq,
|
|
"safety_flagged": flagged,
|
|
"llm_provider": str(stream_meta.get("provider") or engine.engine_mode),
|
|
"model": str(
|
|
stream_meta.get("model")
|
|
or engine.default_model
|
|
or ENGINE_GATEWAY_DEFAULT_MODEL_SENTINEL
|
|
),
|
|
"tokens_in": _safe_int(stream_meta.get("tokens_in")),
|
|
"tokens_out": _safe_int(stream_meta.get("tokens_out")),
|
|
"cost_usd": _safe_float(stream_meta.get("cost_usd")),
|
|
"output_error": output_error,
|
|
},
|
|
)
|
|
except EngineGatewaySseDecodeError as e:
|
|
yield StreamEvent(
|
|
"error",
|
|
{"detail": _safe_engine_error_detail(e, fallback="engine stream decode error")},
|
|
)
|
|
except EngineError as e:
|
|
yield StreamEvent(
|
|
"error",
|
|
{"detail": _safe_engine_error_detail(e, fallback="engine stream error")},
|
|
)
|
|
|
|
|
|
def _optional_str(value: Any) -> Optional[str]:
|
|
if value is None:
|
|
return None
|
|
text = str(value).strip()
|
|
return text or None
|
|
|
|
|
|
def _safe_int(value: Any) -> int:
|
|
try:
|
|
return int(value or 0)
|
|
except (TypeError, ValueError):
|
|
return 0
|
|
|
|
|
|
def _safe_float(value: Any) -> float:
|
|
try:
|
|
return float(value or 0.0)
|
|
except (TypeError, ValueError):
|
|
return 0.0
|
|
|
|
|
|
__all__ = [
|
|
"EvalHook",
|
|
"LlmAuditHook",
|
|
"TurnContext",
|
|
"TurnMemory",
|
|
"TurnResult",
|
|
"StreamEvent",
|
|
"prepare_turn",
|
|
"run_turn_generate",
|
|
"run_turn_stream",
|
|
]
|