vignette/apps/api/app/services/orchestrator.py
Yun Chan 29c406d89f Jev 내담자 평가·표현 v2와 속마음 공개
상담자 발화 판정(A)·감정(B)·표현(C) 20문항 질문 세트, 감쇠 없는 이번 턴 반응과 비대칭 기분 전이, 개방도 게이트로 생성 지시를 만들고 ccd.coping_strategy 전달 누락을 고친다.

trace v2와 고정 문구 속마음 요약(migration 24, AI 경로 차단 RLS)을 같은 트랜잭션에 저장하고 피드백 정책이 켜진 경우에만 done·TurnResponse·음성 reply·리뷰로 노출한다. 회기 화면 속마음 보기 토글, 리뷰 접힘 블록, 관리자 감정 관측 v2 표시를 추가한다.
2026-09-30 13:23:09 +09:00

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