""" 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 {} langfuse.send_later([ langfuse.trace(tid, "chat", userId=user_email, sessionId=session.id, input=state.user_text, output=content, metadata=meta, tags=["codeassist"]), langfuse.generation( tid, "opencode-turn", startTime=datetime.fromtimestamp(started, timezone.utc).isoformat(), 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