"""Small external agent with an SQLite outbox for disconnected execution sites."""
import json
import sqlite3
import uuid
from datetime import datetime
from pathlib import Path
from typing import Any
from urllib.error import URLError
from urllib.request import Request, urlopen
from .models import AgentRegistration, RemoteEvent
from .timeutil import iso, utc_now
class AgentOutbox:
def __init__(self, path: Path) -> None:
path.parent.mkdir(parents=True, exist_ok=True)
self.connection = sqlite3.connect(path, timeout=10.0)
self.connection.execute("PRAGMA journal_mode=WAL")
self.connection.execute("""CREATE TABLE IF NOT EXISTS agent_outbox(
id INTEGER PRIMARY KEY, event_id TEXT NOT NULL UNIQUE, payload_json TEXT NOT NULL,
created_at TEXT NOT NULL, attempts INTEGER NOT NULL DEFAULT 0, last_error TEXT)""")
self.connection.commit()
def enqueue(self, event: RemoteEvent) -> None:
payload = event_payload(event)
self.connection.execute("INSERT OR IGNORE INTO agent_outbox(event_id,payload_json,created_at) VALUES(?,?,?)",
(event.event_id, json.dumps(payload, ensure_ascii=False, default=str), iso(utc_now())))
self.connection.commit()
def pending(self, limit: int = 100) -> list[sqlite3.Row]:
self.connection.row_factory = sqlite3.Row
return list(self.connection.execute("SELECT * FROM agent_outbox ORDER BY id LIMIT ?", (max(1, min(limit, 500)),)))
def acknowledged(self, event_id: str) -> None:
self.connection.execute("DELETE FROM agent_outbox WHERE event_id=?", (event_id,))
self.connection.commit()
def failed(self, event_id: str, message: str) -> None:
self.connection.execute("UPDATE agent_outbox SET attempts=attempts+1,last_error=? WHERE event_id=?", (message, event_id))
self.connection.commit()
def size(self) -> int:
return int(self.connection.execute("SELECT COUNT(*) FROM agent_outbox").fetchone()[0])
def close(self) -> None:
self.connection.close()
def event_payload(event: RemoteEvent) -> dict[str, Any]:
return {"event_id": event.event_id, "agent": {"agent_key": event.agent.agent_key,
"agent_name": event.agent.agent_name, "agent_version": event.agent.agent_version,
"environment": event.agent.environment, "host": event.agent.host, "runtime": event.agent.runtime,
"capabilities": list(event.agent.capabilities), "metadata": event.agent.metadata},
"worker_key": event.worker_key, "job_key": event.job_key, "event_type": event.event_type,
"occurred_at": iso(event.occurred_at), "sequence": event.sequence, "payload": event.payload}
class RemoteAgentClient:
def __init__(self, endpoint: str, registration: AgentRegistration, outbox: AgentOutbox,
token: str | None = None, timeout_seconds: int = 10) -> None:
self.endpoint = endpoint.rstrip("/")
self.registration = registration
self.outbox = outbox
self.token = token
self.timeout_seconds = timeout_seconds
def register(self) -> bool:
return self._post("/api/v1/agents/register", {"agent_key": self.registration.agent_key,
"agent_name": self.registration.agent_name, "agent_version": self.registration.agent_version,
"environment": self.registration.environment, "host": self.registration.host,
"runtime": self.registration.runtime, "capabilities": list(self.registration.capabilities),
"metadata": self.registration.metadata})
def emit(self, worker_key: str, job_key: str, event_type: str, payload: dict[str, Any] | None = None,
occurred_at: datetime | None = None, sequence: int | None = None) -> bool:
event = RemoteEvent(str(uuid.uuid4()), self.registration, worker_key, job_key, event_type,
occurred_at or utc_now(), payload or {}, sequence)
self.outbox.enqueue(event)
return self.flush() > 0
def flush(self, limit: int = 100) -> int:
sent = 0
for row in self.outbox.pending(limit):
payload = json.loads(row["payload_json"])
try:
self._post("/api/v1/agents/events", payload)
except (OSError, URLError, TimeoutError, ValueError) as exc:
self.outbox.failed(row["event_id"], str(exc))
continue
self.outbox.acknowledged(row["event_id"])
sent += 1
return sent
def _post(self, path: str, payload: dict[str, Any]) -> bool:
headers = {"Content-Type": "application/json", "Accept": "application/json"}
if self.token:
headers["Authorization"] = f"Bearer {self.token}"
request = Request(self.endpoint + path, data=json.dumps(payload, default=str).encode("utf-8"),
headers=headers, method="POST")
with urlopen(request, timeout=self.timeout_seconds) as response:
if response.status < 200 or response.status >= 300:
raise OSError(f"remote watcher returned HTTP {response.status}")
return True