""" The agent registry, as this engine sees it. Settings an operator changes in the console (Settings → Skills & Tools) are stored in the Doormile backend's agent registry. This module polls ``GET /api/v1/internal/ai/registry`` and answers, at the moment an agent needs it, three questions: registry.autonomous(agent_id, env_default) may this agent act on its own? registry.model(agent_id) which Claude model, if pinned? registry.skill_enabled(skill_id) is this behaviour switched on? registry.threshold(skill_id, key, default) the tuned value of a knob Precedence — deliberate, and the same for every setting: registry loaded → the registry's value (it is the operator's decision) not yet / never → the env default this engine always had So a backend that is down, unreachable, or not yet deployed leaves the engine exactly as it behaved before the registry existed. Once the registry has been read, the last good copy is kept through later failures: a flapping backend does not flip autonomy back to an env default mid-shift. Polling uses the ETag the backend sends; an unchanged registry costs a 304 with no body. Before Phase 5 of the plan (krow_talent_app/docs/ agent-platform-plan.md) every one of these settings was a module constant read from the environment at import time, so a change needed a redeploy. """ import asyncio import os from typing import Any, Dict, Optional import aiohttp from core.logger import logger POLL_SECONDS = float(os.getenv("REGISTRY_POLL_SECONDS", "30")) class RegistryClient: def __init__(self): self._agents: Dict[str, Dict[str, Any]] = {} self._skills: Dict[str, Dict[str, Any]] = {} self._etag: Optional[str] = None self.loaded = False # ── Reading (sync, cheap, safe to call on every event) ────────────────── def autonomous(self, agent_id: str, env_default: bool) -> bool: if not self.loaded or agent_id not in self._agents: return env_default return bool(self._agents[agent_id].get("autonomous", False)) def model(self, agent_id: str) -> Optional[str]: """A pinned model id, or None for the engine's own LLM_MODEL.""" if not self.loaded: return None return self._agents.get(agent_id, {}).get("model") or None def skill_enabled(self, skill_id: str, default: bool = True) -> bool: if not self.loaded or skill_id not in self._skills: return default return bool(self._skills[skill_id].get("enabled", default)) def threshold(self, skill_id: str, key: str, default: float) -> float: if not self.loaded: return default value = (self._skills.get(skill_id, {}).get("thresholds") or {}).get(key) return value if isinstance(value, (int, float)) and not isinstance(value, bool) else default # ── Loading ────────────────────────────────────────────────────────────── def apply(self, snapshot: Dict[str, Any]) -> None: """Adopt a registry snapshot (the `data` of /internal/ai/registry).""" agents = {a["agentid"]: a for a in snapshot.get("agents") or [] if a.get("agentid")} skills = {s["skillid"]: s for s in snapshot.get("skills") or [] if s.get("skillid")} self._agents, self._skills = agents, skills if not self.loaded: logger.info(f"Agent registry loaded: {len(agents)} agents, {len(skills)} skills") self.loaded = True async def fetch_once(self, session: aiohttp.ClientSession, base_url: str, api_key: str) -> str: """One poll. Returns 'updated', 'unchanged', 'unavailable' or 'refused'.""" headers = {"X-Internal-Key": api_key} if self._etag: headers["If-None-Match"] = self._etag try: async with session.get(f"{base_url}/api/v1/internal/ai/registry", headers=headers, timeout=aiohttp.ClientTimeout(total=10)) as resp: if resp.status == 304: return "unchanged" if resp.status in (401, 403): return "refused" if resp.status != 200: return "unavailable" body = await resp.json() self.apply(body.get("data") or {}) self._etag = resp.headers.get("ETag") return "updated" except Exception as exc: # network, timeout, bad JSON — keep the last good copy logger.debug(f"Agent registry poll failed: {exc}") return "unavailable" async def run(self, base_url: str, api_key: str, poll_seconds: float = POLL_SECONDS) -> None: """Poll forever. Never raises; the engine runs on env defaults meanwhile.""" if not api_key: logger.warning("INTERNAL_API_KEY unset — agent registry not read; running on env defaults") return last = None async with aiohttp.ClientSession() as session: while True: outcome = await self.fetch_once(session, base_url, api_key) if outcome != last and outcome in ("refused", "unavailable"): logger.warning( f"Agent registry {outcome} at {base_url}; " + ("keeping the last loaded copy" if self.loaded else "running on env defaults") ) last = outcome await asyncio.sleep(poll_seconds) registry = RegistryClient()