vignette/apps/api/app/services/orchestrator.py
Yun Chan 24b1b7a6e1 feat: P1 풀빌드 — React 프론트 7화면 + 백엔드 상담루프·평가·음성·RAG
web (Vite+React19+TS, Cloudflare Pages 배포):
- 디자인토큰(세이지틸/테라코타 SSOT), 앱셸, 공통 UI 프리미티브
- 7화면: 로그인/학습자홈/상담세션/회기리뷰/교수자/관리자/설정
- ClientAvatar: SVG 반구상 흉상 4상태 + RMS 립싱크 + 6파라미터 정서
- 회기리뷰는 외부 레퍼런스 디자인을 Vignette 토큰으로 리스킨

api (FastAPI):
- 게이트웨이 /v1/generate·/v1/stream 어댑터(상주풀/EngineSession 보존)
- services: 페르소나 L0~L6 빌더 / 결정론 상태머신 / 가드레일 /
  턴 오케스트레이터 / 회기간 메모리 / 평가AI / 음성 / RAG
- store: DB off 폴백(in-memory), sessions 실구현

검증:
- web: node22 tsc+vite build 통과(node23 segfault 회피), Pages 배포 200
- api: app.main import 통과
- 핫픽스: Topbar initials undefined-safe (undefined.trim 크래시)
- E2E: 서연(P1) 상담 1턴 — 좋은/나쁜 상담에 차등 반응 실증
2026-06-25 23:37:22 +09:00

329 lines
13 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
from dataclasses import dataclass, field
from typing import Any, AsyncIterator, Awaitable, Callable, Optional
from ..engine_client import (
EngineClient,
EngineError,
EngineMessage,
GenerateRequest,
GenerateResponse,
StreamRequest,
)
from . import guardrail, persona, state_machine
from .persona import PersonaCard, PersonaStateContext
from .state_machine import SessionState, Stage
# 평가 훅 타입: U_t(수련생 마스킹 발화) + 내담자응답 + 상태 → 평가 결과(dict)
# Features evaluator 가 이 시그니처에 맞춰 함수를 주입한다(여기선 호출만).
EvalHook = Callable[["TurnContext", str], Awaitable[Optional[dict]]]
# 로깅 훅: TurnContext + 내담자응답 → None (turns insert/임베딩은 주입측 책임)
LogHook = Callable[["TurnContext", str], Awaitable[None]]
@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 = ""
crisis: Optional[guardrail.CrisisResult] = None
state_after: Optional[SessionState] = None
messages: list[EngineMessage] = field(default_factory=list)
# 회상/메모리 주입(memory.RecallContext 에서 옴)
recall_summary: Optional[str] = None
pinned_facts: list[str] = field(default_factory=list)
recent_turns: list[dict[str, str]] = field(default_factory=list)
kb_behavior_cues: list[str] = field(default_factory=list)
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,
)
@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"
# ════════════════════════════════════════════════════════════════════════════
# 1~3단계 — 입력 가드레일 + 상태머신 + 페르소나 컨텍스트 (엔진 호출 전 결정론)
# ════════════════════════════════════════════════════════════════════════════
def prepare_turn(
*,
session_id: str,
case_id: Optional[str],
card: PersonaCard,
state: SessionState,
learner_text: str,
recall_summary: Optional[str] = None,
pinned_facts: Optional[list[str]] = None,
recent_turns: Optional[list[dict[str, str]]] = None,
kb_behavior_cues: Optional[list[str]] = None,
eval_rapport_signal: Optional[float] = None,
) -> TurnContext:
"""엔진 호출 전 결정론 전처리(1~3단계). 순수 — IO/LLM 없음.
eval_rapport_signal 이 주어지면(평가 AI fast-loop 신호) 그걸 쓰고, 없으면
state_machine 의 경량 휴리스틱으로 라포 신호를 추정한다.
"""
ctx = TurnContext(
session_id=session_id,
case_id=case_id,
persona=card,
state_before=state,
learner_text_raw=learner_text,
recall_summary=recall_summary,
pinned_facts=list(pinned_facts or []),
recent_turns=list(recent_turns or []),
kb_behavior_cues=list(kb_behavior_cues or []),
)
# 1) 입력 가드레일 — PII 마스킹 + 위기분류
mask = guardrail.mask_pii(learner_text)
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)
)
ctx.state_after = state_machine.evolve(
state,
rapport_signal=signal,
unlock_rate=card.unlock_rate(),
decay_floor=card.decay_floor(),
)
# 3) 페르소나 컨텍스트 — L0~L6 messages 조립 (CCD 는 행동으로만, L0 가 강제)
ctx.messages = persona.build_turn_messages(
card,
ctx.to_state_context(),
ctx.learner_text_masked,
recall_summary=ctx.recall_summary,
pinned_facts=ctx.pinned_facts,
recent_turns=ctx.recent_turns,
kb_behavior_cues=ctx.kb_behavior_cues,
)
return ctx
# ════════════════════════════════════════════════════════════════════════════
# 4~8단계 — 동기 생성 경로 (폴백/테스트)
# ════════════════════════════════════════════════════════════════════════════
async def run_turn_generate(
ctx: TurnContext,
engine: EngineClient,
*,
eval_hook: Optional[EvalHook] = None,
log_hook: Optional[LogHook] = None,
) -> TurnResult:
"""동기 턴 실행(4~8). 내담자 응답을 한 번에 받아 가드레일·평가·로깅 훅 순차 적용.
eval_hook/log_hook 은 Features 가 주입(없으면 생략). 엔진 장애는 EngineError 전파.
"""
assert ctx.state_after is not None
st = ctx.state_after
# 4) 내담자 AI 생성
req = GenerateRequest(
ai_role="client",
tier="client",
messages=ctx.messages,
session_id=ctx.session_id,
metadata={"stage": st.stage.value},
)
resp: GenerateResponse = await engine.generate(req)
reply = resp.text
# 5) 출력 가드레일 — 수단 차단 + ideation 상한
guard = guardrail.sanitize_client_reply(reply, ideation_stage=st.ideation_stage)
safety_flagged = guard.blocked or (ctx.crisis is not None and ctx.crisis.escalate)
if guard.needs_regeneration:
# 수단정보 누출 → 안전 대체 응답으로 치환(1차). 재생성 루프는 후속.
reply = "…(말을 잇지 못하고 잠시 침묵한다)"
# 6) 평가 훅(주입형) — 평가 AI 4차원 태깅 (Features 소유)
evaluation: Optional[dict] = None
if eval_hook is not None:
try:
evaluation = await eval_hook(ctx, reply)
except Exception:
evaluation = None # 평가 실패가 상담 루프를 막지 않게(비치명적)
# 8) 로깅 훅(주입형) — turns insert + 임베딩
if log_hook is not None:
try:
await log_hook(ctx, reply)
except Exception:
pass
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",
)
# ════════════════════════════════════════════════════════════════════════════
# 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,
*,
log_hook: Optional[LogHook] = None,
) -> AsyncIterator[StreamEvent]:
"""스트리밍 턴 실행(4~8). 게이트웨이 SSE 를 받아 token/done/safety/error 로 재방출.
출력 가드레일은 *누적 텍스트* 기준으로 수단정보를 감지(스트림 중 발견 시 safety 이벤트 +
재생성 신호). 토큰 단위 완벽 차단은 후속(현재는 누적 스캔).
로깅 훅은 done 직전 최종 텍스트로 1회 호출.
"""
assert ctx.state_after is not None
st = ctx.state_after
req = StreamRequest(
ai_role="client",
tier="client",
messages=ctx.messages,
session_id=ctx.session_id,
metadata={"stage": st.stage.value},
)
accumulated = ""
flagged = False
if ctx.crisis is not None and ctx.crisis.escalate:
flagged = True
yield StreamEvent("safety", {"reason": "learner_real_crisis", "level": ctx.crisis.risk_level})
try:
async for raw in engine.stream(req):
# engine_client.stream 은 게이트웨이 SSE 의 *원시 라인*을 그대로 yield 한다.
# 게이트웨이 프레이밍: "event: token\ndata: {\"text\": ...}" 형식.
text_piece = _extract_sse_text(raw)
if text_piece is None:
continue
accumulated += text_piece
# 출력 가드레일(누적 스캔) — 수단정보 발견 시 차단·재생성 신호
guard = guardrail.sanitize_client_reply(accumulated, ideation_stage=st.ideation_stage)
if guard.needs_regeneration and not flagged:
flagged = True
yield StreamEvent("safety", {"reason": "means_info_blocked"})
# 토큰은 더 내보내지 않고 안전 대체로 종결
accumulated = "…(말을 잇지 못하고 잠시 침묵한다)"
break
yield StreamEvent("token", {"text": text_piece})
# 8) 로깅 훅 — 최종 텍스트
if log_hook is not None:
try:
await log_hook(ctx, accumulated)
except Exception:
pass
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,
},
)
except EngineError as e:
yield StreamEvent("error", {"detail": str(e)})
def _extract_sse_text(raw_line: str) -> Optional[str]:
"""게이트웨이 SSE 원시 라인에서 텍스트 델타를 추출.
게이트웨이 /v1/stream 은 'event: token' + 'data: {"text": "..."}' 를 보낸다.
engine_client.stream 은 빈 줄을 필터링하고 비어있지 않은 라인만 흘리므로
여기서 data: 라인의 JSON 만 해석한다. token 이외 이벤트(done/error)는 None.
"""
import json as _json
line = raw_line.strip()
if not line.startswith("data:"):
return None
payload = line[len("data:"):].strip()
if not payload or payload == "[DONE]":
return None
try:
obj = _json.loads(payload)
except _json.JSONDecodeError:
return None
if isinstance(obj, dict) and "text" in obj:
return obj["text"]
return None
__all__ = [
"EvalHook",
"LogHook",
"TurnContext",
"TurnResult",
"StreamEvent",
"prepare_turn",
"run_turn_generate",
"run_turn_stream",
]