Compare commits

..
Author SHA1 Message Date
sysadminandClaude Opus 4.8 b2f6e9a6dc feat(webui): versioned project registry API (Closes #635)
Evolve the MVP project registry (#427) into a versioned, fail-closed
project registry API for the console (Phase 1, read-only).

- Add schema version 2 with project `status`, per-step onboarding
  `state`/`required`, optional redacted `last_seen_health`, and
  `remote_name`. Version 1 files stay loadable and are normalized with
  explicit defaults.
- Serve `/api/v1/projects` and `/api/v1/projects/{project_id}` with API
  provenance (`api_version`, `schema_version`, `source`). `/api/projects`
  is retained as an unversioned Phase 1 alias.
- Replace bare `ValueError` with `RegistryError`, carrying an operator
  `remediation` and `field_path`; invalid registries fail closed as a
  500 JSON payload or a dedicated HTML error page instead of a traceback.
- Reject credential-shaped keys before any DTO is built, reusing
  `registry_safety.is_forbidden_key` as the single source of truth
  shared with the worker registry (#798).
- Render HTML views from `project_to_dict`, so the console and the JSON
  API cannot disagree about status or onboarding progress.
- Document the contract in docs/webui-project-registry-api.md.

Tests: registry load/validate (valid, missing project, schema
validation, credential rejection) and API route coverage.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-22 16:12:07 -04:00
sysadmin 9eb0f29cef Merge pull request 'feat(webui): worker registry and configuration schema (Closes #798)' (#810) from feat/issue-798-worker-registry-schema into master 2026-07-22 06:48:29 -05:00
sysadmin 0b29404031 Merge pull request 'docs(webui): MCP Control Plane Web Console architecture ADR (Closes #632)' (#796) from docs/issue-632-web-console-architecture into master 2026-07-22 06:22:11 -05:00
jcwalker3andClaude Opus 4.8 5463f58933 feat(webui): worker registry and configuration schema (Closes #798)
Add the declarative worker registry that epic #797 makes the source of
truth for the scheduled multi-LLM worker fleet.

Providers and configured workers are modelled as separate entities so a
provider can be listed with no worker configured, and so provider facts
are not copied into every worker record. A worker records provider,
model, project, role, namespace, profile, workflow, schedule, timeout,
enabled state, and scheduler metadata.

Validation fails closed: unknown fields are refused rather than ignored,
so a typo cannot silently disable a timeout; a worker naming an
undeclared provider is rejected; worker ids, provider ids, and
LaunchAgent labels must be unique.

Persistence is atomic (temp file in the same directory, fsync, replace).
Every superseded document is retained as a numbered revision, and
rollback republishes a chosen revision as a new head, so history stays
append-only and a rollback is itself reversible.

The credential-rejection guard is extracted to webui/registry_safety.py
so both registries share one implementation instead of two copies of a
security check; project_registry.py keeps identical behaviour.

Scope: data model, validation, persistence only. No routes, scheduler,
process control, or provider probing - those are #799/#800/#804/#805.
The workers array ships empty because populating it is #808.

Tests: tests/test_webui_worker_registry.py, 44 cases.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-22 06:10:34 -05:00
sysadmin 5032965e3a Merge pull request 'feat: make master-parity live-remote aware so a stale daemon fails closed (Closes #610)' (#788) from feat/issue-610-live-remote-parity into master 2026-07-22 05:51:56 -05:00
jcwalker3andGrok 4.5 57a52b1a99 fix(parity): hermetic live-remote master reads under pytest (Closes #610)
Remediate PR #788 review F1/F2: suite-wide hermetic mode prevents
git ls-remote from running in tests so feature worktrees no longer flip
legacy runtime-context assertions to live_stale, and unit tests stay offline.
Module flag survives patch.dict(clear=True); env override still wins.

Co-Authored-By: Grok 4.5 <[email protected]>
2026-07-22 05:40:15 -05:00
jcwalker3 344dc41ce2 Merge branch 'master' into feat/issue-610-live-remote-parity 2026-07-22 04:30:44 -05:00
jcwalker3 0f19773076 Merge branch 'master' into feat/issue-610-live-remote-parity 2026-07-21 20:48:54 -05:00
sysadminandClaude Opus 4.8 324b4b3e93 feat: make master-parity live-remote aware so a stale daemon fails closed (Closes #610)
Master-parity previously compared only the daemon's startup commit against the
local on-disk HEAD. When the checkout was not pulled, parity reported green even
though the live remote master had advanced, so a stale daemon could claim a
mutation-safe result while running outdated capability gates (observed during
PR #592 recovery, where the resolver correctly required restart but parity said
in_parity=true).

Changes:
- master_parity_gate.assess_master_parity() gains an optional live_remote_head
  and reports the three commits distinctly (daemon_start_head, local_head,
  live_remote_head) plus live_known / live_stale / mutation_safe. A result is
  mutation_safe only when daemon, local checkout, and live remote all agree.
- parity_block_reasons() now blocks mutations on live-staleness too; read-only
  operations remain unblocked (non-goal: never block diagnostics offline).
- New parity_resolver_disagreement(): typed fail-closed blocker naming the
  capability resolver as authoritative when it requires restart but parity
  looks locally green.
- New read_remote_master_head(): best-effort `git ls-remote` for the live
  target, cached with a 60s TTL (bounded offline latency, no network probe per
  gate call); env override GITEA_TEST_LIVE_REMOTE_HEAD keeps tests hermetic.
- Server: _current_master_parity() reads the live remote head;
  gitea_assess_master_parity and runtime_context surface the distinguished
  SHAs, mutation_safe, and resolver-authoritative guidance.

Tests: 13 new cases (live-remote parity, live-stale blocking + typed blocker,
remote-head reader + TTL cache, server wiring). Full suite: 2423 passed; the 8
remaining failures (test_config TestAuthIntegration, test_credentials
TestGetCredentials) are baseline-proven keychain/env failures identical on the
unmodified base.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-09 18:52:59 -04:00
15 changed files with 2992 additions and 143 deletions
+8 -2
View File
@@ -43,6 +43,10 @@ for the console architecture: layer and authority boundaries, the redaction
boundary, `/api/v1/...` versioning, the target page map, and the phase gates boundary, `/api/v1/...` versioning, the target page map, and the phase gates
that govern when a write path may open (#632, epic #631). that govern when a write path may open (#632, epic #631).
See [webui-project-registry-api.md](webui-project-registry-api.md) for the
versioned project registry contract: registry schema versions 1 and 2, project
status, onboarding checklist state, and the fail-closed error payloads (#635).
## Routes (MVP) ## Routes (MVP)
| Path | Description | | Path | Description |
@@ -51,9 +55,11 @@ that govern when a write path may open (#632, epic #631).
| `/health` | JSON liveness (`status`, `service`, `mode`, `timestamp`) | | `/health` | JSON liveness (`status`, `service`, `mode`, `timestamp`) |
| `/queue` | Live PR and issue queue dashboard (#429) | | `/queue` | Live PR and issue queue dashboard (#429) |
| `/api/queue` | JSON queue export with pagination metadata | | `/api/queue` | JSON queue export with pagination metadata |
| `/projects` | Project registry list (#427) | | `/projects` | Project registry list with status and onboarding progress (#427, #635) |
| `/projects/{id}` | Project detail + onboarding checklist | | `/projects/{id}` | Project detail + onboarding checklist |
| `/api/projects` | JSON registry export | | `/api/v1/projects` | Versioned JSON registry export (#635) |
| `/api/v1/projects/{id}` | Versioned JSON project detail (#635) |
| `/api/projects` | JSON registry export — unversioned Phase 1 alias of `/api/v1/projects` |
| `/prompts` | Prompt library with per-prompt copy buttons (#428) | | `/prompts` | Prompt library with per-prompt copy buttons (#428) |
| `/api/prompts` | JSON prompt export with workflow hashes | | `/api/prompts` | JSON prompt export with workflow hashes |
| `/runtime` | MCP runtime health and stale detection (#430) | | `/runtime` | MCP runtime health and stale detection (#430) |
+213
View File
@@ -0,0 +1,213 @@
# Project registry API (#635)
Phase 1 of the [console architecture ADR](architecture/webui-control-plane-console-architecture-adr.md)
gives the project registry a versioned, read-only API. This document is the
field-by-field contract for that API and for the registry file behind it.
Everything here is **read-only**. The console never writes the registry; an
operator edits the JSON file, and an invalid file fails closed rather than
rendering a partial inventory.
## Routes
| Route | Method | Description |
|-------|--------|-------------|
| `/api/v1/projects` | GET | Versioned registry export: all projects, with provenance |
| `/api/v1/projects/{project_id}` | GET | Single project; `404` with `project_not_found` when unknown |
| `/api/projects` | GET | Unversioned MVP alias (#427), retained for all of Phase 1 |
| `/projects` | GET | HTML list — status and onboarding progress per project |
| `/projects/{project_id}` | GET | HTML detail — identity, profiles, paths, checklist |
Per ADR section 6 the unversioned alias may be retired no earlier than Phase 2,
and only after this document and `webui-local-dev.md` record the swap. The alias
returns the same payload as `/api/v1/projects`, including the legacy `version`
and `source_path` keys #427 consumers already read.
The HTML views render from the same DTO the JSON routes serialize
(`project_to_dict`), so the console and the API cannot disagree about a
project's status or onboarding progress.
## Registry file
Default location: `webui/data/projects.registry.json`. Override with the
`WEBUI_PROJECT_REGISTRY` environment variable.
Schema versions: **1** and **2** are accepted; **2** is current. A version 1
file loads unchanged and is normalized with the documented defaults, so an
existing operator registry keeps working without edits.
### Root
| Field | Type | Required | Notes |
|-------|------|----------|-------|
| `version` | int | yes | `1` or `2`. Anything else fails closed |
| `projects` | array | yes | Must be non-empty |
### Project
| Field | Type | Required | Default | Notes |
|-------|------|----------|---------|-------|
| `id` | string | yes | — | Stable registry id used in URLs |
| `repo_name` | string | yes | — | Gitea repository name |
| `gitea_owner` | string | yes | — | Owning org or user |
| `remote_host` | string | yes | — | Instance base URL, no credentials |
| `remote_name` | string | no | `null` | Logical remote label, e.g. `prgs` (v2) |
| `default_branch` | string | yes | — | Stable branch name |
| `local_checkout_path` | string | yes | — | Control checkout path |
| `status` | string | no | `active` | `active`, `onboarding`, `paused`, `archived` (v2) |
| `profiles` | object | yes | — | Must map `author`, `reviewer`, `reconciler` |
| `workflow_paths` | object | yes | — | Non-empty; label to repo-relative path |
| `schema_paths` | object | no | `{}` | Label to repo-relative path |
| `onboarding_checklist` | array | no | `[]` | See below |
| `last_seen_health` | object | no | `null` | Redacted health only (v2) |
### Onboarding step
| Field | Type | Required | Default | Notes |
|-------|------|----------|---------|-------|
| `id` | string | yes | — | Stable step id |
| `title` | string | yes | — | Short operator-facing label |
| `description` | string | yes | — | Self-contained; assumes no chat history |
| `state` | string | no | `pending` | `complete`, `pending`, `blocked`, `not_applicable` (v2) |
| `required` | bool | no | `true` | Optional steps never block readiness (v2) |
### Last-seen health
| Field | Type | Required | Notes |
|-------|------|----------|-------|
| `status` | string | no (default `unknown`) | `healthy`, `degraded`, `unreachable`, `unknown` |
| `checked_at` | string | no | ISO-8601 UTC timestamp, e.g. `2026-01-01T00:00:00Z` |
| `detail` | string | no | Short redacted note |
Health is recorded metadata, not a live probe: Phase 1 performs no outbound
health checks. Endpoints, tokens, and keychain identifiers must never appear
here.
## Response shape
`GET /api/v1/projects`:
```json
{
"api_version": "v1",
"schema_version": 2,
"version": 2,
"source_path": "/path/to/webui/data/projects.registry.json",
"source": {
"kind": "file",
"path": "/path/to/webui/data/projects.registry.json",
"inventory_complete": true
},
"project_count": 1,
"projects": [
{
"id": "example",
"repo_name": "Example",
"gitea_owner": "Org",
"repo_full_name": "Org/Example",
"remote_host": "https://gitea.example.invalid",
"remote_name": "example-remote",
"default_branch": "main",
"local_checkout_path": ".",
"status": "active",
"profiles": {"author": "...", "reviewer": "...", "reconciler": "..."},
"workflow_paths": {"skill": "skills/..."},
"schema_paths": {},
"onboarding_checklist": [
{
"id": "profiles",
"title": "Configure execution profiles",
"description": "...",
"state": "complete",
"required": true
}
],
"onboarding_summary": {
"total": 1,
"complete": 1,
"pending": 0,
"blocked": 0,
"not_applicable": 0,
"required_outstanding": 0,
"onboarding_complete": true
},
"last_seen_health": null
}
]
}
```
`GET /api/v1/projects/{project_id}` returns `api_version`, `schema_version`,
`source`, and a single `project` object with the same fields.
The `source` block satisfies the ADR section 6 provenance rule: every payload
states where the data came from and whether the inventory is complete. A
file-backed registry is always complete — there is no pagination to truncate it.
`onboarding_summary` is derived, never stored. `required_outstanding` counts
steps that are `required` **and** in state `pending` or `blocked`;
`onboarding_complete` is true when that count is zero.
## Fail-closed errors
Validation failures raise `RegistryError`, which routes render instead of a
traceback.
`404` — unknown project id on `/api/v1/projects/{project_id}`:
```json
{
"error": "project_not_found",
"project_id": "not-registered",
"known_project_ids": ["example"],
"remediation": "Request one of the known project ids, or add the project ...",
"source": {"kind": "file", "path": "...", "inventory_complete": true}
}
```
`500` — invalid registry, on both the versioned route and the alias:
```json
{
"error": "registry_invalid",
"detail": "unsupported registry version: 42",
"remediation": "Set 'version' to one of 1, 2 (current schema is 2) ...",
"field_path": "version",
"source_path": "/path/to/registry.json"
}
```
`field_path` points at the offending location (`projects[0].profiles.reconciler`,
`projects[0].onboarding_checklist[2].state`, and so on). The HTML routes render
the same detail, field, source, and remediation on a "Project registry
unavailable" page.
Conditions that fail closed:
* file missing or unreadable;
* invalid JSON (the remediation names line and column);
* root not an object, or `projects` missing/empty;
* unsupported `version`;
* a credential-shaped key anywhere in the file (`token`, `*_secret`, `auth_*`, and similar);
* a project missing a required field, or missing an `author`/`reviewer`/`reconciler` profile;
* an unknown `status`, onboarding `state`, or health `status`.
## Credential rule
The registry stores redacted metadata only. Credential-shaped keys are
rejected at load time, before any DTO is built, consistent with
[safety-model.md](safety-model.md) and
[credential-isolation.md](credential-isolation.md). Tokens live in the keychain
and are resolved server-side by `gitea_auth`.
## Migrating a version 1 registry
1. Set `"version": 2`.
2. Optionally add `"status"` per project (omitted means `active`).
3. Optionally add `"remote_name"` per project.
4. Optionally add `"state"` and `"required"` to each onboarding step (omitted
means `pending` and `true`).
5. Optionally add `"last_seen_health"`.
No step is mandatory: a version 1 file keeps loading. Bumping the version only
declares that the file may use the v2 fields.
+64 -9
View File
@@ -12536,10 +12536,39 @@ def _try_auto_switch_for_operation(op: str, host: str | None = None) -> bool:
return False return False
def _git_default_remote_name(root: str) -> str:
"""First configured git remote name for *root*, defaulting to 'origin'.
Used to resolve the live remote master target for parity (#610). Best
effort: any failure falls back to 'origin' so callers never raise.
"""
try:
res = subprocess.run(
["git", "-C", root, "remote"],
capture_output=True, text=True, check=False,
)
except Exception:
return "origin"
if res.returncode != 0:
return "origin"
names = [n.strip() for n in (res.stdout or "").splitlines() if n.strip()]
return names[0] if names else "origin"
def _current_master_parity() -> dict: def _current_master_parity() -> dict:
"""Assess this process's code against the on-disk master HEAD (#420).""" """Assess this process's code against local and live remote master (#420/#610).
Compares the daemon's startup commit, the on-disk checkout HEAD, and the
live remote master target. A stale daemon relative to live master fails
closed for mutations even when the local checkout HEAD still matches the
startup commit. The live-remote read is best effort: an unresolved live
head leaves read-only diagnostics unblocked but is never mutation-safe.
"""
current_head = master_parity_gate.read_git_head(PROJECT_ROOT) current_head = master_parity_gate.read_git_head(PROJECT_ROOT)
return master_parity_gate.assess_master_parity(_STARTUP_PARITY, current_head) live_head = master_parity_gate.read_remote_master_head(
PROJECT_ROOT, remote=_git_default_remote_name(PROJECT_ROOT))
return master_parity_gate.assess_master_parity(
_STARTUP_PARITY, current_head, live_remote_head=live_head)
def _current_runtime_mode_report(refresh: bool = False) -> dict: def _current_runtime_mode_report(refresh: bool = False) -> dict:
@@ -16480,15 +16509,34 @@ def gitea_get_runtime_context(
"restart_required": parity["restart_required"], "restart_required": parity["restart_required"],
"startup_head": parity["startup_head"], "startup_head": parity["startup_head"],
"current_head": parity["current_head"], "current_head": parity["current_head"],
# #610 distinguished mutation-safety signals:
"daemon_start_head": parity["daemon_start_head"],
"local_head": parity["local_head"],
"live_remote_head": parity["live_remote_head"],
"live_known": parity["live_known"],
"live_stale": parity["live_stale"],
"mutation_safe": parity["mutation_safe"],
"summary": master_parity_gate.format_parity(parity), "summary": master_parity_gate.format_parity(parity),
"mutation_gate_enforced": not master_parity_gate.gate_disabled(), "mutation_gate_enforced": not master_parity_gate.gate_disabled(),
# #610: the capability resolver is authoritative for mutation safety;
# local parity alone must never authorize a mutation.
"resolver_authoritative_for_mutation_safety": True,
} }
if parity["stale"] and not master_parity_gate.gate_disabled(): if parity["restart_required"] and not master_parity_gate.gate_disabled():
safe_next_action = ( if parity["live_stale"]:
"Server code is stale relative to master; restart the Gitea MCP " safe_next_action = (
"server to load current capability gates before mutating. " "Daemon is stale relative to LIVE remote master "
f"({master_parity_gate.format_parity(parity)})" f"(started {parity['startup_head'][:12] if parity['startup_head'] else 'unknown'}, "
) f"live master {parity['live_remote_head'][:12] if parity['live_remote_head'] else 'unknown'}); "
"restart/reconnect the Gitea MCP server before mutating. The "
"capability resolver is authoritative for mutation safety."
)
else:
safe_next_action = (
"Server code is stale relative to master; restart the Gitea MCP "
"server to load current capability gates before mutating. "
f"({master_parity_gate.format_parity(parity)})"
)
result["safe_next_action"] = safe_next_action result["safe_next_action"] = safe_next_action
if reveal and h: if reveal and h:
@@ -16537,6 +16585,13 @@ def gitea_assess_master_parity(
"determinable": parity["determinable"], "determinable": parity["determinable"],
"startup_head": parity["startup_head"], "startup_head": parity["startup_head"],
"current_head": parity["current_head"], "current_head": parity["current_head"],
# #610 distinguished mutation-safety signals:
"daemon_start_head": parity["daemon_start_head"],
"local_head": parity["local_head"],
"live_remote_head": parity["live_remote_head"],
"live_known": parity["live_known"],
"live_stale": parity["live_stale"],
"mutation_safe": parity["mutation_safe"],
"mutation_gate_enforced": enforced, "mutation_gate_enforced": enforced,
"summary": master_parity_gate.format_parity(parity), "summary": master_parity_gate.format_parity(parity),
"reasons": parity["reasons"], "reasons": parity["reasons"],
@@ -16553,7 +16608,7 @@ def gitea_assess_master_parity(
source=canonical_source, source=canonical_source,
), ),
} }
if parity["stale"] and enforced: if parity["restart_required"] and enforced:
out["report"] = master_parity_gate.parity_report(parity) out["report"] = master_parity_gate.parity_report(parity)
return out return out
+213 -17
View File
@@ -24,12 +24,50 @@ from __future__ import annotations
import os import os
import subprocess import subprocess
import time
# Live-remote head cache: the parity gate runs on every mutation and every
# runtime-context read, so the ``git ls-remote`` result is cached briefly to
# avoid a network round-trip per call (#610). Keyed by (root, remote, branch).
_REMOTE_HEAD_CACHE: dict[tuple[str, str, str], tuple[float, str | None]] = {}
_REMOTE_HEAD_TTL = 60.0
# When True, ``read_remote_master_head`` never performs ``git ls-remote`` unless
# ``GITEA_TEST_LIVE_REMOTE_HEAD`` is set. Conftest enables this suite-wide so
# feature worktrees (whose HEAD differs from live master) cannot flip legacy
# runtime-context assertions to live_stale, and so unit tests never depend on
# a live network (PR #788 F1/F2 / issue #610). Module-level (not env-only) so
# ``patch.dict(os.environ, …, clear=True)`` cannot re-enable the probe.
_HERMETIC_TEST_MODE: bool = False
def _clear_remote_head_cache() -> None:
"""Reset the live-remote head cache (test isolation / forced refresh)."""
_REMOTE_HEAD_CACHE.clear()
def set_hermetic_test_mode(enabled: bool) -> None:
"""Enable or disable suite-wide hermetic live-remote reads (tests only)."""
global _HERMETIC_TEST_MODE
_HERMETIC_TEST_MODE = bool(enabled)
_clear_remote_head_cache()
def hermetic_test_mode() -> bool:
"""Return whether hermetic live-remote reads are active."""
return bool(_HERMETIC_TEST_MODE)
# Environment escape hatches (ops + tests): # Environment escape hatches (ops + tests):
# GITEA_MCP_DISABLE_PARITY_GATE -> disable enforcement entirely (fail open). # GITEA_MCP_DISABLE_PARITY_GATE -> disable enforcement entirely (fail open).
# GITEA_TEST_CURRENT_HEAD -> force the "current" HEAD read, for tests. # GITEA_TEST_CURRENT_HEAD -> force the "current" HEAD read, for tests.
ENV_DISABLE = "GITEA_MCP_DISABLE_PARITY_GATE" ENV_DISABLE = "GITEA_MCP_DISABLE_PARITY_GATE"
ENV_TEST_CURRENT_HEAD = "GITEA_TEST_CURRENT_HEAD" ENV_TEST_CURRENT_HEAD = "GITEA_TEST_CURRENT_HEAD"
# GITEA_TEST_LIVE_REMOTE_HEAD -> force the live remote master read, for tests.
ENV_TEST_LIVE_REMOTE_HEAD = "GITEA_TEST_LIVE_REMOTE_HEAD"
# GITEA_TEST_ALLOW_LIVE_REMOTE_PROBE -> opt a single test into a real ls-remote
# even when hermetic mode is on (rare; prefer ENV_TEST_LIVE_REMOTE_HEAD).
ENV_TEST_ALLOW_LIVE_REMOTE_PROBE = "GITEA_TEST_ALLOW_LIVE_REMOTE_PROBE"
def read_git_head(root: str) -> str | None: def read_git_head(root: str) -> str | None:
@@ -58,6 +96,75 @@ def read_git_head(root: str) -> str | None:
return (res.stdout or "").strip() or None return (res.stdout or "").strip() or None
def read_remote_master_head(
root: str,
remote: str = "origin",
branch: str = "master",
ttl: float = _REMOTE_HEAD_TTL,
) -> str | None:
"""Return the live remote ``branch`` commit SHA, or ``None`` (#610).
Resolves the *live* target commit via ``git ls-remote`` so parity can tell
a daemon that is behind the live remote master apart from one whose local
checkout simply hasn't been pulled. ``None`` means the live head could not
be resolved (offline, no such remote, git unavailable, error) -- callers
must treat unknown live state as *not mutation-safe* while never blocking
read-only diagnostics. A ``GITEA_TEST_LIVE_REMOTE_HEAD`` override takes
precedence so the wiring can be exercised deterministically and offline.
The result is cached for *ttl* seconds per (root, remote, branch) so the
gate does not run a network probe on every mutation/read (``ttl=0`` forces
a live probe). Both hits and ``None`` misses are cached to bound offline
latency; the env override bypasses the cache and the subprocess entirely.
Under suite hermetic mode (``set_hermetic_test_mode(True)``, set by
conftest) a missing override returns ``None`` without network I/O so
feature-worktree test runs cannot observe live_stale against real master
(PR #788 F1) and unit tests stay offline (F2). Opt out with an explicit
``GITEA_TEST_LIVE_REMOTE_HEAD`` pin or ``GITEA_TEST_ALLOW_LIVE_REMOTE_PROBE``.
"""
forced = os.environ.get(ENV_TEST_LIVE_REMOTE_HEAD)
if forced is not None:
return forced.strip() or None
if _HERMETIC_TEST_MODE and not (
os.environ.get(ENV_TEST_ALLOW_LIVE_REMOTE_PROBE) or ""
).strip():
# Hermetic default: live head unknown. live_stale stays False;
# mutation_safe is False when live is unknown (documented #610 note).
return None
# Defense in depth: even without the module flag, never probe while pytest
# is running unless the test opted into a real probe or set an override.
if (os.environ.get("PYTEST_CURRENT_TEST") or "").strip() and not (
os.environ.get(ENV_TEST_ALLOW_LIVE_REMOTE_PROBE) or ""
).strip():
return None
if not root:
return None
key = (root, remote, branch)
now = time.monotonic()
if ttl > 0:
cached = _REMOTE_HEAD_CACHE.get(key)
if cached is not None and (now - cached[0]) < ttl:
return cached[1]
sha: str | None = None
try:
res = subprocess.run(
["git", "-C", root, "ls-remote", remote, f"refs/heads/{branch}"],
capture_output=True,
text=True,
check=False,
timeout=5,
)
if res.returncode == 0:
lines = (res.stdout or "").strip().splitlines()
if lines:
sha = lines[0].split("\t", 1)[0].split()[0].strip() or None
except Exception:
sha = None
_REMOTE_HEAD_CACHE[key] = (now, sha)
return sha
def capture_startup_parity(root: str, head: str | None = None) -> dict: def capture_startup_parity(root: str, head: str | None = None) -> dict:
"""Capture the process source-tree baseline once at server startup. """Capture the process source-tree baseline once at server startup.
@@ -72,18 +179,38 @@ def _short(sha: str | None) -> str:
return sha[:12] if sha else "unknown" return sha[:12] if sha else "unknown"
def assess_master_parity(startup: dict | None, current_head: str | None) -> dict: def assess_master_parity(
startup: dict | None,
current_head: str | None,
live_remote_head: str | None = None,
) -> dict:
"""Compare the startup baseline against the current on-disk ``HEAD``. """Compare the startup baseline against the current on-disk ``HEAD``.
Pure: both HEADs are supplied by the caller. Returns a structured result: Pure: all HEADs are supplied by the caller. Returns a structured result:
- ``in_parity`` -- server code matches the on-disk master (or parity - ``in_parity`` -- server code matches the on-disk master (or parity
could not be determined, which is not treated as stale). could not be determined, which is not treated as stale).
- ``stale`` -- the on-disk master has definitively advanced past the - ``stale`` -- the on-disk master has definitively advanced past the
running process. running process.
- ``restart_required`` -- alias of ``stale``; the recovery action. - ``restart_required`` -- ``stale`` or ``live_stale``; the recovery action.
- ``determinable`` -- whether both HEADs were known well enough to compare. - ``determinable`` -- whether both local HEADs were known well enough to
compare.
- ``startup_head`` / ``current_head`` / ``reasons``. - ``startup_head`` / ``current_head`` / ``reasons``.
#610 adds live-remote awareness so a daemon that is stale relative to the
*live* remote master cannot report a mutation-safe result even when the
local checkout HEAD still matches the daemon's startup commit:
- ``daemon_start_head`` -- the commit the running process started at
(alias of ``startup_head``, named for clarity in reports).
- ``local_head`` -- the on-disk checkout HEAD (alias of ``current_head``).
- ``live_remote_head`` -- the live remote target commit, or ``None`` when it
could not be fetched.
- ``live_known`` -- whether the live remote target was resolved.
- ``live_stale`` -- the live remote master has advanced past the running
process (daemon is behind live master) even if local parity is green.
- ``mutation_safe`` -- the daemon code, local checkout, and live remote
target all agree; the only state in which a mutation may rely on parity.
""" """
startup_head = (startup or {}).get("startup_head") startup_head = (startup or {}).get("startup_head")
reasons: list[str] = [] reasons: list[str] = []
@@ -91,32 +218,56 @@ def assess_master_parity(startup: dict | None, current_head: str | None) -> dict
if startup_head is None: if startup_head is None:
reasons.append( reasons.append(
"startup commit was not captured; code parity cannot be enforced") "startup commit was not captured; code parity cannot be enforced")
return _result(True, False, False, startup_head, current_head, reasons) return _result(True, False, False, startup_head, current_head,
live_remote_head, False, reasons)
if current_head is None: if current_head is None:
reasons.append( reasons.append(
"current workspace HEAD could not be read; code parity cannot be " "current workspace HEAD could not be read; code parity cannot be "
"enforced") "enforced")
return _result(True, False, False, startup_head, current_head, reasons) return _result(True, False, False, startup_head, current_head,
live_remote_head, False, reasons)
if startup_head == current_head: local_in_parity = startup_head == current_head
return _result(True, False, True, startup_head, current_head, reasons) local_stale = not local_in_parity
if local_stale:
reasons.append(
f"MCP server started at commit {_short(startup_head)} but the "
f"workspace master is now {_short(current_head)}; restart the "
f"server to load the current capability gates")
reasons.append( live_known = live_remote_head is not None
f"MCP server started at commit {_short(startup_head)} but the workspace " live_stale = live_known and live_remote_head != startup_head
f"master is now {_short(current_head)}; restart the server to load the " if live_stale:
f"current capability gates") reasons.append(
return _result(False, True, True, startup_head, current_head, reasons) f"live remote master is {_short(live_remote_head)} but the MCP "
f"server started at {_short(startup_head)}; the daemon is stale "
f"relative to live master -- restart/reconnect before mutating")
return _result(
local_in_parity, local_stale, True, startup_head, current_head,
live_remote_head, live_stale, reasons)
def _result(in_parity, stale, determinable, startup_head, current_head, reasons): def _result(in_parity, stale, determinable, startup_head, current_head,
live_remote_head, live_stale, reasons):
live_known = live_remote_head is not None
mutation_safe = (
determinable and in_parity and live_known and not live_stale)
return { return {
"in_parity": in_parity, "in_parity": in_parity,
"stale": stale, "stale": stale,
"restart_required": stale, "restart_required": stale or live_stale,
"determinable": determinable, "determinable": determinable,
"startup_head": startup_head, "startup_head": startup_head,
"current_head": current_head, "current_head": current_head,
# #610 distinguished signals:
"daemon_start_head": startup_head,
"local_head": current_head,
"live_remote_head": live_remote_head,
"live_known": live_known,
"live_stale": live_stale,
"mutation_safe": mutation_safe,
"reasons": list(reasons), "reasons": list(reasons),
} }
@@ -130,11 +281,13 @@ def parity_block_reasons(assessment: dict) -> list[str]:
"""Block reasons for a mutation gate (empty when the mutation may proceed). """Block reasons for a mutation gate (empty when the mutation may proceed).
A disabled gate or an in-parity / non-determinable assessment yields no A disabled gate or an in-parity / non-determinable assessment yields no
reasons; only a definitively stale server blocks. reasons. A definitively stale server blocks, and (#610) a daemon that is
stale relative to the *live* remote master blocks even when the local
checkout HEAD still matches the daemon's startup commit.
""" """
if gate_disabled(): if gate_disabled():
return [] return []
if assessment.get("stale"): if assessment.get("stale") or assessment.get("live_stale"):
return list(assessment.get("reasons") or return list(assessment.get("reasons") or
["server code is stale relative to master (fail closed)"]) ["server code is stale relative to master (fail closed)"])
return [] return []
@@ -147,6 +300,10 @@ def parity_report(assessment: dict) -> dict:
"restart_required": True, "restart_required": True,
"startup_head": assessment.get("startup_head"), "startup_head": assessment.get("startup_head"),
"current_head": assessment.get("current_head"), "current_head": assessment.get("current_head"),
# #610: name the live remote target so the report distinguishes a
# local-code stale from a daemon-behind-live-master stale.
"live_remote_head": assessment.get("live_remote_head"),
"live_stale": bool(assessment.get("live_stale")),
"reasons": list(assessment.get("reasons") or []), "reasons": list(assessment.get("reasons") or []),
"recovery": [ "recovery": [
"The running MCP server is executing code older than the current " "The running MCP server is executing code older than the current "
@@ -157,6 +314,45 @@ def parity_report(assessment: dict) -> dict:
} }
def parity_resolver_disagreement(
assessment: dict,
resolver_restart_required: bool,
) -> dict | None:
"""Typed blocker when the resolver requires restart but parity looks green.
The capability resolver (``gitea_resolve_task_capability``) detects stale
runtime authoritatively for mutation safety (#610). When it requires a
restart, local-only parity must never override it: this returns a typed,
fail-closed blocker that names the resolver as authoritative. Returns
``None`` when the resolver does not require a restart.
"""
if not resolver_restart_required:
return None
parity_optimistic = bool(assessment.get("in_parity")) and not (
assessment.get("stale") or assessment.get("live_stale"))
return {
"kind": "parity_resolver_disagreement",
"restart_required": True,
"resolver_authoritative": True,
"parity_optimistic": parity_optimistic,
"daemon_start_head": assessment.get("daemon_start_head"),
"local_head": assessment.get("local_head"),
"live_remote_head": assessment.get("live_remote_head"),
"reasons": [
"The capability resolver requires a restart/reconnect (stale "
"runtime) but master-parity reported local code as in-parity. "
"The resolver is authoritative for mutation safety; do not mutate "
"on local parity alone. Restart/reconnect the Gitea MCP server "
"and re-verify before mutating.",
],
"recovery": [
"Trust the resolver: treat this session as stale.",
"Restart or /mcp reconnect the Gitea MCP namespace so it reloads "
"current master and live target state, then re-run preflight.",
],
}
def format_parity(assessment: dict) -> str: def format_parity(assessment: dict) -> str:
"""One-line human summary for logs / runtime context.""" """One-line human summary for logs / runtime context."""
if assessment.get("stale"): if assessment.get("stale"):
+29
View File
@@ -167,6 +167,35 @@ def _reset_mutation_authority(monkeypatch):
import pytest import pytest
@pytest.fixture(autouse=True)
def _hermetic_live_remote_master_head():
"""#610 / PR #788 F1/F2: keep live-remote parity reads offline in tests.
``read_remote_master_head`` would otherwise ``git ls-remote`` whenever
``GITEA_TEST_LIVE_REMOTE_HEAD`` is unset. Feature worktrees under
``branches/`` always differ from live master, so legacy suites that assert
runtime-context ``safe_next_action`` flip to live_stale. Module-level
hermetic mode survives ``patch.dict(os.environ, …, clear=True)``.
Tests that exercise the real probe path call
``master_parity_gate.set_hermetic_test_mode(False)`` and/or set
``GITEA_TEST_ALLOW_LIVE_REMOTE_PROBE``.
"""
try:
import master_parity_gate as _mpg
_mpg.set_hermetic_test_mode(True)
except Exception:
_mpg = None
try:
yield
finally:
if _mpg is not None:
try:
_mpg.set_hermetic_test_mode(False)
except Exception:
pass
@pytest.fixture(autouse=True) @pytest.fixture(autouse=True)
def _deterministic_workspace_remotes(): def _deterministic_workspace_remotes():
try: try:
+269
View File
@@ -78,6 +78,95 @@ class TestBlockReasonsAndReport(unittest.TestCase):
self.assertTrue(report["recovery"]) self.assertTrue(report["recovery"])
class TestLiveRemoteParity(unittest.TestCase):
"""#610: parity must account for the live remote master, not just local.
The daemon can be stale relative to the live remote target while the local
checkout HEAD still matches the daemon's startup commit, so local parity
reports green even though a mutation would run against outdated code.
"""
SHA_C = "c" * 40
def test_distinguishes_three_shas(self):
res = mp.assess_master_parity(
{"startup_head": SHA_A}, SHA_A, live_remote_head=SHA_B)
self.assertEqual(res["daemon_start_head"], SHA_A)
self.assertEqual(res["local_head"], SHA_A)
self.assertEqual(res["live_remote_head"], SHA_B)
def test_mutation_safe_only_when_all_three_match(self):
res = mp.assess_master_parity(
{"startup_head": SHA_A}, SHA_A, live_remote_head=SHA_A)
self.assertTrue(res["mutation_safe"])
self.assertTrue(res["live_known"])
self.assertFalse(res["live_stale"])
def test_live_stale_when_remote_advanced_past_daemon(self):
# Local checkout still matches the daemon start (local parity green),
# but the live remote master has advanced -> daemon is live-stale.
res = mp.assess_master_parity(
{"startup_head": SHA_A}, SHA_A, live_remote_head=SHA_B)
self.assertTrue(res["in_parity"]) # local parity still green
self.assertTrue(res["live_stale"])
self.assertFalse(res["mutation_safe"])
self.assertTrue(any("live" in r.lower() for r in res["reasons"]))
def test_live_unknown_is_not_mutation_safe_but_not_stale(self):
# Non-goal: unfetchable live remote must not be treated as stale for
# read-only, but a mutation-safe claim fails closed.
res = mp.assess_master_parity(
{"startup_head": SHA_A}, SHA_A, live_remote_head=None)
self.assertFalse(res["live_known"])
self.assertFalse(res["mutation_safe"])
self.assertFalse(res["live_stale"])
self.assertTrue(res["in_parity"])
def test_default_live_remote_preserves_legacy_shape(self):
# Callers that do not supply a live head keep the pre-#610 behavior:
# in-parity, not live-stale, no live-derived block.
res = mp.assess_master_parity({"startup_head": SHA_A}, SHA_A)
self.assertFalse(res["live_stale"])
self.assertEqual(mp.parity_block_reasons(res), [])
class TestLiveStaleBlockAndReport(unittest.TestCase):
"""#610: live-staleness must block mutations and surface a typed blocker."""
def test_live_stale_produces_block_reasons(self):
res = mp.assess_master_parity(
{"startup_head": SHA_A}, SHA_A, live_remote_head=SHA_B)
self.assertTrue(mp.parity_block_reasons(res))
def test_disable_env_suppresses_live_stale_block(self):
res = mp.assess_master_parity(
{"startup_head": SHA_A}, SHA_A, live_remote_head=SHA_B)
with patch.dict(os.environ, {mp.ENV_DISABLE: "1"}):
self.assertEqual(mp.parity_block_reasons(res), [])
def test_resolver_disagreement_returns_typed_blocker(self):
# Parity says local-green, resolver says restart required -> disagreement
# is a typed, fail-closed blocker naming the resolver as authoritative.
res = mp.assess_master_parity({"startup_head": SHA_A}, SHA_A)
blocker = mp.parity_resolver_disagreement(res, resolver_restart_required=True)
self.assertIsNotNone(blocker)
self.assertEqual(blocker["kind"], "parity_resolver_disagreement")
self.assertTrue(blocker["restart_required"])
self.assertTrue(blocker["resolver_authoritative"])
def test_no_disagreement_when_resolver_agrees(self):
res = mp.assess_master_parity({"startup_head": SHA_A}, SHA_A)
self.assertIsNone(
mp.parity_resolver_disagreement(res, resolver_restart_required=False))
def test_live_stale_report_names_live_remote(self):
res = mp.assess_master_parity(
{"startup_head": SHA_A}, SHA_A, live_remote_head=SHA_B)
report = mp.parity_report(res)
self.assertEqual(report["live_remote_head"], SHA_B)
self.assertTrue(report["restart_required"])
class TestReadGitHead(unittest.TestCase): class TestReadGitHead(unittest.TestCase):
def test_test_override_takes_precedence(self): def test_test_override_takes_precedence(self):
with patch.dict(os.environ, {mp.ENV_TEST_CURRENT_HEAD: SHA_B}): with patch.dict(os.environ, {mp.ENV_TEST_CURRENT_HEAD: SHA_B}):
@@ -95,6 +184,149 @@ class TestReadGitHead(unittest.TestCase):
self.assertIsNone(mp.read_git_head("")) self.assertIsNone(mp.read_git_head(""))
class TestReadRemoteMasterHead(unittest.TestCase):
"""#610: live remote master head reader (env-overridable, fails to None)."""
def test_test_override_takes_precedence(self):
with patch.dict(os.environ, {mp.ENV_TEST_LIVE_REMOTE_HEAD: SHA_B}):
self.assertEqual(mp.read_remote_master_head("/nonexistent"), SHA_B)
def test_blank_override_is_none(self):
with patch.dict(os.environ, {mp.ENV_TEST_LIVE_REMOTE_HEAD: " "}):
self.assertIsNone(mp.read_remote_master_head("/nonexistent"))
def test_unfetchable_remote_is_none(self):
# No override; a bogus root/remote must fail closed to None, never raise.
env = {k: v for k, v in os.environ.items()
if k != mp.ENV_TEST_LIVE_REMOTE_HEAD}
with patch.dict(os.environ, env, clear=True):
self.assertIsNone(
mp.read_remote_master_head("/nonexistent", remote="nope"))
class TestRemoteHeadCache(unittest.TestCase):
"""#610: live remote reads are cached with a TTL to stay off the network.
The parity gate runs on every mutation and every runtime-context read, so an
unbounded ``git ls-remote`` per call would be a latency/flakiness regression.
"""
def setUp(self):
# These cases intentionally exercise the subprocess/cache path, so they
# opt out of suite-wide hermetic mode (PR #788 F1).
self._saved_hermetic = mp.hermetic_test_mode()
mp.set_hermetic_test_mode(False)
mp._clear_remote_head_cache()
env = {
k: v for k, v in os.environ.items()
if k not in (mp.ENV_TEST_LIVE_REMOTE_HEAD,
mp.ENV_TEST_ALLOW_LIVE_REMOTE_PROBE,
"PYTEST_CURRENT_TEST")
}
# Allow the probe path under hermetic defenses while still mocking
# subprocess so no real network call runs.
env[mp.ENV_TEST_ALLOW_LIVE_REMOTE_PROBE] = "1"
self._env = patch.dict(os.environ, env, clear=True)
self._env.start()
self.addCleanup(self._env.stop)
self.addCleanup(mp._clear_remote_head_cache)
self.addCleanup(
lambda: mp.set_hermetic_test_mode(self._saved_hermetic)
)
def _fake_run(self, sha):
class _R:
returncode = 0
stdout = f"{sha}\trefs/heads/master\n"
calls = {"n": 0}
def run(*args, **kwargs):
calls["n"] += 1
return _R()
return run, calls
def test_second_call_within_ttl_uses_cache(self):
run, calls = self._fake_run(SHA_B)
with patch.object(mp.subprocess, "run", run):
a = mp.read_remote_master_head("/repo", remote="prgs", ttl=100)
b = mp.read_remote_master_head("/repo", remote="prgs", ttl=100)
self.assertEqual(a, SHA_B)
self.assertEqual(b, SHA_B)
self.assertEqual(calls["n"], 1)
def test_zero_ttl_bypasses_cache(self):
run, calls = self._fake_run(SHA_B)
with patch.object(mp.subprocess, "run", run):
mp.read_remote_master_head("/repo", remote="prgs", ttl=0)
mp.read_remote_master_head("/repo", remote="prgs", ttl=0)
self.assertEqual(calls["n"], 2)
def test_env_override_never_touches_subprocess(self):
run, calls = self._fake_run(SHA_B)
with patch.dict(os.environ, {mp.ENV_TEST_LIVE_REMOTE_HEAD: SHA_A}):
with patch.object(mp.subprocess, "run", run):
self.assertEqual(
mp.read_remote_master_head("/repo", remote="prgs"), SHA_A)
self.assertEqual(calls["n"], 0)
class TestHermeticLiveRemoteReads(unittest.TestCase):
"""#610 / PR #788 F1/F2: suite hermetic mode never hits the network."""
def setUp(self):
self._saved = mp.hermetic_test_mode()
mp.set_hermetic_test_mode(True)
mp._clear_remote_head_cache()
self.addCleanup(lambda: mp.set_hermetic_test_mode(self._saved))
self.addCleanup(mp._clear_remote_head_cache)
def test_hermetic_mode_returns_none_without_subprocess(self):
run_calls = {"n": 0}
def boom(*args, **kwargs):
run_calls["n"] += 1
raise AssertionError("ls-remote must not run under hermetic mode")
env = {
k: v for k, v in os.environ.items()
if k not in (mp.ENV_TEST_LIVE_REMOTE_HEAD,
mp.ENV_TEST_ALLOW_LIVE_REMOTE_PROBE)
}
with patch.dict(os.environ, env, clear=True):
with patch.object(mp.subprocess, "run", boom):
self.assertIsNone(
mp.read_remote_master_head("/repo", remote="prgs")
)
self.assertEqual(run_calls["n"], 0)
def test_hermetic_mode_survives_clear_true_env(self):
"""Module flag, not env pin: clear=True cannot re-enable the probe."""
run_calls = {"n": 0}
def boom(*args, **kwargs):
run_calls["n"] += 1
raise AssertionError("ls-remote must not run after clear=True")
with patch.dict(os.environ, {}, clear=True):
with patch.object(mp.subprocess, "run", boom):
self.assertIsNone(mp.read_remote_master_head("/repo"))
self.assertEqual(run_calls["n"], 0)
def test_explicit_override_still_wins_under_hermetic(self):
run_calls = {"n": 0}
def boom(*args, **kwargs):
run_calls["n"] += 1
raise AssertionError("override must bypass subprocess")
with patch.dict(os.environ, {mp.ENV_TEST_LIVE_REMOTE_HEAD: SHA_B}):
with patch.object(mp.subprocess, "run", boom):
self.assertEqual(
mp.read_remote_master_head("/repo"), SHA_B
)
self.assertEqual(run_calls["n"], 0)
class TestServerWiring(unittest.TestCase): class TestServerWiring(unittest.TestCase):
"""Integration with the gate choke point in the server namespace.""" """Integration with the gate choke point in the server namespace."""
@@ -105,6 +337,13 @@ class TestServerWiring(unittest.TestCase):
self._saved = self.srv._STARTUP_PARITY self._saved = self.srv._STARTUP_PARITY
self.srv._STARTUP_PARITY = {"root": self.srv.PROJECT_ROOT, self.srv._STARTUP_PARITY = {"root": self.srv.PROJECT_ROOT,
"startup_head": SHA_A} "startup_head": SHA_A}
# Keep the live-remote read hermetic (no real ls-remote network call):
# default the live master to the daemon start so parity is fully green
# unless a test overrides the live head explicitly (#610).
self._live_patch = patch.dict(
os.environ, {mp.ENV_TEST_LIVE_REMOTE_HEAD: SHA_A})
self._live_patch.start()
self.addCleanup(self._live_patch.stop)
def tearDown(self): def tearDown(self):
self.srv._STARTUP_PARITY = self._saved self.srv._STARTUP_PARITY = self._saved
@@ -147,6 +386,36 @@ class TestServerWiring(unittest.TestCase):
self.assertTrue(out["in_parity"]) self.assertTrue(out["in_parity"])
self.assertNotIn("report", out) self.assertNotIn("report", out)
# --- #610: live-remote wiring -------------------------------------------
def test_live_stale_blocks_mutation_though_local_green(self):
# Local checkout matches the daemon start (local parity green) but the
# live remote master has advanced -> mutations must fail closed.
with patch.dict(os.environ, {mp.ENV_TEST_CURRENT_HEAD: SHA_A,
mp.ENV_TEST_LIVE_REMOTE_HEAD: SHA_B}):
self.assertEqual(self.srv._master_parity_block("gitea.read"), [])
self.assertTrue(
self.srv._master_parity_block("gitea.pr.create"))
def test_assess_tool_exposes_three_distinct_shas(self):
with patch.dict(os.environ, {mp.ENV_TEST_CURRENT_HEAD: SHA_A,
mp.ENV_TEST_LIVE_REMOTE_HEAD: SHA_B}):
out = self.srv.gitea_assess_master_parity(remote="prgs")
self.assertEqual(out["daemon_start_head"], SHA_A)
self.assertEqual(out["local_head"], SHA_A)
self.assertEqual(out["live_remote_head"], SHA_B)
self.assertTrue(out["live_stale"])
self.assertFalse(out["mutation_safe"])
self.assertIn("report", out)
def test_assess_tool_mutation_safe_when_all_three_match(self):
with patch.dict(os.environ, {mp.ENV_TEST_CURRENT_HEAD: SHA_A,
mp.ENV_TEST_LIVE_REMOTE_HEAD: SHA_A}):
out = self.srv.gitea_assess_master_parity(remote="prgs")
self.assertTrue(out["mutation_safe"])
self.assertFalse(out["live_stale"])
self.assertNotIn("report", out)
if __name__ == "__main__": if __name__ == "__main__":
unittest.main() unittest.main()
+336 -31
View File
@@ -1,4 +1,4 @@
"""Tests for web UI project registry (#427).""" """Tests for web UI project registry (#427) and its API evolution (#635)."""
import json import json
import sys import sys
import tempfile import tempfile
@@ -11,57 +11,246 @@ from starlette.testclient import TestClient
from webui.app import create_app from webui.app import create_app
from webui.project_registry import ( from webui.project_registry import (
CURRENT_SCHEMA_VERSION,
REGISTRY_API_VERSION,
SUPPORTED_SCHEMA_VERSIONS,
RegistryError,
default_registry_path, default_registry_path,
load_registry, load_registry,
onboarding_summary,
project_to_dict, project_to_dict,
) )
from webui.registry_safety import is_forbidden_key
_REPO_ROOT = Path(__file__).resolve().parent.parent
_API_DOC = _REPO_ROOT / "docs" / "webui-project-registry-api.md"
class TestProjectRegistryLoader(unittest.TestCase): def _valid_project(**overrides):
project = {
"id": "example",
"repo_name": "Example",
"gitea_owner": "Org",
"remote_host": "https://gitea.example.invalid",
"default_branch": "main",
"local_checkout_path": ".",
"profiles": {"author": "a", "reviewer": "r", "reconciler": "c"},
"workflow_paths": {"skill": "skills/x.md"},
}
project.update(overrides)
return project
def _write_registry(payload) -> Path:
with tempfile.NamedTemporaryFile("w", suffix=".json", delete=False) as handle:
json.dump(payload, handle)
return Path(handle.name)
class RegistryFileCase(unittest.TestCase):
"""Base class that cleans up temporary registry files."""
def setUp(self):
self._temp_paths: list[Path] = []
def tearDown(self):
for path in self._temp_paths:
path.unlink(missing_ok=True)
def write_registry(self, payload) -> Path:
path = _write_registry(payload)
self._temp_paths.append(path)
return path
class TestProjectRegistryLoader(RegistryFileCase):
def test_default_registry_loads_gitea_tools(self): def test_default_registry_loads_gitea_tools(self):
registry = load_registry() registry = load_registry()
self.assertEqual(registry.version, 1) self.assertEqual(registry.version, CURRENT_SCHEMA_VERSION)
self.assertEqual(registry.schema_version, CURRENT_SCHEMA_VERSION)
self.assertEqual(registry.api_version, REGISTRY_API_VERSION)
self.assertEqual(len(registry.projects), 1) self.assertEqual(len(registry.projects), 1)
project = registry.projects[0] project = registry.projects[0]
self.assertEqual(project.id, "gitea-tools") self.assertEqual(project.id, "gitea-tools")
self.assertEqual(project.repo_name, "Gitea-Tools") self.assertEqual(project.repo_name, "Gitea-Tools")
self.assertEqual(project.gitea_owner, "Scaled-Tech-Consulting") self.assertEqual(project.gitea_owner, "Scaled-Tech-Consulting")
self.assertEqual(project.repo_full_name, "Scaled-Tech-Consulting/Gitea-Tools")
self.assertEqual(project.remote_host, "https://gitea.prgs.cc") self.assertEqual(project.remote_host, "https://gitea.prgs.cc")
self.assertEqual(project.remote_name, "prgs")
self.assertEqual(project.status, "active")
self.assertEqual(project.profiles["author"], "prgs-author") self.assertEqual(project.profiles["author"], "prgs-author")
self.assertEqual(project.profiles["reviewer"], "prgs-reviewer") self.assertEqual(project.profiles["reviewer"], "prgs-reviewer")
self.assertEqual(project.profiles["reconciler"], "prgs-reconciler") self.assertEqual(project.profiles["reconciler"], "prgs-reconciler")
self.assertIn("skill", project.workflow_paths) self.assertIn("skill", project.workflow_paths)
self.assertGreaterEqual(len(project.onboarding_checklist), 4) self.assertGreaterEqual(len(project.onboarding_checklist), 4)
def test_registry_rejects_credential_keys(self): def test_default_registry_onboarding_summary_is_complete(self):
payload = { summary = onboarding_summary(load_registry().projects[0])
self.assertEqual(summary.total, summary.complete)
self.assertEqual(summary.required_outstanding, 0)
self.assertTrue(summary.onboarding_complete)
def test_version_1_registry_still_loads_with_defaults(self):
path = self.write_registry({
"version": 1, "version": 1,
"projects": [ "projects": [
{ _valid_project(
"id": "bad", onboarding_checklist=[
"repo_name": "Bad", {"id": "step", "title": "Step", "description": "Do it"}
"gitea_owner": "Org", ]
"remote_host": "https://gitea.example.invalid", )
"default_branch": "main",
"local_checkout_path": ".",
"profiles": {
"author": "a",
"reviewer": "r",
"reconciler": "c",
},
"workflow_paths": {"skill": "skills/x.md"},
"api_token": "secret",
}
], ],
} })
registry = load_registry(path)
self.assertEqual(registry.schema_version, 1)
self.assertIn(1, SUPPORTED_SCHEMA_VERSIONS)
project = registry.projects[0]
self.assertEqual(project.status, "active")
self.assertIsNone(project.remote_name)
self.assertIsNone(project.last_seen_health)
step = project.onboarding_checklist[0]
self.assertEqual(step.state, "pending")
self.assertTrue(step.required)
self.assertFalse(onboarding_summary(project).onboarding_complete)
def test_onboarding_summary_counts_states(self):
path = self.write_registry({
"version": 2,
"projects": [
_valid_project(
onboarding_checklist=[
{"id": "a", "title": "A", "description": "d", "state": "complete"},
{"id": "b", "title": "B", "description": "d", "state": "blocked"},
{
"id": "c",
"title": "C",
"description": "d",
"state": "pending",
"required": False,
},
{
"id": "d",
"title": "D",
"description": "d",
"state": "not_applicable",
},
]
)
],
})
summary = onboarding_summary(load_registry(path).projects[0])
self.assertEqual(summary.total, 4)
self.assertEqual(summary.complete, 1)
self.assertEqual(summary.blocked, 1)
self.assertEqual(summary.pending, 1)
self.assertEqual(summary.not_applicable, 1)
# Only the blocked step is both required and outstanding.
self.assertEqual(summary.required_outstanding, 1)
self.assertFalse(summary.onboarding_complete)
def test_last_seen_health_is_parsed_when_present(self):
path = self.write_registry({
"version": 2,
"projects": [
_valid_project(
last_seen_health={
"status": "degraded",
"checked_at": "2026-01-01T00:00:00Z",
"detail": "daemon restart pending",
}
)
],
})
health = load_registry(path).projects[0].last_seen_health
self.assertIsNotNone(health)
self.assertEqual(health.status, "degraded")
self.assertEqual(health.checked_at, "2026-01-01T00:00:00Z")
def test_registry_rejects_credential_keys(self):
path = self.write_registry({
"version": 1,
"projects": [_valid_project(id="bad", api_token="redacted-placeholder")],
})
with self.assertRaises(RegistryError) as ctx:
load_registry(path)
self.assertIn("credential", ctx.exception.remediation.lower())
self.assertEqual(ctx.exception.field_path, "projects[0].api_token")
def test_unsupported_version_fails_closed_with_remediation(self):
path = self.write_registry({"version": 99, "projects": [_valid_project()]})
with self.assertRaises(RegistryError) as ctx:
load_registry(path)
self.assertIn("unsupported registry version", ctx.exception.message)
self.assertIn(str(CURRENT_SCHEMA_VERSION), ctx.exception.remediation)
self.assertEqual(ctx.exception.field_path, "version")
def test_missing_required_field_fails_closed(self):
broken = _valid_project()
del broken["default_branch"]
path = self.write_registry({"version": 2, "projects": [broken]})
with self.assertRaises(RegistryError) as ctx:
load_registry(path)
self.assertIn("default_branch", ctx.exception.message)
self.assertEqual(ctx.exception.field_path, "projects[0]")
def test_unknown_status_fails_closed(self):
path = self.write_registry({
"version": 2,
"projects": [_valid_project(status="mystery")],
})
with self.assertRaises(RegistryError) as ctx:
load_registry(path)
self.assertEqual(ctx.exception.field_path, "projects[0].status")
self.assertIn("active", ctx.exception.remediation)
def test_unknown_onboarding_state_fails_closed(self):
path = self.write_registry({
"version": 2,
"projects": [
_valid_project(
onboarding_checklist=[
{"id": "a", "title": "A", "description": "d", "state": "almost"}
]
)
],
})
with self.assertRaises(RegistryError) as ctx:
load_registry(path)
self.assertEqual(
ctx.exception.field_path,
"projects[0].onboarding_checklist[0].state",
)
def test_missing_profile_role_fails_closed(self):
path = self.write_registry({
"version": 2,
"projects": [_valid_project(profiles={"author": "a", "reviewer": "r"})],
})
with self.assertRaises(RegistryError) as ctx:
load_registry(path)
self.assertEqual(ctx.exception.field_path, "projects[0].profiles.reconciler")
def test_empty_projects_fails_closed(self):
path = self.write_registry({"version": 2, "projects": []})
with self.assertRaises(RegistryError) as ctx:
load_registry(path)
self.assertEqual(ctx.exception.field_path, "projects")
def test_invalid_json_fails_closed_with_location(self):
with tempfile.NamedTemporaryFile("w", suffix=".json", delete=False) as handle: with tempfile.NamedTemporaryFile("w", suffix=".json", delete=False) as handle:
json.dump(payload, handle) handle.write("{not json")
path = Path(handle.name) path = Path(handle.name)
try: self._temp_paths.append(path)
with self.assertRaises(ValueError): with self.assertRaises(RegistryError) as ctx:
load_registry(path) load_registry(path)
finally: self.assertIn("not valid JSON", ctx.exception.message)
path.unlink(missing_ok=True) self.assertIn("line", ctx.exception.remediation)
def test_missing_file_fails_closed(self):
missing = Path(tempfile.gettempdir()) / "webui-registry-does-not-exist.json"
with self.assertRaises(RegistryError) as ctx:
load_registry(missing)
self.assertIn("could not be read", ctx.exception.message)
def test_default_registry_path_points_at_packaged_data(self): def test_default_registry_path_points_at_packaged_data(self):
path = default_registry_path() path = default_registry_path()
@@ -81,31 +270,147 @@ class TestProjectRegistryRoutes(unittest.TestCase):
self.assertIn("prgs-author", response.text) self.assertIn("prgs-author", response.text)
self.assertNotIn("child issue", response.text.lower()) self.assertNotIn("child issue", response.text.lower())
def test_projects_page_shows_status_and_progress(self):
response = self.client.get("/projects")
self.assertIn("Status", response.text)
self.assertIn("Onboarding", response.text)
self.assertIn("4/4 complete", response.text)
def test_project_detail_renders_checklist(self): def test_project_detail_renders_checklist(self):
response = self.client.get("/projects/gitea-tools") response = self.client.get("/projects/gitea-tools")
self.assertEqual(response.status_code, 200) self.assertEqual(response.status_code, 200)
self.assertIn("Onboarding checklist", response.text) self.assertIn("Onboarding checklist", response.text)
self.assertIn("Configure execution profiles", response.text) self.assertIn("Configure execution profiles", response.text)
self.assertIn("branches/", response.text) self.assertIn("branches/", response.text)
self.assertIn("Complete", response.text)
self.assertIn("required outstanding 0", response.text)
def test_project_detail_404(self): def test_project_detail_404(self):
response = self.client.get("/projects/unknown-repo") response = self.client.get("/projects/unknown-repo")
self.assertEqual(response.status_code, 404) self.assertEqual(response.status_code, 404)
def test_api_projects_json(self): def test_api_projects_alias_stays_compatible(self):
response = self.client.get("/api/projects") response = self.client.get("/api/projects")
self.assertEqual(response.status_code, 200) self.assertEqual(response.status_code, 200)
data = response.json() data = response.json()
self.assertEqual(data["version"], 1) # #427 consumers keep these keys.
self.assertEqual(data["version"], CURRENT_SCHEMA_VERSION)
self.assertIn("source_path", data)
self.assertEqual(len(data["projects"]), 1) self.assertEqual(len(data["projects"]), 1)
self.assertEqual(data["projects"][0]["id"], "gitea-tools") self.assertEqual(data["projects"][0]["id"], "gitea-tools")
self.assertIn("onboarding_checklist", data["projects"][0]) self.assertIn("onboarding_checklist", data["projects"][0])
def test_api_v1_projects_payload(self):
response = self.client.get("/api/v1/projects")
self.assertEqual(response.status_code, 200)
data = response.json()
self.assertEqual(data["api_version"], REGISTRY_API_VERSION)
self.assertEqual(data["schema_version"], CURRENT_SCHEMA_VERSION)
self.assertEqual(data["project_count"], 1)
self.assertEqual(data["source"]["kind"], "file")
self.assertTrue(data["source"]["inventory_complete"])
project = data["projects"][0]
self.assertEqual(project["status"], "active")
self.assertEqual(project["remote_name"], "prgs")
self.assertEqual(
project["repo_full_name"], "Scaled-Tech-Consulting/Gitea-Tools"
)
self.assertTrue(project["onboarding_summary"]["onboarding_complete"])
self.assertEqual(project["onboarding_checklist"][0]["state"], "complete")
self.assertIsNone(project["last_seen_health"])
def test_api_v1_project_detail(self):
response = self.client.get("/api/v1/projects/gitea-tools")
self.assertEqual(response.status_code, 200)
data = response.json()
self.assertEqual(data["api_version"], REGISTRY_API_VERSION)
self.assertEqual(data["project"]["id"], "gitea-tools")
self.assertEqual(data["source"]["kind"], "file")
def test_api_v1_project_detail_missing_fails_closed(self):
response = self.client.get("/api/v1/projects/not-registered")
self.assertEqual(response.status_code, 404)
data = response.json()
self.assertEqual(data["error"], "project_not_found")
self.assertEqual(data["project_id"], "not-registered")
self.assertIn("gitea-tools", data["known_project_ids"])
self.assertIn("remediation", data)
def test_api_v1_projects_is_read_only(self):
response = self.client.post("/api/v1/projects", json={})
self.assertEqual(response.status_code, 405)
self.assertEqual(response.json()["error"], "read-only-mvp")
def test_project_to_dict_is_json_safe(self): def test_project_to_dict_is_json_safe(self):
registry = load_registry() registry = load_registry()
encoded = json.dumps(project_to_dict(registry.projects[0])) dto = project_to_dict(registry.projects[0])
encoded = json.dumps(dto)
self.assertIn("gitea-tools", encoded) self.assertIn("gitea-tools", encoded)
# Prose may mention tokens; no serialized *key* may look like a secret.
for key in dto:
with self.subTest(key=key):
self.assertFalse(is_forbidden_key(key))
class TestInvalidRegistryFailsClosedOverHttp(RegistryFileCase):
def setUp(self):
super().setUp()
self.path = self.write_registry({"version": 42, "projects": []})
self.client = TestClient(create_app())
def _with_bad_registry(self, url: str):
import os
from unittest import mock
with mock.patch.dict(
os.environ, {"WEBUI_PROJECT_REGISTRY": str(self.path)}, clear=False
):
return self.client.get(url)
def test_api_v1_reports_actionable_error(self):
response = self._with_bad_registry("/api/v1/projects")
self.assertEqual(response.status_code, 500)
data = response.json()
self.assertEqual(data["error"], "registry_invalid")
self.assertIn("unsupported registry version", data["detail"])
self.assertTrue(data["remediation"])
self.assertEqual(data["field_path"], "version")
def test_unversioned_alias_reports_actionable_error(self):
response = self._with_bad_registry("/api/projects")
self.assertEqual(response.status_code, 500)
self.assertEqual(response.json()["error"], "registry_invalid")
def test_html_page_reports_actionable_error(self):
response = self._with_bad_registry("/projects")
self.assertEqual(response.status_code, 500)
self.assertIn("Project registry unavailable", response.text)
self.assertIn("Remediation", response.text)
class TestProjectRegistryApiDocs(unittest.TestCase):
def test_api_contract_is_documented(self):
self.assertTrue(_API_DOC.is_file(), f"missing {_API_DOC}")
text = _API_DOC.read_text(encoding="utf-8")
for token in (
"/api/v1/projects",
"/api/v1/projects/{project_id}",
"/api/projects",
"onboarding_summary",
"last_seen_health",
"registry_invalid",
"#635",
):
with self.subTest(token=token):
self.assertIn(token, text)
def test_route_table_lists_versioned_routes(self):
local_dev = (_REPO_ROOT / "docs" / "webui-local-dev.md").read_text(
encoding="utf-8"
)
self.assertIn("/api/v1/projects", local_dev)
self.assertIn("webui-project-registry-api.md", local_dev)
if __name__ == "__main__": if __name__ == "__main__":
unittest.main() unittest.main()
+458
View File
@@ -0,0 +1,458 @@
"""Tests for the worker registry and configuration schema (#798, epic #797)."""
import json
import sys
import tempfile
import unittest
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from webui.worker_registry import (
ALLOWED_ROLES,
SCHEMA_VERSION,
RegistryValidationError,
WorkerRegistry,
default_registry_path,
find_provider,
find_worker,
history_dir,
list_revisions,
load_registry,
registry_to_dict,
registry_to_document,
rollback_to_revision,
save_registry,
validate_payload,
worker_to_dict,
workers_for_provider,
)
_EXPECTED_PROVIDER_IDS = ("claude", "grok", "codex", "agy", "kimi-k")
def _provider(provider_id: str = "claude", **overrides) -> dict:
payload = {
"id": provider_id,
"display_name": "Claude",
"vendor": "Anthropic",
"executable": "claude",
"available": True,
"models": ["claude-opus-4-8"],
"notes": "",
}
payload.update(overrides)
return payload
def _worker(worker_id: str = "claude-author", **overrides) -> dict:
payload = {
"id": worker_id,
"display_name": "Claude author",
"provider": "claude",
"model": "claude-opus-4-8",
"project": "gitea-tools",
"role": "author",
"namespace": "gitea-author",
"profile": "prgs-author",
"workflow": "skills/llm-project-workflow/workflows/work-issue.md",
"schedule": {"kind": "cron", "expression": "0 * * * *"},
"timeout_seconds": 3600,
"enabled": True,
"scheduler": {"kind": "launchd", "label": "cc.prgs.claude.author"},
"notes": "",
}
payload.update(overrides)
return payload
def _document(providers=None, workers=None, **overrides) -> dict:
payload = {
"version": SCHEMA_VERSION,
"revision": 1,
"updated_at": "2026-07-22T00:00:00Z",
"providers": providers if providers is not None else [_provider()],
"workers": workers if workers is not None else [_worker()],
}
payload.update(overrides)
return payload
class _TempRegistryCase(unittest.TestCase):
"""Base case giving each test an isolated registry file."""
def setUp(self):
self._tmp = tempfile.TemporaryDirectory()
self.addCleanup(self._tmp.cleanup)
self.path = Path(self._tmp.name) / "workers.registry.json"
def write(self, document: dict) -> Path:
self.path.write_text(json.dumps(document, indent=2) + "\n", encoding="utf-8")
return self.path
def parse(self, document: dict) -> WorkerRegistry:
return validate_payload(document, source_path=self.path)
class TestPackagedRegistry(unittest.TestCase):
"""AC: the declarative registry is the source of truth and ships with the app."""
def test_default_path_points_at_packaged_data(self):
path = default_registry_path()
self.assertEqual(path.name, "workers.registry.json")
self.assertEqual(path.parent.name, "data")
def test_packaged_registry_loads_and_validates(self):
registry = load_registry()
self.assertEqual(registry.version, SCHEMA_VERSION)
self.assertGreaterEqual(registry.revision, 1)
def test_packaged_registry_declares_all_five_providers(self):
registry = load_registry()
self.assertEqual(
tuple(provider.id for provider in registry.providers),
_EXPECTED_PROVIDER_IDS,
)
def test_packaged_registry_carries_no_credentials(self):
raw = default_registry_path().read_text(encoding="utf-8").lower()
for marker in ("token", "password", "secret", "api_key", "credential"):
self.assertNotIn(marker, raw)
class TestSeparateEntities(_TempRegistryCase):
"""AC: providers and configured workers are separate entities."""
def test_provider_may_exist_with_no_workers(self):
registry = self.parse(
_document(providers=[_provider("grok", display_name="Grok")], workers=[])
)
self.assertEqual(len(registry.providers), 1)
self.assertEqual(registry.workers, ())
self.assertEqual(workers_for_provider(registry, "grok"), ())
def test_many_workers_may_share_one_provider(self):
registry = self.parse(
_document(
workers=[
_worker("claude-author"),
_worker(
"claude-reviewer",
role="reviewer",
namespace="gitea-reviewer",
profile="prgs-reviewer",
scheduler={"kind": "launchd", "label": "cc.prgs.claude.reviewer"},
),
]
)
)
self.assertEqual(len(workers_for_provider(registry, "claude")), 2)
self.assertEqual(len(registry.providers), 1)
def test_worker_referencing_unknown_provider_is_refused(self):
with self.assertRaises(RegistryValidationError) as ctx:
self.parse(_document(workers=[_worker(provider="mystery")]))
self.assertIn("unknown provider", str(ctx.exception))
def test_lookup_helpers(self):
registry = self.parse(_document())
self.assertIsNotNone(find_worker(registry, "claude-author"))
self.assertIsNone(find_worker(registry, "absent"))
self.assertIsNotNone(find_provider(registry, "claude"))
self.assertIsNone(find_provider(registry, "absent"))
class TestRecordedFields(_TempRegistryCase):
"""AC: records provider, model, project, role, namespace/profile, workflow,
schedule, timeout, enabled state, and scheduler metadata."""
def test_every_required_field_is_recorded(self):
registry = self.parse(_document())
worker = registry.workers[0]
self.assertEqual(worker.provider, "claude")
self.assertEqual(worker.model, "claude-opus-4-8")
self.assertEqual(worker.project, "gitea-tools")
self.assertEqual(worker.role, "author")
self.assertEqual(worker.namespace, "gitea-author")
self.assertEqual(worker.profile, "prgs-author")
self.assertEqual(worker.workflow, "skills/llm-project-workflow/workflows/work-issue.md")
self.assertEqual(worker.schedule.kind, "cron")
self.assertEqual(worker.schedule.expression, "0 * * * *")
self.assertEqual(worker.timeout_seconds, 3600)
self.assertTrue(worker.enabled)
self.assertEqual(worker.scheduler.kind, "launchd")
self.assertEqual(worker.scheduler.label, "cc.prgs.claude.author")
def test_each_required_field_is_individually_required(self):
for field in (
"provider", "model", "project", "role", "namespace",
"profile", "workflow", "schedule", "timeout_seconds",
"enabled", "scheduler", "id", "display_name",
):
with self.subTest(field=field):
worker = _worker()
worker.pop(field)
with self.assertRaises(RegistryValidationError):
self.parse(_document(workers=[worker]))
def test_all_sanctioned_roles_are_accepted(self):
for role in ALLOWED_ROLES:
with self.subTest(role=role):
registry = self.parse(_document(workers=[_worker(role=role)]))
self.assertEqual(registry.workers[0].role, role)
def test_unsanctioned_role_is_refused(self):
with self.assertRaises(RegistryValidationError) as ctx:
self.parse(_document(workers=[_worker(role="admin")]))
self.assertIn("role must be one of", str(ctx.exception))
def test_worker_dict_round_trips_every_field(self):
registry = self.parse(_document())
encoded = worker_to_dict(registry.workers[0])
self.assertEqual(encoded, _worker())
json.dumps(encoded) # must stay JSON-safe for the #799 API
class TestSchemaValidation(_TempRegistryCase):
"""AC: supports schema validation — and fails closed."""
def test_unsupported_version_is_refused(self):
with self.assertRaises(RegistryValidationError):
self.parse(_document(version=2))
def test_root_must_be_an_object(self):
with self.assertRaises(RegistryValidationError):
validate_payload([], source_path=self.path)
def test_providers_must_be_non_empty(self):
with self.assertRaises(RegistryValidationError):
self.parse(_document(providers=[]))
def test_unknown_top_level_field_is_refused(self):
with self.assertRaises(RegistryValidationError) as ctx:
self.parse(_document(fleet=[]))
self.assertIn("unknown fields", str(ctx.exception))
def test_unknown_worker_field_is_refused_not_ignored(self):
# A typo'd field must not be silently dropped: "timeout_second" would
# otherwise read as "no timeout declared".
worker = _worker()
worker["timeout_second"] = 30
with self.assertRaises(RegistryValidationError) as ctx:
self.parse(_document(workers=[worker]))
self.assertIn("timeout_second", str(ctx.exception))
def test_credentials_are_refused_anywhere_in_the_document(self):
for label, mutate in (
("provider.api_token", lambda doc: doc["providers"][0].__setitem__("api_token", "x")),
("worker.password", lambda doc: doc["workers"][0].__setitem__("password", "x")),
("root.secret", lambda doc: doc.__setitem__("secret", "x")),
):
with self.subTest(field=label):
document = _document()
mutate(document)
with self.assertRaises(ValueError) as ctx:
self.parse(document)
self.assertIn("credential", str(ctx.exception).lower())
def test_duplicate_worker_id_is_refused(self):
workers = [_worker("dup"), _worker("dup", scheduler={"kind": "manual"})]
with self.assertRaises(RegistryValidationError) as ctx:
self.parse(_document(workers=workers))
self.assertIn("duplicate worker id", str(ctx.exception))
def test_duplicate_provider_id_is_refused(self):
with self.assertRaises(RegistryValidationError) as ctx:
self.parse(_document(providers=[_provider("claude"), _provider("claude")], workers=[]))
self.assertIn("duplicate provider id", str(ctx.exception))
def test_duplicate_launchagent_label_is_refused(self):
# Two workers sharing a label would silently overwrite each other's agent.
workers = [
_worker("a", scheduler={"kind": "launchd", "label": "cc.prgs.same"}),
_worker("b", scheduler={"kind": "launchd", "label": "cc.prgs.same"}),
]
with self.assertRaises(RegistryValidationError) as ctx:
self.parse(_document(workers=workers))
self.assertIn("duplicate scheduler label", str(ctx.exception))
def test_manual_scheduler_needs_no_label_and_many_may_coexist(self):
workers = [
_worker("a", scheduler={"kind": "manual"}),
_worker("b", scheduler={"kind": "manual"}),
]
registry = self.parse(_document(workers=workers))
self.assertEqual([w.scheduler.label for w in registry.workers], [None, None])
def test_launchd_scheduler_requires_a_label(self):
with self.assertRaises(RegistryValidationError) as ctx:
self.parse(_document(workers=[_worker(scheduler={"kind": "launchd"})]))
self.assertIn("label is required", str(ctx.exception))
def test_unknown_scheduler_kind_is_refused(self):
with self.assertRaises(RegistryValidationError):
self.parse(_document(workers=[_worker(scheduler={"kind": "systemd", "label": "x"})]))
def test_timeout_must_be_a_positive_bounded_integer(self):
for bad in (0, -1, "3600", 1.5, True, 86_401):
with self.subTest(timeout=bad):
with self.assertRaises(RegistryValidationError):
self.parse(_document(workers=[_worker(timeout_seconds=bad)]))
def test_enabled_must_be_a_real_boolean(self):
for bad in ("true", 1, None):
with self.subTest(enabled=bad):
with self.assertRaises(RegistryValidationError):
self.parse(_document(workers=[_worker(enabled=bad)]))
def test_identifier_shape_is_enforced(self):
for bad in ("Claude Author", "-leading", "UPPER", ""):
with self.subTest(worker_id=bad):
with self.assertRaises(RegistryValidationError):
self.parse(_document(workers=[_worker(bad)]))
class TestScheduleValidation(_TempRegistryCase):
"""Schedules are declarations; next-run computation belongs to #803."""
def test_interval_schedule_requires_positive_seconds(self):
registry = self.parse(
_document(workers=[_worker(schedule={"kind": "interval", "seconds": 900})])
)
self.assertEqual(registry.workers[0].schedule.seconds, 900)
with self.assertRaises(RegistryValidationError):
self.parse(_document(workers=[_worker(schedule={"kind": "interval"})]))
with self.assertRaises(RegistryValidationError):
self.parse(_document(workers=[_worker(schedule={"kind": "interval", "seconds": 0})]))
def test_cron_schedule_requires_five_fields(self):
with self.assertRaises(RegistryValidationError) as ctx:
self.parse(_document(workers=[_worker(schedule={"kind": "cron", "expression": "0 *"})]))
self.assertIn("five crontab fields", str(ctx.exception))
def test_manual_schedule_needs_no_timing(self):
registry = self.parse(_document(workers=[_worker(schedule={"kind": "manual"})]))
schedule = registry.workers[0].schedule
self.assertEqual(schedule.kind, "manual")
self.assertIsNone(schedule.seconds)
self.assertIsNone(schedule.expression)
def test_fields_from_the_wrong_kind_are_refused(self):
with self.assertRaises(RegistryValidationError) as ctx:
self.parse(_document(workers=[_worker(schedule={"kind": "manual", "seconds": 60})]))
self.assertIn("not valid for kind", str(ctx.exception))
def test_unknown_schedule_kind_is_refused(self):
with self.assertRaises(RegistryValidationError):
self.parse(_document(workers=[_worker(schedule={"kind": "hourly"})]))
class TestAtomicPersistence(_TempRegistryCase):
"""AC: atomic persistence."""
def test_save_then_load_round_trips(self):
registry = self.parse(_document())
save_registry(registry, self.path)
reloaded = load_registry(self.path)
self.assertEqual(
[worker_to_dict(w) for w in reloaded.workers],
[worker_to_dict(w) for w in registry.workers],
)
def test_save_leaves_no_temp_files_behind(self):
registry = self.parse(_document())
save_registry(registry, self.path)
save_registry(registry, self.path)
leftovers = [p.name for p in self.path.parent.iterdir() if p.name.startswith(".")]
self.assertEqual(leftovers, [])
def test_save_refuses_to_persist_an_invalid_document(self):
registry = self.parse(_document())
broken = WorkerRegistry(
version=registry.version,
revision=registry.revision,
updated_at=registry.updated_at,
providers=registry.providers,
# A worker whose provider is not declared in the registry.
workers=tuple(
type(worker)(**{**worker.__dict__, "provider": "vanished"})
for worker in registry.workers
),
source_path=self.path,
)
with self.assertRaises(RegistryValidationError):
save_registry(broken, self.path)
self.assertFalse(self.path.exists(), "invalid save must not create the file")
def test_document_shape_excludes_local_paths_but_api_shape_includes_it(self):
registry = self.parse(_document())
self.assertNotIn("source_path", registry_to_document(registry))
self.assertEqual(registry_to_dict(registry)["source_path"], str(self.path))
class TestVersioningAndRollback(_TempRegistryCase):
"""AC: versioning and rollback."""
def _seed(self) -> WorkerRegistry:
self.write(_document())
return load_registry(self.path)
def test_revision_increments_on_each_save(self):
registry = self._seed()
self.assertEqual(registry.revision, 1)
second = save_registry(registry, self.path)
self.assertEqual(second.revision, 2)
third = save_registry(second, self.path)
self.assertEqual(third.revision, 3)
def test_updated_at_is_refreshed_and_utc(self):
registry = self._seed()
saved = save_registry(registry, self.path)
self.assertRegex(saved.updated_at, r"^\d{4}-\d{2}-\d{2}T\d{2}:\d{2}:\d{2}Z$")
def test_superseded_revisions_are_retained(self):
registry = self._seed()
second = save_registry(registry, self.path)
save_registry(second, self.path)
self.assertEqual(list_revisions(self.path), (1, 2))
self.assertTrue(history_dir(self.path).is_dir())
def test_rollback_restores_prior_content_as_a_new_revision(self):
self.write(_document(workers=[_worker("original")]))
registry = load_registry(self.path)
changed = WorkerRegistry(
version=registry.version,
revision=registry.revision,
updated_at=registry.updated_at,
providers=registry.providers,
workers=(), # operator deletes every worker
source_path=self.path,
)
save_registry(changed, self.path)
self.assertEqual(load_registry(self.path).workers, ())
restored = rollback_to_revision(1, self.path)
self.assertEqual([w.id for w in restored.workers], ["original"])
# Append-only: the rollback publishes a new head rather than rewinding.
self.assertGreater(restored.revision, 2)
self.assertEqual([w.id for w in load_registry(self.path).workers], ["original"])
def test_rollback_to_unknown_revision_fails_closed(self):
self._seed()
with self.assertRaises(RegistryValidationError) as ctx:
rollback_to_revision(99, self.path)
self.assertIn("not retained", str(ctx.exception))
def test_revision_must_be_a_positive_integer(self):
for bad in (0, -1, "1", None):
with self.subTest(revision=bad):
with self.assertRaises(RegistryValidationError):
self.parse(_document(revision=bad))
def test_history_is_empty_before_any_save(self):
self.write(_document())
self.assertEqual(list_revisions(self.path), ())
if __name__ == "__main__":
unittest.main()
+72 -5
View File
@@ -11,8 +11,20 @@ from starlette.routing import Route
from webui.deployment_boundary import deployment_snapshot from webui.deployment_boundary import deployment_snapshot
from webui.layout import render_page from webui.layout import render_page
from webui.project_registry import find_project, load_registry, registry_to_dict from webui.project_registry import (
from webui.project_views import render_project_detail, render_projects_list ProjectRegistry,
RegistryError,
find_project,
known_project_ids,
load_registry,
project_detail_to_dict,
registry_to_dict,
)
from webui.project_views import (
render_project_detail,
render_projects_list,
render_registry_error,
)
from webui.prompt_library import find_prompt, library_to_dict from webui.prompt_library import find_prompt, library_to_dict
from webui.prompt_views import render_prompt_detail, render_prompts_page from webui.prompt_views import render_prompt_detail, render_prompts_page
from final_report_validator import FINAL_REPORT_TASK_KINDS from final_report_validator import FINAL_REPORT_TASK_KINDS
@@ -81,14 +93,26 @@ async def api_queue(_request: Request) -> JSONResponse:
return JSONResponse(queue_snapshot_to_dict(load_queue_snapshot())) return JSONResponse(queue_snapshot_to_dict(load_queue_snapshot()))
def _load_project_registry() -> tuple[ProjectRegistry | None, RegistryError | None]:
"""Load the registry, converting validation failure into a fail-closed pair."""
try:
return load_registry(), None
except RegistryError as exc:
return None, exc
async def projects(_request: Request) -> HTMLResponse: async def projects(_request: Request) -> HTMLResponse:
registry = load_registry() registry, error = _load_project_registry()
if error is not None:
return HTMLResponse(render_registry_error(error), status_code=500)
return HTMLResponse(render_projects_list(registry)) return HTMLResponse(render_projects_list(registry))
async def project_detail(request: Request) -> HTMLResponse: async def project_detail(request: Request) -> HTMLResponse:
project_id = request.path_params["project_id"] project_id = request.path_params["project_id"]
registry = load_registry() registry, error = _load_project_registry()
if error is not None:
return HTMLResponse(render_registry_error(error), status_code=500)
project = find_project(registry, project_id) project = find_project(registry, project_id)
if project is None: if project is None:
return HTMLResponse( return HTMLResponse(
@@ -106,10 +130,47 @@ async def project_detail(request: Request) -> HTMLResponse:
async def api_projects(_request: Request) -> JSONResponse: async def api_projects(_request: Request) -> JSONResponse:
registry = load_registry() """Unversioned MVP alias, retained through Phase 1 (#632 section 6)."""
registry, error = _load_project_registry()
if error is not None:
return JSONResponse(error.to_dict(), status_code=500)
return JSONResponse(registry_to_dict(registry)) return JSONResponse(registry_to_dict(registry))
async def api_v1_projects(_request: Request) -> JSONResponse:
registry, error = _load_project_registry()
if error is not None:
return JSONResponse(error.to_dict(), status_code=500)
return JSONResponse(registry_to_dict(registry))
async def api_v1_project_detail(request: Request) -> JSONResponse:
project_id = request.path_params["project_id"]
registry, error = _load_project_registry()
if error is not None:
return JSONResponse(error.to_dict(), status_code=500)
project = find_project(registry, project_id)
if project is None:
return JSONResponse(
{
"error": "project_not_found",
"project_id": project_id,
"known_project_ids": known_project_ids(registry),
"remediation": (
"Request one of the known project ids, or add the project to the "
"registry file named in 'source'."
),
"source": {
"kind": "file",
"path": str(registry.source_path),
"inventory_complete": True,
},
},
status_code=404,
)
return JSONResponse(project_detail_to_dict(registry, project))
async def prompts(_request: Request) -> HTMLResponse: async def prompts(_request: Request) -> HTMLResponse:
return HTMLResponse(render_prompts_page()) return HTMLResponse(render_prompts_page())
@@ -268,6 +329,12 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
Route("/projects", projects, methods=["GET"]), Route("/projects", projects, methods=["GET"]),
Route("/projects/{project_id}", project_detail, methods=["GET"]), Route("/projects/{project_id}", project_detail, methods=["GET"]),
Route("/api/projects", api_projects, methods=["GET"]), Route("/api/projects", api_projects, methods=["GET"]),
Route("/api/v1/projects", api_v1_projects, methods=["GET"]),
Route(
"/api/v1/projects/{project_id}",
api_v1_project_detail,
methods=["GET"],
),
Route("/prompts", prompts, methods=["GET"]), Route("/prompts", prompts, methods=["GET"]),
Route("/prompts/{prompt_id}", prompt_detail, methods=["GET"]), Route("/prompts/{prompt_id}", prompt_detail, methods=["GET"]),
Route("/api/prompts", api_prompts, methods=["GET"]), Route("/api/prompts", api_prompts, methods=["GET"]),
+16 -6
View File
@@ -1,13 +1,15 @@
{ {
"version": 1, "version": 2,
"projects": [ "projects": [
{ {
"id": "gitea-tools", "id": "gitea-tools",
"repo_name": "Gitea-Tools", "repo_name": "Gitea-Tools",
"gitea_owner": "Scaled-Tech-Consulting", "gitea_owner": "Scaled-Tech-Consulting",
"remote_name": "prgs",
"remote_host": "https://gitea.prgs.cc", "remote_host": "https://gitea.prgs.cc",
"default_branch": "master", "default_branch": "master",
"local_checkout_path": ".", "local_checkout_path": ".",
"status": "active",
"profiles": { "profiles": {
"author": "prgs-author", "author": "prgs-author",
"reviewer": "prgs-reviewer", "reviewer": "prgs-reviewer",
@@ -26,24 +28,32 @@
{ {
"id": "profiles", "id": "profiles",
"title": "Configure execution profiles", "title": "Configure execution profiles",
"description": "Install author, reviewer, and reconciler MCP profiles (prgs-author, prgs-reviewer, prgs-reconciler) in separate namespaces. Tokens stay in keychain — never in this registry." "description": "Install author, reviewer, and reconciler MCP profiles (prgs-author, prgs-reviewer, prgs-reconciler) in separate namespaces. Tokens stay in keychain — never in this registry.",
"state": "complete",
"required": true
}, },
{ {
"id": "mcp_config", "id": "mcp_config",
"title": "Wire MCP v2 contexts", "title": "Wire MCP v2 contexts",
"description": "Copy and customize gitea-mcp.v2-contexts.example.json for your machine. Map this repo path under projects with default_owner Scaled-Tech-Consulting and default_repo Gitea-Tools." "description": "Copy and customize gitea-mcp.v2-contexts.example.json for your machine. Map this repo path under projects with default_owner Scaled-Tech-Consulting and default_repo Gitea-Tools.",
"state": "complete",
"required": true
}, },
{ {
"id": "wiki_gate", "id": "wiki_gate",
"title": "Wiki publication readiness", "title": "Wiki publication readiness",
"description": "For wiki-tracked work, satisfy the live Gitea Wiki proof gate (#224) before closing issues. See docs/wiki/Safety-and-Gates.md." "description": "For wiki-tracked work, satisfy the live Gitea Wiki proof gate (#224) before closing issues. See docs/wiki/Safety-and-Gates.md.",
"state": "complete",
"required": true
}, },
{ {
"id": "branches_layout", "id": "branches_layout",
"title": "Isolate work under branches/", "title": "Isolate work under branches/",
"description": "All LLM task edits happen in worktrees under branches/. Main checkout stays clean; use skills/llm-project-workflow templates for start-issue and review flows." "description": "All LLM task edits happen in worktrees under branches/. Main checkout stays clean; use skills/llm-project-workflow templates for start-issue and review flows.",
"state": "complete",
"required": true
} }
] ]
} }
] ]
} }
+57
View File
@@ -0,0 +1,57 @@
{
"version": 1,
"revision": 1,
"updated_at": "2026-07-22T00:00:00Z",
"providers": [
{
"id": "claude",
"display_name": "Claude",
"vendor": "Anthropic",
"executable": "claude",
"available": true,
"models": [
"claude-opus-4-8",
"claude-sonnet-5",
"claude-haiku-4-5-20251001"
],
"notes": "Model list is a declaration. Live enumeration and version inspection belong to the provider adapter framework (#800)."
},
{
"id": "grok",
"display_name": "Grok",
"vendor": "xAI",
"executable": "grok",
"available": true,
"models": [],
"notes": "Models enumerated by the provider adapter (#800); not declared here."
},
{
"id": "codex",
"display_name": "Codex",
"vendor": "OpenAI",
"executable": "codex",
"available": true,
"models": [],
"notes": "Models enumerated by the provider adapter (#800); not declared here."
},
{
"id": "agy",
"display_name": "AGY",
"vendor": "Antigravity",
"executable": "agy",
"available": true,
"models": [],
"notes": "MCP allowlist gating applies to this provider; confirm server-side allowlist before configuring a worker."
},
{
"id": "kimi-k",
"display_name": "Kimi K",
"vendor": "Moonshot AI",
"executable": "kimi",
"available": true,
"models": [],
"notes": "Provider id is kimi-k; the executable on PATH is kimi. Models enumerated by the provider adapter (#800)."
}
],
"workers": []
}
+451 -48
View File
@@ -1,4 +1,16 @@
"""Load and validate the web UI project registry (#427).""" """Load and validate the web UI project registry (#427, evolved for #635).
Phase 1 of the console architecture ADR keeps this loader read-only. It owns
the versioned project registry contract served at ``/api/v1/projects``:
* the on-disk file carries a ``version`` (schema version 1 or 2);
* version 1 files stay loadable and are normalized with explicit defaults, so
an operator registry written for #427 keeps working;
* every validation failure raises :class:`RegistryError`, which carries an
actionable ``remediation`` string instead of leaking a traceback;
* serialization never emits credentials — credential-shaped keys are rejected
at load time, before any DTO is built.
"""
from __future__ import annotations from __future__ import annotations
@@ -8,27 +20,59 @@ from dataclasses import dataclass
from pathlib import Path from pathlib import Path
from typing import Any from typing import Any
_FORBIDDEN_EXACT_KEYS = frozenset({ from webui.registry_safety import is_forbidden_key
"token",
"password",
"secret",
"credential",
"auth",
"api_key",
"api-key",
})
_FORBIDDEN_KEY_PREFIXES = ("auth_", "api_key_", "api-key_")
_FORBIDDEN_KEY_SUFFIXES = ("_token", "_secret", "_password", "_credential", "_auth")
#: Version of the JSON contract served under ``/api/v1/...``.
REGISTRY_API_VERSION = "v1"
#: Schema version written by this repository's packaged registry.
CURRENT_SCHEMA_VERSION = 2
#: Schema versions this loader accepts. Version 1 is normalized on load.
SUPPORTED_SCHEMA_VERSIONS = (1, 2)
#: Lifecycle state of a registered project.
PROJECT_STATUSES = ("active", "onboarding", "paused", "archived")
_DEFAULT_PROJECT_STATUS = "active"
#: Completion state of a single onboarding step.
ONBOARDING_STATES = ("complete", "pending", "blocked", "not_applicable")
_DEFAULT_ONBOARDING_STATE = "pending"
#: Redacted, last-seen health of a project's control plane.
HEALTH_STATUSES = ("healthy", "degraded", "unreachable", "unknown")
class RegistryError(ValueError):
"""A registry file could not be loaded or failed validation.
Carries an operator-facing ``remediation`` so routes can fail closed with
an actionable message rather than a stack trace.
"""
def __init__(
self,
message: str,
*,
remediation: str,
source_path: Path | None = None,
field_path: str | None = None,
) -> None:
super().__init__(message)
self.message = message
self.remediation = remediation
self.source_path = source_path
self.field_path = field_path
def to_dict(self) -> dict[str, Any]:
"""Serialize for a fail-closed JSON error response."""
return {
"error": "registry_invalid",
"detail": self.message,
"remediation": self.remediation,
"field_path": self.field_path,
"source_path": str(self.source_path) if self.source_path else None,
}
def _is_forbidden_key(key: str) -> bool:
lowered = key.lower()
if lowered in _FORBIDDEN_EXACT_KEYS:
return True
return (
lowered.startswith(_FORBIDDEN_KEY_PREFIXES)
or lowered.endswith(_FORBIDDEN_KEY_SUFFIXES)
)
_REQUIRED_PROJECT_FIELDS = ( _REQUIRED_PROJECT_FIELDS = (
"id", "id",
@@ -49,6 +93,30 @@ class OnboardingStep:
id: str id: str
title: str title: str
description: str description: str
state: str = _DEFAULT_ONBOARDING_STATE
required: bool = True
@dataclass(frozen=True)
class OnboardingSummary:
"""Aggregate onboarding progress for a single project."""
total: int
complete: int
pending: int
blocked: int
not_applicable: int
required_outstanding: int
onboarding_complete: bool
@dataclass(frozen=True)
class ProjectHealth:
"""Redacted last-seen health. Never carries endpoints or credentials."""
status: str
checked_at: str | None
detail: str | None
@dataclass(frozen=True) @dataclass(frozen=True)
@@ -63,6 +131,13 @@ class ProjectRecord:
workflow_paths: dict[str, str] workflow_paths: dict[str, str]
schema_paths: dict[str, str] schema_paths: dict[str, str]
onboarding_checklist: tuple[OnboardingStep, ...] onboarding_checklist: tuple[OnboardingStep, ...]
status: str = _DEFAULT_PROJECT_STATUS
remote_name: str | None = None
last_seen_health: ProjectHealth | None = None
@property
def repo_full_name(self) -> str:
return f"{self.gitea_owner}/{self.repo_name}"
@dataclass(frozen=True) @dataclass(frozen=True)
@@ -71,6 +146,15 @@ class ProjectRegistry:
projects: tuple[ProjectRecord, ...] projects: tuple[ProjectRecord, ...]
source_path: Path source_path: Path
@property
def schema_version(self) -> int:
"""Alias of :attr:`version` — the schema version read from disk."""
return self.version
@property
def api_version(self) -> str:
return REGISTRY_API_VERSION
def default_registry_path() -> Path: def default_registry_path() -> Path:
override = os.environ.get("WEBUI_PROJECT_REGISTRY", "").strip() override = os.environ.get("WEBUI_PROJECT_REGISTRY", "").strip()
@@ -79,52 +163,221 @@ def default_registry_path() -> Path:
return (Path(__file__).resolve().parent / "data" / "projects.registry.json").resolve() return (Path(__file__).resolve().parent / "data" / "projects.registry.json").resolve()
def _reject_credential_keys(obj: Any, *, path: str = "") -> None: def _reject_credential_keys(obj: Any, *, path: str = "", source: Path | None = None) -> None:
"""Recursive credential-key guard that reports an actionable ``field_path``.
Key *shape* is decided by :func:`webui.registry_safety.is_forbidden_key`, the
single source of truth shared with the worker registry (#798).
"""
if isinstance(obj, dict): if isinstance(obj, dict):
for key, value in obj.items(): for key, value in obj.items():
key_path = f"{path}.{key}" if path else key key_path = f"{path}.{key}" if path else key
if _is_forbidden_key(key): if is_forbidden_key(key):
raise ValueError(f"registry must not store credentials ({key_path})") raise RegistryError(
_reject_credential_keys(value, path=key_path) f"registry must not store credentials ({key_path})",
remediation=(
f"Remove the credential-shaped key '{key_path}' from the registry. "
"Tokens live in the keychain and are resolved server-side by "
"gitea_auth; the registry is redacted metadata only."
),
source_path=source,
field_path=key_path,
)
_reject_credential_keys(value, path=key_path, source=source)
elif isinstance(obj, list): elif isinstance(obj, list):
for index, item in enumerate(obj): for index, item in enumerate(obj):
_reject_credential_keys(item, path=f"{path}[{index}]") _reject_credential_keys(item, path=f"{path}[{index}]", source=source)
def _parse_onboarding(raw: list[dict[str, Any]] | None) -> tuple[OnboardingStep, ...]: def _require_enum(
if not raw: value: Any,
*,
allowed: tuple[str, ...],
field_path: str,
source: Path | None,
) -> str:
text = str(value)
if text not in allowed:
raise RegistryError(
f"{field_path} must be one of {', '.join(allowed)} (got {text!r})",
remediation=(
f"Set {field_path} to one of: {', '.join(allowed)}. "
"Unknown values fail closed so the console never renders an "
"unverified state."
),
source_path=source,
field_path=field_path,
)
return text
def _parse_onboarding(
raw: Any,
*,
project_path: str,
source: Path | None,
) -> tuple[OnboardingStep, ...]:
if raw is None:
return () return ()
if not isinstance(raw, list):
raise RegistryError(
f"{project_path}.onboarding_checklist must be an array",
remediation=(
f"Rewrite {project_path}.onboarding_checklist as a JSON array of "
"steps with id, title, description, and optional state."
),
source_path=source,
field_path=f"{project_path}.onboarding_checklist",
)
steps: list[OnboardingStep] = [] steps: list[OnboardingStep] = []
for item in raw: for index, item in enumerate(raw):
step_path = f"{project_path}.onboarding_checklist[{index}]"
if not isinstance(item, dict):
raise RegistryError(
f"{step_path} must be an object",
remediation=f"Rewrite {step_path} as an object with id, title, description.",
source_path=source,
field_path=step_path,
)
missing = [field for field in ("id", "title", "description") if field not in item]
if missing:
raise RegistryError(
f"{step_path} missing required fields: {', '.join(missing)}",
remediation=(
f"Add {', '.join(missing)} to {step_path}. Every onboarding step "
"must be self-describing for an operator who has no chat history."
),
source_path=source,
field_path=step_path,
)
state = _require_enum(
item.get("state", _DEFAULT_ONBOARDING_STATE),
allowed=ONBOARDING_STATES,
field_path=f"{step_path}.state",
source=source,
)
steps.append( steps.append(
OnboardingStep( OnboardingStep(
id=str(item["id"]), id=str(item["id"]),
title=str(item["title"]), title=str(item["title"]),
description=str(item["description"]), description=str(item["description"]),
state=state,
required=bool(item.get("required", True)),
) )
) )
return tuple(steps) return tuple(steps)
def _parse_project(raw: dict[str, Any]) -> ProjectRecord: def _parse_health(
raw: Any,
*,
project_path: str,
source: Path | None,
) -> ProjectHealth | None:
if raw is None:
return None
if not isinstance(raw, dict):
raise RegistryError(
f"{project_path}.last_seen_health must be an object when present",
remediation=(
f"Rewrite {project_path}.last_seen_health as an object with status "
f"(one of {', '.join(HEALTH_STATUSES)}), optional checked_at and detail, "
"or remove it. Never store endpoints or credentials here."
),
source_path=source,
field_path=f"{project_path}.last_seen_health",
)
status = _require_enum(
raw.get("status", "unknown"),
allowed=HEALTH_STATUSES,
field_path=f"{project_path}.last_seen_health.status",
source=source,
)
checked_at = raw.get("checked_at")
detail = raw.get("detail")
return ProjectHealth(
status=status,
checked_at=str(checked_at) if checked_at is not None else None,
detail=str(detail) if detail is not None else None,
)
def _parse_project(raw: Any, *, index: int, source: Path | None) -> ProjectRecord:
project_path = f"projects[{index}]"
if not isinstance(raw, dict):
raise RegistryError(
f"{project_path} must be an object",
remediation=f"Rewrite {project_path} as a JSON object describing one project.",
source_path=source,
field_path=project_path,
)
missing = [field for field in _REQUIRED_PROJECT_FIELDS if field not in raw] missing = [field for field in _REQUIRED_PROJECT_FIELDS if field not in raw]
if missing: if missing:
raise ValueError(f"project missing required fields: {', '.join(missing)}") raise RegistryError(
f"{project_path} missing required fields: {', '.join(missing)}",
remediation=(
f"Add {', '.join(missing)} to {project_path}. See "
"docs/webui-project-registry-api.md for the field-by-field contract."
),
source_path=source,
field_path=project_path,
)
profiles = raw["profiles"] profiles = raw["profiles"]
if not isinstance(profiles, dict): if not isinstance(profiles, dict):
raise ValueError("profiles must be an object") raise RegistryError(
f"{project_path}.profiles must be an object",
remediation=(
f"Rewrite {project_path}.profiles as an object mapping "
f"{', '.join(_REQUIRED_PROFILE_ROLES)} to MCP profile names."
),
source_path=source,
field_path=f"{project_path}.profiles",
)
for role in _REQUIRED_PROFILE_ROLES: for role in _REQUIRED_PROFILE_ROLES:
if role not in profiles or not profiles[role]: if role not in profiles or not profiles[role]:
raise ValueError(f"profiles.{role} is required") raise RegistryError(
f"{project_path}.profiles.{role} is required",
remediation=(
f"Set {project_path}.profiles.{role} to the configured MCP profile "
"name for that role. Role separation is a workflow-safety invariant."
),
source_path=source,
field_path=f"{project_path}.profiles.{role}",
)
workflow_paths = raw["workflow_paths"] workflow_paths = raw["workflow_paths"]
if not isinstance(workflow_paths, dict) or not workflow_paths: if not isinstance(workflow_paths, dict) or not workflow_paths:
raise ValueError("workflow_paths must be a non-empty object") raise RegistryError(
f"{project_path}.workflow_paths must be a non-empty object",
remediation=(
f"Add at least a 'skill' entry to {project_path}.workflow_paths pointing "
"at the project's canonical workflow skill."
),
source_path=source,
field_path=f"{project_path}.workflow_paths",
)
schema_paths = raw.get("schema_paths") or {} schema_paths = raw.get("schema_paths") or {}
if not isinstance(schema_paths, dict): if not isinstance(schema_paths, dict):
raise ValueError("schema_paths must be an object when present") raise RegistryError(
f"{project_path}.schema_paths must be an object when present",
remediation=(
f"Rewrite {project_path}.schema_paths as an object of label to repo path, "
"or remove it."
),
source_path=source,
field_path=f"{project_path}.schema_paths",
)
status = _require_enum(
raw.get("status", _DEFAULT_PROJECT_STATUS),
allowed=PROJECT_STATUSES,
field_path=f"{project_path}.status",
source=source,
)
remote_name = raw.get("remote_name")
return ProjectRecord( return ProjectRecord(
id=str(raw["id"]), id=str(raw["id"]),
@@ -136,61 +389,211 @@ def _parse_project(raw: dict[str, Any]) -> ProjectRecord:
profiles={role: str(profiles[role]) for role in _REQUIRED_PROFILE_ROLES}, profiles={role: str(profiles[role]) for role in _REQUIRED_PROFILE_ROLES},
workflow_paths={key: str(value) for key, value in workflow_paths.items()}, workflow_paths={key: str(value) for key, value in workflow_paths.items()},
schema_paths={key: str(value) for key, value in schema_paths.items()}, schema_paths={key: str(value) for key, value in schema_paths.items()},
onboarding_checklist=_parse_onboarding(raw.get("onboarding_checklist")), onboarding_checklist=_parse_onboarding(
raw.get("onboarding_checklist"),
project_path=project_path,
source=source,
),
status=status,
remote_name=str(remote_name) if remote_name else None,
last_seen_health=_parse_health(
raw.get("last_seen_health"),
project_path=project_path,
source=source,
),
) )
def load_registry(path: Path | None = None) -> ProjectRegistry: def load_registry(path: Path | None = None) -> ProjectRegistry:
"""Load the versioned project registry from disk.""" """Load the versioned project registry from disk.
Raises:
RegistryError: whenever the file is unreadable, is not valid JSON, or
fails schema validation. The error carries an operator remediation.
"""
source = (path or default_registry_path()).resolve() source = (path or default_registry_path()).resolve()
raw_text = source.read_text(encoding="utf-8") try:
payload = json.loads(raw_text) raw_text = source.read_text(encoding="utf-8")
except OSError as exc:
raise RegistryError(
f"registry file could not be read: {exc.strerror or exc}",
remediation=(
f"Create a readable registry at {source}, or point "
"WEBUI_PROJECT_REGISTRY at an existing file."
),
source_path=source,
) from exc
try:
payload = json.loads(raw_text)
except json.JSONDecodeError as exc:
raise RegistryError(
f"registry is not valid JSON: {exc.msg} (line {exc.lineno}, column {exc.colno})",
remediation=(
f"Fix the JSON syntax in {source} at line {exc.lineno}, column {exc.colno}."
),
source_path=source,
) from exc
if not isinstance(payload, dict): if not isinstance(payload, dict):
raise ValueError("registry root must be an object") raise RegistryError(
"registry root must be an object",
remediation=(
"Wrap the registry in a JSON object with 'version' and 'projects' keys."
),
source_path=source,
)
version = payload.get("version") version = payload.get("version")
if version != 1: if version not in SUPPORTED_SCHEMA_VERSIONS:
raise ValueError(f"unsupported registry version: {version!r}") supported = ", ".join(str(item) for item in SUPPORTED_SCHEMA_VERSIONS)
raise RegistryError(
f"unsupported registry version: {version!r}",
remediation=(
f"Set 'version' to one of {supported} (current schema is "
f"{CURRENT_SCHEMA_VERSION}). Migration notes live in "
"docs/webui-project-registry-api.md."
),
source_path=source,
field_path="version",
)
_reject_credential_keys(payload) _reject_credential_keys(payload, source=source)
projects_raw = payload.get("projects") projects_raw = payload.get("projects")
if not isinstance(projects_raw, list) or not projects_raw: if not isinstance(projects_raw, list) or not projects_raw:
raise ValueError("projects must be a non-empty array") raise RegistryError(
"projects must be a non-empty array",
remediation=(
"Add at least one project object to 'projects'. An empty console "
"registry fails closed rather than rendering a blank inventory."
),
source_path=source,
field_path="projects",
)
projects = tuple(_parse_project(item) for item in projects_raw) projects = tuple(
return ProjectRegistry(version=version, projects=projects, source_path=source) _parse_project(item, index=index, source=source)
for index, item in enumerate(projects_raw)
)
return ProjectRegistry(version=int(version), projects=projects, source_path=source)
def onboarding_summary(project: ProjectRecord) -> OnboardingSummary:
"""Aggregate a project's onboarding checklist state."""
steps = project.onboarding_checklist
counts = {state: 0 for state in ONBOARDING_STATES}
for step in steps:
counts[step.state] += 1
required_outstanding = sum(
1
for step in steps
if step.required and step.state in ("pending", "blocked")
)
return OnboardingSummary(
total=len(steps),
complete=counts["complete"],
pending=counts["pending"],
blocked=counts["blocked"],
not_applicable=counts["not_applicable"],
required_outstanding=required_outstanding,
onboarding_complete=required_outstanding == 0,
)
def project_to_dict(project: ProjectRecord) -> dict[str, Any]: def project_to_dict(project: ProjectRecord) -> dict[str, Any]:
"""Serialize a project for JSON API responses.""" """Serialize a project for JSON API responses and HTML views.
The HTML views render from this same DTO, so the console and the API can
never disagree about a project's status or onboarding progress.
"""
summary = onboarding_summary(project)
health = project.last_seen_health
return { return {
"id": project.id, "id": project.id,
"repo_name": project.repo_name, "repo_name": project.repo_name,
"gitea_owner": project.gitea_owner, "gitea_owner": project.gitea_owner,
"repo_full_name": project.repo_full_name,
"remote_host": project.remote_host, "remote_host": project.remote_host,
"remote_name": project.remote_name,
"default_branch": project.default_branch, "default_branch": project.default_branch,
"local_checkout_path": project.local_checkout_path, "local_checkout_path": project.local_checkout_path,
"status": project.status,
"profiles": dict(project.profiles), "profiles": dict(project.profiles),
"workflow_paths": dict(project.workflow_paths), "workflow_paths": dict(project.workflow_paths),
"schema_paths": dict(project.schema_paths), "schema_paths": dict(project.schema_paths),
"onboarding_checklist": [ "onboarding_checklist": [
{"id": step.id, "title": step.title, "description": step.description} {
"id": step.id,
"title": step.title,
"description": step.description,
"state": step.state,
"required": step.required,
}
for step in project.onboarding_checklist for step in project.onboarding_checklist
], ],
"onboarding_summary": {
"total": summary.total,
"complete": summary.complete,
"pending": summary.pending,
"blocked": summary.blocked,
"not_applicable": summary.not_applicable,
"required_outstanding": summary.required_outstanding,
"onboarding_complete": summary.onboarding_complete,
},
"last_seen_health": (
None
if health is None
else {
"status": health.status,
"checked_at": health.checked_at,
"detail": health.detail,
}
),
} }
def registry_to_dict(registry: ProjectRegistry) -> dict[str, Any]: def registry_to_dict(registry: ProjectRegistry) -> dict[str, Any]:
"""Serialize the whole registry, including API provenance (#632 section 6)."""
return { return {
"api_version": registry.api_version,
"schema_version": registry.schema_version,
# Retained for the unversioned MVP alias consumers (#427).
"version": registry.version, "version": registry.version,
"source_path": str(registry.source_path), "source_path": str(registry.source_path),
"source": {
"kind": "file",
"path": str(registry.source_path),
"inventory_complete": True,
},
"project_count": len(registry.projects),
"projects": [project_to_dict(project) for project in registry.projects], "projects": [project_to_dict(project) for project in registry.projects],
} }
def project_detail_to_dict(
registry: ProjectRegistry,
project: ProjectRecord,
) -> dict[str, Any]:
"""Serialize a single project for ``/api/v1/projects/{project_id}``."""
return {
"api_version": registry.api_version,
"schema_version": registry.schema_version,
"source": {
"kind": "file",
"path": str(registry.source_path),
"inventory_complete": True,
},
"project": project_to_dict(project),
}
def find_project(registry: ProjectRegistry, project_id: str) -> ProjectRecord | None: def find_project(registry: ProjectRegistry, project_id: str) -> ProjectRecord | None:
for project in registry.projects: for project in registry.projects:
if project.id == project_id: if project.id == project_id:
return project return project
return None return None
def known_project_ids(registry: ProjectRegistry) -> list[str]:
return [project.id for project in registry.projects]
+107 -25
View File
@@ -1,34 +1,65 @@
"""HTML views for project registry pages (#427).""" """HTML views for project registry pages (#427, evolved for #635).
Every view renders from :func:`webui.project_registry.project_to_dict`, the
same DTO the ``/api/v1/projects`` JSON responses use, so the HTML console and
the API can never disagree about status or onboarding progress.
"""
from __future__ import annotations from __future__ import annotations
import html import html
from typing import Any
from webui.layout import render_page from webui.layout import render_page
from webui.project_registry import ProjectRecord, ProjectRegistry from webui.project_registry import (
ProjectRecord,
ProjectRegistry,
RegistryError,
project_to_dict,
)
_STATE_LABELS = {
"complete": "Complete",
"pending": "Pending",
"blocked": "Blocked",
"not_applicable": "Not applicable",
}
def _escape(text: str) -> str: def _escape(text: str) -> str:
return html.escape(text, quote=True) return html.escape(text, quote=True)
def _progress_label(summary: dict[str, Any]) -> str:
total = summary["total"]
if not total:
return "no steps"
label = f"{summary['complete']}/{total} complete"
if summary["blocked"]:
label += f", {summary['blocked']} blocked"
return label
def render_projects_list(registry: ProjectRegistry) -> str: def render_projects_list(registry: ProjectRegistry) -> str:
rows = [] rows = []
for project in registry.projects: for project in registry.projects:
dto = project_to_dict(project)
rows.append( rows.append(
"<tr>" "<tr>"
f"<td><a href=\"/projects/{_escape(project.id)}\">{_escape(project.repo_name)}</a></td>" f"<td><a href=\"/projects/{_escape(dto['id'])}\">{_escape(dto['repo_name'])}</a></td>"
f"<td>{_escape(project.gitea_owner)}</td>" f"<td>{_escape(dto['gitea_owner'])}</td>"
f"<td>{_escape(project.remote_host)}</td>" f"<td>{_escape(dto['remote_host'])}</td>"
f"<td>{_escape(project.default_branch)}</td>" f"<td>{_escape(dto['default_branch'])}</td>"
f"<td><code>{_escape(project.profiles['author'])}</code></td>" f"<td><code>{_escape(dto['status'])}</code></td>"
f"<td>{_escape(_progress_label(dto['onboarding_summary']))}</td>"
f"<td><code>{_escape(dto['profiles']['author'])}</code></td>"
"</tr>" "</tr>"
) )
table = ( table = (
"<table class=\"registry\">" "<table class=\"registry\">"
"<thead><tr>" "<thead><tr>"
"<th>Repository</th><th>Owner</th><th>Remote</th>" "<th>Repository</th><th>Owner</th><th>Remote</th>"
"<th>Branch</th><th>Author profile</th>" "<th>Branch</th><th>Status</th><th>Onboarding</th><th>Author profile</th>"
"</tr></thead>" "</tr></thead>"
f"<tbody>{''.join(rows)}</tbody></table>" f"<tbody>{''.join(rows)}</tbody></table>"
) )
@@ -36,32 +67,39 @@ def render_projects_list(registry: ProjectRegistry) -> str:
"<h2>Projects</h2>" "<h2>Projects</h2>"
"<p>Configured repositories managed by the MCP Control Plane.</p>" "<p>Configured repositories managed by the MCP Control Plane.</p>"
f"<p class=\"meta\">Registry: <code>{_escape(str(registry.source_path))}</code> " f"<p class=\"meta\">Registry: <code>{_escape(str(registry.source_path))}</code> "
f"(version {registry.version})</p>" f"(schema version {registry.schema_version}, "
f"API {_escape(registry.api_version)})</p>"
f"{table}" f"{table}"
"<p><a href=\"/api/projects\">JSON API</a></p>" "<p><a href=\"/api/v1/projects\">JSON API</a> "
"(<a href=\"/api/projects\">unversioned alias</a>)</p>"
) )
return render_page(title="Projects", body_html=body) return render_page(title="Projects", body_html=body)
def render_project_detail(project: ProjectRecord) -> str: def render_project_detail(project: ProjectRecord) -> str:
dto = project_to_dict(project)
profile_rows = "".join( profile_rows = "".join(
f"<tr><th>{_escape(role)}</th><td><code>{_escape(name)}</code></td></tr>" f"<tr><th>{_escape(role)}</th><td><code>{_escape(name)}</code></td></tr>"
for role, name in project.profiles.items() for role, name in dto["profiles"].items()
) )
workflow_rows = "".join( workflow_rows = "".join(
f"<tr><th>{_escape(key)}</th><td><code>{_escape(path)}</code></td></tr>" f"<tr><th>{_escape(key)}</th><td><code>{_escape(path)}</code></td></tr>"
for key, path in project.workflow_paths.items() for key, path in dto["workflow_paths"].items()
) )
schema_rows = "".join( schema_rows = "".join(
f"<tr><th>{_escape(key)}</th><td><code>{_escape(path)}</code></td></tr>" f"<tr><th>{_escape(key)}</th><td><code>{_escape(path)}</code></td></tr>"
for key, path in project.schema_paths.items() for key, path in dto["schema_paths"].items()
) )
checklist_items = [] checklist_items = []
for index, step in enumerate(project.onboarding_checklist, start=1): for index, step in enumerate(dto["onboarding_checklist"], start=1):
state_label = _STATE_LABELS.get(step["state"], step["state"])
requirement = "required" if step["required"] else "optional"
checklist_items.append( checklist_items.append(
"<li>" f"<li class=\"step-{_escape(step['state'])}\">"
f"<strong>{index}. {_escape(step.title)}</strong>" f"<strong>{index}. {_escape(step['title'])}</strong>"
f"<p>{_escape(step.description)}</p>" f" <span class=\"badge\">{_escape(state_label)}</span>"
f" <span class=\"meta\">({_escape(requirement)})</span>"
f"<p>{_escape(step['description'])}</p>"
"</li>" "</li>"
) )
checklist_html = ( checklist_html = (
@@ -69,16 +107,38 @@ def render_project_detail(project: ProjectRecord) -> str:
if checklist_items if checklist_items
else "<p>No onboarding steps defined.</p>" else "<p>No onboarding steps defined.</p>"
) )
summary = dto["onboarding_summary"]
summary_html = (
"<p class=\"meta\">Onboarding: "
f"{_escape(_progress_label(summary))}; required outstanding "
f"{summary['required_outstanding']}.</p>"
)
health = dto["last_seen_health"]
health_html = (
"<p class=\"meta\">No health probe recorded (Phase 1 is read-only).</p>"
if health is None
else (
"<table class=\"detail\">"
f"<tr><th>Status</th><td><code>{_escape(health['status'])}</code></td></tr>"
f"<tr><th>Checked at</th><td>{_escape(str(health['checked_at'] or 'unknown'))}</td></tr>"
f"<tr><th>Detail</th><td>{_escape(str(health['detail'] or ''))}</td></tr>"
"</table>"
)
)
remote_name = dto["remote_name"] or "unset"
body = ( body = (
f"<h2>{_escape(project.repo_name)}</h2>" f"<h2>{_escape(dto['repo_name'])}</h2>"
"<p><a href=\"/projects\">← All projects</a></p>" "<p><a href=\"/projects\">← All projects</a></p>"
"<h3>Identity</h3>" "<h3>Identity</h3>"
"<table class=\"detail\">" "<table class=\"detail\">"
f"<tr><th>Registry id</th><td><code>{_escape(project.id)}</code></td></tr>" f"<tr><th>Registry id</th><td><code>{_escape(dto['id'])}</code></td></tr>"
f"<tr><th>Gitea owner</th><td>{_escape(project.gitea_owner)}</td></tr>" f"<tr><th>Status</th><td><code>{_escape(dto['status'])}</code></td></tr>"
f"<tr><th>Remote host</th><td>{_escape(project.remote_host)}</td></tr>" f"<tr><th>Gitea owner</th><td>{_escape(dto['gitea_owner'])}</td></tr>"
f"<tr><th>Default branch</th><td><code>{_escape(project.default_branch)}</code></td></tr>" f"<tr><th>Repository</th><td><code>{_escape(dto['repo_full_name'])}</code></td></tr>"
f"<tr><th>Local checkout</th><td><code>{_escape(project.local_checkout_path)}</code></td></tr>" f"<tr><th>Remote name</th><td><code>{_escape(remote_name)}</code></td></tr>"
f"<tr><th>Remote host</th><td>{_escape(dto['remote_host'])}</td></tr>"
f"<tr><th>Default branch</th><td><code>{_escape(dto['default_branch'])}</code></td></tr>"
f"<tr><th>Local checkout</th><td><code>{_escape(dto['local_checkout_path'])}</code></td></tr>"
"</table>" "</table>"
"<h3>Profiles</h3>" "<h3>Profiles</h3>"
f"<table class=\"detail\">{profile_rows}</table>" f"<table class=\"detail\">{profile_rows}</table>"
@@ -86,8 +146,30 @@ def render_project_detail(project: ProjectRecord) -> str:
f"<table class=\"detail\">{workflow_rows}</table>" f"<table class=\"detail\">{workflow_rows}</table>"
"<h3>Schema paths</h3>" "<h3>Schema paths</h3>"
f"<table class=\"detail\">{schema_rows}</table>" f"<table class=\"detail\">{schema_rows}</table>"
"<h3>Last seen health</h3>"
f"{health_html}"
"<h3>Onboarding checklist</h3>" "<h3>Onboarding checklist</h3>"
"<p class=\"meta\">Read-only MVP — complete these steps outside the UI.</p>" "<p class=\"meta\">Read-only — complete these steps outside the UI.</p>"
f"{summary_html}"
f"{checklist_html}" f"{checklist_html}"
f"<p><a href=\"/api/v1/projects/{_escape(dto['id'])}\">JSON detail</a></p>"
) )
return render_page(title=project.repo_name, body_html=body) return render_page(title=dto["repo_name"], body_html=body)
def render_registry_error(error: RegistryError) -> str:
"""Render a fail-closed page for an invalid registry."""
source = str(error.source_path) if error.source_path else "unknown"
field = error.field_path or "n/a"
body = (
"<h2>Project registry unavailable</h2>"
"<p>The registry failed validation, so the console refuses to render a "
"partial inventory.</p>"
"<table class=\"detail\">"
f"<tr><th>Detail</th><td>{_escape(error.message)}</td></tr>"
f"<tr><th>Field</th><td><code>{_escape(field)}</code></td></tr>"
f"<tr><th>Source</th><td><code>{_escape(source)}</code></td></tr>"
f"<tr><th>Remediation</th><td>{_escape(error.remediation)}</td></tr>"
"</table>"
)
return render_page(title="Project registry unavailable", body_html=body)
+52
View File
@@ -0,0 +1,52 @@
"""Shared credential-rejection guard for web UI registries (#427, #798).
Registries are operator-editable declarative files that the web UI loads and,
for the worker registry, writes back. None of them may ever carry a secret:
credentials belong in the keychain and reach worker processes through
environment injection, never through a file the browser layer can read.
The check is structural rather than value-based on purpose. A value scanner has
to guess what a secret looks like; a key scanner refuses the *shape* of a
credential field, so an operator cannot introduce one by accident and a later
loader cannot silently pass one through.
"""
from __future__ import annotations
from typing import Any
_FORBIDDEN_EXACT_KEYS = frozenset({
"token",
"password",
"secret",
"credential",
"auth",
"api_key",
"api-key",
})
_FORBIDDEN_KEY_PREFIXES = ("auth_", "api_key_", "api-key_")
_FORBIDDEN_KEY_SUFFIXES = ("_token", "_secret", "_password", "_credential", "_auth")
def is_forbidden_key(key: str) -> bool:
"""Return True when *key* names a credential field."""
lowered = key.lower()
if lowered in _FORBIDDEN_EXACT_KEYS:
return True
return (
lowered.startswith(_FORBIDDEN_KEY_PREFIXES)
or lowered.endswith(_FORBIDDEN_KEY_SUFFIXES)
)
def reject_credential_keys(obj: Any, *, path: str = "", subject: str = "registry") -> None:
"""Raise ValueError when *obj* carries a credential-shaped key at any depth."""
if isinstance(obj, dict):
for key, value in obj.items():
key_path = f"{path}.{key}" if path else key
if is_forbidden_key(key):
raise ValueError(f"{subject} must not store credentials ({key_path})")
reject_credential_keys(value, path=key_path, subject=subject)
elif isinstance(obj, list):
for index, item in enumerate(obj):
reject_credential_keys(item, path=f"{path}[{index}]", subject=subject)
+647
View File
@@ -0,0 +1,647 @@
"""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)