"""Declarative worker registry and configuration schema (#798, epic #797). The registry is the single source of truth for the scheduled multi-LLM worker fleet. It is a versioned JSON document holding two *separate* entity kinds: * **Providers** — the LLM runtimes a worker can be built on (Claude, Grok, Codex, AGY, Kimi K). A provider describes the runtime itself: vendor, executable name, models it can serve, and whether it is available on this machine. Providers exist whether or not any worker uses them. * **Workers** — a configured *instance*: one provider, one model, one project, one role, one MCP namespace/profile, one workflow, one schedule. Several workers may share a provider; a worker naming an undeclared provider is refused. Keeping them separate is what lets #799 list all five providers even when a provider currently has no configured worker, and it stops provider facts from being copied into (and drifting across) every worker record. Scope boundary. This module owns the data model, its validation, and its persistence. It does **not** schedule anything, launch anything, probe provider executables, or serve HTTP. Loading a registry never touches a process; the live fields a dashboard wants (PID, elapsed time, next run) are derived elsewhere (#799, #801, #803, #804) from these declarations. Safety invariants: * No credential may be stored (:mod:`webui.registry_safety`), so the registry stays safe to render and to hand to a browser layer. * Validation fails closed. Unknown fields are refused rather than ignored, so a typo cannot silently disable a timeout or a role binding. * Writes are atomic and every superseded document is retained as a numbered revision, so a bad edit is recoverable by rollback rather than hand-repair. """ from __future__ import annotations import json import os import re import tempfile from dataclasses import dataclass from datetime import datetime, timezone from pathlib import Path from typing import Any from webui.registry_safety import reject_credential_keys SCHEMA_VERSION = 1 #: Roles a worker may hold. These mirror the sanctioned MCP role kinds; a #: worker may not invent one, because the role selects the namespace/profile #: whose capability gates constrain it. ALLOWED_ROLES = ("author", "reviewer", "merger", "reconciler", "cleanup") #: Scheduler backends the registry can describe. ``manual`` means the worker is #: only ever started on request and has no recurring trigger. ALLOWED_SCHEDULER_KINDS = ("launchd", "manual") #: Schedule kinds. Next-run computation belongs to #803; this module only #: guarantees the declaration is well formed. ALLOWED_SCHEDULE_KINDS = ("interval", "cron", "manual") _REQUIRED_PROVIDER_FIELDS = ("id", "display_name", "vendor", "executable", "available") _OPTIONAL_PROVIDER_FIELDS = ("models", "notes") _REQUIRED_WORKER_FIELDS = ( "id", "display_name", "provider", "model", "project", "role", "namespace", "profile", "workflow", "schedule", "timeout_seconds", "enabled", "scheduler", ) _OPTIONAL_WORKER_FIELDS = ("notes",) _ID_RE = re.compile(r"^[a-z0-9][a-z0-9._-]*$") #: Guards against an operator writing a timeout that would let a worker hold a #: lease effectively forever. 24h is far above any sanctioned cycle. _MAX_TIMEOUT_SECONDS = 86_400 #: How many superseded revisions to retain beside the live file. _HISTORY_LIMIT = 20 _TOP_LEVEL_FIELDS = frozenset({"version", "revision", "updated_at", "providers", "workers"}) class RegistryValidationError(ValueError): """Raised when a registry document violates the schema.""" @dataclass(frozen=True) class ProviderRecord: id: str display_name: str vendor: str executable: str available: bool models: tuple[str, ...] notes: str @dataclass(frozen=True) class ScheduleSpec: kind: str #: Set for ``interval`` schedules. seconds: int | None #: Set for ``cron`` schedules — a five-field crontab expression. expression: str | None @dataclass(frozen=True) class SchedulerSpec: kind: str #: LaunchAgent label; required for ``launchd``, absent for ``manual``. label: str | None @dataclass(frozen=True) class WorkerRecord: id: str display_name: str provider: str model: str project: str role: str namespace: str profile: str workflow: str schedule: ScheduleSpec timeout_seconds: int enabled: bool scheduler: SchedulerSpec notes: str @dataclass(frozen=True) class WorkerRegistry: version: int revision: int updated_at: str providers: tuple[ProviderRecord, ...] workers: tuple[WorkerRecord, ...] source_path: Path # ── paths ──────────────────────────────────────────────────────────────────── def default_registry_path() -> Path: """Location of the packaged worker registry, overridable for tests/deploys.""" override = os.environ.get("WEBUI_WORKER_REGISTRY", "").strip() if override: return Path(override).expanduser().resolve() return (Path(__file__).resolve().parent / "data" / "workers.registry.json").resolve() def history_dir(path: Path | None = None) -> Path: """Directory holding superseded revisions of *path*.""" source = (path or default_registry_path()).resolve() return source.parent / f"{source.name}.history" # ── field helpers ──────────────────────────────────────────────────────────── def _require_exact_fields( raw: Any, *, required: tuple[str, ...], optional: tuple[str, ...], subject: str, ) -> dict[str, Any]: if not isinstance(raw, dict): raise RegistryValidationError(f"{subject} must be an object") missing = [field for field in required if field not in raw] if missing: raise RegistryValidationError( f"{subject} missing required fields: {', '.join(sorted(missing))}" ) unknown = sorted(set(raw) - set(required) - set(optional)) if unknown: # Fail closed: silently dropping an unrecognized key is how a typo'd # "timeout_second" ends up meaning "no timeout". raise RegistryValidationError(f"{subject} has unknown fields: {', '.join(unknown)}") return raw def _require_identifier(value: Any, *, subject: str) -> str: text = str(value).strip() if not _ID_RE.match(text): raise RegistryValidationError( f"{subject} must be lowercase alphanumeric with '.', '_', or '-' (got {value!r})" ) return text def _require_text(value: Any, *, subject: str) -> str: if not isinstance(value, str): raise RegistryValidationError(f"{subject} must be a string (got {value!r})") text = value.strip() if not text: raise RegistryValidationError(f"{subject} must be a non-empty string") return text def _require_bool(value: Any, *, subject: str) -> bool: if not isinstance(value, bool): raise RegistryValidationError(f"{subject} must be a boolean (got {value!r})") return value def _require_positive_int(value: Any, *, subject: str, maximum: int | None = None) -> int: if isinstance(value, bool) or not isinstance(value, int): raise RegistryValidationError(f"{subject} must be an integer (got {value!r})") if value <= 0: raise RegistryValidationError(f"{subject} must be greater than zero (got {value})") if maximum is not None and value > maximum: raise RegistryValidationError(f"{subject} must not exceed {maximum} (got {value})") return value # ── parsing ────────────────────────────────────────────────────────────────── def _parse_provider(raw: Any) -> ProviderRecord: data = _require_exact_fields( raw, required=_REQUIRED_PROVIDER_FIELDS, optional=_OPTIONAL_PROVIDER_FIELDS, subject="provider", ) provider_id = _require_identifier(data["id"], subject="provider.id") models_raw = data.get("models") or [] if not isinstance(models_raw, list): raise RegistryValidationError(f"provider[{provider_id}].models must be an array") models = tuple( _require_text(item, subject=f"provider[{provider_id}].models[]") for item in models_raw ) return ProviderRecord( id=provider_id, display_name=_require_text( data["display_name"], subject=f"provider[{provider_id}].display_name" ), vendor=_require_text(data["vendor"], subject=f"provider[{provider_id}].vendor"), executable=_require_text(data["executable"], subject=f"provider[{provider_id}].executable"), available=_require_bool(data["available"], subject=f"provider[{provider_id}].available"), models=models, notes=str(data.get("notes") or "").strip(), ) def _parse_schedule(raw: Any, *, subject: str) -> ScheduleSpec: if not isinstance(raw, dict): raise RegistryValidationError(f"{subject} must be an object") kind = _require_text(raw.get("kind"), subject=f"{subject}.kind") if kind not in ALLOWED_SCHEDULE_KINDS: raise RegistryValidationError( f"{subject}.kind must be one of {', '.join(ALLOWED_SCHEDULE_KINDS)} (got {kind!r})" ) seconds: int | None = None expression: str | None = None if kind == "interval": if "seconds" not in raw: raise RegistryValidationError(f"{subject}.seconds is required for interval schedules") seconds = _require_positive_int(raw["seconds"], subject=f"{subject}.seconds") elif kind == "cron": if "expression" not in raw: raise RegistryValidationError(f"{subject}.expression is required for cron schedules") expression = _require_text(raw["expression"], subject=f"{subject}.expression") if len(expression.split()) != 5: raise RegistryValidationError( f"{subject}.expression must have five crontab fields (got {expression!r})" ) allowed = {"kind"} if kind == "interval": allowed.add("seconds") elif kind == "cron": allowed.add("expression") unknown = sorted(set(raw) - allowed) if unknown: raise RegistryValidationError( f"{subject} has fields not valid for kind {kind!r}: {', '.join(unknown)}" ) return ScheduleSpec(kind=kind, seconds=seconds, expression=expression) def _parse_scheduler(raw: Any, *, subject: str) -> SchedulerSpec: if not isinstance(raw, dict): raise RegistryValidationError(f"{subject} must be an object") kind = _require_text(raw.get("kind"), subject=f"{subject}.kind") if kind not in ALLOWED_SCHEDULER_KINDS: raise RegistryValidationError( f"{subject}.kind must be one of {', '.join(ALLOWED_SCHEDULER_KINDS)} (got {kind!r})" ) label: str | None = None if kind == "launchd": if "label" not in raw: raise RegistryValidationError(f"{subject}.label is required for launchd schedulers") label = _require_text(raw["label"], subject=f"{subject}.label") allowed = {"kind"} if kind == "launchd": allowed.add("label") unknown = sorted(set(raw) - allowed) if unknown: raise RegistryValidationError( f"{subject} has fields not valid for kind {kind!r}: {', '.join(unknown)}" ) return SchedulerSpec(kind=kind, label=label) def _parse_worker(raw: Any) -> WorkerRecord: data = _require_exact_fields( raw, required=_REQUIRED_WORKER_FIELDS, optional=_OPTIONAL_WORKER_FIELDS, subject="worker", ) worker_id = _require_identifier(data["id"], subject="worker.id") role = _require_text(data["role"], subject=f"worker[{worker_id}].role") if role not in ALLOWED_ROLES: raise RegistryValidationError( f"worker[{worker_id}].role must be one of {', '.join(ALLOWED_ROLES)} (got {role!r})" ) return WorkerRecord( id=worker_id, display_name=_require_text( data["display_name"], subject=f"worker[{worker_id}].display_name" ), provider=_require_identifier(data["provider"], subject=f"worker[{worker_id}].provider"), model=_require_text(data["model"], subject=f"worker[{worker_id}].model"), project=_require_text(data["project"], subject=f"worker[{worker_id}].project"), role=role, namespace=_require_text(data["namespace"], subject=f"worker[{worker_id}].namespace"), profile=_require_text(data["profile"], subject=f"worker[{worker_id}].profile"), workflow=_require_text(data["workflow"], subject=f"worker[{worker_id}].workflow"), schedule=_parse_schedule(data["schedule"], subject=f"worker[{worker_id}].schedule"), timeout_seconds=_require_positive_int( data["timeout_seconds"], subject=f"worker[{worker_id}].timeout_seconds", maximum=_MAX_TIMEOUT_SECONDS, ), enabled=_require_bool(data["enabled"], subject=f"worker[{worker_id}].enabled"), scheduler=_parse_scheduler(data["scheduler"], subject=f"worker[{worker_id}].scheduler"), notes=str(data.get("notes") or "").strip(), ) def _require_unique(values: list[str], *, subject: str) -> None: seen: set[str] = set() for value in values: if value in seen: raise RegistryValidationError(f"duplicate {subject}: {value}") seen.add(value) def validate_payload(payload: Any, *, source_path: Path) -> WorkerRegistry: """Validate a decoded registry document and return the typed registry. Raises :class:`RegistryValidationError` on any violation; never partially accepts a document. """ if not isinstance(payload, dict): raise RegistryValidationError("registry root must be an object") version = payload.get("version") if version != SCHEMA_VERSION: raise RegistryValidationError(f"unsupported registry version: {version!r}") reject_credential_keys(payload, subject="worker registry") unknown = sorted(set(payload) - _TOP_LEVEL_FIELDS) if unknown: raise RegistryValidationError(f"registry has unknown fields: {', '.join(unknown)}") revision = _require_positive_int(payload.get("revision"), subject="revision") updated_at = _require_text(payload.get("updated_at"), subject="updated_at") providers_raw = payload.get("providers") if not isinstance(providers_raw, list) or not providers_raw: raise RegistryValidationError("providers must be a non-empty array") providers = tuple(_parse_provider(item) for item in providers_raw) _require_unique([provider.id for provider in providers], subject="provider id") workers_raw = payload.get("workers") if not isinstance(workers_raw, list): raise RegistryValidationError("workers must be an array") workers = tuple(_parse_worker(item) for item in workers_raw) _require_unique([worker.id for worker in workers], subject="worker id") # Referential integrity: a worker naming an undeclared provider would look # configured while being unrunnable, which is exactly the ambiguous # ownership the epic requires to fail closed. known_providers = {provider.id for provider in providers} for worker in workers: if worker.provider not in known_providers: raise RegistryValidationError( f"worker[{worker.id}].provider references unknown provider {worker.provider!r}" ) # A LaunchAgent label identifies a job to launchd; two workers sharing one # would silently overwrite each other's agent. _require_unique( [worker.scheduler.label for worker in workers if worker.scheduler.label], subject="scheduler label", ) return WorkerRegistry( version=version, revision=revision, updated_at=updated_at, providers=providers, workers=workers, source_path=source_path, ) def load_registry(path: Path | None = None) -> WorkerRegistry: """Load and validate the worker registry from disk.""" source = (path or default_registry_path()).resolve() payload = json.loads(source.read_text(encoding="utf-8")) return validate_payload(payload, source_path=source) # ── serialization ──────────────────────────────────────────────────────────── def provider_to_dict(provider: ProviderRecord) -> dict[str, Any]: return { "id": provider.id, "display_name": provider.display_name, "vendor": provider.vendor, "executable": provider.executable, "available": provider.available, "models": list(provider.models), "notes": provider.notes, } def _schedule_to_dict(schedule: ScheduleSpec) -> dict[str, Any]: payload: dict[str, Any] = {"kind": schedule.kind} if schedule.kind == "interval": payload["seconds"] = schedule.seconds elif schedule.kind == "cron": payload["expression"] = schedule.expression return payload def _scheduler_to_dict(scheduler: SchedulerSpec) -> dict[str, Any]: payload: dict[str, Any] = {"kind": scheduler.kind} if scheduler.kind == "launchd": payload["label"] = scheduler.label return payload def worker_to_dict(worker: WorkerRecord) -> dict[str, Any]: return { "id": worker.id, "display_name": worker.display_name, "provider": worker.provider, "model": worker.model, "project": worker.project, "role": worker.role, "namespace": worker.namespace, "profile": worker.profile, "workflow": worker.workflow, "schedule": _schedule_to_dict(worker.schedule), "timeout_seconds": worker.timeout_seconds, "enabled": worker.enabled, "scheduler": _scheduler_to_dict(worker.scheduler), "notes": worker.notes, } def registry_to_document(registry: WorkerRegistry) -> dict[str, Any]: """Serialize to the on-disk document shape (no local paths embedded).""" return { "version": registry.version, "revision": registry.revision, "updated_at": registry.updated_at, "providers": [provider_to_dict(provider) for provider in registry.providers], "workers": [worker_to_dict(worker) for worker in registry.workers], } def registry_to_dict(registry: WorkerRegistry) -> dict[str, Any]: """Serialize for JSON API responses (adds the resolved source path).""" document = registry_to_document(registry) document["source_path"] = str(registry.source_path) return document def find_worker(registry: WorkerRegistry, worker_id: str) -> WorkerRecord | None: for worker in registry.workers: if worker.id == worker_id: return worker return None def find_provider(registry: WorkerRegistry, provider_id: str) -> ProviderRecord | None: for provider in registry.providers: if provider.id == provider_id: return provider return None def workers_for_provider(registry: WorkerRegistry, provider_id: str) -> tuple[WorkerRecord, ...]: return tuple(worker for worker in registry.workers if worker.provider == provider_id) # ── persistence ────────────────────────────────────────────────────────────── def _utc_now() -> str: return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") def _atomic_write(path: Path, payload: str) -> None: """Write *payload* to *path* atomically: temp file in the same dir, fsync, replace.""" parent = path.parent parent.mkdir(parents=True, exist_ok=True) fd, temp_path = tempfile.mkstemp(prefix=f".{path.name}-", suffix=".tmp", dir=parent) try: with os.fdopen(fd, "w", encoding="utf-8") as handle: handle.write(payload) handle.flush() os.fsync(handle.fileno()) os.replace(temp_path, path) finally: if os.path.exists(temp_path): try: os.remove(temp_path) except OSError: pass def _revision_path(directory: Path, revision: int) -> Path: return directory / f"rev-{revision:06d}.json" def _prune_history(path: Path) -> None: directory = history_dir(path) revisions = list_revisions(path) excess = len(revisions) - _HISTORY_LIMIT for revision in revisions[: max(0, excess)]: _revision_path(directory, revision).unlink(missing_ok=True) def _archive_current(path: Path) -> int | None: """Copy the live document into the history dir under its own revision number.""" if not path.exists(): return None try: existing = json.loads(path.read_text(encoding="utf-8")) revision = int(existing.get("revision", 0)) except (json.JSONDecodeError, TypeError, ValueError, AttributeError): # An unreadable live file has no trustworthy revision number to file it # under, so it cannot join the history chain. return None if revision <= 0: return None _atomic_write( _revision_path(history_dir(path), revision), json.dumps(existing, indent=2, sort_keys=True) + "\n", ) _prune_history(path) return revision def list_revisions(path: Path | None = None) -> tuple[int, ...]: """Revision numbers retained in history for *path*, oldest first.""" directory = history_dir(path) if not directory.is_dir(): return () revisions: list[int] = [] for entry in directory.glob("rev-*.json"): try: revisions.append(int(entry.stem.split("-", 1)[1])) except (IndexError, ValueError): continue return tuple(sorted(revisions)) def save_registry( registry: WorkerRegistry, path: Path | None = None, *, updated_at: str | None = None, ) -> WorkerRegistry: """Validate, archive the superseded revision, then atomically persist a new one. The stored revision is always the previous revision plus one, so a reader can tell two documents apart even when their content is otherwise equal. Returns the registry exactly as persisted. """ target = (path or registry.source_path or default_registry_path()).resolve() document = registry_to_document(registry) # Re-validate before writing: a registry assembled in memory has not # necessarily been through the loader. validate_payload(document, source_path=target) archived = _archive_current(target) document["revision"] = (archived + 1) if archived is not None else registry.revision document["updated_at"] = updated_at or _utc_now() persisted = validate_payload(document, source_path=target) _atomic_write(target, json.dumps(document, indent=2, sort_keys=True) + "\n") return persisted def rollback_to_revision(revision: int, path: Path | None = None) -> WorkerRegistry: """Restore a retained *revision* as a new head revision. History is append-only: rolling back does not delete the revisions in between, it republishes the chosen one under the next revision number, so a rollback is itself reversible. """ target = (path or default_registry_path()).resolve() snapshot_path = _revision_path(history_dir(target), revision) if not snapshot_path.exists(): available = ", ".join(str(item) for item in list_revisions(target)) or "(none)" raise RegistryValidationError( f"revision {revision} is not retained for {target.name}; available: {available}" ) payload = json.loads(snapshot_path.read_text(encoding="utf-8")) restored = validate_payload(payload, source_path=target) return save_registry(restored, target)