Merge branch 'master' into feat/issue-628-autonomous-handoffs-orchestration

This commit is contained in:
2026-07-23 01:12:59 -05:00
37 changed files with 9998 additions and 200 deletions
+29
View File
@@ -167,6 +167,35 @@ def _reset_mutation_authority(monkeypatch):
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)
def _deterministic_workspace_remotes():
try:
+572
View File
@@ -0,0 +1,572 @@
"""Executable acceptance tests for ARCH-01 Slice A (#822).
Each acceptance criterion (#822 §12) and named test (#822 §13) is exercised
against a real SQLite database. The migration runs on a fresh DB in ``setUp``;
the test-run output is the durable evidence the issue requires (§14).
Enforcement being proven:
* ``[TRUSTED-SERVICE]`` — the ``cp_*`` actor functions exist only on the
trusted kernel connection; a raw connection cannot satisfy the triggers.
* ``[SCHEMA]`` — fail-closed aborts, exact dominance set, NOT-NULL class,
immutability, and the last-active-grant floor are enforced by
CHECK/FK/trigger, verified here including raw-write bypass and concurrency.
"""
from __future__ import annotations
import os
import sqlite3
import tempfile
import threading
import unittest
from concurrent.futures import ThreadPoolExecutor
import arch01_platform as ap
from arch01_platform import (
ALREADY_INSTALLED,
AUTHORIZATION_DENIED,
CONCURRENT_INSTALLATION_LOST,
DISTINGUISHED_ISSUER_ID,
DOMINANCE_SET_MISMATCH,
DOMINANCE_TUPLES,
INSTALLED,
INVALID_ACTOR_CONTEXT,
INVALID_BOOTSTRAP_STATE,
PlatformKernel,
)
INSTALLER = "platform.installer"
_BOOTSTRAP_TABLES = (
"principal_equivalence_classes",
"principals",
"authoritative_issuers",
"authority_dominance",
"platform_bootstrap_seed",
"platform_bootstrap_grants",
"platform_active_invariant",
"install_state",
)
def _count(kernel: PlatformKernel, table: str) -> int:
return kernel._conn.execute(f"SELECT COUNT(*) FROM {table}").fetchone()[0]
def _count_where(kernel: PlatformKernel, table: str, where: str) -> int:
return kernel._conn.execute(f"SELECT COUNT(*) FROM {table} WHERE {where}").fetchone()[0]
def _all_bootstrap_empty(kernel: PlatformKernel) -> bool:
return all(_count(kernel, t) == 0 for t in _BOOTSTRAP_TABLES)
class Arch01MemoryTest(unittest.TestCase):
"""Single-connection behavior on an in-memory database."""
def setUp(self) -> None:
self.kernel = PlatformKernel(":memory:")
def tearDown(self) -> None:
self.kernel.close()
# -- AC1 -------------------------------------------------------------- #
def test_install_clean(self) -> None: # t_install_clean(+)
res = self.kernel.install_platform(INSTALLER)
self.assertEqual(res.code, INSTALLED)
self.assertTrue(self.kernel.is_installed())
self.assertEqual(_count(self.kernel, "install_state"), 1)
self.assertEqual(self.kernel.active_grant_count(), 1)
self.assertIn(ap.EVT_PLATFORM_INSTALLED, self.kernel.audit_events())
rows = set(
self.kernel._conn.execute(
"SELECT dominant, subordinate FROM authority_dominance"
).fetchall()
)
self.assertEqual(rows, set(DOMINANCE_TUPLES))
issuer_ref = self.kernel._conn.execute(
"SELECT i.issuer_ref FROM principals p JOIN authoritative_issuers i "
"ON p.issuer_id = i.issuer_id WHERE p.principal_id = ?",
(INSTALLER,),
).fetchone()
self.assertEqual(issuer_ref[0], DISTINGUISHED_ISSUER_ID)
# -- AC2 -------------------------------------------------------------- #
def test_install_twice(self) -> None: # t_install_twice(-)
self.assertEqual(self.kernel.install_platform(INSTALLER).code, INSTALLED)
res2 = self.kernel.install_platform(INSTALLER)
self.assertEqual(res2.code, ALREADY_INSTALLED)
self.assertEqual(_count(self.kernel, "principals"), 1)
self.assertEqual(_count(self.kernel, "platform_bootstrap_grants"), 1)
self.assertEqual(_count(self.kernel, "install_state"), 1)
# -- AC3 / AC5 -------------------------------------------------------- #
def test_install_stage_rollback(self) -> None: # t_install_stage_rollback
for stop in range(1, 9):
with self.subTest(stages=stop):
k = PlatformKernel(":memory:")
try:
self._partial_bootstrap_then_rollback(k, stop)
self.assertTrue(
_all_bootstrap_empty(k),
f"partial rows survived rollback at stage {stop}",
)
self.assertFalse(k.is_installed())
finally:
k.close()
def test_no_partial_after_rollback(self) -> None: # t_no_partial_after_rollback
k = PlatformKernel(":memory:")
try:
code = self._seed_bootstrap_and_mark(k, dominance=DOMINANCE_TUPLES[:-1])
self.assertEqual(code, DOMINANCE_SET_MISMATCH)
self.assertTrue(_all_bootstrap_empty(k))
self.assertFalse(k.is_installed())
finally:
k.close()
# -- AC4 -------------------------------------------------------------- #
def test_dominance_missing(self) -> None: # t_dominance_missing(-)
k = PlatformKernel(":memory:")
try:
self.assertEqual(
self._seed_bootstrap_and_mark(k, dominance=DOMINANCE_TUPLES[:-1]),
DOMINANCE_SET_MISMATCH,
)
self.assertFalse(k.is_installed())
finally:
k.close()
def test_dominance_extra(self) -> None: # t_dominance_extra(-)
k = PlatformKernel(":memory:")
try:
extra = DOMINANCE_TUPLES + (("platform.bootstrap", "rogue.extra"),)
self.assertEqual(
self._seed_bootstrap_and_mark(k, dominance=extra),
DOMINANCE_SET_MISMATCH,
)
self.assertFalse(k.is_installed())
finally:
k.close()
def test_dominance_malformed(self) -> None: # t_dominance_malformed(-)
k = PlatformKernel(":memory:")
try:
malformed = DOMINANCE_TUPLES[:-1] + (("supervisor.root", "WRONG.subordinate"),)
self.assertEqual(
self._seed_bootstrap_and_mark(k, dominance=malformed),
DOMINANCE_SET_MISMATCH,
)
self.assertFalse(k.is_installed())
finally:
k.close()
# -- AC6 -------------------------------------------------------------- #
def test_principal_no_class(self) -> None: # t_principal_no_class(-)
with self.kernel.actor_context("op", "operator", "install"):
with self.assertRaises(sqlite3.IntegrityError):
self.kernel._conn.execute(
"INSERT INTO principals"
"(principal_id, actor_kind, current_class_id, issuer_id, registered_by, created_at) "
"VALUES ('x', 'operator', NULL, NULL, NULL, '2026-01-01T00:00:00Z')"
)
# -- AC7 -------------------------------------------------------------- #
def test_noninstaller_null_issuer(self) -> None: # t_nonobstaller_null_issuer(-)
self.assertEqual(self.kernel.install_platform(INSTALLER).code, INSTALLED)
with self.kernel.actor_context("op", "operator", "normal"):
cur = self.kernel._conn.execute(
"INSERT INTO principal_equivalence_classes(created_at) VALUES ('2026-01-01T00:00:00Z')"
)
class_id = cur.lastrowid
with self.assertRaises(sqlite3.IntegrityError) as ctx:
self.kernel._conn.execute(
"INSERT INTO principals"
"(principal_id, actor_kind, current_class_id, issuer_id, registered_by, created_at) "
"VALUES ('rogue', 'operator', ?, NULL, NULL, '2026-01-01T00:00:00Z')",
(class_id,),
)
self.assertIn("INVALID_BOOTSTRAP_STATE", str(ctx.exception))
def test_installer_null_issuer_only_during_install(self) -> None:
self.assertEqual(self.kernel.install_platform(INSTALLER).code, INSTALLED)
with self.kernel.actor_context("i2", "installer", "install"):
cur = self.kernel._conn.execute(
"INSERT INTO principal_equivalence_classes(created_at) VALUES ('2026-01-01T00:00:00Z')"
)
class_id = cur.lastrowid
with self.assertRaises(sqlite3.IntegrityError):
self.kernel._conn.execute(
"INSERT INTO principals"
"(principal_id, actor_kind, current_class_id, issuer_id, registered_by, created_at) "
"VALUES ('i2', 'installer', ?, NULL, NULL, '2026-01-01T00:00:00Z')",
(class_id,),
)
# -- AC8 -------------------------------------------------------------- #
def test_context_missing(self) -> None: # t_context_missing(-)
self.assertIsNone(self.kernel._ctx)
with self.assertRaises(sqlite3.IntegrityError) as ctx:
self.kernel._conn.execute(
"INSERT INTO principal_equivalence_classes(created_at) VALUES ('2026-01-01T00:00:00Z')"
)
self.assertIn("INVALID_ACTOR_CONTEXT", str(ctx.exception))
def test_context_stale(self) -> None: # t_context_stale(-)
with self.kernel.actor_context("op", "operator", "normal"):
self.kernel._ctx.expired = True
with self.assertRaises(sqlite3.IntegrityError) as ctx:
self.kernel._conn.execute(
"INSERT INTO principal_equivalence_classes(created_at) VALUES ('2026-01-01T00:00:00Z')"
)
self.assertIn("INVALID_ACTOR_CONTEXT", str(ctx.exception))
def test_context_epoch_shift(self) -> None: # t_context_epoch_shift(-)
with self.kernel.actor_context("op", "operator", "normal"):
self.kernel._ctx.live_epoch = self.kernel._ctx.bound_epoch + 99
with self.assertRaises(sqlite3.IntegrityError) as ctx:
self.kernel._conn.execute(
"INSERT INTO principal_equivalence_classes(created_at) VALUES ('2026-01-01T00:00:00Z')"
)
self.assertIn("INVALID_ACTOR_CONTEXT", str(ctx.exception))
def test_bad_actor_kind_or_mode_rejected(self) -> None:
for kind, mode in (("intruder", "normal"), ("operator", "sabotage")):
with self.subTest(kind=kind, mode=mode):
with self.kernel.actor_context("op", kind, mode):
with self.assertRaises(sqlite3.IntegrityError):
self.kernel._conn.execute(
"INSERT INTO principal_equivalence_classes(created_at) "
"VALUES ('2026-01-01T00:00:00Z')"
)
# -- AC9 -------------------------------------------------------------- #
def test_bootstrap_immutable_update(self) -> None: # t_bootstrap_immutable_{update}
self.assertEqual(self.kernel.install_platform(INSTALLER).code, INSTALLED)
cases = [
("UPDATE install_state SET installed_at = 'x' WHERE id = 1", "IMMUTABLE_INSTALL_STATE"),
("UPDATE platform_bootstrap_seed SET created_at = 'x' WHERE seed_id = 1", "IMMUTABLE_SEED"),
("UPDATE authority_dominance SET subordinate = 'x' WHERE dominant = 'supervisor.root'", "IMMUTABLE_DOMINANCE"),
(f"UPDATE authoritative_issuers SET issuer_ref = 'x' WHERE issuer_ref = '{DISTINGUISHED_ISSUER_ID}'", "IMMUTABLE_ISSUER"),
(f"UPDATE principals SET actor_kind = 'operator' WHERE principal_id = '{INSTALLER}'", "IMMUTABLE_PRINCIPAL"),
]
for sql, tag in cases:
with self.subTest(sql=sql):
with self.kernel.actor_context("op", "operator", "normal"):
with self.assertRaises(sqlite3.IntegrityError) as ctx:
self.kernel._conn.execute(sql)
self.assertIn(tag, str(ctx.exception))
def test_bootstrap_immutable_delete(self) -> None: # t_bootstrap_immutable_{delete}
self.assertEqual(self.kernel.install_platform(INSTALLER).code, INSTALLED)
cases = [
("DELETE FROM install_state WHERE id = 1", "IMMUTABLE_INSTALL_STATE"),
("DELETE FROM platform_bootstrap_seed WHERE seed_id = 1", "IMMUTABLE_SEED"),
("DELETE FROM authority_dominance", "IMMUTABLE_DOMINANCE"),
("DELETE FROM authoritative_issuers", "IMMUTABLE_ISSUER"),
(f"DELETE FROM principals WHERE principal_id = '{INSTALLER}'", "IMMUTABLE_PRINCIPAL"),
("DELETE FROM platform_bootstrap_grants", "IMMUTABLE_GRANT"),
]
for sql, tag in cases:
with self.subTest(sql=sql):
with self.kernel.actor_context("op", "operator", "normal"):
with self.assertRaises(sqlite3.IntegrityError) as ctx:
self.kernel._conn.execute(sql)
self.assertIn(tag, str(ctx.exception))
def test_grant_reactivation_rejected(self) -> None:
self.assertEqual(self.kernel.install_platform(INSTALLER).code, INSTALLED)
self.kernel.register_principal(
"op1", "operator", DISTINGUISHED_ISSUER_ID, actor_principal=INSTALLER
)
self.assertEqual(
self.kernel.grant_platform_bootstrap("op1", INSTALLER).code, INSTALLED
)
gid = self.kernel._conn.execute(
"SELECT grant_id FROM platform_bootstrap_grants WHERE grantee_principal_id = 'op1'"
).fetchone()[0]
self.assertEqual(
self.kernel.revoke_platform_bootstrap(gid, actor_principal=INSTALLER).code,
INSTALLED,
)
with self.kernel.actor_context("op", "operator", "normal"):
with self.assertRaises(sqlite3.IntegrityError) as ctx:
self.kernel._conn.execute(
"UPDATE platform_bootstrap_grants SET active = 1 WHERE grant_id = ?",
(gid,),
)
self.assertIn("IMMUTABLE_GRANT", str(ctx.exception))
# -- AC12 ------------------------------------------------------------- #
def test_raw_write_bypass(self) -> None: # t_raw_write_bypass(raw-bypass)
with tempfile.TemporaryDirectory() as tmp:
path = os.path.join(tmp, "p.sqlite3")
k = PlatformKernel(path)
self.assertEqual(k.install_platform(INSTALLER).code, INSTALLED)
k.close()
raw = sqlite3.connect(path)
raw.execute("PRAGMA foreign_keys = ON")
try:
with self.assertRaises(sqlite3.Error):
raw.execute(
"INSERT INTO audit_records(event, created_at) "
"VALUES ('forged', '2026-01-01T00:00:00Z')"
)
raw.commit()
with self.assertRaises(sqlite3.Error):
raw.execute("UPDATE install_state SET installed_at = 'x' WHERE id = 1")
raw.commit()
with self.assertRaises(sqlite3.Error):
raw.execute("DELETE FROM platform_bootstrap_grants")
raw.commit()
finally:
raw.close()
# -- AC13 ------------------------------------------------------------- #
def test_audit_created(self) -> None: # t_audit_created(+)
self.assertEqual(self.kernel.install_platform(INSTALLER).code, INSTALLED)
self.kernel.register_principal(
"op1", "operator", DISTINGUISHED_ISSUER_ID, actor_principal=INSTALLER
)
self.assertEqual(
self.kernel.grant_platform_bootstrap("op1", INSTALLER).code, INSTALLED
)
gid = self.kernel._conn.execute(
"SELECT grant_id FROM platform_bootstrap_grants WHERE grantee_principal_id = 'op1'"
).fetchone()[0]
self.assertEqual(
self.kernel.revoke_platform_bootstrap(gid, actor_principal=INSTALLER).code,
INSTALLED,
)
events = self.kernel.audit_events()
for evt in (
ap.EVT_PLATFORM_INSTALLED,
ap.EVT_GRANT_CREATED,
ap.EVT_GRANT_REVOKED,
ap.EVT_PRINCIPAL_REGISTERED,
):
self.assertIn(evt, events)
# -- AC14 ------------------------------------------------------------- #
def test_audit_immutable(self) -> None: # t_audit_immutable(raw-bypass)
self.assertEqual(self.kernel.install_platform(INSTALLER).code, INSTALLED)
with self.kernel.actor_context("op", "operator", "normal"):
with self.assertRaises(sqlite3.IntegrityError) as up:
self.kernel._conn.execute("UPDATE audit_records SET event = 'x' WHERE audit_id = 1")
self.assertIn("IMMUTABLE_AUDIT", str(up.exception))
with self.assertRaises(sqlite3.IntegrityError) as dl:
self.kernel._conn.execute("DELETE FROM audit_records WHERE audit_id = 1")
self.assertIn("IMMUTABLE_AUDIT", str(dl.exception))
# -- meta ------------------------------------------------------------- #
def test_schema_meta(self) -> None:
rows = dict(self.kernel._conn.execute("SELECT key, value FROM arch01_meta").fetchall())
self.assertEqual(rows["schema_version"], str(ap.SCHEMA_VERSION))
self.assertIn("disabled by default", rows["architecture"])
def test_register_principal_creates_class_first(self) -> None:
self.assertEqual(self.kernel.install_platform(INSTALLER).code, INSTALLED)
res = self.kernel.register_principal(
"svc1", "service", DISTINGUISHED_ISSUER_ID, actor_principal=INSTALLER
)
self.assertEqual(res.code, INSTALLED)
row = self.kernel._conn.execute(
"SELECT current_class_id FROM principals WHERE principal_id = 'svc1'"
).fetchone()
self.assertIsNotNone(row[0])
# -- helpers ---------------------------------------------------------- #
def _partial_bootstrap_then_rollback(self, k: PlatformKernel, stop: int) -> None:
"""Execute the first ``stop`` bootstrap statements, then ROLLBACK."""
now = "2026-01-01T00:00:00Z"
k._conn.execute("BEGIN IMMEDIATE")
class_id = None
issuer_id = None
try:
with k.actor_context(INSTALLER, "installer", "install"):
c = k._conn
if stop >= 1:
class_id = c.execute(
"INSERT INTO principal_equivalence_classes(created_at) VALUES (?)", (now,)
).lastrowid
if stop >= 2:
c.execute(
"INSERT INTO principals(principal_id, actor_kind, current_class_id, issuer_id, registered_by, created_at) "
"VALUES (?, 'installer', ?, NULL, ?, ?)",
(INSTALLER, class_id, INSTALLER, now),
)
if stop >= 3:
issuer_id = c.execute(
"INSERT INTO authoritative_issuers(issuer_kind, issuer_ref, created_at) VALUES ('operator-key', ?, ?)",
(DISTINGUISHED_ISSUER_ID, now),
).lastrowid
if stop >= 4:
c.execute(
"UPDATE principals SET issuer_id = ? WHERE principal_id = ?",
(issuer_id, INSTALLER),
)
if stop >= 5:
c.executemany(
"INSERT INTO authority_dominance(dominant, subordinate) VALUES (?, ?)",
DOMINANCE_TUPLES,
)
if stop >= 6:
c.execute(
"INSERT INTO platform_bootstrap_seed(seed_id, installer_principal_id, created_at) VALUES (1, ?, ?)",
(INSTALLER, now),
)
if stop >= 7:
c.execute(
"INSERT INTO platform_bootstrap_grants(grantee_principal_id, granted_by, active, created_at) VALUES (?, NULL, 1, ?)",
(INSTALLER, now),
)
if stop >= 8:
c.execute("INSERT INTO platform_active_invariant(id, active_count) VALUES (1, 1)")
finally:
k._conn.execute("ROLLBACK")
def _seed_bootstrap_and_mark(self, k: PlatformKernel, dominance) -> str:
"""Seed a full bootstrap with a caller-supplied dominance set, then
attempt the marker insert. Returns the classified failure code (or
INSTALLED). Rolls back on failure so no partial rows remain."""
now = "2026-01-01T00:00:00Z"
k._conn.execute("BEGIN IMMEDIATE")
try:
with k.actor_context(INSTALLER, "installer", "install"):
c = k._conn
class_id = c.execute(
"INSERT INTO principal_equivalence_classes(created_at) VALUES (?)", (now,)
).lastrowid
c.execute(
"INSERT INTO principals(principal_id, actor_kind, current_class_id, issuer_id, registered_by, created_at) "
"VALUES (?, 'installer', ?, NULL, ?, ?)",
(INSTALLER, class_id, INSTALLER, now),
)
issuer_id = c.execute(
"INSERT INTO authoritative_issuers(issuer_kind, issuer_ref, created_at) VALUES ('operator-key', ?, ?)",
(DISTINGUISHED_ISSUER_ID, now),
).lastrowid
c.execute(
"UPDATE principals SET issuer_id = ? WHERE principal_id = ?",
(issuer_id, INSTALLER),
)
c.executemany(
"INSERT INTO authority_dominance(dominant, subordinate) VALUES (?, ?)",
dominance,
)
c.execute(
"INSERT INTO platform_bootstrap_seed(seed_id, installer_principal_id, created_at) VALUES (1, ?, ?)",
(INSTALLER, now),
)
c.execute(
"INSERT INTO platform_bootstrap_grants(grantee_principal_id, granted_by, active, created_at) VALUES (?, NULL, 1, ?)",
(INSTALLER, now),
)
c.execute("INSERT INTO platform_active_invariant(id, active_count) VALUES (1, 1)")
c.execute(
"INSERT INTO install_state(id, marker, installed_at) VALUES (1, 'installed', ?)",
(now,),
)
k._conn.execute("COMMIT")
return INSTALLED
except sqlite3.Error as exc:
k._safe_rollback()
return PlatformKernel._classify(exc)
class Arch01ConcurrencyTest(unittest.TestCase):
"""Concurrency invariants require file-backed DBs and independent connections."""
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.path = os.path.join(self._tmp.name, "p.sqlite3")
def tearDown(self) -> None:
self._tmp.cleanup()
# -- AC10 ------------------------------------------------------------- #
def test_concurrent_install(self) -> None: # t_concurrent_install(concurrency)
k1 = PlatformKernel(self.path, busy_timeout_ms=0)
k2 = PlatformKernel(self.path, busy_timeout_ms=0)
barrier = threading.Barrier(2)
results = {}
def _install(name, kernel):
barrier.wait()
results[name] = kernel.install_platform(INSTALLER).code
try:
with ThreadPoolExecutor(max_workers=2) as ex:
f1 = ex.submit(_install, "a", k1)
f2 = ex.submit(_install, "b", k2)
f1.result()
f2.result()
codes = sorted(results.values())
self.assertEqual(codes.count(INSTALLED), 1, f"exactly one install expected: {results}")
other = [c for c in results.values() if c != INSTALLED][0]
self.assertIn(other, (ALREADY_INSTALLED, CONCURRENT_INSTALLATION_LOST))
self.assertTrue(k1.is_installed())
self.assertEqual(_count(k1, "install_state"), 1)
self.assertEqual(_count(k1, "principals"), 1)
finally:
k1.close()
k2.close()
# -- AC11 ------------------------------------------------------------- #
def test_concurrent_last_grant_revoke(self) -> None: # t_concurrent_last_grant_revoke
setup = PlatformKernel(self.path)
self.assertEqual(setup.install_platform(INSTALLER).code, INSTALLED)
setup.register_principal("op1", "operator", DISTINGUISHED_ISSUER_ID, actor_principal=INSTALLER)
self.assertEqual(setup.grant_platform_bootstrap("op1", INSTALLER).code, INSTALLED)
self.assertEqual(setup.active_grant_count(), 2)
gids = [
r[0]
for r in setup._conn.execute(
"SELECT grant_id FROM platform_bootstrap_grants WHERE active = 1 ORDER BY grant_id"
).fetchall()
]
setup.close()
self.assertEqual(len(gids), 2)
k1 = PlatformKernel(self.path, busy_timeout_ms=3000)
k2 = PlatformKernel(self.path, busy_timeout_ms=3000)
barrier = threading.Barrier(2)
results = {}
def _revoke(name, kernel, gid):
barrier.wait()
results[name] = kernel.revoke_platform_bootstrap(gid, actor_principal=INSTALLER).code
try:
with ThreadPoolExecutor(max_workers=2) as ex:
f1 = ex.submit(_revoke, "a", k1, gids[0])
f2 = ex.submit(_revoke, "b", k2, gids[1])
f1.result()
f2.result()
codes = list(results.values())
self.assertEqual(codes.count(INSTALLED), 1, f"exactly one revoke should win: {results}")
self.assertEqual(codes.count(AUTHORIZATION_DENIED), 1, f"one revoke must be denied: {results}")
self.assertEqual(k1.active_grant_count(), 1)
self.assertEqual(_count_where(k1, "platform_bootstrap_grants", "active = 1"), 1)
finally:
k1.close()
k2.close()
def test_revoke_final_grant_denied(self) -> None:
k = PlatformKernel(self.path)
try:
self.assertEqual(k.install_platform(INSTALLER).code, INSTALLED)
gid = k._conn.execute(
"SELECT grant_id FROM platform_bootstrap_grants WHERE active = 1"
).fetchone()[0]
res = k.revoke_platform_bootstrap(gid, actor_principal=INSTALLER)
self.assertEqual(res.code, AUTHORIZATION_DENIED)
self.assertEqual(k.active_grant_count(), 1)
self.assertEqual(_count_where(k, "platform_bootstrap_grants", "active = 1"), 1)
finally:
k.close()
if __name__ == "__main__":
unittest.main()
+581
View File
@@ -0,0 +1,581 @@
"""Authoritative controller cross-role generic queue allocation (#840)."""
from __future__ import annotations
import os
import tempfile
import unittest
from unittest.mock import patch
from allocator_service import (
ALLOCATION_MODE_CROSS_ROLE,
ALLOCATION_MODE_ROLE_SCOPED,
OUTCOME_NO_SAFE,
OUTCOME_PREVIEW,
OUTCOME_WAIT,
ROLE_AUTHOR,
ROLE_CONTROLLER,
ROLE_MERGER,
ROLE_RECONCILER,
ROLE_REVIEWER,
WorkCandidate,
allocate_next_work,
build_selection_dict,
classify_skip,
required_namespace_for_role,
required_profile_for_role,
resolve_allocation_mode,
selected_action_for_candidate,
)
from control_plane_db import ControlPlaneDB
import role_session_router
from role_session_router import (
ROUTE_ALLOWED,
ROUTE_AMBIGUOUS,
ROUTE_WRONG_ROLE,
route_task_session,
)
import namespace_workspace_binding as nwb
import task_capability_map
class CrossRoleAllocationModeTest(unittest.TestCase):
def test_controller_defaults_to_cross_role(self) -> None:
self.assertEqual(
resolve_allocation_mode(ROLE_CONTROLLER),
ALLOCATION_MODE_CROSS_ROLE,
)
def test_worker_defaults_to_role_scoped(self) -> None:
for role in (ROLE_AUTHOR, ROLE_REVIEWER, ROLE_MERGER, ROLE_RECONCILER):
self.assertEqual(
resolve_allocation_mode(role),
ALLOCATION_MODE_ROLE_SCOPED,
)
def test_explicit_modes(self) -> None:
self.assertEqual(
resolve_allocation_mode(ROLE_CONTROLLER, "role_scoped"),
ALLOCATION_MODE_ROLE_SCOPED,
)
self.assertEqual(
resolve_allocation_mode(ROLE_AUTHOR, "cross_role"),
ALLOCATION_MODE_CROSS_ROLE,
)
class CrossRoleSelectionPayloadTest(unittest.TestCase):
def test_selection_contains_required_fields(self) -> None:
c = WorkCandidate(
kind="issue",
number=840,
labels=("status:ready",),
title="cross-role",
priority=20,
)
sel = build_selection_dict(
c,
active_role=ROLE_CONTROLLER,
required_role=ROLE_AUTHOR,
profile_name="prgs-controller",
allocation_mode=ALLOCATION_MODE_CROSS_ROLE,
)
self.assertEqual(sel["number"], 840)
self.assertEqual(sel["kind"], "issue")
self.assertEqual(sel["required_role"], ROLE_AUTHOR)
self.assertEqual(sel["selected_action"], "implement")
self.assertEqual(sel["action"], "implement")
self.assertEqual(sel["required_profile"], "prgs-author")
self.assertEqual(sel["required_namespace"], "gitea-author")
self.assertEqual(sel["pinned"]["number"], 840)
self.assertIsNone(sel["pinned"]["head_sha"])
def test_profile_prefix_preserved(self) -> None:
self.assertEqual(
required_profile_for_role(ROLE_REVIEWER, profile_name="dadeschools-controller"),
"dadeschools-reviewer",
)
self.assertEqual(
required_namespace_for_role(ROLE_MERGER),
"gitea-merger",
)
def test_selected_actions_per_role(self) -> None:
issue = WorkCandidate(kind="issue", number=1, labels=("status:ready",))
pr_review = WorkCandidate(kind="pr", number=2, head_sha="a" * 40)
pr_rc = WorkCandidate(
kind="pr",
number=3,
head_sha="b" * 40,
request_changes_current_head=True,
)
pr_merge = WorkCandidate(
kind="pr",
number=4,
head_sha="c" * 40,
approval_on_current_head=True,
mergeable=True,
)
pr_recon = WorkCandidate(
kind="pr",
number=5,
head_sha="d" * 40,
approval_contaminated=True,
)
self.assertEqual(selected_action_for_candidate(issue, ROLE_AUTHOR), "implement")
self.assertEqual(
selected_action_for_candidate(pr_rc, ROLE_AUTHOR),
"address_pr_change_requests",
)
self.assertEqual(
selected_action_for_candidate(pr_review, ROLE_REVIEWER), "review"
)
self.assertEqual(selected_action_for_candidate(pr_merge, ROLE_MERGER), "merge")
self.assertEqual(
selected_action_for_candidate(pr_recon, ROLE_RECONCILER),
"reconcile_contaminated_approval",
)
class CrossRoleAllocateServiceTest(unittest.TestCase):
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.db = ControlPlaneDB(os.path.join(self._tmp.name, "cp.sqlite3"))
def tearDown(self) -> None:
self._tmp.cleanup()
def _alloc(self, **kwargs):
defaults = dict(
db=self.db,
session_id="ctrl-session",
role=ROLE_CONTROLLER,
remote="prgs",
org="org",
repo="repo",
candidates=[],
apply=False,
profile_name="prgs-controller",
username="controller-bot",
controller_instance_id="ctrl-1",
)
defaults.update(kwargs)
return allocate_next_work(**defaults)
def test_eligible_author_work(self) -> None:
cands = [
WorkCandidate(
kind="issue",
number=100,
labels=("status:ready",),
title="author work",
priority=20,
),
]
res = self._alloc(candidates=cands)
self.assertTrue(res["success"])
self.assertEqual(res["outcome"], OUTCOME_PREVIEW)
self.assertEqual(res["allocation_mode"], ALLOCATION_MODE_CROSS_ROLE)
self.assertIsNotNone(res["selected"])
self.assertEqual(res["selected"]["number"], 100)
self.assertEqual(res["required_role"], ROLE_AUTHOR)
self.assertEqual(res["selected_action"], "implement")
self.assertEqual(res["required_profile"], "prgs-author")
self.assertEqual(res["required_namespace"], "gitea-author")
self.assertIn("allocate", res["controller_allowed_actions"])
self.assertIn("merge", res["controller_forbidden_actions"])
self.assertFalse(res["allocation_evidence"]["lease_created"])
def test_eligible_reviewer_work(self) -> None:
cands = [
WorkCandidate(
kind="pr",
number=200,
head_sha="e" * 40,
title="needs review",
priority=30,
),
]
res = self._alloc(candidates=cands)
self.assertEqual(res["selected"]["number"], 200)
self.assertEqual(res["required_role"], ROLE_REVIEWER)
self.assertEqual(res["selected_action"], "review")
self.assertEqual(res["required_profile"], "prgs-reviewer")
self.assertEqual(res["selected"]["pinned"]["head_sha"], "e" * 40)
def test_eligible_merger_work(self) -> None:
cands = [
WorkCandidate(
kind="pr",
number=300,
head_sha="f" * 40,
approval_on_current_head=True,
mergeable=True,
priority=40,
),
]
res = self._alloc(candidates=cands)
self.assertEqual(res["selected"]["number"], 300)
self.assertEqual(res["required_role"], ROLE_MERGER)
self.assertEqual(res["selected_action"], "merge")
def test_eligible_reconciler_work(self) -> None:
cands = [
WorkCandidate(
kind="pr",
number=400,
head_sha="1" * 40,
approval_contaminated=True,
priority=50,
),
]
res = self._alloc(candidates=cands)
self.assertEqual(res["selected"]["number"], 400)
self.assertEqual(res["required_role"], ROLE_RECONCILER)
self.assertIn("reconcile", res["selected_action"])
def test_no_eligible_work(self) -> None:
cands = [
WorkCandidate(
kind="issue",
number=10,
labels=("status:blocked",),
blocked=True,
priority=99,
),
WorkCandidate(
kind="issue",
number=11,
labels=("status:ready",),
dependency_unmet=True,
dependency_reason="blocked by #10",
priority=98,
),
]
res = self._alloc(candidates=cands)
self.assertTrue(res["success"])
self.assertEqual(res["outcome"], OUTCOME_NO_SAFE)
self.assertIsNone(res["selected"])
self.assertEqual(res["allocation_mode"], ALLOCATION_MODE_CROSS_ROLE)
def test_leased_work_skipped(self) -> None:
cands = [
WorkCandidate(
kind="issue",
number=50,
labels=("status:ready",),
priority=20,
),
WorkCandidate(
kind="issue",
number=51,
labels=("status:ready",),
priority=10,
),
]
# Seed a foreign lease on issue 50 via assign_and_lease under another session.
other = allocate_next_work(
self.db,
session_id="other-worker",
role=ROLE_AUTHOR,
remote="prgs",
org="org",
repo="repo",
candidates=cands[:1],
apply=True,
profile_name="prgs-author",
controller_instance_id="other-ctrl",
)
self.assertEqual(other["outcome"], "assigned_work")
res = self._alloc(candidates=cands)
self.assertIsNotNone(res["selected"])
self.assertEqual(res["selected"]["number"], 51)
self.assertTrue(any(s["number"] == 50 for s in res["skipped"]))
self.assertTrue(res["claims_excluded"])
def test_dependencies_skipped(self) -> None:
cands = [
WorkCandidate(
kind="issue",
number=1,
labels=("status:ready",),
priority=99,
dependency_unmet=True,
dependency_reason="needs #2",
),
WorkCandidate(
kind="issue",
number=2,
labels=("status:ready",),
priority=1,
),
]
res = self._alloc(candidates=cands)
self.assertEqual(res["selected"]["number"], 2)
skipped = {s["number"]: s["reason"] for s in res["skipped"]}
self.assertIn(1, skipped)
self.assertIn("needs #2", skipped[1])
def test_pagination_limit_only_truncates_skip_report(self) -> None:
"""Ranking uses full inventory; reporting limit is MCP-layer only.
Service ranks all candidates; prove higher-priority eligible item
wins even when many skipped precede it.
"""
cands = []
for n in range(1, 30):
cands.append(
WorkCandidate(
kind="issue",
number=n,
labels=("status:ready",),
priority=100 - n,
dependency_unmet=True,
dependency_reason=f"dep {n}",
)
)
cands.append(
WorkCandidate(
kind="issue",
number=999,
labels=("status:ready",),
priority=1,
)
)
res = self._alloc(candidates=cands)
self.assertEqual(res["selected"]["number"], 999)
self.assertGreaterEqual(len(res["skipped"]), 29)
def test_role_scoped_controller_legacy_still_restricts(self) -> None:
"""role_scoped controller only takes reconciler-needed items."""
cands = [
WorkCandidate(
kind="issue",
number=1,
labels=("status:ready",),
priority=50,
),
WorkCandidate(
kind="pr",
number=2,
head_sha="a" * 40,
approval_contaminated=True,
priority=1,
),
]
res = self._alloc(
candidates=cands,
allocation_mode=ALLOCATION_MODE_ROLE_SCOPED,
)
self.assertEqual(res["allocation_mode"], ALLOCATION_MODE_ROLE_SCOPED)
self.assertEqual(res["selected"]["number"], 2)
self.assertEqual(res["required_role"], ROLE_RECONCILER)
def test_cross_role_prefers_highest_priority_across_roles(self) -> None:
cands = [
WorkCandidate(
kind="issue",
number=10,
labels=("status:ready",),
priority=10,
),
WorkCandidate(
kind="pr",
number=20,
head_sha="b" * 40,
priority=50,
),
WorkCandidate(
kind="pr",
number=30,
head_sha="c" * 40,
approval_on_current_head=True,
mergeable=True,
priority=20,
),
]
res = self._alloc(candidates=cands)
# PR #20 highest priority → reviewer
self.assertEqual(res["selected"]["number"], 20)
self.assertEqual(res["required_role"], ROLE_REVIEWER)
def test_apply_creates_lease_evidence_for_required_role(self) -> None:
cands = [
WorkCandidate(
kind="issue",
number=777,
labels=("status:ready",),
priority=20,
),
]
res = self._alloc(candidates=cands, apply=True)
self.assertEqual(res["outcome"], "assigned_work")
self.assertTrue(res["allocation_evidence"]["lease_created"])
self.assertEqual(res["allocation_evidence"]["lease_role"], ROLE_AUTHOR)
proof = res["lease_proof"]
self.assertIsNotNone(proof["lease_id"])
self.assertEqual(proof["lease_role"], ROLE_AUTHOR)
self.assertIn("implement", proof["allowed_actions"])
# Controller isolation: controller still forbids merge/push/create_pr
self.assertIn("merge", res["controller_forbidden_actions"])
self.assertIn("push", res["controller_forbidden_actions"])
def test_metadata_consistency_role_is_controller(self) -> None:
cands = [
WorkCandidate(
kind="issue",
number=1,
labels=("status:ready",),
),
]
res = self._alloc(candidates=cands)
self.assertEqual(res["role"], ROLE_CONTROLLER)
self.assertEqual(res["routing_role"], ROLE_CONTROLLER)
self.assertEqual(res["required_role"], ROLE_AUTHOR)
class ProcessWorkQueueRouterTest(unittest.TestCase):
def tearDown(self) -> None:
role_session_router.clear_route_state()
def test_process_work_queue_allowed_for_controller(self) -> None:
res = route_task_session(
"process_work_queue",
active_profile="prgs-controller",
active_role_kind="controller",
allowed_in_current_session=True,
)
self.assertEqual(res["route_result"], ROUTE_ALLOWED)
self.assertEqual(res["required_role"], "controller")
self.assertTrue(res["downstream_allowed"])
def test_process_work_queue_hyphen_alias(self) -> None:
res = route_task_session(
"process-work-queue",
active_profile="prgs-controller",
active_role_kind="controller",
allowed_in_current_session=True,
)
self.assertEqual(res["route_result"], ROUTE_ALLOWED)
def test_process_work_queue_wrong_role_for_author(self) -> None:
res = route_task_session(
"process_work_queue",
active_profile="prgs-author",
active_role_kind="author",
allowed_in_current_session=False,
)
self.assertEqual(res["route_result"], ROUTE_WRONG_ROLE)
self.assertEqual(res["required_role"], "controller")
self.assertFalse(res["downstream_allowed"])
def test_unknown_still_ambiguous(self) -> None:
res = route_task_session(
"not_a_real_task",
active_profile="prgs-controller",
active_role_kind="controller",
allowed_in_current_session=False,
)
self.assertEqual(res["route_result"], ROUTE_AMBIGUOUS)
def test_capability_map_process_work_queue_is_controller(self) -> None:
self.assertEqual(
task_capability_map.required_role("process_work_queue"),
"controller",
)
self.assertEqual(
task_capability_map.required_permission("process_work_queue"),
"gitea.read",
)
class ControllerRoleMetadataTest(unittest.TestCase):
def test_normalize_role_kind_controller(self) -> None:
self.assertEqual(
nwb.normalize_role_kind("controller"),
"controller",
)
self.assertEqual(
nwb.normalize_role_kind("author", profile_name="prgs-controller"),
"controller",
)
self.assertEqual(
nwb.normalize_role_kind("reconciler", profile_name="prgs-controller"),
"controller",
)
def test_profile_role_kind_prefers_declared_controller(self) -> None:
# Import from worktree package path via sys.path already set by pytest.
import gitea_mcp_server as mcp
profile = {
"profile_name": "prgs-controller",
"role": "controller",
"allowed_operations": [
"gitea.read",
"gitea.issue.comment",
"gitea.pr.close",
],
"forbidden_operations": [
"gitea.pr.approve",
"gitea.pr.merge",
"gitea.pr.create",
"gitea.branch.push",
],
}
# Declared role wins even if permissions look reconciler-like.
self.assertEqual(mcp._profile_role_kind(profile), "controller")
# Name-based fallback.
profile_no_role = dict(profile)
profile_no_role["role"] = None
profile_no_role["role_kind"] = None
self.assertEqual(mcp._profile_role_kind(profile_no_role), "controller")
def test_permission_inference_without_controller_name_stays_reconciler(self) -> None:
import gitea_mcp_server as mcp
# Pure permission inference still may return reconciler when no controller
# declaration exists — that is intentional for reconciler profiles.
role = mcp._role_kind(
["gitea.read", "gitea.pr.close", "gitea.issue.comment"],
["gitea.pr.approve", "gitea.pr.merge", "gitea.pr.create", "gitea.branch.push"],
)
self.assertEqual(role, "reconciler")
class DashboardRemainsExplanatoryTest(unittest.TestCase):
def test_dashboard_prompt_points_at_allocator_not_self_select(self) -> None:
import workflow_dashboard as wd
self.assertIn("gitea_allocate_next_work", wd.PROMPT_CONTROLLER)
self.assertIn("process_work_queue", wd.PROMPT_CONTROLLER)
self.assertIn("never replaces allocator", wd.PROMPT_CONTROLLER.lower())
self.assertNotIn("self-select", wd.PROMPT_CONTROLLER.lower())
class ClassifySkipCrossRoleTest(unittest.TestCase):
def test_controller_cross_role_accepts_author_issue(self) -> None:
c = WorkCandidate(kind="issue", number=1, labels=("status:ready",))
self.assertIsNone(
classify_skip(
c,
role=ROLE_CONTROLLER,
terminal_pr=None,
allocation_mode=ALLOCATION_MODE_CROSS_ROLE,
)
)
def test_legacy_controller_skips_author_issue(self) -> None:
c = WorkCandidate(kind="issue", number=1, labels=("status:ready",))
reason = classify_skip(
c,
role=ROLE_CONTROLLER,
terminal_pr=None,
allocation_mode=ALLOCATION_MODE_ROLE_SCOPED,
)
self.assertIsNotNone(reason)
self.assertIn("does not require controller", reason or "")
if __name__ == "__main__":
unittest.main()
@@ -0,0 +1,650 @@
"""Publication of an unpublished local commit (#812 AC20).
Entry point B of #812: a registered worktree, clean, on its issue branch,
holding a local commit that has never been published. Exact-owner lease renewal
refuses such a claim for want of an observable remote head, and every existing
publication path is lock-derived, so the two predicates close a cycle around
work that is otherwise complete.
These tests exercise the disposition through its *evidence*, never through any
particular issue number: every case uses an arbitrary issue number against a
synthetic repository, and the same assertions hold for any other. Nothing here
reads, writes, or references the live protected worktree named in #812 AC17 —
that content is preserved evidence for the duration of this work, so the
fixtures below build their own repositories from scratch.
The remote is a local bare repository, so publication and read-after-write
verification are genuinely executed rather than mocked.
"""
from __future__ import annotations
import os
import subprocess
import sys
import tempfile
import unittest
from datetime import datetime, timedelta, timezone
from unittest.mock import patch
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
import branch_publish # noqa: E402
import issue_lock_provenance # noqa: E402
import issue_lock_renewal # noqa: E402
import issue_lock_store # noqa: E402
import mcp_server # noqa: E402
from mutation_profile_fixture import shared_mutation_env # noqa: E402
ISSUE = 9812
BRANCH = f"feat/issue-{ISSUE}-publish-fixture"
IDENTITY = "example-user"
PROFILE = "test-author-prgs"
ORG = "Scaled-Tech-Consulting"
REPO = "Gitea-Tools"
GIT_REMOTE = "prgs"
def _ts(hours: int) -> str:
return (
(datetime.now(timezone.utc) + timedelta(hours=hours))
.isoformat()
.replace("+00:00", "Z")
)
class _PublishBase(unittest.TestCase):
"""Real git repo + real bare remote + durable lock naming the caller.
The recorded owner pid is deliberately **this live process**. That mirrors
the production shape #812 documents, where the pid belongs to a long-running
MCP daemon rather than to a dead author client, and it proves publication
never depends on a dead process (#812 AC24).
"""
def setUp(self):
self.lock_dir = tempfile.TemporaryDirectory()
self.addCleanup(self.lock_dir.cleanup)
self.origin = tempfile.mkdtemp(prefix="issue812-origin-")
self.repo = tempfile.mkdtemp(prefix="issue812-work-")
for path in (self.origin, self.repo):
self.addCleanup(
lambda p=path: subprocess.run(["rm", "-rf", p], check=False)
)
self._init_repos()
self.remotes = patch.dict(
mcp_server.REMOTES,
{"prgs": {"host": "gitea.prgs.cc", "org": ORG, "repo": REPO}},
)
self.remotes.start()
self.addCleanup(patch.stopall)
mcp_server._IDENTITY_CACHE.clear()
# ── fixture construction ─────────────────────────────────────────────
def _git(self, *args, cwd=None):
return subprocess.run(
["git", "-C", cwd or self.repo, *args],
capture_output=True,
text=True,
check=True,
)
def _init_repos(self):
subprocess.run(
["git", "init", "-q", "--bare", "-b", "master", self.origin], check=True
)
self._git("init", "-q", "-b", "master")
self._git("config", "user.email", "[email protected]")
self._git("config", "user.name", "Test")
self._git("remote", "add", GIT_REMOTE, self.origin)
with open(os.path.join(self.repo, "seed.txt"), "w") as fh:
fh.write("seed\n")
self._git("add", "seed.txt")
self._git("commit", "-q", "-m", "seed")
self.base_sha = self._git("rev-parse", "HEAD").stdout.strip()
self._git("push", "-q", GIT_REMOTE, "master")
self._git("checkout", "-q", "-b", BRANCH)
with open(os.path.join(self.repo, "work.txt"), "w") as fh:
fh.write("unpublished implementation\n")
self._git("add", "work.txt")
self._git("commit", "-q", "-m", "unpublished implementation")
self.head_sha = self._git("rev-parse", "HEAD").stdout.strip()
self.worktree = os.path.realpath(self.repo)
def lock_path(self):
return issue_lock_store.lock_file_path(
remote="prgs", org=ORG, repo=REPO, issue_number=ISSUE,
lock_dir=self.lock_dir.name,
)
def write_lock(self, **overrides):
path = self.lock_path()
claimant = overrides.pop(
"claimant", {"username": IDENTITY, "profile": PROFILE}
)
pid = overrides.pop("session_pid", os.getpid())
lease = {
"operation_type": issue_lock_store.AUTHOR_ISSUE_WORK_LEASE,
"issue_number": ISSUE,
"pr_number": None,
"branch": overrides.get("branch_name", BRANCH),
"worktree_path": overrides.get("worktree_path", self.worktree),
"claimant": claimant,
"created_at": _ts(-2),
"last_heartbeat_at": _ts(-2),
# Expired: entry point B's lease has lapsed, which is precisely why
# renewal — and therefore a published head — is needed.
"expires_at": _ts(-1),
}
lease.update(overrides.pop("work_lease", {}))
data = {
"issue_number": ISSUE,
"branch_name": BRANCH,
"remote": "prgs",
"org": ORG,
"repo": REPO,
"worktree_path": self.worktree,
"session_pid": pid,
"pid": pid,
"lock_generation": 1,
"work_lease": lease,
"lock_provenance": issue_lock_provenance.build_sanctioned_lock_provenance(
tool="gitea_lock_issue", claimant=claimant
),
}
data.update(overrides)
data["lock_file_path"] = path
issue_lock_store.save_lock_file(path, data)
return path
def _tool_env(self):
env = shared_mutation_env(
PROFILE, include_example_repo=True,
GITEA_ISSUE_LOCK_DIR=self.lock_dir.name,
)
env["GITEA_ISSUE_LOCK_DIR"] = self.lock_dir.name
# These tests repoint PROJECT_ROOT at a synthetic repository so the
# registered-worktree proof runs for real. Pin the parity gate to the
# server's own startup head so that repointing does not read as a stale
# daemon; the gate itself stays live and enforced.
startup_head = mcp_server._STARTUP_PARITY.get("startup_head") or ""
env["GITEA_TEST_CURRENT_HEAD"] = startup_head
env["GITEA_TEST_LIVE_REMOTE_HEAD"] = startup_head
return env
# ── tool driver ──────────────────────────────────────────────────────
def run_publish(self, *, open_prs=None, expected_head=None, **kwargs):
"""Drive the public publication tool against the synthetic fixture."""
env = self._tool_env()
with patch(
"mcp_server._list_open_pulls", return_value=list(open_prs or [])
), patch(
"mcp_server._auth", return_value="token x"
), patch(
"mcp_server.get_auth_header", return_value="token x"
), patch(
"mcp_server._work_lease_claimant",
return_value={"username": IDENTITY, "profile": PROFILE},
), patch.object(
mcp_server, "PROJECT_ROOT", self.repo
), patch.dict(os.environ, env, clear=True):
os.environ["GITEA_ISSUE_LOCK_DIR"] = self.lock_dir.name
return mcp_server.gitea_publish_unpublished_issue_branch(
issue_number=kwargs.pop("issue_number", ISSUE),
branch_name=kwargs.pop("branch_name", BRANCH),
worktree_path=kwargs.pop("worktree_path", self.worktree),
expected_head=expected_head or self.head_sha,
remote="prgs",
git_remote_name=kwargs.pop("git_remote_name", GIT_REMOTE),
**kwargs,
)
def remote_head(self, branch=BRANCH):
res = subprocess.run(
["git", "-C", self.origin, "rev-parse", "--verify", "--quiet", branch],
capture_output=True, text=True, check=False,
)
return (res.stdout or "").strip() or None
class TestSuccessfulPublication(_PublishBase):
"""AC20 — the branch becomes observable and is verified after the write."""
def test_publishes_clean_unpublished_commit(self):
self.write_lock()
self.assertIsNone(self.remote_head(), "fixture must start unpublished")
result = self.run_publish()
self.assertTrue(result["success"], result.get("reasons"))
self.assertTrue(result["performed"])
self.assertTrue(result["published"])
self.assertTrue(result["verified"], "read-after-write must be proven")
self.assertEqual(result["remote_head_sha"], self.head_sha)
self.assertEqual(self.remote_head(), self.head_sha)
def test_publication_does_not_rewrite_the_commit(self):
self.write_lock()
self.run_publish()
# The published object is the same commit, not a copy or a rewrite.
self.assertEqual(self.remote_head(), self.head_sha)
self.assertEqual(
self._git("rev-parse", "HEAD").stdout.strip(), self.head_sha
)
def test_exact_next_action_names_the_lock_call(self):
self.write_lock()
result = self.run_publish()
self.assertIn("gitea_lock_issue", result["exact_next_action"])
class TestFailsClosed(_PublishBase):
"""AC20/AC9 — each refusal reason, exercised independently."""
def test_changed_local_head_refuses(self):
self.write_lock()
stale = self.base_sha # a real commit, but not the declared head
result = self.run_publish(expected_head=stale)
self.assertFalse(result["success"])
self.assertTrue(
any("local commit changed" in r for r in result["reasons"]),
result["reasons"],
)
self.assertIsNone(self.remote_head(), "refusal must not publish")
def test_abbreviated_sha_refuses(self):
self.write_lock()
result = self.run_publish(expected_head=self.head_sha[:8])
self.assertFalse(result["success"])
self.assertTrue(
any("40-character" in r for r in result["reasons"]), result["reasons"]
)
def test_dirty_tracked_worktree_refuses(self):
self.write_lock()
with open(os.path.join(self.repo, "work.txt"), "a") as fh:
fh.write("uncommitted edit\n")
result = self.run_publish()
self.assertFalse(result["success"])
self.assertTrue(
any("dirty tracked files" in r for r in result["reasons"]),
result["reasons"],
)
self.assertIn("work.txt", result["evidence"]["dirty_tracked_files"])
self.assertIsNone(self.remote_head())
def test_untracked_file_refuses(self):
self.write_lock()
with open(os.path.join(self.repo, "stray.txt"), "w") as fh:
fh.write("not committed\n")
result = self.run_publish()
self.assertFalse(result["success"])
self.assertTrue(
any("untracked files" in r for r in result["reasons"]), result["reasons"]
)
self.assertIn("stray.txt", result["evidence"]["untracked_files"])
self.assertIsNone(self.remote_head())
def test_unexpected_remote_head_refuses(self):
"""A remote head that is not an ancestor must never be overwritten."""
self.write_lock()
# Publish a divergent commit to the branch from a separate line.
self._git("checkout", "-q", "-b", "divergent", self.base_sha)
with open(os.path.join(self.repo, "other.txt"), "w") as fh:
fh.write("someone else's work\n")
self._git("add", "other.txt")
self._git("commit", "-q", "-m", "divergent")
divergent = self._git("rev-parse", "HEAD").stdout.strip()
self._git("push", "-q", GIT_REMOTE, f"{divergent}:refs/heads/{BRANCH}")
self._git("checkout", "-q", BRANCH)
result = self.run_publish()
self.assertFalse(result["success"])
self.assertTrue(
any("not an ancestor" in r for r in result["reasons"]), result["reasons"]
)
self.assertEqual(
self.remote_head(), divergent, "the other head must survive intact"
)
def test_fast_forward_remote_head_is_allowed(self):
"""An ancestor head is an honest fast-forward, not a conflict."""
self.write_lock()
self._git("push", "-q", GIT_REMOTE, f"{self.base_sha}:refs/heads/{BRANCH}")
result = self.run_publish()
self.assertTrue(result["success"], result.get("reasons"))
self.assertTrue(result["evidence"]["fast_forward_from_remote"])
self.assertEqual(self.remote_head(), self.head_sha)
def test_content_hash_mismatch_refuses(self):
self.write_lock()
wrong = {"work.txt": "0" * 64}
result = self.run_publish(expected_file_hashes=wrong)
self.assertFalse(result["success"])
self.assertTrue(
any("declared content hashes" in r for r in result["reasons"]),
result["reasons"],
)
self.assertFalse(result["evidence"]["file_hashes_verified"])
self.assertIsNone(self.remote_head())
def test_matching_content_hashes_publish(self):
self.write_lock()
digests = branch_publish.hash_worktree_files(self.worktree, ["work.txt"])
result = self.run_publish(expected_file_hashes=digests)
self.assertTrue(result["success"], result.get("reasons"))
self.assertTrue(result["evidence"]["file_hashes_verified"])
def test_missing_declared_file_refuses(self):
self.write_lock()
result = self.run_publish(expected_file_hashes={"absent.txt": "0" * 64})
self.assertFalse(result["success"])
self.assertTrue(
any("missing or unreadable" in r for r in result["reasons"]),
result["reasons"],
)
def test_foreign_claimant_refuses(self):
"""Ownership comes from the durable record, not from the caller."""
self.write_lock(claimant={"username": "someone-else", "profile": PROFILE})
result = self.run_publish()
self.assertFalse(result["success"])
self.assertTrue(
any("foreign claim" in r for r in result["reasons"]), result["reasons"]
)
self.assertIsNone(self.remote_head())
def test_foreign_profile_refuses(self):
self.write_lock(
claimant={"username": IDENTITY, "profile": "test-reviewer-prgs"}
)
result = self.run_publish()
self.assertFalse(result["success"])
self.assertTrue(
any("claimant profile" in r for r in result["reasons"]), result["reasons"]
)
def test_absent_lock_record_refuses(self):
"""No recorded claim means this cannot be used to bypass the lock."""
result = self.run_publish() # no write_lock()
self.assertFalse(result["success"])
self.assertTrue(
any("no durable issue-lock record" in r for r in result["reasons"]),
result["reasons"],
)
self.assertIsNone(self.remote_head())
def test_branch_mismatch_against_lock_refuses(self):
self.write_lock(branch_name=f"feat/issue-{ISSUE}-different")
result = self.run_publish()
self.assertFalse(result["success"])
self.assertTrue(
any("records branch" in r for r in result["reasons"]), result["reasons"]
)
def test_worktree_mismatch_against_lock_refuses(self):
self.write_lock(worktree_path="/tmp/some/other/worktree")
result = self.run_publish()
self.assertFalse(result["success"])
self.assertTrue(
any("records worktree" in r for r in result["reasons"]), result["reasons"]
)
def test_competing_open_pr_on_another_branch_refuses(self):
self.write_lock()
competing = [{"number": 4242, "head": {"ref": f"fix/issue-{ISSUE}-rival"}}]
result = self.run_publish(open_prs=competing)
self.assertFalse(result["success"])
self.assertTrue(
any("already claim issue" in r for r in result["reasons"]),
result["reasons"],
)
self.assertIsNone(self.remote_head())
def test_open_pr_on_the_same_branch_is_not_competing(self):
"""This branch's own PR is not a rival claim against itself."""
self.write_lock()
own = [{"number": 77, "head": {"ref": BRANCH}}]
result = self.run_publish(open_prs=own)
self.assertTrue(result["success"], result.get("reasons"))
class TestGuardStrictnessPreserved(_PublishBase):
"""AC15 — publication is an operation, never a weakening of the guards."""
def test_non_issue_branch_refuses(self):
self._git("checkout", "-q", "-b", "scratch/not-issue-linked")
self.write_lock(branch_name="scratch/not-issue-linked")
result = self.run_publish(branch_name="scratch/not-issue-linked")
self.assertFalse(result["success"])
self.assertTrue(
any("issue-linked" in r for r in result["reasons"]), result["reasons"]
)
def test_stable_branch_refuses(self):
self.write_lock(branch_name="master")
result = self.run_publish(branch_name="master")
self.assertFalse(result["success"])
self.assertTrue(
any("issue-linked" in r or "stable branch" in r for r in result["reasons"]),
result["reasons"],
)
def test_branch_number_must_match_the_issue(self):
other = "feat/issue-7777-mismatched"
self._git("checkout", "-q", "-b", other)
self.write_lock(branch_name=other)
result = self.run_publish(branch_name=other)
self.assertFalse(result["success"])
self.assertTrue(
any("does not carry issue number" in r for r in result["reasons"]),
result["reasons"],
)
def test_unregistered_worktree_refuses(self):
"""#713 — an improvised directory is not a registered worktree."""
path = self.write_lock()
assessment = branch_publish.assess_unpublished_commit_publication(
issue_lock_store.read_lock_file(path),
issue_number=ISSUE, branch_name=BRANCH, worktree_path=self.worktree,
expected_head=self.head_sha, remote="prgs", org=ORG, repo=REPO,
identity=IDENTITY, profile=PROFILE,
worktree_state={
"current_branch": BRANCH, "porcelain_status": "",
"head_sha": self.head_sha,
},
worktree_registered=False,
remote_probe={"probe_ok": True, "remote_branch_exists": False},
)
self.assertEqual(assessment["outcome"], branch_publish.REFUSED)
self.assertTrue(
any("not listed in git worktree list" in r
for r in assessment["reasons"]),
assessment["reasons"],
)
def test_unobservable_remote_refuses(self):
"""An unknown remote state must not be mistaken for an absent branch."""
self.write_lock()
result = self.run_publish(git_remote_name="no-such-remote")
self.assertFalse(result["success"])
self.assertTrue(
any("could not be observed" in r for r in result["reasons"]),
result["reasons"],
)
class TestRecordSeparation(_PublishBase):
"""AC23 — the durable issue lock and the workflow lease are distinct."""
def test_publication_leaves_the_issue_lock_byte_identical(self):
path = self.write_lock()
with open(path, "rb") as fh:
before = fh.read()
result = self.run_publish()
self.assertTrue(result["success"], result.get("reasons"))
with open(path, "rb") as fh:
after = fh.read()
self.assertEqual(before, after, "publication must not mutate the lock record")
self.assertFalse(result["issue_lock_record_mutated"])
self.assertFalse(result["workflow_lease_touched"])
def test_refusal_also_reports_untouched_records(self):
result = self.run_publish() # refuses: no lock record
self.assertFalse(result["issue_lock_record_mutated"])
self.assertFalse(result["workflow_lease_touched"])
def test_lock_generation_is_not_advanced(self):
path = self.write_lock()
self.run_publish()
lock = issue_lock_store.read_lock_file(path)
self.assertEqual(lock["lock_generation"], 1)
class TestTruthfulProcessEvidence(_PublishBase):
"""AC24 — a live daemon pid is never represented as a dead process."""
def test_live_recorded_pid_does_not_block_publication(self):
# The recorded pid is this live process, standing in for the live MCP
# daemon. Reclaim would refuse here; publication legitimately does not.
path = self.write_lock(session_pid=os.getpid())
lock = issue_lock_store.read_lock_file(path)
self.assertEqual(lock["pid"], os.getpid())
result = self.run_publish()
self.assertTrue(result["success"], result.get("reasons"))
self.assertEqual(self.remote_head(), self.head_sha)
def test_liveness_is_not_consulted_as_evidence(self):
self.write_lock(session_pid=os.getpid())
result = self.run_publish()
self.assertFalse(result["evidence"]["owner_pid_liveness_consulted"])
def test_reclaim_still_refuses_for_the_same_live_pid(self):
"""Publication does not soften the reclaim predicate it routes around."""
path = self.write_lock(session_pid=os.getpid())
lock = issue_lock_store.read_lock_file(path)
reclaim = issue_lock_store.assess_expired_lock_reclaim(lock)
self.assertFalse(reclaim["reclaim_allowed"])
class TestIdempotentRetry(_PublishBase):
"""AC20 — retry is safe and read-after-write is proven every time."""
def test_second_publication_reports_already_published(self):
self.write_lock()
first = self.run_publish()
self.assertTrue(first["performed"])
second = self.run_publish()
self.assertTrue(second["success"], second.get("reasons"))
self.assertFalse(second["performed"], "no second push is needed")
self.assertTrue(second["published"])
self.assertTrue(second["verified"])
self.assertEqual(second["outcome"], branch_publish.ALREADY_PUBLISHED)
self.assertEqual(self.remote_head(), self.head_sha)
class TestDryRun(_PublishBase):
"""AC12 — dry run reports the decision and mutates nothing."""
def test_dry_run_reports_intent_without_publishing(self):
self.write_lock()
result = self.run_publish(dry_run=True)
self.assertTrue(result["success"])
self.assertTrue(result["dry_run"])
self.assertTrue(result["would_publish"])
self.assertFalse(result["performed"])
self.assertIsNone(self.remote_head(), "dry run must not publish")
def test_dry_run_and_apply_agree_on_a_refusal(self):
"""AC11 — the reported decision does not depend on which mode ran."""
self.write_lock(claimant={"username": "someone-else", "profile": PROFILE})
dry = self.run_publish(dry_run=True)
applied = self.run_publish()
self.assertFalse(dry["success"])
self.assertFalse(applied["success"])
self.assertEqual(dry["reasons"], applied["reasons"])
class TestRenewalUnblocked(_PublishBase):
"""AC20/AC21 — renewal is permitted only after verified publication."""
def _renewal(self, remote_head):
return issue_lock_renewal.assess_exact_owner_lease_renewal(
issue_lock_store.read_lock_file(self.lock_path()),
issue_number=ISSUE, branch_name=BRANCH, worktree_path=self.worktree,
remote="prgs", org=ORG, repo=REPO,
identity=IDENTITY, profile=PROFILE,
current_branch=BRANCH, porcelain_status="", worktree_exists=True,
head_sha=self.head_sha, remote_head_sha=remote_head,
)
def test_renewal_refuses_before_publication(self):
self.write_lock()
decision = self._renewal(None)
self.assertFalse(decision["renewal_sanctioned"])
self.assertTrue(
any("unpublished branch" in r for r in decision["reasons"]),
decision["reasons"],
)
def test_renewal_is_sanctioned_after_publication(self):
self.write_lock()
result = self.run_publish()
self.assertTrue(result["verified"], result.get("reasons"))
decision = self._renewal(self.remote_head())
self.assertTrue(decision["renewal_sanctioned"], decision["reasons"])
class TestProtectedAssetUntouched(unittest.TestCase):
"""AC17 — no test or fixture may reference the protected worktree."""
def test_no_reference_to_the_protected_worktree(self):
here = os.path.dirname(os.path.abspath(__file__))
root = os.path.dirname(here)
needle = "issue-635-project-registry" + "-api"
for path in (
os.path.join(here, "test_issue_812_publish_unpublished_commit.py"),
os.path.join(root, "branch_publish.py"),
):
with open(path, "r", encoding="utf-8") as fh:
body = fh.read()
self.assertNotIn(needle, body)
if __name__ == "__main__": # pragma: no cover
unittest.main()
@@ -0,0 +1,622 @@
"""Publication preflight must receive the caller's worktree (#815).
``gitea_publish_unpublished_issue_branch`` takes a **required** ``worktree_path``
but resolved it only *after* ``verify_preflight_purity`` had already run. Every
workspace-resolution layer behind that preflight — canonical root, root checkout,
create-issue bootstrap, the #618 branches-only guard, issue scope, and anti-stomp
— therefore received ``None`` and fell back to the MCP process root. A daemon
rooted at the stable control checkout refused a valid registered issue worktree
that the caller had explicitly supplied, before the publication assessor ever ran.
The #812 suite could not see this. Its fixture sets ``self.worktree =
os.path.realpath(self.repo)`` and patches ``PROJECT_ROOT`` to that same path, so
the fallback resolved to the very worktree the argument named. The production
topology — control checkout on a stable branch, issue worktree somewhere else —
was never constructed, and preflight additionally no-ops under pytest unless
production guards are forced on.
These tests build that topology honestly:
* ``PROJECT_ROOT`` is a control checkout sitting on ``master``;
* the registered issue worktree is a genuinely separate path under ``branches/``;
* ``GITEA_TEST_FORCE_PRODUCTION_GUARDS`` is set so the #618 guard really runs;
* no patch makes the issue worktree appear to be ``PROJECT_ROOT``.
Nothing here reads, writes, or references the protected worktree named in #812
AC17 and #815 AC9. Every fixture is built from scratch against a local bare
remote, so publication and read-after-write verification genuinely execute.
"""
from __future__ import annotations
import os
import subprocess
import sys
import tempfile
import unittest
from datetime import datetime, timedelta, timezone
from unittest.mock import patch
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
import issue_lock_provenance # noqa: E402
import issue_lock_store # noqa: E402
import mcp_server # noqa: E402
from mutation_profile_fixture import shared_mutation_env # noqa: E402
ISSUE = 9815
BRANCH = f"feat/issue-{ISSUE}-forwarding-fixture"
WORKTREE_DIRNAME = BRANCH.replace("/", "-")
IDENTITY = "example-user"
PROFILE = "test-author-prgs"
ORG = "Scaled-Tech-Consulting"
REPO = "Gitea-Tools"
GIT_REMOTE = "prgs"
AUTHOR_PROFILE = {
"profile_name": "prgs-author",
"role": "author",
"allowed_operations": [
"gitea.read", "gitea.issue.create", "gitea.issue.comment",
"gitea.pr.create", "gitea.repo.commit", "gitea.branch.push",
],
"forbidden_operations": [],
"audit_label": "prgs-author",
}
def _ts(hours: int) -> str:
return (
(datetime.now(timezone.utc) + timedelta(hours=hours))
.isoformat()
.replace("+00:00", "Z")
)
class TestPreflightReceivesTheWorktree(unittest.TestCase):
"""AC1 — the supplied path reaches ``verify_preflight_purity`` itself.
Follows the #735 capture pattern: replace preflight with a recorder that
raises, so the argument can be proven forwarded without performing the
mutation. This is the direct unit-level statement of the defect.
"""
def _capture_preflight(self, **kwargs):
captured: dict = {}
def _capture(*a, **kw):
captured.update(kw)
captured["_args"] = a
raise RuntimeError("capture-only")
with patch.object(
mcp_server, "verify_preflight_purity", side_effect=_capture
), patch.object(
mcp_server, "get_profile", return_value=AUTHOR_PROFILE
), patch.object(
mcp_server, "_resolve",
return_value=("gitea.prgs.cc", ORG, REPO),
), patch.object(
mcp_server, "_auth", return_value="token fake",
), patch.object(
mcp_server.role_session_router,
"check_author_mutation_after_reviewer_stop",
return_value=(True, []),
), patch.object(
mcp_server, "_namespace_mutation_block", return_value=None
), patch.object(
mcp_server, "_profile_permission_block", return_value=None
):
try:
mcp_server.gitea_publish_unpublished_issue_branch(**kwargs)
except RuntimeError as exc:
if "capture-only" not in str(exc) and not captured:
raise
self.assertTrue(
captured,
"gitea_publish_unpublished_issue_branch never called "
"verify_preflight_purity",
)
return captured
def _base_kwargs(self, **overrides):
kwargs = {
"issue_number": ISSUE,
"branch_name": BRANCH,
"worktree_path": "/tmp/issue-815-explicit-worktree",
"expected_head": "a" * 40,
"remote": "prgs",
"org": ORG,
"repo": REPO,
"git_remote_name": GIT_REMOTE,
}
kwargs.update(overrides)
return kwargs
def test_explicit_worktree_path_reaches_preflight(self):
captured = self._capture_preflight(**self._base_kwargs())
self.assertEqual(
captured.get("worktree_path"),
os.path.realpath(os.path.abspath("/tmp/issue-815-explicit-worktree")),
"the authoritative worktree_path must be forwarded into preflight",
)
def test_forwarded_path_is_the_one_publication_uses(self):
"""AC4 — preflight and publication must judge the same resolved path."""
raw = "/tmp/issue-815-explicit-worktree/./"
captured = self._capture_preflight(**self._base_kwargs(worktree_path=raw))
expected = os.path.realpath(os.path.abspath(raw.strip()))
self.assertEqual(captured.get("worktree_path"), expected)
def test_blank_worktree_path_forwards_none(self):
"""AC5/AC8 — nothing usable supplied keeps the fail-closed fallback."""
for blank in ("", " "):
with self.subTest(blank=repr(blank)):
captured = self._capture_preflight(
**self._base_kwargs(worktree_path=blank)
)
self.assertIsNone(
captured.get("worktree_path"),
"a blank worktree must not resolve to the process cwd",
)
def test_org_repo_and_task_forwarding_are_not_regressed(self):
"""AC6 — #735's org/repo forwarding and the task name still hold."""
captured = self._capture_preflight(**self._base_kwargs())
self.assertEqual(captured.get("org"), ORG)
self.assertEqual(captured.get("repo"), REPO)
self.assertEqual(captured.get("task"), "publish_unpublished_branch")
class _ProductionTopologyBase(unittest.TestCase):
"""Control checkout on master + a distinct registered issue worktree.
This is the shape the production daemon runs in and the shape the #812
fixture never built. ``PROJECT_ROOT`` is the control checkout; the issue
worktree is a real registered worktree at a different path; production
guards are forced on so the #618 branches-only guard genuinely evaluates.
"""
def setUp(self):
self.lock_dir = tempfile.TemporaryDirectory()
self.addCleanup(self.lock_dir.cleanup)
self.origin = tempfile.mkdtemp(prefix="issue815-origin-")
self.control = tempfile.mkdtemp(prefix="issue815-control-")
for path in (self.origin, self.control):
self.addCleanup(
lambda p=path: subprocess.run(["rm", "-rf", p], check=False)
)
self._init_repos()
self.remotes = patch.dict(
mcp_server.REMOTES,
{"prgs": {"host": "gitea.prgs.cc", "org": ORG, "repo": REPO}},
)
self.remotes.start()
self.addCleanup(patch.stopall)
mcp_server._IDENTITY_CACHE.clear()
def _git(self, *args, cwd=None):
return subprocess.run(
["git", "-C", cwd or self.control, *args],
capture_output=True, text=True, check=True,
)
def _init_repos(self):
subprocess.run(
["git", "init", "-q", "--bare", "-b", "master", self.origin], check=True
)
self._git("init", "-q", "-b", "master")
self._git("config", "user.email", "[email protected]")
self._git("config", "user.name", "Test")
self._git("remote", "add", GIT_REMOTE, self.origin)
with open(os.path.join(self.control, "seed.txt"), "w") as fh:
fh.write("seed\n")
# The real repository gitignores branches/, so a registered worktree
# living there does not dirty the stable control checkout. Mirror that,
# or the #615 dirty-runtime block fires on the worktree we just created.
with open(os.path.join(self.control, ".gitignore"), "w") as fh:
fh.write("branches/\n")
self._git("add", "seed.txt", ".gitignore")
self._git("commit", "-q", "-m", "seed")
self.base_sha = self._git("rev-parse", "HEAD").stdout.strip()
self._git("push", "-q", GIT_REMOTE, "master")
# The control checkout STAYS on master. This is the whole point: the
# daemon's process root is the stable control checkout, never the
# worktree the publication targets.
self.worktree = os.path.realpath(
os.path.join(self.control, "branches", WORKTREE_DIRNAME)
)
self._git("worktree", "add", "-q", "-b", BRANCH, self.worktree, "master")
with open(os.path.join(self.worktree, "work.txt"), "w") as fh:
fh.write("unpublished implementation\n")
self._git("add", "work.txt", cwd=self.worktree)
self._git("commit", "-q", "-m", "unpublished implementation", cwd=self.worktree)
self.head_sha = self._git("rev-parse", "HEAD", cwd=self.worktree).stdout.strip()
self.control_branch = self._git(
"rev-parse", "--abbrev-ref", "HEAD"
).stdout.strip()
# ── durable lock naming the caller and the issue worktree ────────────
def lock_path(self):
return issue_lock_store.lock_file_path(
remote="prgs", org=ORG, repo=REPO, issue_number=ISSUE,
lock_dir=self.lock_dir.name,
)
def write_lock(self, *, bind_session=True, **overrides):
path = self.lock_path()
claimant = overrides.pop(
"claimant", {"username": IDENTITY, "profile": PROFILE}
)
pid = overrides.pop("session_pid", os.getpid())
lease = {
"operation_type": issue_lock_store.AUTHOR_ISSUE_WORK_LEASE,
"issue_number": ISSUE,
"pr_number": None,
"branch": overrides.get("branch_name", BRANCH),
"worktree_path": overrides.get("worktree_path", self.worktree),
"claimant": claimant,
"created_at": _ts(-2),
"last_heartbeat_at": _ts(-2),
"expires_at": _ts(-1),
}
lease.update(overrides.pop("work_lease", {}))
data = {
"issue_number": ISSUE,
"branch_name": BRANCH,
"remote": "prgs",
"org": ORG,
"repo": REPO,
"worktree_path": self.worktree,
"session_pid": pid,
"pid": pid,
"lock_generation": 1,
"work_lease": lease,
"lock_provenance": issue_lock_provenance.build_sanctioned_lock_provenance(
tool="gitea_lock_issue", claimant=claimant
),
}
data.update(overrides)
data["lock_file_path"] = path
issue_lock_store.save_lock_file(path, data)
# Bind the session pointer so the #683 issue-scope guard resolves an
# owning issue for this author session. In real production the publish
# task does not require a session lock — require_author_lock is keyed on
# the test-only production_guards_forced() flag, which this suite must
# set to make preflight run at all — so this pointer is fixture
# scaffolding to clear a guard production would not apply here, never a
# softening of the worktree-forwarding behaviour under test. The
# preflight-negative cases below leave it unbound precisely so the #618
# guard is reached with no session fallback to rescue a bad worktree.
if bind_session:
pointer = {
"pid": os.getpid(),
"lock_file_path": path,
"issue_number": ISSUE,
"branch_name": data["branch_name"],
"remote": "prgs",
"org": ORG,
"repo": REPO,
}
issue_lock_store.save_lock_file(
issue_lock_store.session_pointer_path(self.lock_dir.name), pointer
)
return path
def _tool_env(self):
env = shared_mutation_env(
PROFILE, include_example_repo=True,
GITEA_ISSUE_LOCK_DIR=self.lock_dir.name,
)
env["GITEA_ISSUE_LOCK_DIR"] = self.lock_dir.name
# The defect only exists where preflight actually runs. Under pytest the
# production root/branches/scope guards are skipped unless forced on, so
# force them: this test exists to exercise the #618 guard, not to bypass
# it. Parity is pinned to the server's own startup head so repointing
# PROJECT_ROOT does not read as a stale daemon.
env["GITEA_TEST_FORCE_PRODUCTION_GUARDS"] = "1"
# Production is a promoted stable-control runtime. The pytest process
# itself runs from a branches/ worktree, which the #615 runtime-mode
# gate correctly classifies as dev-test; declaring the sanctioned mode
# models the production daemon rather than defeating the gate. Without
# this, forcing production guards on would trip the *runtime-mode* block
# for a reason unrelated to the #815 worktree-forwarding defect.
env["GITEA_MCP_RUNTIME_MODE"] = "stable-control"
startup_head = mcp_server._STARTUP_PARITY.get("startup_head") or ""
env["GITEA_TEST_CURRENT_HEAD"] = startup_head
env["GITEA_TEST_LIVE_REMOTE_HEAD"] = startup_head
return env
def run_publish(self, *, open_prs=None, expected_head=None, **kwargs):
"""Drive the public tool with PROJECT_ROOT pinned to the CONTROL checkout."""
env = self._tool_env()
with patch(
"mcp_server._list_open_pulls", return_value=list(open_prs or [])
), patch(
"mcp_server._auth", return_value="token x"
), patch(
"mcp_server.get_auth_header", return_value="token x"
), patch(
"mcp_server._work_lease_claimant",
return_value={"username": IDENTITY, "profile": PROFILE},
), patch.object(
# NOTE: the control checkout — deliberately NOT self.worktree.
mcp_server, "PROJECT_ROOT", self.control
), patch.dict(os.environ, env, clear=True):
os.environ["GITEA_ISSUE_LOCK_DIR"] = self.lock_dir.name
return mcp_server.gitea_publish_unpublished_issue_branch(
issue_number=kwargs.pop("issue_number", ISSUE),
branch_name=kwargs.pop("branch_name", BRANCH),
worktree_path=kwargs.pop("worktree_path", self.worktree),
expected_head=expected_head or self.head_sha,
remote="prgs",
git_remote_name=kwargs.pop("git_remote_name", GIT_REMOTE),
**kwargs,
)
def remote_head(self, branch=BRANCH):
res = subprocess.run(
["git", "-C", self.origin, "rev-parse", "--verify", "--quiet", branch],
capture_output=True, text=True, check=False,
)
return (res.stdout or "").strip() or None
class TestForwardingClearsThe618Guard(_ProductionTopologyBase):
"""AC2 — the faithful production reproduction, and the sharpest fix proof.
The production recovery worker had **no** session issue lock — acquiring one
was the very thing the deadlock prevented — so preflight had nothing but the
explicit ``worktree_path`` argument to resolve the workspace from. This class
reproduces exactly that: no session pointer is bound, so there is no
author-lock fallback to rescue a dropped argument.
With the argument forwarded (fixed source) the #618 branches-only guard
accepts the registered issue worktree and the call advances to the next
guard. With the argument dropped (the buggy source this issue reports)
preflight falls back to ``PROJECT_ROOT`` — the stable control checkout — and
the #618 guard traps the call there. The two outcomes are told apart by the
guard that fired, on its own error text.
This test therefore *fails* against the unpatched source (the call is trapped
at #618 instead of clearing it), which is what makes it a regression rather
than a smoke test.
"""
_CONTROL_CHECKOUT_MARKERS = ("stable control checkout", "#618")
def test_explicit_worktree_clears_618_without_a_session_lock(self):
# No write_lock(): the session is deliberately unbound, as in production.
with self.assertRaises(RuntimeError) as ctx:
self.run_publish()
message = str(ctx.exception)
# The workspace guard is satisfied — the failure is the *later* scope
# guard (no owning issue), never the control-checkout refusal. If the
# argument were dropped, this call would be trapped at #618 instead.
for marker in self._CONTROL_CHECKOUT_MARKERS:
self.assertNotIn(
marker, message,
f"the explicit worktree must clear #618; got a control-checkout "
f"refusal instead: {message}",
)
self.assertIn(
"owning issue", message,
f"expected the downstream scope guard to fire, got: {message}",
)
self.assertIsNone(self.remote_head())
def test_dropped_argument_would_be_trapped_at_618(self):
# Simulate the buggy call shape directly: no session lock, and preflight
# given no worktree, exactly as the unpatched source left it. This pins
# the control-checkout refusal that the fix eliminates, so the pair of
# tests brackets the defect from both sides regardless of which source
# version is loaded.
env = self._tool_env()
with patch(
"mcp_server._list_open_pulls", return_value=[]
), patch(
"mcp_server._auth", return_value="token x"
), patch(
"mcp_server.get_auth_header", return_value="token x"
), patch(
"mcp_server._work_lease_claimant",
return_value={"username": IDENTITY, "profile": PROFILE},
), patch.object(
mcp_server, "PROJECT_ROOT", self.control
), patch.dict(os.environ, env, clear=True):
os.environ["GITEA_ISSUE_LOCK_DIR"] = self.lock_dir.name
with self.assertRaises(RuntimeError) as ctx:
# Drive verify_preflight_purity the way the buggy body did:
# no worktree_path forwarded at all.
mcp_server.verify_preflight_purity(
"prgs",
task="publish_unpublished_branch",
org=ORG,
repo=REPO,
)
message = str(ctx.exception)
self.assertTrue(
any(m in message for m in self._CONTROL_CHECKOUT_MARKERS),
f"a dropped worktree must trap at the control checkout: {message}",
)
self.assertIsNone(self.remote_head())
class TestProductionTopologyPublishes(_ProductionTopologyBase):
"""AC2/AC4/AC7 — the explicit registered worktree is what preflight validates."""
def test_fixture_is_genuinely_the_production_topology(self):
"""Guard the guard: if this drifts, the regression stops meaning anything."""
self.assertNotEqual(
os.path.realpath(self.control), self.worktree,
"the issue worktree must not be PROJECT_ROOT",
)
self.assertEqual(
self.control_branch, "master",
"the control checkout must sit on a stable branch",
)
self.assertTrue(
os.path.realpath(self.worktree).startswith(
os.path.realpath(os.path.join(self.control, "branches")) + os.sep
),
"the issue worktree must live under branches/",
)
listed = subprocess.run(
["git", "-C", self.control, "worktree", "list"],
capture_output=True, text=True, check=True,
).stdout
self.assertIn(
self.worktree, listed, "the issue worktree must be genuinely registered"
)
def test_publishes_from_a_control_rooted_daemon(self):
"""The exact production failure: this refused with #618 before the fix."""
self.write_lock()
self.assertIsNone(self.remote_head(), "fixture must start unpublished")
res = self.run_publish()
self.assertTrue(res.get("success"), res)
self.assertTrue(res.get("performed"), res)
self.assertEqual(self.remote_head(), self.head_sha)
def test_dry_run_uses_the_explicit_worktree(self):
"""AC4 — dry-run reaches the same decision without publishing."""
self.write_lock()
res = self.run_publish(dry_run=True)
self.assertTrue(res.get("success"), res)
self.assertFalse(res.get("performed"), res)
self.assertTrue(res.get("would_publish"), res)
self.assertIsNone(self.remote_head(), "dry-run must not publish")
def test_dry_run_and_apply_agree_on_the_same_worktree(self):
"""AC4 — both paths resolve the same workspace, so both succeed."""
self.write_lock()
dry = self.run_publish(dry_run=True)
self.assertTrue(dry.get("would_publish"), dry)
applied = self.run_publish()
self.assertTrue(applied.get("performed"), applied)
self.assertEqual(self.remote_head(), self.head_sha)
def test_read_after_write_verification_still_runs(self):
"""AC6 — PR #814's post-publication verification is unchanged."""
self.write_lock()
res = self.run_publish()
self.assertTrue(res.get("verified"), res)
self.assertEqual(res.get("remote_head_sha"), self.head_sha)
class TestProductionTopologyFailsClosed(_ProductionTopologyBase):
"""AC3/AC5/AC8 — the fix does not weaken any refusal.
A refusal reaches the caller by one of two mechanisms, and this class holds
them apart deliberately. A bad *workspace* is caught by the #618 preflight
guard, which raises before the assessor is built. A bad *content/ownership*
fact passes preflight (the worktree itself is fine) and is then refused by
the publication assessor, which returns ``success: False``. Both are
fail-closed; asserting the wrong mechanism would hide a regression.
"""
# ── #618 preflight refusals: no session lock, so nothing rescues a bad
# workspace and the guard fires exactly as it does in production ──────
def _assert_preflight_raises(self, **kwargs):
with self.assertRaises(RuntimeError) as ctx:
self.run_publish(**kwargs)
self.assertIsNone(
self.remote_head(), "a blocked publication must not reach the remote"
)
return str(ctx.exception)
def test_blank_worktree_path_fails_closed_via_618(self):
"""AC5 — a blank path forwards None, so preflight sees the control root."""
for blank in ("", " "):
with self.subTest(blank=repr(blank)):
message = self._assert_preflight_raises(worktree_path=blank)
self.assertIn("618", message)
def test_control_checkout_as_worktree_fails_closed_via_618(self):
"""AC5 — naming the stable control checkout explicitly is still refused."""
message = self._assert_preflight_raises(worktree_path=self.control)
self.assertIn("618", message)
def test_unregistered_directory_fails_closed(self):
"""AC3 — a plain directory under branches/ is not a registered worktree."""
bogus = os.path.join(self.control, "branches", "not-a-worktree")
os.makedirs(bogus, exist_ok=True)
self._assert_preflight_raises(worktree_path=bogus)
def test_missing_worktree_path_fails_closed(self):
"""AC3 — a path that does not exist is refused, not silently replaced."""
missing = os.path.join(self.control, "branches", "absent-worktree")
self._assert_preflight_raises(worktree_path=missing)
# ── assessor refusals: preflight passes on a valid worktree, then the
# publication assessor refuses on content/ownership evidence ──────────
def _assert_assessor_refuses(self, **kwargs):
res = self.run_publish(**kwargs)
self.assertFalse(res.get("success"), res)
self.assertFalse(res.get("performed"), res)
self.assertIsNone(self.remote_head())
return res
def test_changed_local_head_still_refuses(self):
"""AC6 — the declared expected_head remains authoritative."""
self.write_lock()
self._assert_assessor_refuses(expected_head="b" * 40)
def test_foreign_claimant_still_refuses(self):
"""AC6 — ownership still comes from the durable lock record."""
self.write_lock(claimant={"username": "someone-else", "profile": PROFILE})
self._assert_assessor_refuses()
def test_dirty_worktree_still_refuses(self):
"""AC6 — cleanliness enforcement survives the forwarding change."""
self.write_lock()
with open(os.path.join(self.worktree, "work.txt"), "a") as fh:
fh.write("uncommitted drift\n")
self._assert_assessor_refuses()
def test_competing_open_pr_still_refuses(self):
"""AC6 — a rival claim on another branch still blocks."""
self.write_lock()
self._assert_assessor_refuses(
open_prs=[{"number": 4242, "head": {"ref": f"fix/issue-{ISSUE}-rival"}}]
)
def test_issue_lock_record_is_not_mutated_by_a_refusal(self):
"""AC6 — record separation (#812 AC23) is unaffected by this change."""
path = self.write_lock()
with open(path, "rb") as fh:
before = fh.read()
self._assert_assessor_refuses(expected_head="c" * 40)
with open(path, "rb") as fh:
self.assertEqual(before, fh.read())
class TestProtectedFixtureNotReferenced(unittest.TestCase):
"""AC9 — this regression never names the protected #635 fixture.
The forbidden tokens are reconstructed from fragments so this assertion
file does not itself contain them and produce a false positive.
"""
def test_no_reference_to_the_protected_worktree(self):
forbidden = [
"issue-635-" + "project-registry-api",
"b2f6e9a6dc40e9651ef8" + "76f322dd0a68bddebfd8",
]
here = os.path.abspath(__file__)
with open(here, "r", encoding="utf-8") as fh:
text = fh.read()
for token in forbidden:
self.assertNotIn(
token, text,
f"the protected #635 fixture must not be referenced: {token}",
)
if __name__ == "__main__":
unittest.main()
+269
View File
@@ -78,6 +78,95 @@ class TestBlockReasonsAndReport(unittest.TestCase):
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):
def test_test_override_takes_precedence(self):
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(""))
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):
"""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.srv._STARTUP_PARITY = {"root": self.srv.PROJECT_ROOT,
"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):
self.srv._STARTUP_PARITY = self._saved
@@ -147,6 +386,36 @@ class TestServerWiring(unittest.TestCase):
self.assertTrue(out["in_parity"])
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__":
unittest.main()
@@ -139,6 +139,9 @@ EXPECTED_ROLE_EXCLUSIVE_TASKS = frozenset(
"gitea_release_merger_pr_lease",
"create_branch",
"push_branch",
# #812 AC20: publishing an unpublished local head is author-only for the
# same reason every other push is — it writes a branch to the remote.
"publish_unpublished_branch",
"create_pr",
"commit_files",
"gitea_commit_files",
+149
View File
@@ -0,0 +1,149 @@
"""Documentation acceptance for the web console architecture ADR (#632 / epic #631).
Enforces the acceptance criteria of issue #632:
* AC1 — the ADR exists and covers layers, authority, phases, API versioning,
and a page map.
* AC2 — every #631 child (#632#651) maps to at least one architectural
component.
* AC3 — the closed MVP (#425#436) is stated as foundation, not recreated.
* AC4 — forbidden paths are explicit: raw provider incidents as work,
browser-held tokens, process-kill recovery.
* AC5 — a controller can approve the document without reading chat history.
Plus the linkage requirement: ``docs/webui-local-dev.md`` cross-links the ADR.
"""
from pathlib import Path
REPO_ROOT = Path(__file__).resolve().parent.parent
ADR = (
REPO_ROOT
/ "docs"
/ "architecture"
/ "webui-control-plane-console-architecture-adr.md"
)
ADR_BASENAME = "webui-control-plane-console-architecture-adr.md"
LOCAL_DEV = REPO_ROOT / "docs" / "webui-local-dev.md"
# Epic #631 children, phases 1-4 (twenty capability areas).
EPIC_CHILDREN = tuple(f"#{number}" for number in range(632, 652))
def _read(path: Path) -> str:
assert path.is_file(), f"missing {path.relative_to(REPO_ROOT)}"
return path.read_text(encoding="utf-8")
def test_ac1_adr_exists_with_required_sections():
text = _read(ADR)
lower = text.lower()
assert text.lstrip().startswith("#"), "ADR lacks a title"
assert "#631" in text and "#632" in text
for heading in (
"## 2. Decision summary",
"## 4. Authority boundaries",
"## 5. Request flow and the redaction boundary",
"## 6. API naming and versioning",
"## 7. Page map",
"## 8. Component ownership",
"## 9. Phase gates",
"## 11. Forbidden paths",
):
assert heading in text, f"ADR must contain section {heading!r}"
assert "browser ui" in lower and "domain loader" in lower
assert "control-plane db" in lower and "capability gate" in lower
def test_ac1_api_versioning_is_decided_including_legacy_routes():
text = _read(ADR)
assert "/api/v1/" in text, "ADR must decide the versioned API prefix"
assert "/api/v2/" in text, "ADR must state how breaking changes are handled"
lower = text.lower()
assert "compatibility alias" in lower, (
"ADR must say what happens to the existing unversioned MVP exports"
)
def test_ac1_page_map_covers_mvp_routes():
text = _read(ADR)
for route in ("`/`", "`/health`", "`/projects`", "`/prompts`", "`/runtime`",
"`/audit`", "`/actions`"):
assert route in text, f"page map must account for MVP route {route}"
def test_ac2_every_epic_child_maps_to_a_component():
text = _read(ADR)
ownership = text.split("## 8. Component ownership", 1)[-1].split("## 9.", 1)[0]
missing = [child for child in EPIC_CHILDREN if child not in ownership]
assert not missing, (
f"epic #631 children without an architectural component: {missing}"
)
def test_ac2_every_child_row_declares_a_phase():
text = _read(ADR)
ownership = text.split("## 8. Component ownership", 1)[-1].split("## 9.", 1)[0]
for child in EPIC_CHILDREN:
row = next(
(line for line in ownership.splitlines() if line.startswith(f"| {child} ")),
None,
)
assert row is not None, f"no ownership row for {child}"
assert row.rstrip().endswith(("| 1 |", "| 2 |", "| 3 |", "| 4 |")), (
f"ownership row for {child} must end with its phase: {row!r}"
)
def test_ac3_mvp_is_foundation_not_recreated():
text = _read(ADR)
assert "#425" in text and "#436" in text
lower = text.lower()
assert "do not recreate" in lower or "recreating mvp scope" in lower
assert "retained and evolved" in lower
def test_ac4_forbidden_paths_are_explicit():
text = _read(ADR)
forbidden = text.split("## 11. Forbidden paths", 1)[-1].split("## 12.", 1)[0]
lower = forbidden.lower()
assert "raw provider incidents" in lower and "#612" in forbidden
assert "browser-held tokens" in lower
assert "process-kill recovery" in lower and "#630" in forbidden
assert "ungated browser mutations" in lower
def test_ac5_approval_checklist_is_self_contained():
text = _read(ADR)
assert "## 12. Approval checklist" in text
checklist = text.split("## 12. Approval checklist", 1)[-1].split("## 13.", 1)[0]
for marker in ("1.", "2.", "3.", "4.", "5.", "6."):
assert marker in checklist, f"approval checklist missing item {marker}"
def test_adr_states_the_two_boundary_invariants():
text = _read(ADR)
lower = text.lower()
assert "no secrets to the browser" in lower
assert "no ungated mutations" in lower
def test_open_questions_are_recorded_not_implied():
text = _read(ADR)
assert "## 13. Open questions and follow-ups" in text
section = text.split("## 13. Open questions and follow-ups", 1)[-1]
assert "#633" in section, "deferred authorization work must name its issue"
def test_local_dev_doc_cross_links_the_adr():
text = _read(LOCAL_DEV)
assert ADR_BASENAME in text, (
"docs/webui-local-dev.md must cross-link the console architecture ADR "
"(issue #632 scope)"
)
def test_docs_do_not_embed_secrets():
for path in (ADR, LOCAL_DEV):
text = _read(path)
for marker in ("ghp_", "BEGIN PRIVATE KEY", "Authorization: Bearer"):
assert marker not in text, f"{path.name} contains {marker!r}"
+703
View File
@@ -0,0 +1,703 @@
"""Console authorization, redaction, and audit model tests (#633).
Covers each acceptance criterion and each required test named in the issue:
* AC1 — RBAC matrix and privileged-action list.
* AC2 — redaction rules, unit-tested against sample payloads.
* AC3 — audit event schema with required fields and retention defaults.
* AC4 — Phase 2 integration points.
* AC5 — local-dev mode with explicit insecurity warnings.
Required tests: redaction units (token, keychain, password patterns),
default-deny for unauthenticated write stubs, and audit record creation for a
simulated privileged preview.
"""
from __future__ import annotations
import datetime
import json
import os
import pathlib
import sys
import tempfile
import unittest
from starlette.testclient import TestClient
sys.path.insert(0, str(pathlib.Path(__file__).resolve().parents[1]))
from task_capability_map import TASK_CAPABILITY_MAP # noqa: E402
from webui import console_audit, console_authz # noqa: E402
from webui.app import create_app # noqa: E402
from webui.console_redaction import ( # noqa: E402
REDACTED,
redact_payload,
redact_text,
redaction_policy,
scan_for_secrets,
)
DOCS = pathlib.Path(__file__).resolve().parents[1] / "docs"
AUTHZ_DOC = DOCS / "webui-authz-audit.md"
def _principal(role: str) -> console_authz.Principal:
return console_authz.Principal(
subject=f"{role}@example.com",
role=role,
identity_source=console_authz.IDENTITY_ACCESS_PROXY,
authenticated=True,
)
class TestRoleMatrix(unittest.TestCase):
"""AC1 — the written RBAC matrix and privileged-action list."""
def test_roles_are_ordered_least_to_most_authority(self):
self.assertEqual(
console_authz.ROLE_ORDER,
("viewer", "operator", "controller", "admin"),
)
def test_every_role_has_a_description(self):
for role in console_authz.ROLE_ORDER:
with self.subTest(role=role):
self.assertTrue(console_authz.ROLE_DESCRIPTIONS[role].strip())
def test_higher_roles_inherit_lower_role_actions(self):
matrix = {
entry["role"]: set(entry["permitted_actions"])
for entry in console_authz.rbac_matrix()["roles"]
}
for lower, higher in zip(
console_authz.ROLE_ORDER, console_authz.ROLE_ORDER[1:]
):
with self.subTest(lower=lower, higher=higher):
self.assertTrue(matrix[lower].issubset(matrix[higher]))
def test_viewer_holds_no_write_action(self):
matrix = {
entry["role"]: set(entry["permitted_actions"])
for entry in console_authz.rbac_matrix()["roles"]
}
self.assertEqual(matrix["viewer"], set())
def test_privileged_action_list_is_non_empty_and_classified(self):
privileged = console_authz.privileged_actions()
self.assertTrue(privileged)
ids = {action.action_id for action in privileged}
# Merge and branch deletion are the canonical privileged pair.
self.assertIn("merge_pr", ids)
self.assertIn("delete_branch", ids)
def test_merge_and_delete_require_dual_control_and_break_glass(self):
for action_id in ("merge_pr", "delete_branch"):
with self.subTest(action=action_id):
action = console_authz.get_action(action_id)
self.assertTrue(action.dual_control)
self.assertTrue(action.break_glass)
self.assertTrue(action.requires_confirmation)
def test_every_write_action_requires_confirmation(self):
for action in console_authz.ACTIONS.values():
with self.subTest(action=action.action_id):
self.assertTrue(action.requires_confirmation)
def test_delete_branch_is_admin_only(self):
self.assertEqual(
console_authz.get_action("delete_branch").minimum_role,
console_authz.ADMIN,
)
def test_actions_map_to_real_mcp_capability_vocabulary(self):
"""The console must not invent an authority the MCP layer lacks."""
for action in console_authz.ACTIONS.values():
with self.subTest(action=action.action_id):
self.assertIn(action.task_key, TASK_CAPABILITY_MAP)
self.assertEqual(
action.mcp_permission,
TASK_CAPABILITY_MAP[action.task_key]["permission"],
)
self.assertEqual(
action.mcp_role,
TASK_CAPABILITY_MAP[action.task_key]["role"],
)
def test_matrix_declares_deny_by_default_and_execution_disabled(self):
matrix = console_authz.rbac_matrix()
self.assertEqual(matrix["default_decision"], "deny")
self.assertFalse(matrix["execution_enabled"])
class TestAuthorizeDefaultDeny(unittest.TestCase):
"""Fail-closed behaviour of the authorization decision."""
def test_anonymous_is_denied_every_action(self):
for action_id in console_authz.ACTIONS:
with self.subTest(action=action_id):
decision = console_authz.authorize(action_id)
self.assertFalse(decision.allowed)
self.assertEqual(
decision.reason_code, console_authz.DENY_UNAUTHENTICATED
)
def test_unknown_action_is_denied(self):
decision = console_authz.authorize(
"not_a_real_action", _principal("admin")
)
self.assertFalse(decision.allowed)
self.assertEqual(decision.reason_code, console_authz.DENY_UNKNOWN_ACTION)
def test_unknown_role_is_denied(self):
rogue = console_authz.Principal(
subject="[email protected]",
role="superuser",
identity_source=console_authz.IDENTITY_ACCESS_PROXY,
authenticated=True,
)
decision = console_authz.authorize("comment_issue", rogue)
self.assertFalse(decision.allowed)
self.assertEqual(decision.reason_code, console_authz.DENY_UNKNOWN_ROLE)
def test_insufficient_role_is_denied(self):
decision = console_authz.authorize("merge_pr", _principal("operator"))
self.assertFalse(decision.allowed)
self.assertEqual(
decision.reason_code, console_authz.DENY_INSUFFICIENT_ROLE
)
def test_sufficient_role_allows_preview_only(self):
decision = console_authz.authorize("merge_pr", _principal("controller"))
self.assertTrue(decision.allowed)
self.assertFalse(decision.execution_enabled)
def test_execution_is_refused_while_phase_is_not_active(self):
decision = console_authz.authorize(
"merge_pr", _principal("controller"), for_execution=True
)
self.assertFalse(decision.allowed)
self.assertEqual(
decision.reason_code, console_authz.DENY_PHASE_NOT_ACTIVE
)
def test_allowed_decision_never_reports_execution_enabled(self):
for action_id in console_authz.ACTIONS:
with self.subTest(action=action_id):
decision = console_authz.authorize(
action_id, _principal("admin")
)
self.assertFalse(decision.execution_enabled)
class TestIdentityResolution(unittest.TestCase):
"""AC5 — identity sources, including the insecure local-dev mode."""
def test_no_auth_mode_yields_anonymous_viewer(self):
principal = console_authz.resolve_principal(env={})
self.assertFalse(principal.authenticated)
self.assertEqual(principal.role, console_authz.VIEWER)
self.assertEqual(principal.identity_source, console_authz.IDENTITY_NONE)
def test_local_dev_mode_warns_that_identity_is_unverified(self):
principal = console_authz.resolve_principal(
env={
console_authz.AUTH_MODE_ENV: "local-dev",
console_authz.DEV_SUBJECT_ENV: "[email protected]",
console_authz.DEV_ROLE_ENV: "admin",
}
)
self.assertTrue(principal.authenticated)
self.assertEqual(principal.role, "admin")
self.assertTrue(principal.warnings)
self.assertIn("asserted", " ".join(principal.warnings).lower())
def test_local_dev_without_subject_falls_back_to_anonymous(self):
principal = console_authz.resolve_principal(
env={console_authz.AUTH_MODE_ENV: "local-dev"}
)
self.assertFalse(principal.authenticated)
def test_local_dev_unknown_role_degrades_to_viewer(self):
principal = console_authz.resolve_principal(
env={
console_authz.AUTH_MODE_ENV: "local_dev",
console_authz.DEV_SUBJECT_ENV: "[email protected]",
console_authz.DEV_ROLE_ENV: "root",
}
)
self.assertEqual(principal.role, console_authz.VIEWER)
def test_access_proxy_without_header_fails_closed(self):
"""A proxy-mode request that did not traverse the proxy is anonymous."""
principal = console_authz.resolve_principal(
headers={},
env={console_authz.AUTH_MODE_ENV: "access_proxy"},
)
self.assertFalse(principal.authenticated)
def test_access_proxy_role_comes_from_server_config_not_client(self):
env = {
console_authz.AUTH_MODE_ENV: "access_proxy",
console_authz.ROLE_MAP_ENV: json.dumps(
{"[email protected]": "controller"}
),
}
principal = console_authz.resolve_principal(
headers={
console_authz.ACCESS_SUBJECT_HEADER: "[email protected]",
"x-role": "admin", # client-supplied role must be ignored
},
env=env,
)
self.assertEqual(principal.role, "controller")
def test_access_proxy_unmapped_subject_defaults_to_viewer(self):
principal = console_authz.resolve_principal(
headers={
console_authz.ACCESS_SUBJECT_HEADER: "[email protected]"
},
env={console_authz.AUTH_MODE_ENV: "access_proxy"},
)
self.assertEqual(principal.role, console_authz.VIEWER)
def test_malformed_role_map_does_not_raise_and_denies(self):
principal = console_authz.resolve_principal(
headers={console_authz.ACCESS_SUBJECT_HEADER: "[email protected]"},
env={
console_authz.AUTH_MODE_ENV: "access_proxy",
console_authz.ROLE_MAP_ENV: "{not json",
},
)
self.assertEqual(principal.role, console_authz.VIEWER)
def test_probe_auth_is_opt_in(self):
self.assertFalse(console_authz.probe_auth_required(env={}))
self.assertTrue(
console_authz.probe_auth_required(
env={console_authz.REQUIRE_PROBE_AUTH_ENV: "1"}
)
)
def test_probe_auth_is_declared_but_not_yet_enforced(self):
"""Phase 1 declares the probe-auth policy; no route enforces it yet.
The flag exists so the Phase 2 action framework has a declared policy
to honour instead of inventing a second one. Pinning the current
not-enforced status here means wiring it later is a deliberate change
that updates this test and the documentation together, rather than a
silent behaviour shift. The documentation must say so plainly, because
an operator who sets the variable believing it protects a probe is
worse off than one who knows it does not.
"""
import inspect
from webui import app as webui_app
source = inspect.getsource(webui_app)
self.assertNotIn(
"probe_auth_required",
source,
msg=(
"webui.app now consults probe_auth_required, so probe auth is "
"no longer merely declared. Update the 'Probe authentication' "
"section of docs/webui-authz-audit.md, which states it "
"enforces nothing, and replace this test with real "
"enforcement coverage."
),
)
self.assertIn(
"enforces nothing today",
AUTHZ_DOC.read_text(encoding="utf-8"),
)
class TestRedaction(unittest.TestCase):
"""AC2 — required redaction units: token, keychain, password patterns."""
def test_token_assignment_is_redacted(self):
out = redact_text("GITEA_TOKEN=abcd1234efgh5678ijkl")
self.assertIn(REDACTED, out)
self.assertNotIn("abcd1234efgh5678ijkl", out)
def test_password_assignment_is_redacted(self):
out = redact_text("password: hunter2supersecret")
self.assertIn(REDACTED, out)
self.assertNotIn("hunter2supersecret", out)
def test_keychain_reference_is_redacted(self):
out = redact_text("keychain:gitea-prgs-token")
self.assertIn(REDACTED, out)
self.assertNotIn("gitea-prgs-token", out)
def test_keychain_command_is_redacted(self):
out = redact_text("security find-generic-password -s gitea -w")
self.assertIn(REDACTED, out)
self.assertNotIn("find-generic-password -s gitea", out)
def test_bearer_credential_is_redacted(self):
out = redact_text("Authorization: Bearer abcdef1234567890abcdef")
self.assertNotIn("abcdef1234567890abcdef", out)
def test_jwt_is_redacted(self):
token = "eyJhbGciOiJIUzI1NiJ9.eyJzdWIiOiIxIn0.abcdefghijklmnop"
out = redact_text(f"session={token}")
self.assertNotIn(token, out)
def test_private_key_block_is_redacted(self):
pem = (
"-----BEGIN RSA PRIVATE KEY-----\n"
"MIIEowIBAAKCAQEAsecretmaterial\n"
"-----END RSA PRIVATE KEY-----"
)
out = redact_text(pem)
self.assertNotIn("MIIEowIBAAKCAQEAsecretmaterial", out)
def test_api_key_assignment_is_redacted(self):
out = redact_text('api_key = "sk-live-9f8e7d6c5b4a3210"')
self.assertNotIn("sk-live-9f8e7d6c5b4a3210", out)
def test_nested_payload_is_redacted_recursively(self):
payload = {
"token": "abc123456789",
"nested": {"note": "password=letmein12345"},
"list": ["keychain:some-entry"],
"safe": "plain text",
}
out = redact_payload(payload)
self.assertEqual(out["token"], REDACTED)
self.assertNotIn("letmein12345", json.dumps(out))
self.assertNotIn("some-entry", json.dumps(out))
self.assertEqual(out["safe"], "plain text")
def test_scan_reports_findings_before_and_none_after(self):
dirty = "password: hunter2supersecret"
self.assertTrue(scan_for_secrets(dirty))
self.assertEqual(scan_for_secrets(redact_text(dirty)), [])
def test_non_strings_pass_through_untouched(self):
self.assertEqual(redact_text(42), 42)
self.assertEqual(
redact_payload({"n": 1, "b": True}), {"n": 1, "b": True}
)
def test_policy_is_documented_and_declares_redact_before_persist(self):
policy = redaction_policy()
self.assertTrue(policy["redact_before_persist"])
self.assertIn("audit_records", policy["applies_to"])
self.assertTrue(policy["console_rules"])
def test_policy_statement_contains_no_secret_material(self):
self.assertEqual(scan_for_secrets(redaction_policy()), [])
class TestAuditSchema(unittest.TestCase):
"""AC3 — audit event schema, required fields, and retention defaults."""
def _event(self, action_id="merge_pr", **kwargs):
return console_audit.build_event(
action_id=action_id,
result=console_audit.RESULT_DENIED,
decision=console_authz.authorize(action_id, _principal("operator")),
target={"kind": "pr", "ref": "#123"},
request_id="req-test",
**kwargs,
)
def test_every_required_field_is_present(self):
event = self._event()
for field in console_audit.REQUIRED_FIELDS:
with self.subTest(field=field):
self.assertIn(field, event)
def test_actor_carries_who_and_how_they_were_identified(self):
event = self._event()
for field in console_audit.REQUIRED_ACTOR_FIELDS:
with self.subTest(field=field):
self.assertIn(field, event["actor"])
def test_correlation_ids_are_present(self):
event = self._event()
for field in console_audit.REQUIRED_CORRELATION_FIELDS:
with self.subTest(field=field):
self.assertIn(field, event["correlation"])
self.assertEqual(event["correlation"]["request_id"], "req-test")
self.assertEqual(event["correlation"]["mcp_task"], "merge_pr")
def test_timestamp_is_timezone_aware_utc_iso8601(self):
now = datetime.datetime(
2026, 7, 22, 10, 16, 42, tzinfo=datetime.timezone.utc
)
event = self._event(now=now)
self.assertEqual(event["timestamp"], "2026-07-22T10:16:42+00:00")
parsed = datetime.datetime.fromisoformat(event["timestamp"])
self.assertIsNotNone(parsed.tzinfo)
def test_retention_defaults_by_class(self):
self.assertEqual(
console_audit.RETENTION_DAYS[console_audit.RETENTION_STANDARD], 90
)
self.assertEqual(
console_audit.RETENTION_DAYS[console_audit.RETENTION_PRIVILEGED],
365,
)
self.assertEqual(
console_audit.RETENTION_DAYS[console_audit.RETENTION_BREAK_GLASS],
730,
)
def test_break_glass_action_retains_longest(self):
event = self._event("merge_pr")
self.assertEqual(
event["retention"]["class"], console_audit.RETENTION_BREAK_GLASS
)
def test_routine_write_uses_standard_retention(self):
event = self._event("comment_issue")
self.assertEqual(
event["retention"]["class"], console_audit.RETENTION_STANDARD
)
def test_unknown_action_retains_as_privileged_not_standard(self):
"""Conservative direction: keep an unclassifiable record longer."""
self.assertEqual(
console_audit.retention_class_for(None),
console_audit.RETENTION_PRIVILEGED,
)
def test_retention_expiry_matches_declared_days(self):
now = datetime.datetime(2026, 7, 22, tzinfo=datetime.timezone.utc)
event = self._event("comment_issue", now=now)
expires = datetime.datetime.fromisoformat(
event["retention"]["expires_at"]
)
self.assertEqual((expires - now).days, 90)
def test_invalid_result_degrades_to_failed(self):
event = console_audit.build_event(action_id="merge_pr", result="banana")
self.assertEqual(event["result"], console_audit.RESULT_FAILED)
def test_denied_result_is_representable(self):
"""An authorization denial has no MCP-side mutation record."""
self.assertIn(console_audit.RESULT_DENIED, console_audit.RESULTS)
def test_event_is_redacted_before_it_is_returned(self):
event = console_audit.build_event(
action_id="merge_pr",
result=console_audit.RESULT_DENIED,
detail="failed with token=abcdef1234567890",
metadata={"password": "hunter2supersecret"},
)
serialized = json.dumps(event)
self.assertNotIn("abcdef1234567890", serialized)
self.assertNotIn("hunter2supersecret", serialized)
self.assertTrue(event["redacted"])
def test_audit_policy_reports_schema_and_retention(self):
policy = console_audit.audit_policy()
self.assertTrue(policy["append_only"])
self.assertTrue(policy["redact_before_persist"])
self.assertEqual(
policy["retention_defaults_days"], console_audit.RETENTION_DAYS
)
class TestAuditSink(unittest.TestCase):
"""Append-only persistence behaviour."""
def test_write_is_a_noop_when_sink_is_unconfigured(self):
saved = os.environ.pop(console_audit.AUDIT_LOG_ENV, None)
try:
self.assertFalse(console_audit.audit_enabled())
self.assertFalse(console_audit.write_event({"schema_version": 1}))
finally:
if saved is not None:
os.environ[console_audit.AUDIT_LOG_ENV] = saved
def test_records_append_one_json_line_each(self):
with tempfile.TemporaryDirectory() as tmp:
sink = os.path.join(tmp, "console-audit.jsonl")
for _ in range(3):
event = console_audit.build_event(
action_id="merge_pr", result=console_audit.RESULT_DENIED
)
self.assertTrue(console_audit.write_event(event, path=sink))
with open(sink, encoding="utf-8") as handle:
lines = [json.loads(line) for line in handle if line.strip()]
self.assertEqual(len(lines), 3)
self.assertEqual(len({line["event_id"] for line in lines}), 3)
def test_a_record_that_still_carries_a_secret_is_not_persisted(self):
with tempfile.TemporaryDirectory() as tmp:
sink = os.path.join(tmp, "console-audit.jsonl")
leaky = {
"schema_version": 1,
"detail": "password: hunter2supersecret",
}
self.assertFalse(console_audit.write_event(leaky, path=sink))
self.assertFalse(os.path.exists(sink))
def test_write_never_raises_on_a_bad_path(self):
self.assertFalse(
console_audit.write_event(
{"schema_version": 1}, path="/nonexistent-dir/audit.jsonl"
)
)
def test_simulated_privileged_preview_creates_an_audit_record(self):
"""Required test: audit record creation for a privileged preview."""
with tempfile.TemporaryDirectory() as tmp:
sink = os.path.join(tmp, "console-audit.jsonl")
os.environ[console_audit.AUDIT_LOG_ENV] = sink
try:
decision = console_authz.authorize(
"merge_pr", _principal("controller")
)
outcome = console_audit.record_event(
action_id="merge_pr",
result=console_audit.RESULT_PREVIEWED,
decision=decision,
target={"kind": "pr", "ref": "#123"},
request_id="req-preview",
)
finally:
os.environ.pop(console_audit.AUDIT_LOG_ENV, None)
self.assertTrue(outcome["written"])
with open(sink, encoding="utf-8") as handle:
record = json.loads(handle.read().strip())
self.assertEqual(record["action"], "merge_pr")
self.assertEqual(record["result"], console_audit.RESULT_PREVIEWED)
self.assertEqual(record["action_class"], "privileged")
self.assertTrue(record["decision"]["allowed"])
self.assertFalse(record["decision"]["execution_enabled"])
self.assertEqual(record["actor"]["role"], "controller")
def test_decision_block_survives_redaction(self):
"""Regression: naming it 'authorization' collided with a secret hint.
``gitea_audit._SECRET_KEY_HINTS`` contains "authorization" (for the
HTTP header), so a block under that key was replaced wholesale by the
placeholder and the record lost its decision entirely.
"""
event = console_audit.build_event(
action_id="merge_pr",
result=console_audit.RESULT_DENIED,
decision=console_authz.authorize("merge_pr", _principal("admin")),
)
self.assertIsInstance(event["decision"], dict)
self.assertIn("allowed", event["decision"])
class TestConsoleRoutes(unittest.TestCase):
"""AC4 — the wired Phase 2 integration points, still fail-closed."""
def setUp(self):
self.client = TestClient(create_app(bind_host="127.0.0.1"))
def test_unauthenticated_write_stub_is_denied(self):
"""Required test: default-deny for unauthenticated write stubs."""
response = self.client.post(
"/api/actions/merge_pr/attempt", json={"pr_number": 99}
)
self.assertEqual(response.status_code, 403)
body = response.json()
self.assertFalse(body["success"])
authorization = body["authorization"]
self.assertFalse(authorization["allowed"])
self.assertEqual(
authorization["reason_code"], console_authz.DENY_UNAUTHENTICATED
)
self.assertFalse(authorization["execution_enabled"])
def test_preview_reports_an_authorization_decision(self):
response = self.client.get("/api/actions/merge_pr/preview?pr_number=7")
self.assertEqual(response.status_code, 200)
authorization = response.json()["authorization"]
self.assertFalse(authorization["allowed"])
self.assertTrue(authorization["dual_control"])
self.assertEqual(authorization["required_role"], "controller")
def test_unknown_action_preview_still_404s(self):
response = self.client.get("/api/actions/no_such_action/preview")
self.assertEqual(response.status_code, 404)
def test_security_model_endpoint_publishes_all_three_policies(self):
response = self.client.get("/api/console/security-model")
self.assertEqual(response.status_code, 200)
body = response.json()
self.assertIn("rbac", body)
self.assertIn("redaction", body)
self.assertIn("audit", body)
self.assertEqual(body["rbac"]["default_decision"], "deny")
def test_security_model_endpoint_leaks_no_secrets(self):
response = self.client.get("/api/console/security-model")
self.assertEqual(scan_for_secrets(response.json()), [])
def test_security_model_rejects_writes(self):
response = self.client.post("/api/console/security-model", json={})
self.assertEqual(response.status_code, 405)
def test_existing_read_routes_are_unaffected(self):
for path in ("/", "/health", "/actions", "/api/actions"):
with self.subTest(path=path):
self.assertEqual(self.client.get(path).status_code, 200)
class TestAuthzAuditDoc(unittest.TestCase):
"""The model must be written down, not only coded."""
@classmethod
def setUpClass(cls):
cls.text = (
AUTHZ_DOC.read_text(encoding="utf-8") if AUTHZ_DOC.exists() else ""
)
def test_doc_exists(self):
self.assertTrue(AUTHZ_DOC.exists(), f"missing {AUTHZ_DOC}")
def test_doc_covers_each_required_section(self):
for heading in (
"Identity sources",
"Role matrix",
"Privileged actions",
"Secret redaction",
"Audit event schema",
"Retention",
"Phase 2 integration",
"Local-dev mode",
):
with self.subTest(heading=heading):
self.assertIn(heading, self.text)
def test_doc_names_every_role(self):
for role in console_authz.ROLE_ORDER:
with self.subTest(role=role):
self.assertIn(role, self.text)
def test_doc_names_every_console_action(self):
for action_id in console_authz.ACTIONS:
with self.subTest(action=action_id):
self.assertIn(action_id, self.text)
def test_doc_states_retention_defaults(self):
for days in console_audit.RETENTION_DAYS.values():
with self.subTest(days=days):
self.assertIn(str(days), self.text)
def test_doc_warns_local_dev_is_insecure(self):
self.assertIn("INSECURE", self.text.upper())
def test_doc_states_default_deny(self):
self.assertIn("deny", self.text.lower())
def test_doc_contains_no_secret_material(self):
self.assertEqual(scan_for_secrets(self.text), [])
def test_deployment_doc_links_to_the_model(self):
deployment = (DOCS / "webui-deployment.md").read_text(encoding="utf-8")
self.assertIn("webui-authz-audit", deployment)
if __name__ == "__main__": # pragma: no cover
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 sys
import tempfile
@@ -11,57 +11,246 @@ from starlette.testclient import TestClient
from webui.app import create_app
from webui.project_registry import (
CURRENT_SCHEMA_VERSION,
REGISTRY_API_VERSION,
SUPPORTED_SCHEMA_VERSIONS,
RegistryError,
default_registry_path,
load_registry,
onboarding_summary,
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):
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)
project = registry.projects[0]
self.assertEqual(project.id, "gitea-tools")
self.assertEqual(project.repo_name, "Gitea-Tools")
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_name, "prgs")
self.assertEqual(project.status, "active")
self.assertEqual(project.profiles["author"], "prgs-author")
self.assertEqual(project.profiles["reviewer"], "prgs-reviewer")
self.assertEqual(project.profiles["reconciler"], "prgs-reconciler")
self.assertIn("skill", project.workflow_paths)
self.assertGreaterEqual(len(project.onboarding_checklist), 4)
def test_registry_rejects_credential_keys(self):
payload = {
def test_default_registry_onboarding_summary_is_complete(self):
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,
"projects": [
{
"id": "bad",
"repo_name": "Bad",
"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",
}
_valid_project(
onboarding_checklist=[
{"id": "step", "title": "Step", "description": "Do it"}
]
)
],
}
})
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:
json.dump(payload, handle)
handle.write("{not json")
path = Path(handle.name)
try:
with self.assertRaises(ValueError):
load_registry(path)
finally:
path.unlink(missing_ok=True)
self._temp_paths.append(path)
with self.assertRaises(RegistryError) as ctx:
load_registry(path)
self.assertIn("not valid JSON", ctx.exception.message)
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):
path = default_registry_path()
@@ -81,31 +270,147 @@ class TestProjectRegistryRoutes(unittest.TestCase):
self.assertIn("prgs-author", response.text)
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):
response = self.client.get("/projects/gitea-tools")
self.assertEqual(response.status_code, 200)
self.assertIn("Onboarding checklist", response.text)
self.assertIn("Configure execution profiles", 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):
response = self.client.get("/projects/unknown-repo")
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")
self.assertEqual(response.status_code, 200)
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(data["projects"][0]["id"], "gitea-tools")
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):
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)
# 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__":
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()