Files
CODE_ASSISTANT/5_django_backend/apps/chat/stream.py
T

419 lines
19 KiB
Python
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
"""
POST /api/v1/chat/stream — 프론트 `lib/streaming/streamLLM.ts` SSE 계약.
요청 {sessionId, content, images?[{mediaType,data}], forcedSkill?, explain?}
이벤트 token{delta} ×N → usage{used,limit,ratio,elapsed_ms} → done{}
중간 title{title}, 실패 error{message,code}
흐름
1. 인증(Bearer) → 본인 세션 → 생성 중이면 409
2. 미러: user 메시지 저장 + is_generating=true
3. 이벤트 버스 구독 → OpenCode prompt_async (블록 없음)
4. 턴 감시 태스크(run_turn)가 /event 를 우리 이벤트로 바꿔 큐에 넣음.
클라이언트가 끊겨도 태스크는 idle 까지 돌고 DB 를 마무리함 → 프론트 isGenerating 폴링이 복구.
5. session.idle 이면 OpenCode 에서 최종 메시지·tokens·cost 가져와 저장 → usage → done
순수 Django async 뷰 — DRF 는 async 스트리밍을 못 함. 반드시 ASGI(uvicorn).
"""
import asyncio
import json
import logging
import time
from datetime import datetime, timezone
from django.conf import settings
from django.http import StreamingHttpResponse
from django.views.decorators.csrf import csrf_exempt
from apps.accounts.authentication import user_from_token
from apps.gateway import langfuse
from asgiref.sync import sync_to_async
from common.envelope import json_error
from common.opencode_service import opencode_service
from .events import event_bus
from .models import ChatMessage, ChatSession
log = logging.getLogger(__name__)
ALLOWED_IMAGE_TYPES = {"image/png", "image/jpeg", "image/webp"}
MAX_IMAGES = 4
# 타임아웃은 .env 로 조정 (settings.STREAM_*). Gemma4 가 큐 밀리면 첫 토큰도 늦어서 넉넉히.
TURN_TIMEOUT_S = settings.STREAM_TURN_TIMEOUT_S # OpenCode 답변 상한 (기본 600)
# prompt 를 받았는데 이 세션 이벤트가 하나도 안 오면 OpenCode 가 조용히 실패한 것
# (예: 다른 프로젝트의 세션 — 로그에만 "prompt_async failed" 남고 session.error 안 옴)
FIRST_EVENT_TIMEOUT_S = settings.STREAM_FIRST_EVENT_TIMEOUT_S # 기본 300
KEEPALIVE_S = 15
_END = object() # 큐 종료 표시
# ── SSE 프레임 ──────────────────────────────────────────────────
def sse(event: str, data) -> bytes:
return f"event: {event}\ndata: {json.dumps(data, ensure_ascii=False)}\n\n".encode()
def _token_from(request) -> str:
auth = request.headers.get("Authorization", "")
return auth[7:] if auth.startswith("Bearer ") else request.GET.get("token", "")
def _parse_body(request) -> tuple[dict | None, str | None]:
try:
body = json.loads(request.body or b"{}")
except ValueError:
return None, "JSON 이 아니야"
if not isinstance(body, dict) or not body.get("sessionId"):
return None, "sessionId 가 없어"
content = body.get("content") or ""
images = body.get("images") or []
if not isinstance(images, list) or len(images) > MAX_IMAGES:
return None, f"이미지는 최대 {MAX_IMAGES}장"
for img in images:
if not isinstance(img, dict) or img.get("mediaType") not in ALLOWED_IMAGE_TYPES or not img.get("data"):
return None, "이미지 형식이 잘못됐어 (png/jpeg/webp data URL)"
if not content.strip() and not images:
return None, "content 나 images 중 하나는 있어야 해"
return {"sessionId": body["sessionId"], "content": content, "images": images, "explain": bool(body.get("explain"))}, None
def _prompt_payload(req: dict) -> dict:
text = req["content"]
if req["explain"]:
text = f"[설명 모드] {text}".strip()
parts: list[dict] = []
if text.strip():
parts.append({"type": "text", "text": text})
for i, img in enumerate(req["images"]):
ext = img["mediaType"].split("/")[1]
parts.append({"type": "file", "mime": img["mediaType"], "filename": f"image-{i + 1}.{ext}", "url": img["data"]})
payload: dict = {"agent": settings.OPENCODE_AGENT, "parts": parts}
# 이미지 있으면 vision 모델로 — 텍스트는 기본 모델(빠름). 고객사 FabriX 는 게이트웨이가 따로 분기.
if req["images"] and settings.OPENCODE_VISION_MODEL and "/" in settings.OPENCODE_VISION_MODEL:
provider, model = settings.OPENCODE_VISION_MODEL.split("/", 1)
payload["model"] = {"providerID": provider, "modelID": model}
return payload
# ── OpenCode 이벤트 → 우리 이벤트 ────────────────────────────────
class TurnState:
"""한 턴 동안 delta 계산에 필요한 것."""
def __init__(self, session_id: str, user_text: str):
self.session_id = session_id
self.user_text = user_text
self.user_message_ids: set[str] = set()
self.assistant_message_ids: set[str] = set()
self.seen_len: dict[str, int] = {} # partID → 지금까지 보낸 글자 수
self.part_types: dict[str, str] = {} # partID → "text"/"reasoning"/… (part.updated 가 delta 보다 먼저 옴)
self.text_by_part: dict[str, str] = {}
self.part_order: list[str] = []
self.title: str | None = None
def accumulated(self) -> str:
return "".join(self.text_by_part.get(p, "") for p in self.part_order)
def handle(self, event: dict) -> list[tuple[str, dict]]:
"""이벤트 하나 → 프론트로 보낼 (event, data) 목록. idle/error 는 호출자가 봄."""
out: list[tuple[str, dict]] = []
etype = event.get("type")
props = event.get("properties") or {}
if etype == "message.updated":
info = props.get("info") or {}
if info.get("role") == "user":
self.user_message_ids.add(info.get("id", ""))
elif info.get("role") == "assistant":
self.assistant_message_ids.add(info.get("id", ""))
return out
if etype == "message.part.delta":
# 실서버(1.18.6)의 글자 단위 스트리밍은 이 이벤트로 옴 — SDK 타입엔 없지만 실제로 174개/답변.
# {sessionID, messageID, partID, field:"text", delta}. 파트 타입은 안 실려서 part.updated 로 알아둔 것으로 거름.
if props.get("field") != "text" or not props.get("delta"):
return out
mid = props.get("messageID", "")
if mid in self.user_message_ids and mid not in self.assistant_message_ids:
return out
pid = props.get("partID", "")
ptype = self.part_types.get(pid)
if ptype == "reasoning":
out.append(("step", {"id": pid, "kind": "reasoning", "text": props["delta"]}))
return out
if ptype != "text":
return out # tool 파트, 또는 아직 타입 모름(최종 스냅샷이 메워줌)
delta = props["delta"]
if pid not in self.text_by_part:
self.part_order.append(pid)
self.text_by_part[pid] = self.text_by_part.get(pid, "") + delta
self.seen_len[pid] = len(self.text_by_part[pid])
out.append(("token", {"delta": delta}))
return out
if etype == "message.part.updated":
part = props.get("part") or {}
if part.get("id"):
self.part_types[part["id"]] = part.get("type", "")
if part.get("type") == "tool":
# 도구 호출 진행 — 화면 "생각 과정" 에 이름·상태만. 입력/출력은 안 보냄(크고 사용자 관심 밖)
st = part.get("state") or {}
out.append(("step", {"id": part.get("id", ""), "kind": "tool", "tool": part.get("tool", ""), "status": st.get("status", ""), "title": st.get("title") or ""}))
return out
if part.get("type") != "text" or part.get("synthetic") or part.get("ignored"):
return out
mid = part.get("messageID", "")
# 사용자 메시지 파트(우리가 보낸 질문의 echo)는 건너뜀
if mid in self.user_message_ids and mid not in self.assistant_message_ids:
return out
pid = part.get("id", "")
text = part.get("text") or ""
delta = props.get("delta")
if pid not in self.text_by_part:
self.part_order.append(pid)
if delta:
self.text_by_part[pid] = self.text_by_part.get(pid, "") + delta
self.seen_len[pid] = len(self.text_by_part[pid])
out.append(("token", {"delta": delta}))
else:
seen = self.seen_len.get(pid, 0)
if len(text) > seen and text.startswith(self.text_by_part.get(pid, "")):
out.append(("token", {"delta": text[seen:]}))
self.seen_len[pid] = len(text)
self.text_by_part[pid] = text if len(text) >= seen else self.text_by_part.get(pid, "")
return out
if etype == "session.updated":
title = _real_title((props.get("info") or {}).get("title"))
if title and title != self.title:
self.title = title
out.append(("title", {"title": title}))
return out
return out
def _real_title(raw) -> str | None:
"""OpenCode 기본 제목("New session - 2026-…")은 제목이 아님 → None."""
t = (raw or "").strip()
if not t or t.startswith("New session"):
return None
return t
def _error_message(err: dict | None) -> str:
if not err:
return "OpenCode 오류"
data = err.get("data") or {}
return data.get("message") or err.get("name") or "OpenCode 오류"
async def _finalize(session: ChatSession, state: TurnState, started: float, *, failed: str | None = None):
"""idle/에러/타임아웃 후 미러 DB 마무리. 실패해도 is_generating 은 반드시 풀림."""
content = state.accumulated()
usage: dict | None = None
title = state.title
try:
msgs = await opencode_service.list_messages_a(session.id) or []
assistant = [m for m in msgs if (m.get("info") or {}).get("role") == "assistant"]
if assistant:
info = assistant[-1]["info"]
parts = assistant[-1].get("parts") or []
full = "".join(p.get("text", "") for p in parts if p.get("type") == "text" and not p.get("synthetic"))
if full:
content = full
tokens = info.get("tokens") or {}
t = info.get("time") or {}
elapsed = (t.get("completed") or int(time.time() * 1000)) - (t.get("created") or int(started * 1000))
usage = {
"input": int(tokens.get("input") or 0),
"output": int(tokens.get("output") or 0) + int(tokens.get("reasoning") or 0),
"cost": float(info.get("cost") or 0.0),
"elapsed_ms": max(0, int(elapsed)),
}
if not title:
info = await opencode_service.get_session_a(session.id) or {}
title = _real_title(info.get("title"))
except Exception: # noqa: BLE001
log.warning("턴 마무리 중 OpenCode 조회 실패 — 누적 텍스트로 저장: %s", session.id)
if content or failed is None:
await ChatMessage.objects.acreate(
session=session,
role="assistant",
content=content or (f"(오류: {failed})" if failed else ""),
input_tokens=usage["input"] if usage else None,
output_tokens=usage["output"] if usage else None,
cost_usd=usage["cost"] if usage else None,
elapsed_ms=usage["elapsed_ms"] if usage else None,
)
session.is_generating = False
if title:
session.title_llm = title[:200]
await session.asave(update_fields=["is_generating", "title_llm", "updated_at"])
# 관측: 사용자 단위 trace 하나 + 답변 generation 하나. 실패해도 여기까진 이미 저장됨.
if langfuse.enabled():
tid = f"{session.id}-{int(started * 1000)}"
user_email = await sync_to_async(lambda: session.user.email)()
meta = {"failed": failed} if failed else {}
t0 = datetime.fromtimestamp(started, timezone.utc).isoformat()
root = langfuse.trace(tid, "chat", userId=user_email, sessionId=session.id, input=state.user_text,
output=content, metadata=meta, tags=["codeassist"], start=t0, end=langfuse.now_iso())
langfuse.send_later([
root,
langfuse.generation(
tid, "opencode-turn", parent=root,
startTime=t0, endTime=langfuse.now_iso(),
input=state.user_text, output=content,
usage=langfuse.usage_of(usage["input"], usage["output"]) if usage else None,
level="ERROR" if failed else "DEFAULT", statusMessage=failed or "",
metadata={"elapsed_ms": usage["elapsed_ms"]} if usage else {},
),
])
return usage, title
TITLE_WAIT_S = 20 # OpenCode 제목 생성은 별도 LLM 호출이라 idle 뒤에 오기도 함
async def _wait_late_title(session: ChatSession, q: asyncio.Queue) -> None:
"""idle 뒤 늦게 오는 session.updated 제목을 미러에만 반영 (목록 재조회 때 보임)."""
deadline = time.time() + TITLE_WAIT_S
while (remaining := deadline - time.time()) > 0:
try:
event = await asyncio.wait_for(q.get(), timeout=remaining)
except asyncio.TimeoutError:
return
if not event or event.get("type") != "session.updated":
continue
title = _real_title(((event.get("properties") or {}).get("info") or {}).get("title"))
if title:
session.title_llm = title[:200]
await session.asave(update_fields=["title_llm", "updated_at"])
return
async def run_turn(session: ChatSession, req: dict, out: asyncio.Queue) -> None:
"""한 턴 감시 — 클라이언트와 무관하게 idle 까지 돌고 DB 마무리. 결과는 out 큐로."""
state = TurnState(session.id, req["content"])
started = time.time()
q = event_bus.subscribe(session.id)
try:
await opencode_service.prompt_async(session.id, _prompt_payload(req))
except Exception as e: # noqa: BLE001
event_bus.unsubscribe(session.id, q)
log.warning("prompt_async 실패: %s", e)
await _finalize(session, state, started, failed=str(e))
out.put_nowait(("error", {"message": "OpenCode 에 질문을 못 보냈어", "code": "LLM_ERROR"}))
out.put_nowait(_END)
return
try:
deadline = started + TURN_TIMEOUT_S
got_any = False
while True:
now = time.time()
remaining = deadline - now
silent_too_long = not got_any and now - started > FIRST_EVENT_TIMEOUT_S
if remaining <= 0 or silent_too_long:
try:
await opencode_service.abort_session_a(session.id)
except Exception: # noqa: BLE001
pass
reason = "no-events" if silent_too_long else "timeout"
await _finalize(session, state, started, failed=reason)
msg = "OpenCode 가 응답을 시작하지 않았어" if silent_too_long else "답변 시간이 너무 길어서 중단했어"
out.put_nowait(("error", {"message": msg, "code": "LLM_ERROR"}))
break
try:
event = await asyncio.wait_for(q.get(), timeout=min(KEEPALIVE_S, remaining))
except asyncio.TimeoutError:
out.put_nowait(None) # keepalive
continue
if event is None:
out.put_nowait(None)
continue
got_any = True
etype = event.get("type")
props = event.get("properties") or {}
if etype == "session.error" or (
etype == "message.updated"
and (props.get("info") or {}).get("role") == "assistant"
and (props.get("info") or {}).get("error")
):
err = props.get("error") or (props.get("info") or {}).get("error")
msg = _error_message(err)
await _finalize(session, state, started, failed=msg)
code = "LLM_ABORTED" if (err or {}).get("name") == "MessageAbortedError" else "LLM_ERROR"
out.put_nowait(("error", {"message": msg, "code": code}))
break
for ev in state.handle(event):
out.put_nowait(ev)
if etype == "session.idle":
usage, title = await _finalize(session, state, started)
if title and title != state.title:
out.put_nowait(("title", {"title": title}))
if usage:
used = usage["input"] + usage["output"]
limit = settings.CONTEXT_LIMIT_TOKENS
out.put_nowait(
("usage", {"used": used, "limit": limit, "ratio": round(used / limit, 4) if limit else 0, "elapsed_ms": usage["elapsed_ms"]})
)
out.put_nowait(("done", {}))
out.put_nowait(_END) # 클라이언트는 여기서 끝. 제목은 뒤에서 조용히 기다림
if not title:
await _wait_late_title(session, q)
break
except Exception as e: # noqa: BLE001
log.exception("턴 감시 실패: %s", session.id)
await _finalize(session, state, started, failed=str(e))
out.put_nowait(("error", {"message": "서버 오류로 답변을 못 받았어", "code": "LLM_ERROR"}))
finally:
event_bus.unsubscribe(session.id, q)
out.put_nowait(_END)
# ── 뷰 ──────────────────────────────────────────────────────────
@csrf_exempt
async def stream_chat(request):
if request.method != "POST":
return json_error(405, "METHOD_NOT_ALLOWED", "POST 만 받아")
user = await sync_to_async(user_from_token)(_token_from(request))
if user is None:
return json_error(401, "UNAUTHORIZED", "로그인이 필요해")
req, err = _parse_body(request)
if err:
return json_error(400, "VALIDATION_ERROR", err)
try:
session = await ChatSession.objects.aget(id=req["sessionId"], user=user)
except ChatSession.DoesNotExist:
return json_error(404, "NOT_FOUND", "세션이 없어")
if session.is_generating:
return json_error(409, "CHAT_GENERATION_IN_PROGRESS", "아직 답변 생성 중이야")
session.is_generating = True
await session.asave(update_fields=["is_generating", "updated_at"])
await ChatMessage.objects.acreate(session=session, role="user", content=req["content"])
out: asyncio.Queue = asyncio.Queue()
asyncio.get_running_loop().create_task(run_turn(session, req, out), name=f"turn-{session.id}")
async def gen():
while True:
item = await out.get()
if item is _END:
break
if item is None:
yield b": keepalive\n\n"
continue
yield sse(*item)
resp = StreamingHttpResponse(gen(), content_type="text/event-stream")
resp["Cache-Control"] = "no-cache"
resp["X-Accel-Buffering"] = "no"
return resp