|
# retoor <retoor@molodetz.nl>
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
import logging
|
|
from typing import Any
|
|
|
|
import uuid_utils
|
|
|
|
from ..actions import Dispatcher
|
|
from ..agent import Agent
|
|
from ..agentic import AgenticController, LessonStore
|
|
from ..avatar import AvatarController
|
|
from ..behavior import BehaviorController
|
|
from ..chunks import ChunkStore
|
|
from ..client import ClientController
|
|
from ..config import Settings
|
|
from ..cost import CostTracker
|
|
from ..cost.tracker import Pricing
|
|
from ..http_client import PlatformClient
|
|
from ..llm import LLMClient
|
|
from ..registry import CATALOG
|
|
from ..tasks import Scheduler, TaskController, TaskStore
|
|
from ..virtual_tools import VirtualToolController, VirtualToolStore
|
|
from ._helpers import _format_offset, _now_iso, _repair_history
|
|
from .prompts import (
|
|
BEHAVIOR_HEADER,
|
|
DOCS_GREETING,
|
|
DOCS_SYSTEM_PROMPT,
|
|
DOCS_TOOLS,
|
|
LOGIN_REQUEST,
|
|
TOPIC_CLASSIFIER_PROMPT,
|
|
_parse_topic,
|
|
_system_prompt_for,
|
|
)
|
|
|
|
logger = logging.getLogger("devii.session")
|
|
|
|
AVATAR_QUERY_TIMEOUT_SECONDS = 30.0
|
|
BROWSER_REQUEST_DEADLINE_SECONDS = 60.0
|
|
RESULT_NOTICE_CHARS = 20000
|
|
INTERNAL_PREFIXES = (
|
|
"[memory]",
|
|
"[reflection-trigger]",
|
|
"[verification-gate]",
|
|
"[stopped]",
|
|
"Protocol violation:",
|
|
)
|
|
BUFFERED_TYPES = ("task", "error", "reply", "status")
|
|
|
|
|
|
class DeviiSession:
|
|
def __init__(
|
|
self,
|
|
owner_kind: str,
|
|
owner_id: str,
|
|
username: str,
|
|
settings: Settings,
|
|
llm: LLMClient,
|
|
lessons: LessonStore,
|
|
pricing: Pricing,
|
|
task_store: TaskStore,
|
|
virtual_tool_store: VirtualToolStore,
|
|
behavior_store: Any,
|
|
stores: dict[str, Any],
|
|
is_admin: bool = False,
|
|
is_primary_admin: bool = False,
|
|
channel: str = "main",
|
|
) -> None:
|
|
self.owner_kind = owner_kind
|
|
self.owner_id = owner_id
|
|
self.channel = channel
|
|
self.username = username
|
|
self.settings = settings
|
|
self.is_admin = is_admin
|
|
self.is_primary_admin = is_primary_admin
|
|
self.persist_conversation = owner_kind == "user"
|
|
self._timezone: str = ""
|
|
self._tz_offset_minutes: int | None = None
|
|
self._llm = llm
|
|
self._lessons = lessons
|
|
self.client = PlatformClient(
|
|
settings.base_url, settings.timeout_seconds, settings.platform_api_key
|
|
)
|
|
self.avatar = AvatarController()
|
|
self.browser = ClientController(
|
|
allow_eval=settings.allow_eval, allow_control=settings.browser_control
|
|
)
|
|
self.store = task_store
|
|
self.task_controller = TaskController(self.store)
|
|
self.agentic = AgenticController(lessons, settings)
|
|
self.cost = CostTracker(pricing=pricing)
|
|
self.chunks = ChunkStore()
|
|
self._virtual_tool_store = virtual_tool_store
|
|
self.virtual_tools = VirtualToolController(
|
|
virtual_tool_store, self.agentic.run_subagent, set(CATALOG.by_name())
|
|
)
|
|
self._behavior_store = behavior_store
|
|
self.behavior = BehaviorController(behavior_store)
|
|
self.dispatcher = Dispatcher(
|
|
CATALOG,
|
|
self.client,
|
|
settings,
|
|
self.task_controller,
|
|
self.agentic,
|
|
avatar=self.avatar,
|
|
browser=self.browser,
|
|
is_admin=is_admin,
|
|
is_primary_admin=is_primary_admin,
|
|
quota_provider=self._quota_snapshot,
|
|
owner_kind=owner_kind,
|
|
owner_id=owner_id,
|
|
virtual_tools=self.virtual_tools,
|
|
behavior=self.behavior,
|
|
)
|
|
self.tools = self._builtin_tools()
|
|
self._system_prompt = (
|
|
DOCS_SYSTEM_PROMPT if channel == "docs" else _system_prompt_for(is_admin)
|
|
)
|
|
self.agentic.bind(
|
|
llm=llm,
|
|
dispatcher=self.dispatcher,
|
|
tools=self.tools,
|
|
on_trace=self._trace,
|
|
cost_tracker=self.cost,
|
|
chunk_store=self.chunks,
|
|
)
|
|
self.agent = Agent(
|
|
settings,
|
|
llm,
|
|
self.dispatcher,
|
|
self.tools,
|
|
lessons=lessons,
|
|
on_trace=self._trace,
|
|
cost_tracker=self.cost,
|
|
chunk_store=self.chunks,
|
|
system_prompt=self._compose_system_prompt(),
|
|
)
|
|
self.scheduler = Scheduler(
|
|
self.store,
|
|
self._make_executor(),
|
|
self._task_event,
|
|
settings.scheduler_tick_seconds,
|
|
)
|
|
self._conv = stores["conversations"]
|
|
self._ledger = stores["ledger"]
|
|
self._audit = stores["turns"]
|
|
self._conns: set[Any] = set()
|
|
self._conn_meta: dict[Any, dict[str, Any]] = {}
|
|
self._attach_seq = 0
|
|
self._lock = asyncio.Lock()
|
|
self._send_lock = asyncio.Lock()
|
|
self._pending: dict[str, asyncio.Future] = {}
|
|
self._query_seq = 0
|
|
self._connected = asyncio.Event()
|
|
self._disconnected = asyncio.Event()
|
|
self._disconnected.set()
|
|
self._started = False
|
|
self._buffer: list[dict[str, Any]] = []
|
|
self._turn_tool_calls = 0
|
|
self._pending_session: str | None = None
|
|
self._turns: set[asyncio.Task] = set()
|
|
self._turn_epoch = 0
|
|
|
|
def restore_history(self, messages: list[dict[str, Any]]) -> None:
|
|
if messages:
|
|
self.agent._messages = messages
|
|
|
|
def history(self) -> list[dict[str, Any]]:
|
|
visible: list[dict[str, Any]] = []
|
|
for message in self.agent._messages:
|
|
role = message.get("role")
|
|
content = message.get("content")
|
|
if role not in ("user", "assistant") or not content:
|
|
continue
|
|
if (
|
|
role == "user"
|
|
and isinstance(content, str)
|
|
and content.startswith(INTERNAL_PREFIXES)
|
|
):
|
|
continue
|
|
visible.append({"role": role, "content": content})
|
|
return visible
|
|
|
|
def _make_executor(self) -> Any:
|
|
async def execute(prompt: str) -> str:
|
|
async with self._lock:
|
|
turn_id = uuid_utils.uuid7().hex
|
|
started_at = _now_iso()
|
|
self._turn_tool_calls = 0
|
|
before = self._cost_snapshot()
|
|
reply = ""
|
|
error = ""
|
|
worker = Agent(
|
|
self.settings,
|
|
self._llm,
|
|
self.dispatcher,
|
|
self.tools,
|
|
lessons=self._lessons,
|
|
on_trace=self._trace,
|
|
cost_tracker=self.cost,
|
|
chunk_store=self.chunks,
|
|
system_prompt=self._compose_system_prompt(),
|
|
)
|
|
try:
|
|
reply = await worker.respond(prompt)
|
|
return reply
|
|
except Exception as exc: # noqa: BLE001 - recorded, then re-raised to the scheduler
|
|
error = str(exc)
|
|
raise
|
|
finally:
|
|
self._record_spend(
|
|
turn_id, started_at, prompt, reply, error, before
|
|
)
|
|
|
|
return execute
|
|
|
|
def attach(self, ws: Any) -> None:
|
|
self._conns.add(ws)
|
|
self._attach_seq += 1
|
|
self._conn_meta[ws] = {
|
|
"seq": self._attach_seq,
|
|
"visible": True,
|
|
"focused": False,
|
|
}
|
|
self._connected.set()
|
|
self._disconnected.clear()
|
|
self.avatar.bind(self._avatar_request)
|
|
self.browser.bind(self._client_request)
|
|
self.ensure_scheduler_started()
|
|
if self._buffer:
|
|
pending = self._buffer
|
|
self._buffer = []
|
|
asyncio.create_task(self._flush(pending))
|
|
logger.info(
|
|
"Session %s/%s attached (%d conns)",
|
|
self.owner_kind,
|
|
self.owner_id,
|
|
len(self._conns),
|
|
)
|
|
|
|
def detach(self, ws: Any) -> None:
|
|
self._conns.discard(ws)
|
|
self._conn_meta.pop(ws, None)
|
|
if not self._conns:
|
|
self._connected.clear()
|
|
self._disconnected.set()
|
|
self.avatar.unbind()
|
|
self.browser.unbind()
|
|
logger.info(
|
|
"Session %s/%s detached (%d conns)",
|
|
self.owner_kind,
|
|
self.owner_id,
|
|
len(self._conns),
|
|
)
|
|
|
|
def ensure_scheduler_started(self) -> None:
|
|
if not self._started and self.channel == "main":
|
|
self.scheduler.start()
|
|
self._started = True
|
|
|
|
def has_pending_tasks(self) -> bool:
|
|
if self.channel != "main":
|
|
return False
|
|
try:
|
|
return self.store.has_pending()
|
|
except Exception: # noqa: BLE001 - never block GC on a store read
|
|
return False
|
|
|
|
def set_clientinfo(self, timezone_name: str, offset_minutes: int | None) -> None:
|
|
if timezone_name:
|
|
self._timezone = timezone_name
|
|
if offset_minutes is not None:
|
|
self._tz_offset_minutes = int(offset_minutes)
|
|
if self.owner_kind == "user" and timezone_name:
|
|
from devplacepy.database import set_user_timezone
|
|
|
|
try:
|
|
set_user_timezone(self.owner_id, timezone_name)
|
|
except Exception: # noqa: BLE001 - persistence must not break the socket
|
|
logger.exception(
|
|
"Failed to persist timezone for %s/%s",
|
|
self.owner_kind,
|
|
self.owner_id,
|
|
)
|
|
|
|
def set_timezone(self, timezone_name: str) -> None:
|
|
if timezone_name:
|
|
self._timezone = timezone_name
|
|
|
|
def set_visibility(self, ws: Any, visible: bool, focused: bool) -> None:
|
|
meta = self._conn_meta.get(ws)
|
|
if meta is not None:
|
|
meta["visible"] = visible
|
|
meta["focused"] = focused
|
|
|
|
def _pick_target(self) -> Any:
|
|
best = None
|
|
highest_rank = None
|
|
for ws in self._conns:
|
|
meta = self._conn_meta.get(ws)
|
|
if meta is None:
|
|
continue
|
|
rank = (meta["focused"], meta["visible"], meta["seq"])
|
|
if highest_rank is None or rank > highest_rank:
|
|
highest_rank = rank
|
|
best = ws
|
|
return best
|
|
|
|
@property
|
|
def connection_count(self) -> int:
|
|
return len(self._conns)
|
|
|
|
async def aclose(self) -> None:
|
|
await self.scheduler.stop()
|
|
await self.client.aclose()
|
|
await self._llm.aclose()
|
|
|
|
async def bootstrap_greeting(self) -> str:
|
|
if self.channel == "docs":
|
|
return DOCS_GREETING
|
|
if self.client.authenticated:
|
|
return f"Signed in as {self.username}. Devii is operating your DevPlace account. How can I help?"
|
|
return LOGIN_REQUEST
|
|
|
|
async def send_bootstrap(self, ws: Any) -> None:
|
|
await self._send_to(ws, {"type": "status", "text": "Devii ready"})
|
|
history = self.history()
|
|
if history:
|
|
await self._send_to(ws, {"type": "history", "messages": history})
|
|
else:
|
|
await self._send_to(
|
|
ws, {"type": "reply", "text": await self.bootstrap_greeting()}
|
|
)
|
|
|
|
async def reset_conversation(self) -> None:
|
|
async with self._lock:
|
|
messages = self.agent._messages
|
|
if messages and messages[0].get("role") == "system":
|
|
system = messages[0]
|
|
else:
|
|
system = {"role": "system", "content": self._system_prompt}
|
|
self.agent._messages = [system]
|
|
self._buffer = []
|
|
if self.persist_conversation:
|
|
try:
|
|
self._conv.clear(self.owner_kind, self.owner_id, self.channel)
|
|
except Exception: # noqa: BLE001 - clearing storage must not break the socket
|
|
logger.exception(
|
|
"Failed to clear conversation for %s/%s",
|
|
self.owner_kind,
|
|
self.owner_id,
|
|
)
|
|
await self._emit({"type": "clear"}, buffer=False)
|
|
|
|
def spawn_turn(self, content: Any, audit_text: str | None = None) -> None:
|
|
epoch = self._turn_epoch
|
|
task = asyncio.create_task(self._run_turn(content, epoch, audit_text))
|
|
self._turns.add(task)
|
|
task.add_done_callback(self._turns.discard)
|
|
|
|
async def cancel_turns(self) -> None:
|
|
self._turn_epoch += 1
|
|
turns = [task for task in self._turns if not task.done()]
|
|
if not turns:
|
|
return
|
|
for task in turns:
|
|
task.cancel()
|
|
try:
|
|
await asyncio.wait_for(
|
|
asyncio.gather(*turns, return_exceptions=True), timeout=10
|
|
)
|
|
except asyncio.TimeoutError:
|
|
logger.warning(
|
|
"cancel_turns: a turn did not unwind within 10s for %s/%s (epoch-guarded)",
|
|
self.owner_kind,
|
|
self.owner_id,
|
|
)
|
|
|
|
async def reset(self) -> None:
|
|
await self.cancel_turns()
|
|
await self.reset_conversation()
|
|
await self._emit(
|
|
{"type": "status", "text": "Reset. Conversation cleared."}, buffer=False
|
|
)
|
|
|
|
async def stop(self) -> None:
|
|
await self.cancel_turns()
|
|
async with self._lock:
|
|
_repair_history(self.agent._messages)
|
|
if self.persist_conversation:
|
|
try:
|
|
self._conv.save(
|
|
self.owner_kind,
|
|
self.owner_id,
|
|
self.agent._messages,
|
|
self.channel,
|
|
)
|
|
except Exception: # noqa: BLE001 - persistence must not break the socket
|
|
logger.exception(
|
|
"Failed to persist history after stop for %s/%s",
|
|
self.owner_kind,
|
|
self.owner_id,
|
|
)
|
|
await self._emit({"type": "status", "text": "Stopped."}, buffer=False)
|
|
|
|
async def _run_turn(
|
|
self, content: Any, epoch: int, audit_text: str | None = None
|
|
) -> None:
|
|
turn_id = uuid_utils.uuid7().hex
|
|
started_at = _now_iso()
|
|
before = self._cost_snapshot()
|
|
self._turn_tool_calls = 0
|
|
reply = ""
|
|
error = ""
|
|
cancelled = False
|
|
prompt_text = audit_text if audit_text is not None else (
|
|
content if isinstance(content, str) else "[image]"
|
|
)
|
|
try:
|
|
async with self._lock:
|
|
self._refresh_tools()
|
|
self._refresh_system_prompt()
|
|
if self.channel == "docs":
|
|
await self._docs_topic_gate(prompt_text)
|
|
before = self._cost_snapshot()
|
|
reply = await self.agent.respond(content)
|
|
if not isinstance(content, str):
|
|
self._redact_last_user_message(prompt_text)
|
|
if epoch == self._turn_epoch:
|
|
await self._emit({"type": "reply", "text": reply}, buffer=True)
|
|
except asyncio.CancelledError:
|
|
cancelled = True
|
|
raise
|
|
except Exception as exc: # noqa: BLE001 - reported to the browser, recorded for audit
|
|
error = str(exc)
|
|
logger.exception("Turn failed for %s/%s", self.owner_kind, self.owner_id)
|
|
if epoch == self._turn_epoch:
|
|
await self._emit({"type": "error", "text": error}, buffer=True)
|
|
finally:
|
|
self._record_spend(turn_id, started_at, prompt_text, reply, error, before)
|
|
if not cancelled and epoch == self._turn_epoch:
|
|
self._persist_history()
|
|
if self.owner_kind == "user" and self.channel != "docs" and not error:
|
|
from devplacepy.utils import track_action
|
|
|
|
track_action(self.owner_id, "devii")
|
|
|
|
def _builtin_tools(self) -> list[dict[str, Any]]:
|
|
schemas = CATALOG.tool_schemas_for(
|
|
self.client.authenticated, self.is_admin, self.is_primary_admin
|
|
)
|
|
if self.channel == "docs":
|
|
return [
|
|
s
|
|
for s in schemas
|
|
if s.get("function", {}).get("name") in DOCS_TOOLS
|
|
]
|
|
return schemas
|
|
|
|
def _refresh_tools(self) -> None:
|
|
if self.channel == "docs":
|
|
self.tools[:] = self._builtin_tools()
|
|
return
|
|
virtual = self._virtual_tool_store.tool_schemas()
|
|
self.tools[:] = self._builtin_tools() + virtual
|
|
|
|
def _compose_system_prompt(self) -> str:
|
|
if self.channel == "docs":
|
|
return self._system_prompt
|
|
body = self._behavior_store.text().strip()
|
|
section = BEHAVIOR_HEADER if not body else f"{BEHAVIOR_HEADER}\n{body}"
|
|
return f"{self._system_prompt}\n\n{self._clock_line()}\n\n{section}"
|
|
|
|
def _clock_line(self) -> str:
|
|
from datetime import datetime, timedelta, timezone
|
|
|
|
now = datetime.now(timezone.utc).replace(microsecond=0)
|
|
tz_name = self._timezone or self._stored_timezone()
|
|
local = ""
|
|
if tz_name:
|
|
try:
|
|
from zoneinfo import ZoneInfo
|
|
|
|
local_now = now.astimezone(ZoneInfo(tz_name))
|
|
offset = local_now.utcoffset()
|
|
offset_text = _format_offset(offset)
|
|
local = (
|
|
f" The user's local time is {local_now.strftime('%Y-%m-%d %H:%M:%S')} "
|
|
f"({tz_name}, UTC{offset_text})."
|
|
)
|
|
except Exception: # noqa: BLE001 - unknown tz name: fall back to UTC only
|
|
local = ""
|
|
if not local and self._tz_offset_minutes is not None:
|
|
offset = timedelta(minutes=self._tz_offset_minutes)
|
|
local = (
|
|
f" The user's local time is "
|
|
f"{(now + offset).strftime('%Y-%m-%d %H:%M:%S')} "
|
|
f"(UTC{_format_offset(offset)})."
|
|
)
|
|
return (
|
|
f"# CURRENT TIME\n"
|
|
f"The current UTC time is {now.strftime('%Y-%m-%dT%H:%M:%S')}Z.{local} "
|
|
"When the user gives a wall-clock time (for example '3pm' or 'tomorrow at 09:00'), "
|
|
"interpret it in the user's local timezone and convert it to UTC for the run_at field. "
|
|
"For a relative request (for example 'in 40 seconds' or 'in 2 hours'), use delay_seconds "
|
|
"instead and do not compute an absolute time."
|
|
)
|
|
|
|
def _stored_timezone(self) -> str:
|
|
if self.owner_kind != "user":
|
|
return ""
|
|
try:
|
|
from devplacepy.database import db
|
|
|
|
if "users" not in db.tables:
|
|
return ""
|
|
row = db["users"].find_one(uid=self.owner_id)
|
|
return (row or {}).get("timezone") or ""
|
|
except Exception: # noqa: BLE001 - clock line is best-effort
|
|
return ""
|
|
|
|
def _refresh_system_prompt(self) -> None:
|
|
messages = self.agent._messages
|
|
if messages and messages[0].get("role") == "system":
|
|
messages[0]["content"] = self._compose_system_prompt()
|
|
|
|
async def _docs_topic_gate(self, text: str) -> None:
|
|
prior = [
|
|
m
|
|
for m in self.agent._messages
|
|
if m.get("role") in ("user", "assistant") and m.get("content")
|
|
]
|
|
if not prior:
|
|
return
|
|
decision, reason = await self._classify_topic(prior, text)
|
|
if decision == "new":
|
|
messages = self.agent._messages
|
|
if messages and messages[0].get("role") == "system":
|
|
system = messages[0]
|
|
else:
|
|
system = {"role": "system", "content": self._compose_system_prompt()}
|
|
self.agent._messages = [system]
|
|
self._buffer = []
|
|
if self.persist_conversation:
|
|
try:
|
|
self._conv.clear(self.owner_kind, self.owner_id, self.channel)
|
|
except Exception: # noqa: BLE001 - clearing storage must not break the turn
|
|
logger.exception(
|
|
"Failed to clear docs conversation for %s/%s",
|
|
self.owner_kind,
|
|
self.owner_id,
|
|
)
|
|
await self._emit(
|
|
{"type": "topic", "decision": decision, "reason": reason}, buffer=False
|
|
)
|
|
|
|
async def _classify_topic(
|
|
self, prior: list[dict[str, Any]], text: str
|
|
) -> tuple[str, str]:
|
|
recent = prior[-6:]
|
|
convo = "\n".join(
|
|
f"{m['role']}: {str(m['content'])[:300]}" for m in recent
|
|
)
|
|
messages = [
|
|
{"role": "system", "content": TOPIC_CLASSIFIER_PROMPT},
|
|
{
|
|
"role": "user",
|
|
"content": f"Prior conversation:\n{convo}\n\nNew message:\n{text}\n\nClassify.",
|
|
},
|
|
]
|
|
try:
|
|
raw = await self._llm.complete_text(messages)
|
|
except Exception: # noqa: BLE001 - on any failure, keep context (treat as follow-up)
|
|
return "follow_up", ""
|
|
return _parse_topic(raw)
|
|
|
|
def _quota_snapshot(self) -> dict[str, Any]:
|
|
spent = self._ledger.spent_24h(self.owner_kind, self.owner_id)
|
|
turns = self._ledger.turns_24h(self.owner_kind, self.owner_id)
|
|
limit = self.settings.daily_limit_usd
|
|
used_pct = round(min(100.0, spent / limit * 100), 1) if limit > 0 else 0.0
|
|
return {
|
|
"used_pct": used_pct,
|
|
"turns_today": turns,
|
|
"limit_reached": limit > 0 and spent >= limit,
|
|
}
|
|
|
|
def _cost_snapshot(self) -> dict[str, Any]:
|
|
return {
|
|
"requests": self.cost.requests,
|
|
"prompt_tokens": self.cost.prompt_tokens,
|
|
"completion_tokens": self.cost.completion_tokens,
|
|
"cache_hit_tokens": self.cost.cache_hit_tokens,
|
|
"cache_miss_tokens": self.cost.cache_miss_tokens,
|
|
"cost_usd": self.cost.cost_usd()["total"],
|
|
}
|
|
|
|
def _persist_history(self) -> None:
|
|
if not self.persist_conversation:
|
|
return
|
|
try:
|
|
self._conv.save(
|
|
self.owner_kind,
|
|
self.owner_id,
|
|
self.agent._messages,
|
|
self.channel,
|
|
)
|
|
except Exception: # noqa: BLE001 - persistence must never break a turn
|
|
logger.exception(
|
|
"Failed to persist conversation for %s/%s",
|
|
self.owner_kind,
|
|
self.owner_id,
|
|
)
|
|
|
|
def _redact_last_user_message(self, text: str) -> None:
|
|
messages = self.agent._messages
|
|
for message in reversed(messages):
|
|
if message.get("role") == "user":
|
|
message["content"] = text
|
|
return
|
|
|
|
def _record_spend(self, turn_id, started_at, prompt, reply, error, before) -> None:
|
|
after = self._cost_snapshot()
|
|
usage = {
|
|
"prompt_tokens": after["prompt_tokens"] - before["prompt_tokens"],
|
|
"completion_tokens": after["completion_tokens"]
|
|
- before["completion_tokens"],
|
|
"cache_hit_tokens": after["cache_hit_tokens"] - before["cache_hit_tokens"],
|
|
"cache_miss_tokens": after["cache_miss_tokens"]
|
|
- before["cache_miss_tokens"],
|
|
}
|
|
cost_delta = round(after["cost_usd"] - before["cost_usd"], 8)
|
|
if (
|
|
cost_delta <= 0
|
|
and usage["prompt_tokens"] <= 0
|
|
and usage["completion_tokens"] <= 0
|
|
):
|
|
return
|
|
try:
|
|
self._ledger.record(
|
|
self.owner_kind,
|
|
self.owner_id,
|
|
turn_id,
|
|
usage,
|
|
cost_delta,
|
|
self.settings.ai_model,
|
|
)
|
|
self._audit.record(
|
|
{
|
|
"turn_id": turn_id,
|
|
"owner_kind": self.owner_kind,
|
|
"owner_id": self.owner_id,
|
|
"username": self.username,
|
|
"started_at": started_at,
|
|
"ended_at": _now_iso(),
|
|
"prompt": prompt,
|
|
"reply": reply,
|
|
"iterations": after["requests"] - before["requests"],
|
|
"tool_calls": self._turn_tool_calls,
|
|
"cost_usd": cost_delta,
|
|
"error": error,
|
|
}
|
|
)
|
|
from devplacepy.services.audit import record as audit
|
|
|
|
audit.record_system(
|
|
"devii.turn",
|
|
actor_kind="user" if self.owner_kind == "user" else self.owner_kind,
|
|
actor_uid=self.owner_id if self.owner_kind == "user" else None,
|
|
actor_username=self.username,
|
|
actor_role="user" if self.owner_kind == "user" else self.owner_kind,
|
|
origin="devii",
|
|
via_agent=1,
|
|
result="failure" if error else "success",
|
|
summary=f"Devii processed a turn for {self.username or self.owner_id}",
|
|
metadata={
|
|
"iterations": after["requests"] - before["requests"],
|
|
"tool_calls": self._turn_tool_calls,
|
|
"cost_usd": cost_delta,
|
|
"prompt_tokens": usage["prompt_tokens"],
|
|
"completion_tokens": usage["completion_tokens"],
|
|
},
|
|
)
|
|
except Exception: # noqa: BLE001 - persistence must never break a turn
|
|
logger.exception(
|
|
"Failed to persist turn for %s/%s", self.owner_kind, self.owner_id
|
|
)
|
|
|
|
def resolve_query(self, query_id: str, payload: Any) -> None:
|
|
future = self._pending.pop(query_id, None)
|
|
if future is not None and not future.done():
|
|
future.set_result(payload)
|
|
|
|
async def _flush(self, frames: list[dict[str, Any]]) -> None:
|
|
for frame in frames:
|
|
await self._emit(frame, buffer=True)
|
|
|
|
def _trace(self, event: str, name: str = "", detail: str = "") -> None:
|
|
if event == "call":
|
|
self._turn_tool_calls += 1
|
|
elif event == "ok":
|
|
self._sync_auth(name)
|
|
payload = {"type": "trace", "event": event, "name": name, "detail": detail}
|
|
asyncio.create_task(self._emit(payload, buffer=False))
|
|
|
|
def _sync_auth(self, name: str) -> None:
|
|
# When the agent changes who is logged in, propagate it to the browser so the terminal and
|
|
# the browser share one session. Login mints a real platform session in the client's cookie
|
|
# jar; hand that token to the browser via /devii/adopt. Logout reuses /auth/logout.
|
|
if name in ("login", "signup"):
|
|
token = self.client.session_cookie()
|
|
if token:
|
|
self._pending_session = token
|
|
asyncio.create_task(
|
|
self._emit({"type": "auth", "action": "adopt"}, buffer=False)
|
|
)
|
|
elif name == "logout":
|
|
self._pending_session = None
|
|
asyncio.create_task(
|
|
self._emit({"type": "auth", "action": "logout"}, buffer=False)
|
|
)
|
|
|
|
def take_pending_session(self) -> str | None:
|
|
token = self._pending_session
|
|
self._pending_session = None
|
|
return token
|
|
|
|
async def _send_to(self, ws: Any, payload: dict[str, Any]) -> None:
|
|
async with self._send_lock:
|
|
await ws.send_json(payload)
|
|
|
|
async def _emit(self, payload: dict[str, Any], buffer: bool = True) -> None:
|
|
delivered = False
|
|
for ws in list(self._conns):
|
|
try:
|
|
async with self._send_lock:
|
|
await ws.send_json(payload)
|
|
delivered = True
|
|
except Exception: # noqa: BLE001 - drop dead sockets, keep serving the rest
|
|
self._conns.discard(ws)
|
|
if not delivered and buffer and payload.get("type") in BUFFERED_TYPES:
|
|
self._buffer.append(payload)
|
|
if len(self._buffer) > 100:
|
|
self._buffer = self._buffer[-100:]
|
|
|
|
async def _avatar_request(self, action: str, args: dict[str, Any]) -> Any:
|
|
return await self._browser_request("avatar", action, args)
|
|
|
|
async def _client_request(self, action: str, args: dict[str, Any]) -> Any:
|
|
return await self._browser_request("client", action, args)
|
|
|
|
async def _browser_request(
|
|
self, channel: str, action: str, args: dict[str, Any]
|
|
) -> Any:
|
|
loop = asyncio.get_event_loop()
|
|
deadline = loop.time() + BROWSER_REQUEST_DEADLINE_SECONDS
|
|
self._query_seq += 1
|
|
request_id = f"{self.owner_id}:{self._query_seq}"
|
|
future: asyncio.Future = loop.create_future()
|
|
self._pending[request_id] = future
|
|
frame = {"type": channel, "id": request_id, "action": action, "args": args}
|
|
try:
|
|
while not future.done():
|
|
remaining = deadline - loop.time()
|
|
if remaining <= 0:
|
|
raise RuntimeError(f"Browser did not respond to '{action}' in time")
|
|
try:
|
|
await asyncio.wait_for(self._connected.wait(), timeout=remaining)
|
|
except asyncio.TimeoutError as exc:
|
|
raise RuntimeError(
|
|
f"No browser connected to handle '{action}'"
|
|
) from exc
|
|
ws = self._pick_target()
|
|
if ws is None:
|
|
continue
|
|
try:
|
|
await self._send_to(ws, frame)
|
|
except Exception: # noqa: BLE001 - socket dying mid-send; retry on reconnect
|
|
await asyncio.sleep(0.1)
|
|
continue
|
|
drop = asyncio.create_task(self._disconnected.wait())
|
|
try:
|
|
await asyncio.wait(
|
|
{future, drop},
|
|
timeout=max(0.0, deadline - loop.time()),
|
|
return_when=asyncio.FIRST_COMPLETED,
|
|
)
|
|
finally:
|
|
drop.cancel()
|
|
# future done -> loop exits and returns; disconnected -> re-send after
|
|
# reconnect; otherwise the deadline check at the top raises.
|
|
return future.result()
|
|
finally:
|
|
self._pending.pop(request_id, None)
|
|
|
|
def _task_event(self, kind: str, row: dict[str, Any], payload: str) -> None:
|
|
label = row.get("label") or row.get("uid", "task")
|
|
message = {
|
|
"type": "task",
|
|
"kind": kind,
|
|
"label": label,
|
|
"payload": payload[:RESULT_NOTICE_CHARS],
|
|
}
|
|
asyncio.create_task(self._emit(message, buffer=True))
|
|
if kind in ("done", "finished") and row.get("notify"):
|
|
self._deliver_reminder(row, payload)
|
|
|
|
def _deliver_reminder(self, row: dict[str, Any], payload: str) -> None:
|
|
if self.owner_kind != "user" or not self.owner_id:
|
|
return
|
|
text = (payload or "").strip() or (
|
|
row.get("label") or "Your scheduled reminder fired."
|
|
)
|
|
try:
|
|
from devplacepy.utils import create_notification
|
|
|
|
create_notification(
|
|
self.owner_id,
|
|
"reminder",
|
|
text[:500],
|
|
row.get("uid", ""),
|
|
target_url="/devii",
|
|
)
|
|
except Exception: # noqa: BLE001 - a reminder notification must never crash the scheduler
|
|
logger.exception(
|
|
"Failed to deliver reminder notification for %s/%s",
|
|
self.owner_kind,
|
|
self.owner_id,
|
|
)
|