"""Shared original Agents SDK conversation, tool events and approval/resume engine. No intent router or answer templates: the agent chooses tools and writes answers. """ from __future__ import annotations import asyncio import json import logging import queue import threading import uuid from agents import Runner, RunState, SQLiteSession, set_tracing_disabled from agents.items import ItemHelpers from openai.types.responses import (ResponseCompletedEvent, ResponseCreatedEvent, ResponseReasoningSummaryTextDeltaEvent, ResponseReasoningTextDeltaEvent, ResponseTextDeltaEvent) set_tracing_disabled(True) log = logging.getLogger("ontic.agent") class SessionRuntime: def __init__(self, stream: bool = True): self.id = uuid.uuid4().hex self.events: queue.Queue = queue.Queue() self.thread: threading.Thread | None = None self.running = False self.stream = stream self.memory = SQLiteSession(self.id) # 대화 이력 (프로세스 메모리; 파일 경로를 주면 영속) self.state: RunState | None = None # 승인 대기 중인 런 self.usage = None self.mcp = self.make_mcp() # 세션당 하나 (런마다 connect/cleanup; 이벤트 루프가 런마다 새로 생기므로 접속은 그때) steps = {} own_tools = set() def make_mcp(self): return None def _agent(self, mcp_servers=()): raise NotImplementedError def command(self, name, args): return "" def _emit(self, **e): if e["type"] not in ("token", "reasoning"): log.info("session %s %s %s", self.id, e["type"], e.get("tool") or e.get("state") or "") self.events.put(e) @staticmethod def _args(raw: str | None) -> dict: try: return {k: v for k, v in json.loads(raw or "{}").items() if k != "summary"} except ValueError: return {} @property def status(self) -> str: return "running" if self.busy() else "waiting_for_confirmation" if self.state is not None else "idle" def pending(self) -> list[dict]: return [{"tool": i.tool_name, "args": self._args(i.arguments)} for i in (self.state.get_interruptions() if self.state else [])] async def _arun(self, inp): return await self._stream(inp, ()) async def _stream(self, inp, servers): result = Runner.run_streamed(self._agent(servers), inp, context=self, session=self.memory, max_turns=30) calls: dict[str, str] = {} n_llm, summary_seen = 0, False async for ev in result.stream_events(): if ev.type == "raw_response_event": # LLM 호출마다: 시작 → 생각(reasoning 델타) → 끝. 기다림이 보이게 한다 d = ev.data if isinstance(d, ResponseCreatedEvent): n_llm += 1; summary_seen = False; self._emit(type="llm", state="start", n=n_llm) elif isinstance(d, ResponseCompletedEvent): self._emit(type="llm", state="end", n=n_llm, tokens=getattr(getattr(d.response, "usage", None), "output_tokens", None)) elif isinstance(d, ResponseReasoningSummaryTextDeltaEvent) and d.delta: # LiteLLM 은 같은 생각을 reasoning_content 와 reasoning 두 필드로 보내고 summary_seen = True; self._emit(type="reasoning", text=d.delta) # SDK 는 각각 summary/text 델타로 바꾼다 → 한쪽만 쓴다 (둘 다 쓰면 글자가 겹침) elif isinstance(d, ResponseReasoningTextDeltaEvent) and d.delta and not summary_seen: self._emit(type="reasoning", text=d.delta) elif self.stream and isinstance(d, ResponseTextDeltaEvent) and d.delta: self._emit(type="token", text=d.delta) elif ev.type == "run_item_stream_event": it = ev.item if it.type == "tool_call_item": raw = json.loads(getattr(it.raw_item, "arguments", None) or "{}") self._emit(type="progress", step=self.steps.get(it.raw_item.name, 0)) # D14: SSE 7종 + progress · llm · reasoning args = {k: v for k, v in raw.items() if k != "summary"} self._emit(type="action", tool=it.raw_item.name, args=args, summary=raw.get("summary", ""), cmd=self.command(it.raw_item.name, args)) # ⑤ 단계 카드의 명령줄 calls[getattr(it.raw_item, "call_id", "")] = it.raw_item.name elif it.type == "tool_call_output_item": # MCP 도구는 _run 을 안 거치므로 여기서 관측 이벤트 ri = it.raw_item; cid = ri.get("call_id") if isinstance(ri, dict) else getattr(ri, "call_id", "") name = calls.get(cid or "") if name and name not in self.own_tools: self._emit(type="observation", tool=name, data={"result": str(it.output)[:6000]}) elif it.type == "message_output_item": text = ItemHelpers.text_message_output(it) if text.strip(): self._emit(type="message", text=text) self.usage = result.context_wrapper.usage if result.interruptions: self.state = result.to_state() return dict(type="approval", actions=self.pending()) self.state = None return dict(type="status", state="finished") def _run(self, inp): try: final = asyncio.run(self._arun(inp)) except Exception as e: final = dict(type="error", text=f"{type(e).__name__}: {str(e)[:300]}") self.running = False # 끝 이벤트보다 먼저 내려야 받은 쪽이 바로 confirm/send 할 수 있다 self._emit(**final) def _start(self, inp): self.running = True self._emit(type="status", state="running") self.thread = threading.Thread(target=self._run, args=(inp,), daemon=True) self.thread.start() def busy(self) -> bool: return self.running def send(self, text: str): if self.busy(): raise RuntimeError("agent is still running") if self.state is not None: # 승인 카드를 두고 말을 걸면 = 이유 있는 거부 (프롬프트 3번) return self.confirm(False, text) self._start(text) def confirm(self, approve: bool, reason: str = ""): if self.busy(): raise RuntimeError("agent is still running") if self.state is None: raise RuntimeError("nothing to confirm") for it in self.state.get_interruptions(): if approve: self.state.approve(it) else: self.state.reject(it, rejection_message=reason or "User rejected the action") self._start(self.state)