- engine_gateway/gateway.py: 회기당 claude -p 상주 프로세스(stream-json), 턴 직렬, budget 제한 - 세션 생성/턴/종료 HTTP API(FastAPI, :9099) - 검증: 멀티턴 컨텍스트 유지 + prompt caching 재사용(턴2 +$0.07) 실동작 확인
188 lines
8 KiB
Python
188 lines
8 KiB
Python
"""상담 세션 라우트 — 시작 / 턴 / 종료 + SSE 스트림 스텁.
|
||
|
||
흐름 (설계서 §2 회기 라이프사이클 + 마스터플랜 §2.2 턴 사이클):
|
||
POST /sessions — 회기 시작 (case_profile/summary 회상 + session_state 초기화)
|
||
POST /sessions/{id}/turn — 수련생 발화 1턴 (가드레일→상태머신→내담자AI→평가)
|
||
GET /sessions/{id}/stream — 내담자 AI 응답 SSE 스트림 (Cloudflare 우회 heartbeat)
|
||
POST /sessions/{id}/end — 회기 종료 (무손실 carry-over + LLM 압축 트리거)
|
||
|
||
상태머신(라포→탐색→개입→정리)은 백엔드가 결정론적으로 소유(LLM 아님, 마스터플랜 §0).
|
||
이 파일은 핸들러 시그니처 + 계약 + TODO. 실제 상태머신/가드레일/압축은 Phase 1~2a 트랙 A.
|
||
"""
|
||
|
||
from __future__ import annotations
|
||
|
||
import asyncio
|
||
import json
|
||
from typing import Annotated, Literal, Optional
|
||
from uuid import UUID, uuid4
|
||
|
||
from fastapi import APIRouter, Depends, HTTPException, status
|
||
from pydantic import BaseModel, Field
|
||
from sse_starlette.sse import EventSourceResponse
|
||
|
||
from ..config import settings
|
||
from ..deps import CurrentPrincipal, HumanDB
|
||
from ..engine_client import EngineMessage, StreamRequest, engine_client, EngineError
|
||
|
||
router = APIRouter(prefix="/sessions", tags=["sessions"])
|
||
|
||
Stage = Literal["라포", "탐색", "개입", "정리"]
|
||
|
||
|
||
# ── 요청/응답 모델 ──────────────────────────────────────
|
||
class SessionStartRequest(BaseModel):
|
||
persona_code: str = Field(..., examples=["P1"]) # 시드 페르소나 (P1/P2/P3)
|
||
theory_mode: Literal["humanistic", "cbt", "integrative"] = "humanistic"
|
||
|
||
|
||
class SessionStartResponse(BaseModel):
|
||
session_id: UUID
|
||
case_id: UUID
|
||
session_no: int
|
||
stage: Stage
|
||
# 회기 시작 회상 요약 (큰그림→세부, UI 카드용. CCD/정답은 절대 미포함)
|
||
recall_summary: Optional[str] = None
|
||
|
||
|
||
class TurnRequest(BaseModel):
|
||
text: str = Field(..., min_length=1) # 수련생 발화 (저장 전 PII 마스킹)
|
||
|
||
|
||
class TurnResponse(BaseModel):
|
||
turn_seq: int
|
||
stage: Stage
|
||
effective_openness: float
|
||
# 내담자 응답은 스트림(GET /stream)으로 받는 게 기본. 동기 응답은 폴백/테스트용.
|
||
client_reply: Optional[str] = None
|
||
safety_flagged: bool = False
|
||
|
||
|
||
class SessionEndResponse(BaseModel):
|
||
session_id: UUID
|
||
session_no: int
|
||
digest_pending: bool # 압축은 비동기 비블로킹 (설계서 §2-C)
|
||
|
||
|
||
# ── 핸들러 ──────────────────────────────────────────────
|
||
@router.post("", response_model=SessionStartResponse, status_code=status.HTTP_201_CREATED)
|
||
async def start_session(
|
||
body: SessionStartRequest,
|
||
principal: CurrentPrincipal,
|
||
conn: HumanDB,
|
||
) -> SessionStartResponse:
|
||
"""회기 시작 — 회상 + 상태 복원 (설계서 §2-A).
|
||
|
||
절차:
|
||
1. persona_card(approved) 조회 + (persona_id, learner_id) -> case_profile upsert
|
||
2. case_digest + 직전 session_summary + episodic recall (Phase 2a, 1차는 단일회기)
|
||
3. session_state 초기화: stage='라포', carry-over (rapport×0.7, ideation 보수적 유지) [P2]
|
||
4. sessions 행 insert
|
||
TODO: persona 조회/회상/상태머신 init 구현 (트랙 A). 현재 스텁 응답.
|
||
"""
|
||
# TODO: SELECT persona_id FROM app.persona_card WHERE code=$1 AND status='approved'
|
||
# TODO: init_session_state_from_history() — 결정론 carry-over
|
||
session_id = uuid4()
|
||
case_id = uuid4()
|
||
return SessionStartResponse(
|
||
session_id=session_id,
|
||
case_id=case_id,
|
||
session_no=1,
|
||
stage="라포",
|
||
recall_summary=None, # Phase 2a 회상 채움
|
||
)
|
||
|
||
|
||
@router.post("/{session_id}/turn", response_model=TurnResponse)
|
||
async def submit_turn(
|
||
session_id: UUID,
|
||
body: TurnRequest,
|
||
principal: CurrentPrincipal,
|
||
conn: HumanDB,
|
||
) -> TurnResponse:
|
||
"""수련생 발화 1턴 (마스터플랜 §2.2 / 설계서 §2-B).
|
||
|
||
파이프라인 (전부 백엔드 결정론 게이트):
|
||
1. [입력 가드레일] Presidio PII 마스킹 + 위기분류(실제위기 vs 페르소나 연기) [R7/F-03]
|
||
2. [상태머신] effective_openness = clamp(stage.openness
|
||
+ rapport_credit*unlock_rate - resistance*decay, 0, 1) [P2, 결정론]
|
||
3. [모순 검사] pinned_fact locked 모순 -> 차단·재생성 (설계서 §2-B)
|
||
4. [내담자 AI] engine_client.stream/generate (CCD 직접노출 금지, Structured Outputs)
|
||
5. [출력 가드레일] 자살수단 차단, ideation_stage <= 3 상한 [R5]
|
||
6. [평가 AI] fast-loop 4차원 태깅 (deep-loop 은 단계전환/회기말)
|
||
7. [working 갱신] session_state UPSERT (체크포인트)
|
||
8. [로깅] turns insert + 임베딩 (재귀학습 원천)
|
||
TODO: 1~8 구현 (트랙 A). 현재 스텁: 발화 검증만.
|
||
"""
|
||
# TODO: load session_state, run guardrail + state machine deterministically
|
||
# 동기 응답은 폴백. 기본 UX 는 GET /stream 으로 토큰 스트리밍.
|
||
return TurnResponse(
|
||
turn_seq=0,
|
||
stage="라포",
|
||
effective_openness=0.15,
|
||
client_reply=None,
|
||
safety_flagged=False,
|
||
)
|
||
|
||
|
||
@router.get("/{session_id}/stream")
|
||
async def stream_client_reply(
|
||
session_id: UUID,
|
||
principal: CurrentPrincipal,
|
||
):
|
||
"""내담자 AI 응답 SSE 스트림 (마스터플랜 §1.1 SSE 분리경로).
|
||
|
||
- Cloudflare 100초 timeout 회피: settings.sse_heartbeat_seconds 마다 ping 이벤트 [R2]
|
||
- 게이트웨이 SSE(engine_client.stream)를 프록시해 토큰을 재방출
|
||
- 이벤트: {event: "token"|"done"|"safety"|"ping", data: ...}
|
||
|
||
TODO: 게이트웨이와 StreamRequest 바디 결합(현재 stage 회상 컨텍스트 없이 placeholder).
|
||
상태머신 컨텍스트(L3 stage/openness) + 마스킹된 최근 N턴 주입.
|
||
"""
|
||
|
||
async def event_generator():
|
||
# heartbeat 와 엔진 스트림을 병행 (Cloudflare 버퍼링/타임아웃 회피)
|
||
last_beat = asyncio.get_event_loop().time()
|
||
|
||
# TODO: 실제 StreamRequest 조립 — session_state 에서 stage/openness/최근턴 로드
|
||
req = StreamRequest(
|
||
ai_role="client",
|
||
tier="client",
|
||
messages=[
|
||
EngineMessage(role="system", content="<persona L0~L2 cache_control 주입 TODO>", cache=True),
|
||
EngineMessage(role="user", content="<masked latest learner turn TODO>"),
|
||
],
|
||
)
|
||
|
||
try:
|
||
async for chunk in engine_client.stream(req):
|
||
yield {"event": "token", "data": chunk}
|
||
now = asyncio.get_event_loop().time()
|
||
if now - last_beat >= settings.sse_heartbeat_seconds:
|
||
yield {"event": "ping", "data": "{}"}
|
||
last_beat = now
|
||
yield {"event": "done", "data": json.dumps({"session_id": str(session_id)})}
|
||
except EngineError as e:
|
||
yield {"event": "error", "data": json.dumps({"detail": str(e)})}
|
||
|
||
return EventSourceResponse(event_generator())
|
||
|
||
|
||
@router.post("/{session_id}/end", response_model=SessionEndResponse)
|
||
async def end_session(
|
||
session_id: UUID,
|
||
principal: CurrentPrincipal,
|
||
conn: HumanDB,
|
||
) -> SessionEndResponse:
|
||
"""회기 종료 — carry-over + 압축 트리거 (설계서 §2-C, 비동기 비블로킹).
|
||
|
||
절차:
|
||
(A) 무손실 carry-over: end_state = session_state 종료 snapshot (코드 복사, LLM 미경유) [P4]
|
||
(B) salience 산출 -> 망각/유지
|
||
(C) narrative 압축 (LLM, 상주 claude -p 재사용) — 비동기 큐
|
||
(D~G) digest 임베딩 / case_profile 병합 / pinned 모순처리 / RAG 동기화
|
||
+ 상주 프로세스 회수 (말투표류 회기경계 차단)
|
||
TODO: (A) 동기 수행 후 (B~G)는 BackgroundTasks/큐로 비블로킹. 현재 스텁.
|
||
"""
|
||
# TODO: UPDATE app.sessions SET ended_at=now(); copy end_state; enqueue compression
|
||
return SessionEndResponse(session_id=session_id, session_no=1, digest_pending=True)
|