ontic-tabicl-api / agent_runtime.py
fbdeme's picture
Restore ONTIC tool agent and original approval runtime
df42e17 verified
Raw History Blame Contribute Delete
7 kB
"""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)