Compare commits

..
Author SHA1 Message Date
jcwalker3 f0c9ffb25e Merge branch 'master' into fix/issue-843-cross-role-allocation-handoff 2026-07-23 12:31:10 -05:00
sysadmin 1c455b6ec0 Merge pull request 'feat(webui): read-only system-health API (Closes #634)' (#813) from feat/issue-634-readonly-system-health-api into master 2026-07-23 04:14:33 -05:00
sysadmin 5eb89f8830 fix: bind cross-role handoff consume role to authenticated profile (Closes #843)
Review #515 F1: gitea_adopt_workflow_lease trusted a caller-supplied role
((role or active_role)), so any namespace holding gitea.read could consume
an author-only cross-role handoff by passing role="author".

- Derive the adopter role authoritatively from the active profile; reject
  any supplied role that does not exactly match (no silent accept).
- Pass the profile-derived role and authoritative profile/namespace context
  to lease_lifecycle.adopt_lease; validate handoff provenance
  required_profile/required_namespace against it (fail closed).
- Fail closed when the profile role cannot be derived (no author default).
- Add MCP-boundary regression tests: reviewer/merger profiles cannot
  consume an author handoff via role="author"; the legitimate author
  profile still consumes; foreign required_profile rejected.
2026-07-23 03:41:28 -04:00
sysadminandClaude Opus 4.8 a6c15afec1 fix: make cross-role allocations consumable by independent workers (Closes #843)
Controller-created role=author allocations were owned by the allocating
controller session with no authorized consume path for independent author
workers. When the controller exited, the lease became stale_dead_process
and required abandon/reassign instead of a usable handoff.

- Mark cross-role apply with durable handoff provenance (pending)
- Allow gitea_adopt_workflow_lease to consume pending handoffs by the
  required role without sharing controller session identity or requiring
  the controller process to remain alive
- Atomically transfer assignment+lease ownership and set
  adopted_by_session_id with read-after-write evidence
- Reject wrong-role, second, and terminal adoptions
- Surface consume_allocation identifiers in process_work_queue results
- Preserve same-role allocation and genuine abandon recovery behavior

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-23 02:38:19 -04:00
jcwalker3 6868b345ee Merge branch 'master' into feat/issue-634-readonly-system-health-api 2026-07-23 01:12:52 -05:00
sysadmin 4f3a464a90 Merge pull request 'fix: authoritative cross-role generic queue allocation (Closes #840)' (#841) from fix/issue-840-cross-role-queue-allocation into master 2026-07-23 00:37:55 -05:00
jcwalker3 da6a864463 Merge branch 'master' into feat/issue-634-readonly-system-health-api 2026-07-23 00:06:06 -05:00
sysadmin 9468dd624d merge master into fix/issue-840-cross-role-queue-allocation 2026-07-23 00:06:09 -04:00
sysadminandClaude Opus 4.8 648d9464ba fix: authoritative cross-role generic queue allocation (Closes #840)
Add controller-owned cross_role allocation mode that inspects the full
queue and returns one selection with required role/profile/action and
lease evidence. Document process_work_queue routing, normalize
controller role metadata, and keep the dashboard explanatory only.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-23 00:05:35 -04:00
jcwalker3andsysadmin caaae9b6ee feat(arch01): atomic platform install + authority kernel (Closes #822) (#839)
Co-authored-by: jcwalker3 <[email protected]>
2026-07-22 23:00:24 -05:00
sysadmin 689c60fc7c Merge pull request 'feat(webui): console authorization, RBAC, redaction, and audit model (Closes #633)' (#811) from feat/issue-633-console-authz-audit-model into master 2026-07-22 22:16:38 -05:00
jcwalker3andClaude Opus 4.8 5494696227 feat(webui): read-only system-health API (Closes #634)
Adds `GET /api/v1/system/health`, a structured read-only health surface for
automated readiness checks, and keeps `/health` as the cheap liveness probe.

webui/system_health.py composes a DTO from fail-soft dependency probes: the
control-plane database, the local checkout, and — opt-in via `?deep=1` — live
Gitea reachability, each carrying status, reason, and probe latency. Required
probes drive readiness; the optional Gitea probe can only degrade overall
status, because local inventory stays serveable when the remote is
unreachable. A probe that did not run leaves readiness incomplete rather than
silently passing.

Read-only throughout: the control-plane database is opened through a `mode=ro`
URI because `ControlPlaneDB.__init__` creates directories and runs migrations,
which a health check must never do. No restart or reload control is exposed;
those are Phase 2 and #630 forbids process-kill recovery.

No unproven claims: `stale_runtime.mutation_safe` is true only when the
runtime, checkout, and remote commits are all known and equal, and MCP
namespaces always report `unproven` because a web process cannot exercise the
IDE-managed client path (#543). Probe details are redacted at the browser
boundary — URLs lose userinfo and query strings, credential-shaped text is
masked.

`/health` is expanded additively: every MVP key is retained, plus `started_at`,
`uptime_seconds`, and a pointer to the versioned API. The versioned route
returns 503 when not ready so automation can branch on the status code alone.

Verified at master 9eb0f29: focused file 40 passed / 11 subtests; `-k "webui or
health"` 230 passed / 159 subtests; full suite 4358 passed with the 11
pre-existing master-drift failures unchanged from the clean-master baseline.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-22 15:59:30 -05:00
16 changed files with 4940 additions and 82 deletions
+364 -41
View File
@@ -78,6 +78,33 @@ VALID_ROLES = frozenset(
{ROLE_AUTHOR, ROLE_REVIEWER, ROLE_MERGER, ROLE_RECONCILER, ROLE_CONTROLLER} {ROLE_AUTHOR, ROLE_REVIEWER, ROLE_MERGER, ROLE_RECONCILER, ROLE_CONTROLLER}
) )
# Allocation modes (#840).
# role_scoped: only candidates whose expected role matches the caller role.
# cross_role: controller-owned generic queue selection — inspect full queue,
# rank/eligibility canonically, return one selection naming the required
# downstream role/profile. Controller routes; it does not perform mutations.
ALLOCATION_MODE_ROLE_SCOPED = "role_scoped"
ALLOCATION_MODE_CROSS_ROLE = "cross_role"
VALID_ALLOCATION_MODES = frozenset(
{ALLOCATION_MODE_ROLE_SCOPED, ALLOCATION_MODE_CROSS_ROLE}
)
# Default execution-profile / MCP-namespace names for each role.
DEFAULT_ROLE_PROFILES: dict[str, str] = {
ROLE_AUTHOR: "prgs-author",
ROLE_REVIEWER: "prgs-reviewer",
ROLE_MERGER: "prgs-merger",
ROLE_RECONCILER: "prgs-reconciler",
ROLE_CONTROLLER: "prgs-controller",
}
DEFAULT_ROLE_NAMESPACES: dict[str, str] = {
ROLE_AUTHOR: "gitea-author",
ROLE_REVIEWER: "gitea-reviewer",
ROLE_MERGER: "gitea-merger",
ROLE_RECONCILER: "gitea-reconciler",
ROLE_CONTROLLER: "gitea-controller",
}
# Default action matrices by role (mutation gate will re-check). # Default action matrices by role (mutation gate will re-check).
ROLE_ACTIONS: dict[str, tuple[tuple[str, ...], tuple[str, ...]]] = { ROLE_ACTIONS: dict[str, tuple[tuple[str, ...], tuple[str, ...]]] = {
ROLE_AUTHOR: ( ROLE_AUTHOR: (
@@ -259,6 +286,126 @@ def normalize_role(role: str | None, *, profile_name: str | None = None) -> str:
) )
def resolve_allocation_mode(
role: str,
allocation_mode: str | None = None,
) -> str:
"""Resolve allocation mode; controller defaults to cross_role (#840)."""
raw = (allocation_mode or "").strip().lower()
if raw:
if raw not in VALID_ALLOCATION_MODES:
raise ControlPlaneError(
f"unknown allocation_mode {allocation_mode!r}; expected one of "
f"{sorted(VALID_ALLOCATION_MODES)}"
)
return raw
if role == ROLE_CONTROLLER:
return ALLOCATION_MODE_CROSS_ROLE
return ALLOCATION_MODE_ROLE_SCOPED
def required_profile_for_role(
role: str,
*,
profile_name: str | None = None,
) -> str:
"""Map a required role to the canonical execution profile name."""
role_norm = (role or "").strip().lower()
# Preserve remote/env prefix from the active profile when present
# (e.g. dadeschools-author → dadeschools-reviewer).
active = (profile_name or "").strip()
if active:
lower = active.lower()
for token in ("author", "reviewer", "merger", "reconciler", "controller"):
if lower.endswith(f"-{token}") or lower == token:
prefix = active[: -len(token)].rstrip("-")
if prefix:
return f"{prefix}-{role_norm}"
return role_norm
return DEFAULT_ROLE_PROFILES.get(role_norm, f"prgs-{role_norm}")
def required_namespace_for_role(
role: str,
*,
profile_name: str | None = None,
) -> str:
"""Map a required role to the canonical MCP namespace name."""
role_norm = (role or "").strip().lower()
profile = required_profile_for_role(role_norm, profile_name=profile_name)
# Namespace is typically gitea-<role>; keep stable mapping when profile is
# non-prgs (still gitea-<role> for isolation).
return DEFAULT_ROLE_NAMESPACES.get(role_norm, f"gitea-{role_norm}")
def selected_action_for_candidate(c: WorkCandidate, required_role: str) -> str:
"""Canonical next action for the selected work under *required_role*."""
role = (required_role or "").strip().lower()
if role == ROLE_AUTHOR:
if c.kind == "pr" and c.request_changes_current_head:
return "address_pr_change_requests"
if c.kind == "pr":
return "update_pr"
return "implement"
if role == ROLE_REVIEWER:
if c.approval_stale:
return "re_review"
return "review"
if role == ROLE_MERGER:
return "merge"
if role == ROLE_RECONCILER:
if c.approval_contaminated:
return "reconcile_contaminated_approval"
return "reconcile"
if role == ROLE_CONTROLLER:
return "diagnose"
return "process"
def build_selection_dict(
selected: WorkCandidate,
*,
active_role: str,
required_role: str,
profile_name: str | None = None,
allocation_mode: str,
) -> dict[str, Any]:
"""Authoritative single selection payload for allocator results (#840)."""
action = selected_action_for_candidate(selected, required_role)
req_profile = required_profile_for_role(
required_role, profile_name=profile_name
)
req_ns = required_namespace_for_role(
required_role, profile_name=profile_name
)
return {
"kind": selected.kind,
"number": selected.number,
"title": selected.title,
"labels": list(selected.labels),
"head_sha": selected.head_sha,
"priority": selected.priority,
"expected_role_next": required_role,
"required_role": required_role,
"selected_action": action,
"action": action,
"required_profile": req_profile,
"required_namespace": req_ns,
"pinned": {
"kind": selected.kind,
"number": selected.number,
"head_sha": selected.head_sha,
"issue_number": selected.number if selected.kind == "issue" else None,
"pr_number": selected.number if selected.kind == "pr" else None,
},
"reason_selected": (
f"highest-priority eligible candidate under allocation_mode="
f"'{allocation_mode}' (active_role={active_role}, "
f"required_role={required_role}, action={action})"
),
}
def expected_role_for_candidate(c: WorkCandidate) -> str: def expected_role_for_candidate(c: WorkCandidate) -> str:
"""ADR §5.3 routing: which role should take this work next.""" """ADR §5.3 routing: which role should take this work next."""
if c.kind == "pr": if c.kind == "pr":
@@ -289,6 +436,7 @@ def classify_skip(
role: str, role: str,
terminal_pr: int | None, terminal_pr: int | None,
claim_ownership: str | None = None, claim_ownership: str | None = None,
allocation_mode: str | None = None,
) -> str | None: ) -> str | None:
"""Return skip reason, or None if candidate is selectable for *role*. """Return skip reason, or None if candidate is selectable for *role*.
@@ -297,7 +445,12 @@ def classify_skip(
and unknown claims are excluded so one session's in-progress task can never and unknown claims are excluded so one session's in-progress task can never
blockade the queue for a different controller; ``own`` stays selectable so blockade the queue for a different controller; ``own`` stays selectable so
a controller can resume its own work. a controller can resume its own work.
*allocation_mode* (#840): ``cross_role`` (controller default) ranks the full
queue and selects the highest-priority eligible item for any downstream
role. ``role_scoped`` retains prior role-match filtering.
""" """
mode = resolve_allocation_mode(role, allocation_mode)
if c.state in ("merged", "closed"): if c.state in ("merged", "closed"):
return f"{c.kind}#{c.number} is {c.state}; never assign" return f"{c.kind}#{c.number} is {c.state}; never assign"
if c.blocked or "status:blocked" in c.labels: if c.blocked or "status:blocked" in c.labels:
@@ -322,34 +475,58 @@ def classify_skip(
if c.kind == "pr" and not (c.head_sha or "").strip(): if c.kind == "pr" and not (c.head_sha or "").strip():
return f"pr#{c.number} missing head_sha pin" return f"pr#{c.number} missing head_sha pin"
expected = expected_role_for_candidate(c)
# Terminal path first: when an active terminal PR exists, only that PR # Terminal path first: when an active terminal PR exists, only that PR
# (or controller diagnosis) is assignable for review-path roles. # is assignable for review-path roles (or for work whose expected role is
# review/merge under cross_role selection).
if terminal_pr is not None and c.kind == "pr" and c.number != terminal_pr: if terminal_pr is not None and c.kind == "pr" and c.number != terminal_pr:
if role in (ROLE_REVIEWER, ROLE_MERGER): terminal_roles = (ROLE_REVIEWER, ROLE_MERGER)
if mode == ALLOCATION_MODE_CROSS_ROLE:
if expected in terminal_roles:
return (
f"pr#{c.number} skipped: active terminal-review lock on "
f"PR #{terminal_pr} must be resolved first"
)
elif role in terminal_roles:
return ( return (
f"pr#{c.number} skipped: active terminal-review lock on " f"pr#{c.number} skipped: active terminal-review lock on "
f"PR #{terminal_pr} must be resolved first" f"PR #{terminal_pr} must be resolved first"
) )
expected = expected_role_for_candidate(c) if mode == ALLOCATION_MODE_CROSS_ROLE:
if role == ROLE_CONTROLLER: # Cross-role controller selection: eligibility only — no active-role
# Controller may inspect anything but only assigns diagnosis targets # match filter. The selection payload names required_role.
# when contaminated / blocked. pass
elif role == ROLE_CONTROLLER:
# Legacy diagnosis-only controller path (role_scoped): only reconciler-
# needed targets. Prefer cross_role for generic queue allocation.
if expected == ROLE_RECONCILER or c.blocked: if expected == ROLE_RECONCILER or c.blocked:
return None return None
return f"{c.kind}#{c.number} does not require controller (expected {expected})" return (
f"{c.kind}#{c.number} does not require controller "
if role != expected: f"(expected {expected})"
)
elif role != expected:
return ( return (
f"{c.kind}#{c.number} expects role '{expected}', active role is '{role}'" f"{c.kind}#{c.number} expects role '{expected}', active role is '{role}'"
) )
# Ready-gate for issues: prefer status:ready when labels present. # Ready-gate for issues: prefer status:ready when labels present.
# Applies for author-bound work in both modes (cross_role only gates
# author-expected issues so reconciler/reviewer PRs stay selectable).
if c.kind == "issue" and c.labels: if c.kind == "issue" and c.labels:
if "status:ready" not in c.labels and "status:in-progress" not in c.labels: gate_role = expected if mode == ALLOCATION_MODE_CROSS_ROLE else role
# Allow unlabeled open issues; only skip explicit non-ready states. if gate_role in (ROLE_AUTHOR, ROLE_CONTROLLER):
if (
"status:ready" not in c.labels
and "status:in-progress" not in c.labels
):
if any(l.startswith("status:") for l in c.labels): if any(l.startswith("status:") for l in c.labels):
return f"issue#{c.number} not status:ready ({','.join(c.labels)})" return (
f"issue#{c.number} not status:ready "
f"({','.join(c.labels)})"
)
return None return None
@@ -501,12 +678,19 @@ def allocate_next_work(
claims: Mapping[tuple[str, int], dict[str, Any]] | None = None, claims: Mapping[tuple[str, int], dict[str, Any]] | None = None,
exclude_issue_numbers: Sequence[int] | None = None, exclude_issue_numbers: Sequence[int] | None = None,
expected_candidate_set_fingerprint: str | None = None, expected_candidate_set_fingerprint: str | None = None,
allocation_mode: str | None = None,
) -> dict[str, Any]: ) -> dict[str, Any]:
"""Select and optionally reserve the next work unit via control-plane DB. """Select and optionally reserve the next work unit via control-plane DB.
*apply=False* (default): dry-run selection only — no lease/assignment. *apply=False* (default): dry-run selection only — no lease/assignment.
*apply=True*: atomic ``assign_and_lease`` for the selected candidate. *apply=True*: atomic ``assign_and_lease`` for the selected candidate.
*allocation_mode* (#840): ``cross_role`` (default for controller) inspects
the complete queue and returns one authoritative selection naming the
required downstream role/profile/action. ``role_scoped`` keeps prior
per-role filtering. Controller routes only — never grants author/reviewer/
merger/reconciler mutation rights to the controller session.
*exclude_issue_numbers* (#776): numbers removed before ranking. Omitted / *exclude_issue_numbers* (#776): numbers removed before ranking. Omitted /
empty preserves prior behavior. empty preserves prior behavior.
@@ -541,6 +725,19 @@ def allocate_next_work(
"substrate": "control_plane_db", "substrate": "control_plane_db",
} }
try:
mode = resolve_allocation_mode(role_norm, allocation_mode)
except ControlPlaneError as exc:
return {
"success": False,
"outcome": OUTCOME_ROLE_INELIGIBLE,
"reasons": [str(exc)],
"skipped": [],
"assignment": None,
"substrate": "control_plane_db",
"allocation_mode": (allocation_mode or "").strip() or None,
}
session_id = (session_id or "").strip() or f"alloc-{uuid.uuid4().hex[:12]}" session_id = (session_id or "").strip() or f"alloc-{uuid.uuid4().hex[:12]}"
try: try:
db.upsert_session( db.upsert_session(
@@ -748,6 +945,7 @@ def allocate_next_work(
role=role_norm, role=role_norm,
terminal_pr=terminal_pr, terminal_pr=terminal_pr,
claim_ownership=ownership, claim_ownership=ownership,
allocation_mode=mode,
) )
if reason: if reason:
is_claim_skip = SKIP_CLAIMED_BY_OTHER_SESSION in reason is_claim_skip = SKIP_CLAIMED_BY_OTHER_SESSION in reason
@@ -828,6 +1026,10 @@ def allocate_next_work(
"outcome": outcome, "outcome": outcome,
"apply": bool(apply), "apply": bool(apply),
"role": role_norm, "role": role_norm,
"allocation_mode": mode,
"routing_role": role_norm,
"required_role": None,
"selected_action": None,
"profile_name": profile_name, "profile_name": profile_name,
"username": username, "username": username,
"session_id": session_id, "session_id": session_id,
@@ -840,6 +1042,12 @@ def allocate_next_work(
"skipped": [s.as_dict() for s in skipped], "skipped": [s.as_dict() for s in skipped],
"terminal_pr": terminal_pr, "terminal_pr": terminal_pr,
"assignment": None, "assignment": None,
"allocation_evidence": {
"mode": "empty",
"allocation_mode": mode,
"lease_created": False,
"selection_policy": SELECTION_POLICY,
},
"substrate": "control_plane_db", "substrate": "control_plane_db",
"file_lock_only": False, "file_lock_only": False,
"comment_lease_only": False, "comment_lease_only": False,
@@ -852,25 +1060,37 @@ def allocate_next_work(
"owner_session_id": owner_session_id, "owner_session_id": owner_session_id,
"downstream_note": ( "downstream_note": (
"#612 incident bridge remains downstream of #600; " "#612 incident bridge remains downstream of #600; "
"allocator never assigns raw monitoring incidents" "allocator never assigns raw monitoring incidents; "
"controller routes only under cross_role (#840)"
), ),
} }
expected_role = expected_role_for_candidate(selected) expected_role = expected_role_for_candidate(selected)
allowed, forbidden = role_actions(role_norm) # Cross-role: lease/action matrix follows the required downstream role so
selection = { # evidence names the worker that must act. Controller session still owns
"kind": selected.kind, # the routing decision; mutation isolation is enforced by role gates on
"number": selected.number, # mutation tools (controller profile lacks author/review/merge ops).
"title": selected.title, lease_role = (
"labels": list(selected.labels), expected_role if mode == ALLOCATION_MODE_CROSS_ROLE else role_norm
"head_sha": selected.head_sha, )
"priority": selected.priority, allowed, forbidden = role_actions(lease_role)
"expected_role_next": expected_role, # Controller must never receive mutation-class rights via cross-role apply.
"reason_selected": ( if role_norm == ROLE_CONTROLLER:
f"highest-priority candidate for role '{role_norm}' " ctrl_allowed, ctrl_forbidden = role_actions(ROLE_CONTROLLER)
f"(expected_role={expected_role})" # Keep controller session capability evidence separate from lease_role.
), controller_allowed_actions = ctrl_allowed
} controller_forbidden_actions = ctrl_forbidden
else:
controller_allowed_actions = allowed
controller_forbidden_actions = forbidden
selection = build_selection_dict(
selected,
active_role=role_norm,
required_role=expected_role,
profile_name=profile_name,
allocation_mode=mode,
)
if not apply: if not apply:
return { return {
@@ -878,6 +1098,12 @@ def allocate_next_work(
"outcome": OUTCOME_PREVIEW, "outcome": OUTCOME_PREVIEW,
"apply": False, "apply": False,
"role": role_norm, "role": role_norm,
"allocation_mode": mode,
"routing_role": role_norm,
"required_role": expected_role,
"selected_action": selection["selected_action"],
"required_profile": selection["required_profile"],
"required_namespace": selection["required_namespace"],
"profile_name": profile_name, "profile_name": profile_name,
"username": username, "username": username,
"session_id": session_id, "session_id": session_id,
@@ -889,10 +1115,22 @@ def allocate_next_work(
"reasons": [ "reasons": [
"dry-run only (apply=false); no assignment/lease created — " "dry-run only (apply=false); no assignment/lease created — "
"call again with apply=true to reserve via control-plane DB" "call again with apply=true to reserve via control-plane DB"
+ (
"; after apply, the required-role worker consumes via "
"gitea_adopt_workflow_lease (#843)"
if mode == ALLOCATION_MODE_CROSS_ROLE and expected_role != role_norm
else ""
)
], ],
"skipped": [s.as_dict() for s in skipped], "skipped": [s.as_dict() for s in skipped],
"terminal_pr": terminal_pr, "terminal_pr": terminal_pr,
"assignment": None, "assignment": None,
"allocation_evidence": {
"mode": "preview",
"allocation_mode": mode,
"lease_created": False,
"selection_policy": SELECTION_POLICY,
},
"substrate": "control_plane_db", "substrate": "control_plane_db",
"file_lock_only": False, "file_lock_only": False,
"comment_lease_only": False, "comment_lease_only": False,
@@ -902,18 +1140,24 @@ def allocate_next_work(
"controller_excluded": list(controller_excluded), "controller_excluded": list(controller_excluded),
"exclude_issue_numbers": list(exclude_nums), "exclude_issue_numbers": list(exclude_nums),
"candidate_set_fingerprint": cas_fp, "candidate_set_fingerprint": cas_fp,
"controller_allowed_actions": list(controller_allowed_actions),
"controller_forbidden_actions": list(controller_forbidden_actions),
"downstream_note": ( "downstream_note": (
"#612 incident bridge remains downstream of #600; " "#612 incident bridge remains downstream of #600; "
"allocator never assigns raw monitoring incidents" "allocator never assigns raw monitoring incidents; "
"controller routes only under cross_role (#840)"
), ),
} }
# Atomic reserve via #613 substrate. # Atomic reserve via #613 substrate.
ttl = lease_ttl_seconds if lease_ttl_seconds is not None else None ttl = lease_ttl_seconds if lease_ttl_seconds is not None else None
try: try:
cross_role_handoff = (
mode == ALLOCATION_MODE_CROSS_ROLE and lease_role != role_norm
)
kwargs: dict[str, Any] = { kwargs: dict[str, Any] = {
"session_id": session_id, "session_id": session_id,
"role": role_norm, "role": lease_role,
"remote": remote, "remote": remote,
"org": org, "org": org,
"repo": repo, "repo": repo,
@@ -922,7 +1166,8 @@ def allocate_next_work(
"expected_head_sha": selected.head_sha, "expected_head_sha": selected.head_sha,
"allowed_actions": allowed, "allowed_actions": allowed,
"forbidden_actions": forbidden, "forbidden_actions": forbidden,
"phase": "allocated", # #843: mark cross-role allocations as awaiting independent consume
"phase": "awaiting_handoff" if cross_role_handoff else "allocated",
} }
if ttl is not None: if ttl is not None:
kwargs["lease_ttl_seconds"] = int(ttl) kwargs["lease_ttl_seconds"] = int(ttl)
@@ -992,11 +1237,72 @@ def allocate_next_work(
} }
# assigned # assigned
return { lease_proof = {
"assignment_id": result.assignment_id,
"lease_id": result.lease_id,
"expires_at": result.expires_at,
"expected_head_sha": result.expected_head_sha,
"allowed_actions": list(result.allowed_actions),
"forbidden_actions": list(result.forbidden_actions),
"lease_role": lease_role,
"source": "control_plane_db.assign_and_lease",
}
consume_allocation = None
if cross_role_handoff and result.lease_id:
# Durable handoff marker so independent required-role workers can
# consume without sharing the controller session (#843).
handoff_prov = {
"cross_role_handoff": True,
"handoff_status": "pending",
"allocating_session_id": session_id,
"allocating_role": role_norm,
"required_role": expected_role,
"required_profile": selection["required_profile"],
"required_namespace": selection["required_namespace"],
"assignment_id": result.assignment_id,
"lease_id": result.lease_id,
"allocation_mode": mode,
"adopted_by_session_id": None,
}
try:
db.attach_lease_provenance(result.lease_id, handoff_prov)
except ControlPlaneError:
# Still return assignment evidence; consume path may be unavailable
handoff_prov["attach_failed"] = True
consume_allocation = {
"tool": "gitea_adopt_workflow_lease",
"lease_id": result.lease_id,
"assignment_id": result.assignment_id,
"required_role": expected_role,
"required_profile": selection["required_profile"],
"required_namespace": selection["required_namespace"],
"handoff_status": "pending",
"controller_session_required": False,
"instructions": (
f"From an independent {expected_role} session "
f"({selection['required_namespace']} / "
f"{selection['required_profile']}), call "
f"gitea_adopt_workflow_lease(lease_id={result.lease_id!r}) "
"to consume this controller allocation. The allocating "
"controller process does not need to remain alive. Wrong-role "
"and second-adoption attempts fail closed."
),
}
lease_proof["cross_role_handoff"] = True
lease_proof["handoff_status"] = "pending"
lease_proof["consume_tool"] = "gitea_adopt_workflow_lease"
out = {
"success": True, "success": True,
"outcome": OUTCOME_ASSIGNED, "outcome": OUTCOME_ASSIGNED,
"apply": True, "apply": True,
"role": role_norm, "role": role_norm,
"allocation_mode": mode,
"routing_role": role_norm,
"required_role": expected_role,
"selected_action": selection["selected_action"],
"required_profile": selection["required_profile"],
"required_namespace": selection["required_namespace"],
"profile_name": profile_name, "profile_name": profile_name,
"username": username, "username": username,
"session_id": session_id, "session_id": session_id,
@@ -1012,16 +1318,25 @@ def allocate_next_work(
"skipped": [s.as_dict() for s in skipped], "skipped": [s.as_dict() for s in skipped],
"terminal_pr": terminal_pr, "terminal_pr": terminal_pr,
"assignment": result.as_dict(), "assignment": result.as_dict(),
"lease_proof": { "lease_proof": lease_proof,
"assignment_id": result.assignment_id, "allocation_evidence": {
"lease_id": result.lease_id, "mode": "assigned",
"expires_at": result.expires_at, "allocation_mode": mode,
"expected_head_sha": result.expected_head_sha, "lease_created": True,
"allowed_actions": list(result.allowed_actions), "lease_role": lease_role,
"forbidden_actions": list(result.forbidden_actions), "lease_proof": lease_proof,
"source": "control_plane_db.assign_and_lease", "selection_policy": SELECTION_POLICY,
"cross_role_handoff": bool(cross_role_handoff),
}, },
"next_valid_command": _next_command(role_norm, selected), "next_valid_command": (
(
f"consume lease {result.lease_id} via gitea_adopt_workflow_lease "
f"as {expected_role}, then "
)
+ _next_command(lease_role, selected)
if cross_role_handoff
else _next_command(lease_role, selected)
),
"substrate": "control_plane_db", "substrate": "control_plane_db",
"file_lock_only": False, "file_lock_only": False,
"comment_lease_only": False, "comment_lease_only": False,
@@ -1031,11 +1346,19 @@ def allocate_next_work(
"controller_excluded": list(controller_excluded), "controller_excluded": list(controller_excluded),
"exclude_issue_numbers": list(exclude_nums), "exclude_issue_numbers": list(exclude_nums),
"candidate_set_fingerprint": cas_fp, "candidate_set_fingerprint": cas_fp,
"controller_allowed_actions": list(controller_allowed_actions),
"controller_forbidden_actions": list(controller_forbidden_actions),
"downstream_note": ( "downstream_note": (
"#612 incident bridge remains downstream of #600; " "#612 incident bridge remains downstream of #600; "
"allocator never assigns raw monitoring incidents" "allocator never assigns raw monitoring incidents; "
"controller routes only under cross_role (#840); "
"cross-role assignments are consumable by independent "
"required-role workers via gitea_adopt_workflow_lease (#843)"
), ),
} }
if consume_allocation is not None:
out["consume_allocation"] = consume_allocation
return out
def _next_command(role: str, c: WorkCandidate) -> str: def _next_command(role: str, c: WorkCandidate) -> str:
+892
View File
@@ -0,0 +1,892 @@
"""ARCH-01 Foundation Slice A — atomic platform installation + authority kernel (#822).
Parents: #820, #821. **First implementation leaf of the ARCH-01 program.**
This module implements the smallest executable ARCH-01 foundation:
* a connection-bound authenticated actor context (``cp_actor_*`` /
``cp_operation_mode`` / ``cp_context_epoch`` SQLite scalar functions that SQL
may *read* but can never *set* — ``[TRUSTED-SERVICE]`` authenticity);
* an immutable authority-dominance lattice with an exact seeded tuple set
(``[SCHEMA]``);
* the principal-equivalence root (a class exists *before* its first principal;
``principals.current_class_id`` is ``NOT NULL``; ``[SCHEMA]``);
* a single-transaction platform installation that seeds the initial
``platform.bootstrap`` grant and an immutable ``installed`` marker, validated
by a fail-closed ``install_state`` ``BEFORE INSERT`` trigger (``[SCHEMA]``).
Everything else in the ARCH-01/02/04 program (evidence stores, repository
bindings, workspaces, PostgreSQL parity, full grant succession, full principal
merge) is out of scope here and tracked in its own issue — see #822 §5/§17.
**Readiness / production posture.** This subsystem is *disabled by default*.
Nothing in the running MCP server imports or enables it. It becomes a security
boundary only once its readiness checks (the ACs in #822) pass in the target
environment. Instantiating :class:`PlatformKernel` creates an isolated SQLite
database and never touches the operational control-plane store.
Enforcement classification (per #820 vocabulary):
* ``[TRUSTED-SERVICE]`` — actor-context authenticity: the scalar functions are
registered by the trusted Python process; SQL cannot define or redefine them.
* ``[SCHEMA]`` — fail-closed aborts, the dominance/immutability/NOT-NULL-class/
last-active-grant invariants, enforced by CHECK/FK/trigger.
* ``[RUNTIME-ADAPTER]`` — *none* in this slice.
SQLite-first. ``BEGIN IMMEDIATE`` serializes concurrent installs and concurrent
grant/revoke on the singleton invariant row. PostgreSQL parity is a distinct
issue (#827); this module does **not** claim it.
"""
from __future__ import annotations
import os
import sqlite3
import threading
from contextlib import contextmanager
from dataclasses import dataclass
from datetime import datetime, timezone
from typing import Iterator, Optional
# --------------------------------------------------------------------------- #
# Closed enumerations (#822 §4).
# --------------------------------------------------------------------------- #
ACTOR_KINDS = ("operator", "supervisor", "service", "installer")
OPERATION_MODES = ("normal", "install", "merge", "internal_service")
# Exact seeded authority-dominance tuple set (#822 §4). This set is normative:
# the install-state trigger rejects any missing, additional, or malformed tuple.
DOMINANCE_TUPLES = (
("platform.bootstrap", "platform.bootstrap"),
("platform.bootstrap", "project.admin"),
("platform.bootstrap", "supervisor.root.establish"),
("supervisor.root", "supervisor.register"),
("supervisor.root", "supervisor.verify"),
("supervisor.root", "supervisor.recover"),
)
# The distinguished operator-key issuer seeded during install.
DISTINGUISHED_ISSUER_KIND = "operator-key"
DISTINGUISHED_ISSUER_ID = "platform.bootstrap.operator-key"
# Structured result codes (#822 §10).
INSTALLED = "INSTALLED"
ALREADY_INSTALLED = "ALREADY_INSTALLED"
INVALID_ACTOR_CONTEXT = "INVALID_ACTOR_CONTEXT"
INVALID_BOOTSTRAP_STATE = "INVALID_BOOTSTRAP_STATE"
DOMINANCE_SET_MISMATCH = "DOMINANCE_SET_MISMATCH"
AUTHORIZATION_DENIED = "AUTHORIZATION_DENIED"
CONCURRENT_INSTALLATION_LOST = "CONCURRENT_INSTALLATION_LOST"
# Required audit events (#822 §14).
EVT_PLATFORM_INSTALLED = "platform_installed"
EVT_GRANT_CREATED = "platform_grant_created"
EVT_GRANT_REVOKED = "platform_grant_revoked"
EVT_PRINCIPAL_REGISTERED = "principal_registered"
SCHEMA_VERSION = 1
DB_PATH_ENV = "ARCH01_PLATFORM_DB"
class PlatformKernelError(RuntimeError):
"""Base class for structured, code-bearing kernel failures."""
def __init__(self, code: str, message: str = "") -> None:
super().__init__(message or code)
self.code = code
class ActorContextError(PlatformKernelError):
"""Raised when a mutation is attempted without a valid actor context."""
# --------------------------------------------------------------------------- #
# Schema (#822 §6). Tables + fail-closed triggers.
#
# Every *mutating* trigger opens with the actor protocol: read the context
# epoch, read the actor fields, and abort unless the context is present,
# non-null, mode/kind well-formed, and epoch-consistent with the active
# transaction. The scalar functions ``cp_*`` are registered from Python only;
# SQL has no statement that can set them, which is the trusted-service boundary.
# --------------------------------------------------------------------------- #
_ACTOR_KINDS_SQL = ", ".join("'%s'" % k for k in ACTOR_KINDS)
_OP_MODES_SQL = ", ".join("'%s'" % m for m in OPERATION_MODES)
# Actor-protocol predicate: TRUE when the context is INVALID and the trigger
# must abort. ``cp_actor_context_valid()`` folds "present + non-expired +
# live-epoch == bound-epoch" (the read/re-read epoch equality of #822 §4) into
# one trusted-service answer; the remaining reads assert field well-formedness.
_INVALID_ACTOR = (
"cp_actor_context_valid() IS NOT 1 "
"OR cp_context_epoch() IS NULL "
"OR cp_actor_principal() IS NULL "
"OR cp_actor_kind() NOT IN (%s) "
"OR cp_operation_mode() NOT IN (%s)" % (_ACTOR_KINDS_SQL, _OP_MODES_SQL)
)
_ACTOR_GUARD = (
"SELECT CASE WHEN (%s) "
"THEN RAISE(ABORT, 'INVALID_ACTOR_CONTEXT') END;" % _INVALID_ACTOR
)
# require_installed: abort a privileged mutation when there is no install
# marker and we are not currently installing (#822 §4).
_REQUIRE_INSTALLED = (
"SELECT CASE WHEN ((SELECT COUNT(*) FROM install_state) = 0 "
"AND cp_operation_mode() <> 'install') "
"THEN RAISE(ABORT, 'NOT_INSTALLED') END;"
)
_SCHEMA_SQL = f"""
PRAGMA foreign_keys = ON;
CREATE TABLE IF NOT EXISTS arch01_meta (
key TEXT PRIMARY KEY,
value TEXT NOT NULL
);
-- Equivalence classes are created BEFORE their first principal (#822 §4).
CREATE TABLE IF NOT EXISTS principal_equivalence_classes (
class_id INTEGER PRIMARY KEY AUTOINCREMENT,
created_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS authoritative_issuers (
issuer_id INTEGER PRIMARY KEY AUTOINCREMENT,
issuer_kind TEXT NOT NULL,
issuer_ref TEXT NOT NULL,
created_at TEXT NOT NULL,
UNIQUE (issuer_kind, issuer_ref)
);
-- current_class_id is NOT NULL: a principal cannot exist without a class
-- (#822 AC6). issuer_id is nullable ONLY for the installer during install
-- (#822 AC7), enforced by trg_principals_null_issuer below.
CREATE TABLE IF NOT EXISTS principals (
principal_id TEXT PRIMARY KEY,
actor_kind TEXT NOT NULL CHECK (actor_kind IN ({_ACTOR_KINDS_SQL})),
current_class_id INTEGER NOT NULL REFERENCES principal_equivalence_classes(class_id),
issuer_id INTEGER REFERENCES authoritative_issuers(issuer_id),
registered_by TEXT REFERENCES principals(principal_id),
created_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS authority_dominance (
dominant TEXT NOT NULL,
subordinate TEXT NOT NULL,
PRIMARY KEY (dominant, subordinate)
);
CREATE TABLE IF NOT EXISTS platform_bootstrap_seed (
seed_id INTEGER PRIMARY KEY CHECK (seed_id = 1),
installer_principal_id TEXT NOT NULL REFERENCES principals(principal_id),
created_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS platform_bootstrap_grants (
grant_id INTEGER PRIMARY KEY AUTOINCREMENT,
grantee_principal_id TEXT NOT NULL REFERENCES principals(principal_id),
granted_by TEXT REFERENCES principals(principal_id),
active INTEGER NOT NULL DEFAULT 1 CHECK (active IN (0, 1)),
created_at TEXT NOT NULL,
revoked_at TEXT
);
-- Singleton row; active_count floored at 1 by CHECK so the last active grant
-- can never be revoked (#822 AC11).
CREATE TABLE IF NOT EXISTS platform_active_invariant (
id INTEGER PRIMARY KEY CHECK (id = 1),
active_count INTEGER NOT NULL CHECK (active_count >= 1)
);
-- The immutable install marker; inserted LAST in the install transaction.
CREATE TABLE IF NOT EXISTS install_state (
id INTEGER PRIMARY KEY CHECK (id = 1),
marker TEXT NOT NULL CHECK (marker = 'installed'),
installed_at TEXT NOT NULL
);
-- Append-only (#822 AC14).
CREATE TABLE IF NOT EXISTS audit_records (
audit_id INTEGER PRIMARY KEY AUTOINCREMENT,
event TEXT NOT NULL,
principal_id TEXT,
detail TEXT,
created_at TEXT NOT NULL
);
-- ------------------------------------------------------------------------- --
-- Actor protocol on every mutating trigger (#822 §4, [SCHEMA] fail-closed).
-- ------------------------------------------------------------------------- --
CREATE TRIGGER IF NOT EXISTS trg_classes_actor
BEFORE INSERT ON principal_equivalence_classes
BEGIN
{_ACTOR_GUARD}
END;
CREATE TRIGGER IF NOT EXISTS trg_issuers_actor
BEFORE INSERT ON authoritative_issuers
BEGIN
{_ACTOR_GUARD}
END;
CREATE TRIGGER IF NOT EXISTS trg_principals_actor
BEFORE INSERT ON principals
BEGIN
{_ACTOR_GUARD}
END;
CREATE TRIGGER IF NOT EXISTS trg_dominance_actor
BEFORE INSERT ON authority_dominance
BEGIN
{_ACTOR_GUARD}
END;
CREATE TRIGGER IF NOT EXISTS trg_seed_actor
BEFORE INSERT ON platform_bootstrap_seed
BEGIN
{_ACTOR_GUARD}
END;
CREATE TRIGGER IF NOT EXISTS trg_grants_actor_insert
BEFORE INSERT ON platform_bootstrap_grants
BEGIN
{_ACTOR_GUARD}
{_REQUIRE_INSTALLED}
END;
CREATE TRIGGER IF NOT EXISTS trg_grants_actor_update
BEFORE UPDATE ON platform_bootstrap_grants
BEGIN
{_ACTOR_GUARD}
END;
CREATE TRIGGER IF NOT EXISTS trg_invariant_actor_insert
BEFORE INSERT ON platform_active_invariant
BEGIN
{_ACTOR_GUARD}
END;
CREATE TRIGGER IF NOT EXISTS trg_invariant_actor_update
BEFORE UPDATE ON platform_active_invariant
BEGIN
{_ACTOR_GUARD}
END;
CREATE TRIGGER IF NOT EXISTS trg_audit_actor
BEFORE INSERT ON audit_records
BEGIN
{_ACTOR_GUARD}
END;
-- ------------------------------------------------------------------------- --
-- NOT-NULL-issuer exception for the installer only (#822 AC7).
-- A NULL issuer_id is accepted solely for an installer principal during
-- install mode, before the marker exists; any other NULL-issuer principal is
-- rejected. install-time issuer linkage (installer -> distinguished issuer)
-- is applied by a later UPDATE, permitted while no marker exists.
-- ------------------------------------------------------------------------- --
CREATE TRIGGER IF NOT EXISTS trg_principals_null_issuer
BEFORE INSERT ON principals
WHEN NEW.issuer_id IS NULL
BEGIN
SELECT CASE WHEN NOT (
NEW.actor_kind = 'installer'
AND cp_operation_mode() = 'install'
AND (SELECT COUNT(*) FROM install_state) = 0
AND (SELECT COUNT(*) FROM principals WHERE issuer_id IS NULL) = 0
) THEN RAISE(ABORT, 'INVALID_BOOTSTRAP_STATE') END;
END;
-- ------------------------------------------------------------------------- --
-- Post-install immutability of the authority root (#822 §4, AC9).
-- Registration fields freeze only AFTER the marker exists, so the install
-- transaction's own installer issuer-linkage UPDATE is permitted.
-- ------------------------------------------------------------------------- --
CREATE TRIGGER IF NOT EXISTS trg_principals_frozen_update
BEFORE UPDATE ON principals
WHEN (SELECT COUNT(*) FROM install_state) > 0
BEGIN
SELECT RAISE(ABORT, 'IMMUTABLE_PRINCIPAL');
END;
CREATE TRIGGER IF NOT EXISTS trg_principals_frozen_delete
BEFORE DELETE ON principals
BEGIN
SELECT RAISE(ABORT, 'IMMUTABLE_PRINCIPAL');
END;
-- Distinguished issuer identity is immutable once written.
CREATE TRIGGER IF NOT EXISTS trg_issuers_immutable_update
BEFORE UPDATE ON authoritative_issuers
BEGIN
SELECT RAISE(ABORT, 'IMMUTABLE_ISSUER');
END;
CREATE TRIGGER IF NOT EXISTS trg_issuers_immutable_delete
BEFORE DELETE ON authoritative_issuers
BEGIN
SELECT RAISE(ABORT, 'IMMUTABLE_ISSUER');
END;
-- The dominance lattice is immutable once seeded.
CREATE TRIGGER IF NOT EXISTS trg_dominance_immutable_update
BEFORE UPDATE ON authority_dominance
BEGIN
SELECT RAISE(ABORT, 'IMMUTABLE_DOMINANCE');
END;
CREATE TRIGGER IF NOT EXISTS trg_dominance_immutable_delete
BEFORE DELETE ON authority_dominance
BEGIN
SELECT RAISE(ABORT, 'IMMUTABLE_DOMINANCE');
END;
-- The bootstrap seed is immutable once written.
CREATE TRIGGER IF NOT EXISTS trg_seed_immutable_update
BEFORE UPDATE ON platform_bootstrap_seed
BEGIN
SELECT RAISE(ABORT, 'IMMUTABLE_SEED');
END;
CREATE TRIGGER IF NOT EXISTS trg_seed_immutable_delete
BEFORE DELETE ON platform_bootstrap_seed
BEGIN
SELECT RAISE(ABORT, 'IMMUTABLE_SEED');
END;
-- The install marker is immutable once written.
CREATE TRIGGER IF NOT EXISTS trg_install_state_immutable_update
BEFORE UPDATE ON install_state
BEGIN
SELECT RAISE(ABORT, 'IMMUTABLE_INSTALL_STATE');
END;
CREATE TRIGGER IF NOT EXISTS trg_install_state_immutable_delete
BEFORE DELETE ON install_state
BEGIN
SELECT RAISE(ABORT, 'IMMUTABLE_INSTALL_STATE');
END;
-- Grants: identity is immutable; the ONLY permitted mutation is a single
-- active 1 -> 0 revocation (#822 §4 initial-grant identity immutability +
-- grant/revoke). Reactivation and identity edits are rejected.
CREATE TRIGGER IF NOT EXISTS trg_grants_identity_frozen
BEFORE UPDATE ON platform_bootstrap_grants
WHEN NOT (
NEW.grant_id = OLD.grant_id
AND NEW.grantee_principal_id = OLD.grantee_principal_id
AND NEW.granted_by IS OLD.granted_by
AND NEW.created_at = OLD.created_at
AND OLD.active = 1
AND NEW.active = 0
)
BEGIN
SELECT RAISE(ABORT, 'IMMUTABLE_GRANT');
END;
CREATE TRIGGER IF NOT EXISTS trg_grants_no_delete
BEFORE DELETE ON platform_bootstrap_grants
BEGIN
SELECT RAISE(ABORT, 'IMMUTABLE_GRANT');
END;
-- audit_records is append-only.
CREATE TRIGGER IF NOT EXISTS trg_audit_immutable_update
BEFORE UPDATE ON audit_records
BEGIN
SELECT RAISE(ABORT, 'IMMUTABLE_AUDIT');
END;
CREATE TRIGGER IF NOT EXISTS trg_audit_immutable_delete
BEFORE DELETE ON audit_records
BEGIN
SELECT RAISE(ABORT, 'IMMUTABLE_AUDIT');
END;
-- ------------------------------------------------------------------------- --
-- install_state BEFORE INSERT: validate the whole bootstrap atomically
-- (#822 §4, AC4). Each dominance tuple is checked individually; a missing,
-- additional, or malformed tuple -> DOMINANCE_SET_MISMATCH. The seed<->installer
-- link, the single active NULL-grantor installer grant, the installer's
-- non-NULL issuer, the active invariant, and "no extra principal created under
-- the NULL-issuer exception" -> INVALID_BOOTSTRAP_STATE.
-- ------------------------------------------------------------------------- --
CREATE TRIGGER IF NOT EXISTS trg_install_state_validate
BEFORE INSERT ON install_state
BEGIN
SELECT CASE WHEN NOT (
(SELECT COUNT(*) FROM authority_dominance) = {len(DOMINANCE_TUPLES)}
AND EXISTS (SELECT 1 FROM authority_dominance WHERE dominant='platform.bootstrap' AND subordinate='platform.bootstrap')
AND EXISTS (SELECT 1 FROM authority_dominance WHERE dominant='platform.bootstrap' AND subordinate='project.admin')
AND EXISTS (SELECT 1 FROM authority_dominance WHERE dominant='platform.bootstrap' AND subordinate='supervisor.root.establish')
AND EXISTS (SELECT 1 FROM authority_dominance WHERE dominant='supervisor.root' AND subordinate='supervisor.register')
AND EXISTS (SELECT 1 FROM authority_dominance WHERE dominant='supervisor.root' AND subordinate='supervisor.verify')
AND EXISTS (SELECT 1 FROM authority_dominance WHERE dominant='supervisor.root' AND subordinate='supervisor.recover')
) THEN RAISE(ABORT, 'DOMINANCE_SET_MISMATCH') END;
SELECT CASE WHEN NOT (
(SELECT COUNT(*) FROM platform_bootstrap_seed) = 1
AND (SELECT COUNT(*) FROM principals) = 1
AND (SELECT actor_kind FROM principals
WHERE principal_id = (SELECT installer_principal_id FROM platform_bootstrap_seed WHERE seed_id = 1)
) = 'installer'
AND (SELECT issuer_id FROM principals
WHERE principal_id = (SELECT installer_principal_id FROM platform_bootstrap_seed WHERE seed_id = 1)
) IS NOT NULL
AND (SELECT COUNT(*) FROM platform_bootstrap_grants
WHERE granted_by IS NULL AND active = 1
AND grantee_principal_id = (SELECT installer_principal_id FROM platform_bootstrap_seed WHERE seed_id = 1)
) = 1
AND (SELECT COUNT(*) FROM platform_bootstrap_grants) = 1
AND (SELECT active_count FROM platform_active_invariant WHERE id = 1) = 1
) THEN RAISE(ABORT, 'INVALID_BOOTSTRAP_STATE') END;
END;
"""
def default_db_path() -> str:
return os.environ.get(
DB_PATH_ENV,
os.path.expanduser("~/.cache/gitea-tools/arch01/platform.sqlite3"),
)
def _utc_now_iso() -> str:
return datetime.now(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
@dataclass(frozen=True)
class OperationResult:
"""Structured result of a kernel operation (#822 §10)."""
code: str
detail: str = ""
@property
def ok(self) -> bool:
return self.code in (INSTALLED, ALREADY_INSTALLED)
@dataclass
class _ActorContext:
principal: str
kind: str
mode: str
session: Optional[str]
bound_epoch: int
live_epoch: int
expired: bool = False
class PlatformKernel:
"""ARCH-01 authority kernel over a single SQLite connection.
The connection carries the trusted-service actor context: the ``cp_*``
scalar functions read the context this object holds. Only Python code here
can bind or clear it, so no SQL statement can assert an actor identity — the
trusted-service authenticity boundary of #822 §4.
"""
def __init__(self, db_path: Optional[str] = None, *, busy_timeout_ms: int = 5000) -> None:
self.db_path = db_path or default_db_path()
if self.db_path != ":memory:":
parent = os.path.dirname(self.db_path)
if parent:
os.makedirs(parent, exist_ok=True)
self._ctx: Optional[_ActorContext] = None
self._epoch_seq = 0
self._lock = threading.Lock()
# check_same_thread=False is safe: every mutation path is serialized
# by self._lock, so the connection is never used concurrently even when
# callers drive the kernel from different threads (concurrency tests).
self._conn = sqlite3.connect(
self.db_path, isolation_level=None, check_same_thread=False
)
self._conn.execute("PRAGMA foreign_keys = ON")
self._conn.execute(f"PRAGMA busy_timeout = {int(busy_timeout_ms)}")
self._register_actor_functions()
self._migrate()
# -- trusted-service actor functions ---------------------------------- #
def _register_actor_functions(self) -> None:
c = self._conn
c.create_function("cp_actor_principal", 0, lambda: self._ctx.principal if self._ctx else None)
c.create_function("cp_actor_kind", 0, lambda: self._ctx.kind if self._ctx else None)
c.create_function("cp_operation_mode", 0, lambda: self._ctx.mode if self._ctx else None)
c.create_function("cp_service_session", 0, lambda: self._ctx.session if self._ctx else None)
c.create_function("cp_context_epoch", 0, self._fn_context_epoch)
# Trusted-service helper: folds present + non-expired + epoch-consistent
# into the read/re-read epoch equality of #822 §4.
c.create_function("cp_actor_context_valid", 0, self._fn_context_valid)
def _fn_context_epoch(self) -> Optional[int]:
if self._ctx is None or self._ctx.expired:
return None
return self._ctx.live_epoch
def _fn_context_valid(self) -> int:
ctx = self._ctx
if ctx is None or ctx.expired:
return 0
# read/re-read epoch equality: a context whose live epoch has drifted
# from the epoch it was bound to (a stale/replaced connection context)
# is not bound to the active transaction and fails closed.
if ctx.live_epoch != ctx.bound_epoch:
return 0
if ctx.principal is None:
return 0
if ctx.kind not in ACTOR_KINDS or ctx.mode not in OPERATION_MODES:
return 0
return 1
# -- context lifecycle ------------------------------------------------ #
@contextmanager
def actor_context(
self, principal: str, kind: str, mode: str, session: Optional[str] = None
) -> Iterator[None]:
"""Bind a trusted actor context for the duration of the block."""
prev = self._ctx
self._epoch_seq += 1
epoch = self._epoch_seq
self._ctx = _ActorContext(
principal=principal, kind=kind, mode=mode, session=session,
bound_epoch=epoch, live_epoch=epoch,
)
try:
yield
finally:
self._ctx = prev
def _clear_context(self) -> None:
self._ctx = None
# -- migration -------------------------------------------------------- #
def _migrate(self) -> None:
self._conn.executescript(_SCHEMA_SQL)
self._conn.execute(
"INSERT OR IGNORE INTO arch01_meta(key, value) VALUES ('schema_version', ?)",
(str(SCHEMA_VERSION),),
)
self._conn.execute(
"INSERT OR IGNORE INTO arch01_meta(key, value) VALUES "
"('architecture', 'ARCH-01 Slice A: atomic install + authority kernel (#822); "
"disabled by default until readiness checks pass')"
)
# -- introspection ---------------------------------------------------- #
def is_installed(self) -> bool:
row = self._conn.execute("SELECT COUNT(*) FROM install_state").fetchone()
return bool(row[0])
def active_grant_count(self) -> int:
row = self._conn.execute(
"SELECT active_count FROM platform_active_invariant WHERE id = 1"
).fetchone()
return int(row[0]) if row else 0
def audit_events(self) -> list[str]:
return [
r[0]
for r in self._conn.execute(
"SELECT event FROM audit_records ORDER BY audit_id"
).fetchall()
]
def close(self) -> None:
self._conn.close()
# -- operations ------------------------------------------------------- #
def install_platform(
self,
installer_principal_id: str = "platform.installer",
*,
session: Optional[str] = None,
) -> OperationResult:
"""Single atomic install transaction (#822 §4/§7).
``BEGIN IMMEDIATE`` serializes concurrent installs; the loser rechecks
the marker and returns ``ALREADY_INSTALLED``, or — if it never acquires
the write lock — ``CONCURRENT_INSTALLATION_LOST``. On any stage failure
the whole transaction rolls back leaving no partial rows (AC3/AC5).
"""
now = _utc_now_iso()
with self._lock:
try:
self._conn.execute("BEGIN IMMEDIATE")
except sqlite3.OperationalError as exc:
if "locked" in str(exc).lower() or "busy" in str(exc).lower():
return OperationResult(CONCURRENT_INSTALLATION_LOST, str(exc))
raise
try:
if self.is_installed():
self._conn.execute("ROLLBACK")
return OperationResult(ALREADY_INSTALLED, "install marker already present")
with self.actor_context(installer_principal_id, "installer", "install", session):
c = self._conn
# class -> installer principal (temporary NULL issuer)
cur = c.execute(
"INSERT INTO principal_equivalence_classes(created_at) VALUES (?)",
(now,),
)
class_id = cur.lastrowid
c.execute(
"INSERT INTO principals"
"(principal_id, actor_kind, current_class_id, issuer_id, registered_by, created_at) "
"VALUES (?, 'installer', ?, NULL, ?, ?)",
(installer_principal_id, class_id, installer_principal_id, now),
)
# distinguished operator-key issuer
cur = c.execute(
"INSERT INTO authoritative_issuers(issuer_kind, issuer_ref, created_at) "
"VALUES (?, ?, ?)",
(DISTINGUISHED_ISSUER_KIND, DISTINGUISHED_ISSUER_ID, now),
)
issuer_id = cur.lastrowid
# link installer -> issuer (permitted pre-marker)
c.execute(
"UPDATE principals SET issuer_id = ? WHERE principal_id = ?",
(issuer_id, installer_principal_id),
)
# dominance tuples
c.executemany(
"INSERT INTO authority_dominance(dominant, subordinate) VALUES (?, ?)",
DOMINANCE_TUPLES,
)
# seed
c.execute(
"INSERT INTO platform_bootstrap_seed(seed_id, installer_principal_id, created_at) "
"VALUES (1, ?, ?)",
(installer_principal_id, now),
)
# initial grant (granted_by NULL, active)
c.execute(
"INSERT INTO platform_bootstrap_grants"
"(grantee_principal_id, granted_by, active, created_at) "
"VALUES (?, NULL, 1, ?)",
(installer_principal_id, now),
)
# active invariant
c.execute(
"INSERT INTO platform_active_invariant(id, active_count) VALUES (1, 1)"
)
# audit rows for the security-sensitive operation
c.execute(
"INSERT INTO audit_records(event, principal_id, detail, created_at) "
"VALUES (?, ?, ?, ?)",
(EVT_PRINCIPAL_REGISTERED, installer_principal_id, "installer", now),
)
c.execute(
"INSERT INTO audit_records(event, principal_id, detail, created_at) "
"VALUES (?, ?, ?, ?)",
(EVT_GRANT_CREATED, installer_principal_id, "initial platform.bootstrap grant", now),
)
# install marker LAST -> fires the whole-bootstrap validator
c.execute(
"INSERT INTO install_state(id, marker, installed_at) VALUES (1, 'installed', ?)",
(now,),
)
c.execute(
"INSERT INTO audit_records(event, principal_id, detail, created_at) "
"VALUES (?, ?, ?, ?)",
(EVT_PLATFORM_INSTALLED, installer_principal_id, "platform installed", now),
)
self._conn.execute("COMMIT")
return OperationResult(INSTALLED, "platform installed")
except sqlite3.Error as exc:
self._safe_rollback()
return OperationResult(self._classify(exc), str(exc))
def register_principal(
self,
principal_id: str,
actor_kind: str,
issuer_ref: str,
*,
actor_principal: str,
actor_kind_ctx: str = "operator",
session: Optional[str] = None,
) -> OperationResult:
"""Atomically create an equivalence class and its first principal.
The class is inserted *before* the principal, and ``current_class_id``
is ``NOT NULL`` (#822 AC6): a principal can never exist classless.
The principal references an existing issuer (non-NULL); the temporary
NULL-issuer exception is reserved for the installer during install
(AC7).
"""
if actor_kind not in ACTOR_KINDS:
return OperationResult(INVALID_BOOTSTRAP_STATE, f"bad actor_kind {actor_kind!r}")
now = _utc_now_iso()
with self._lock:
try:
self._conn.execute("BEGIN IMMEDIATE")
except sqlite3.OperationalError as exc:
return OperationResult(AUTHORIZATION_DENIED, str(exc))
try:
if not self.is_installed():
self._conn.execute("ROLLBACK")
return OperationResult(INVALID_BOOTSTRAP_STATE, "platform not installed")
row = self._conn.execute(
"SELECT issuer_id FROM authoritative_issuers WHERE issuer_ref = ?",
(issuer_ref,),
).fetchone()
if row is None:
self._conn.execute("ROLLBACK")
return OperationResult(INVALID_BOOTSTRAP_STATE, f"unknown issuer {issuer_ref!r}")
issuer_id = row[0]
with self.actor_context(actor_principal, actor_kind_ctx, "normal", session):
cur = self._conn.execute(
"INSERT INTO principal_equivalence_classes(created_at) VALUES (?)",
(now,),
)
class_id = cur.lastrowid
self._conn.execute(
"INSERT INTO principals"
"(principal_id, actor_kind, current_class_id, issuer_id, registered_by, created_at) "
"VALUES (?, ?, ?, ?, ?, ?)",
(principal_id, actor_kind, class_id, issuer_id, actor_principal, now),
)
self._conn.execute(
"INSERT INTO audit_records(event, principal_id, detail, created_at) "
"VALUES (?, ?, ?, ?)",
(EVT_PRINCIPAL_REGISTERED, principal_id, actor_kind, now),
)
self._conn.execute("COMMIT")
return OperationResult(INSTALLED, f"registered {principal_id}")
except sqlite3.Error as exc:
self._safe_rollback()
return OperationResult(self._classify(exc), str(exc))
def grant_platform_bootstrap(
self,
grantee_principal_id: str,
granted_by: str,
*,
actor_kind_ctx: str = "operator",
session: Optional[str] = None,
) -> OperationResult:
"""Create an additional active platform.bootstrap grant.
Serialized on the singleton invariant row via ``BEGIN IMMEDIATE``.
"""
now = _utc_now_iso()
with self._lock:
try:
self._conn.execute("BEGIN IMMEDIATE")
except sqlite3.OperationalError as exc:
return OperationResult(AUTHORIZATION_DENIED, str(exc))
try:
if not self.is_installed():
self._conn.execute("ROLLBACK")
return OperationResult(INVALID_BOOTSTRAP_STATE, "platform not installed")
with self.actor_context(granted_by, actor_kind_ctx, "normal", session):
self._conn.execute(
"INSERT INTO platform_bootstrap_grants"
"(grantee_principal_id, granted_by, active, created_at) "
"VALUES (?, ?, 1, ?)",
(grantee_principal_id, granted_by, now),
)
self._conn.execute(
"UPDATE platform_active_invariant SET active_count = active_count + 1 WHERE id = 1"
)
self._conn.execute(
"INSERT INTO audit_records(event, principal_id, detail, created_at) "
"VALUES (?, ?, ?, ?)",
(EVT_GRANT_CREATED, grantee_principal_id, f"granted_by={granted_by}", now),
)
self._conn.execute("COMMIT")
return OperationResult(INSTALLED, f"granted to {grantee_principal_id}")
except sqlite3.Error as exc:
self._safe_rollback()
return OperationResult(self._classify(exc), str(exc))
def revoke_platform_bootstrap(
self,
grant_id: int,
*,
actor_principal: str,
actor_kind_ctx: str = "operator",
session: Optional[str] = None,
) -> OperationResult:
"""Revoke an active grant, floored so the last one can never drop.
The ``active_count >= 1`` CHECK plus ``BEGIN IMMEDIATE`` serialization
make two concurrent revocations unable to remove the final active grant
(#822 AC11): the decrement that would reach zero fails and rolls back.
"""
now = _utc_now_iso()
with self._lock:
try:
self._conn.execute("BEGIN IMMEDIATE")
except sqlite3.OperationalError as exc:
return OperationResult(AUTHORIZATION_DENIED, str(exc))
try:
if not self.is_installed():
self._conn.execute("ROLLBACK")
return OperationResult(INVALID_BOOTSTRAP_STATE, "platform not installed")
row = self._conn.execute(
"SELECT active, grantee_principal_id FROM platform_bootstrap_grants WHERE grant_id = ?",
(grant_id,),
).fetchone()
if row is None or row[0] != 1:
self._conn.execute("ROLLBACK")
return OperationResult(AUTHORIZATION_DENIED, "grant absent or already inactive")
grantee = row[1]
with self.actor_context(actor_principal, actor_kind_ctx, "normal", session):
# Decrement first: the CHECK floor rejects dropping below 1,
# aborting the whole revoke before the grant flips inactive.
self._conn.execute(
"UPDATE platform_active_invariant SET active_count = active_count - 1 WHERE id = 1"
)
self._conn.execute(
"UPDATE platform_bootstrap_grants SET active = 0, revoked_at = ? WHERE grant_id = ?",
(now, grant_id),
)
self._conn.execute(
"INSERT INTO audit_records(event, principal_id, detail, created_at) "
"VALUES (?, ?, ?, ?)",
(EVT_GRANT_REVOKED, grantee, f"grant_id={grant_id}", now),
)
self._conn.execute("COMMIT")
return OperationResult(INSTALLED, f"revoked grant {grant_id}")
except sqlite3.Error as exc:
self._safe_rollback()
return OperationResult(self._classify(exc), str(exc))
# -- helpers ---------------------------------------------------------- #
def _safe_rollback(self) -> None:
try:
self._conn.execute("ROLLBACK")
except sqlite3.Error:
pass
@staticmethod
def _classify(exc: sqlite3.Error) -> str:
msg = str(exc)
if "INVALID_ACTOR_CONTEXT" in msg:
return INVALID_ACTOR_CONTEXT
if "DOMINANCE_SET_MISMATCH" in msg:
return DOMINANCE_SET_MISMATCH
if "active_count" in msg or "CHECK constraint failed: platform_active_invariant" in msg:
# last-active-grant floor tripped
return AUTHORIZATION_DENIED
if any(tag in msg for tag in (
"INVALID_BOOTSTRAP_STATE", "IMMUTABLE_", "NOT_INSTALLED",
)):
return INVALID_BOOTSTRAP_STATE
return INVALID_BOOTSTRAP_STATE
+169 -3
View File
@@ -1637,11 +1637,13 @@ class ControlPlaneDB:
provenance: dict[str, Any] | None = None, provenance: dict[str, Any] | None = None,
lease_ttl_seconds: int = DEFAULT_LEASE_TTL_SECONDS, lease_ttl_seconds: int = DEFAULT_LEASE_TTL_SECONDS,
) -> dict[str, Any]: ) -> dict[str, Any]:
"""Transfer or refresh a lease with provenance (#601). """Transfer or refresh a lease with provenance (#601 / #843).
* Same owner + active → refresh (owner-resume). * Same owner + active → refresh (owner-resume).
* Cross-role handoff pending + matching required role → atomic consume
(even while the allocating controller session still "owns" the lease).
* Expired/abandoned/released → create new assignment+lease with provenance. * Expired/abandoned/released → create new assignment+lease with provenance.
* Active foreign → raise ForeignLeaseError (never silent steal). * Active foreign (non-handoff) → raise ForeignLeaseError (never silent steal).
""" """
now = _utc_now() now = _utc_now()
now_s = _ts(now) now_s = _ts(now)
@@ -1677,7 +1679,35 @@ class ControlPlaneDB:
status = "expired" status = "expired"
owner = lease["session_id"] owner = lease["session_id"]
if status == "active" and owner != adopter_session_id: # Parse durable provenance for cross-role handoff consume (#843).
lease_prov: dict[str, Any] = {}
if "provenance_json" in lease.keys() and lease["provenance_json"]:
try:
loaded = json.loads(lease["provenance_json"])
if isinstance(loaded, dict):
lease_prov = loaded
except (TypeError, json.JSONDecodeError):
lease_prov = {}
handoff_pending = bool(lease_prov.get("cross_role_handoff")) and (
str(lease_prov.get("handoff_status") or "pending").strip().lower()
== "pending"
)
already_adopted = bool(
(lease["adopted_by_session_id"] if "adopted_by_session_id" in lease.keys() else None)
or lease_prov.get("adopted_by_session_id")
)
required_role = str(
lease_prov.get("required_role") or lease["role"] or ""
).strip().lower()
adopter_role = (role or "").strip().lower()
cross_role_consume = (
handoff_pending
and not already_adopted
and status == "active"
and owner != adopter_session_id
)
if status == "active" and owner != adopter_session_id and not cross_role_consume:
raise ForeignLeaseError( raise ForeignLeaseError(
f"cannot adopt active foreign lease {lease_id} owned by {owner}" f"cannot adopt active foreign lease {lease_id} owned by {owner}"
) )
@@ -1761,6 +1791,142 @@ class ControlPlaneDB:
"reasons": ["owner-resume: refreshed lease with provenance"], "reasons": ["owner-resume: refreshed lease with provenance"],
} }
# #843: controller→required-role handoff consume (atomic, same lease_id)
if cross_role_consume:
if not required_role:
raise ControlPlaneError(
f"cross-role handoff lease {lease_id} missing required_role"
)
if adopter_role != required_role:
raise ForeignLeaseError(
f"wrong role for cross-role handoff consume: "
f"required={required_role} adopter={adopter_role or 'none'} "
f"(fail closed)"
)
# CAS: only transfer if still owned by allocating session and unadopted
cols = self._lease_columns(conn)
adopted_col_null = (
"(adopted_by_session_id IS NULL OR adopted_by_session_id = '')"
if "adopted_by_session_id" in cols
else "1=1"
)
cas = conn.execute(
f"""
UPDATE leases
SET session_id = ?,
heartbeat_at = ?,
expires_at = ?,
phase = ?,
role = ?
WHERE lease_id = ?
AND status = 'active'
AND session_id = ?
AND {adopted_col_null}
""",
(
adopter_session_id,
now_s,
expires,
"adopted",
required_role,
lease_id,
owner,
),
)
if cas.rowcount != 1:
raise ForeignLeaseError(
f"cross-role handoff consume lost race for lease {lease_id} "
"(already adopted or no longer pending; fail closed)"
)
if "adopted_from_session_id" in cols:
conn.execute(
"""
UPDATE leases
SET adopted_from_session_id = ?, adopted_by_session_id = ?
WHERE lease_id = ?
""",
(owner, adopter_session_id, lease_id),
)
if "worktree_path" in cols and worktree_path:
conn.execute(
"UPDATE leases SET worktree_path = ? WHERE lease_id = ?",
(worktree_path, lease_id),
)
if "owner_pid" in cols and owner_pid is not None:
conn.execute(
"UPDATE leases SET owner_pid = ? WHERE lease_id = ?",
(owner_pid, lease_id),
)
if "expected_head_sha" in cols and expected_head_sha:
conn.execute(
"UPDATE leases SET expected_head_sha = ? WHERE lease_id = ?",
(expected_head_sha, lease_id),
)
# Merge handoff provenance + caller provenance
merged = dict(lease_prov)
merged.update(provenance or {})
merged["cross_role_handoff"] = True
merged["handoff_status"] = "adopted"
merged["adopted_from_session_id"] = owner
merged["adopted_by_session_id"] = adopter_session_id
merged["required_role"] = required_role
if "provenance_json" in cols:
conn.execute(
"UPDATE leases SET provenance_json = ? WHERE lease_id = ?",
(json.dumps(merged), lease_id),
)
# Transfer active assignment ownership atomically
asn_cas = conn.execute(
"""
UPDATE assignments
SET session_id = ?, role = ?
WHERE lease_id = ? AND status = 'active' AND session_id = ?
""",
(adopter_session_id, required_role, lease_id, owner),
)
if asn_cas.rowcount < 1:
# Fail closed: assignment must move with the lease
raise ControlPlaneError(
f"cross-role handoff: no active assignment for lease {lease_id} "
f"owned by {owner}"
)
lease2 = conn.execute(
"SELECT * FROM leases WHERE lease_id = ?", (lease_id,)
).fetchone()
asn = conn.execute(
"""
SELECT * FROM assignments
WHERE lease_id = ? AND status = 'active'
ORDER BY created_at DESC LIMIT 1
""",
(lease_id,),
).fetchone()
conn.execute(
"""
INSERT INTO events(work_item_id, event_type, message, created_at)
VALUES (?, 'lease_adopted', ?, ?)
""",
(
lease["work_item_id"],
f"cross-role handoff: {adopter_session_id} consumed "
f"{lease_id} from {owner} as {required_role}",
now_s,
),
)
return {
"outcome": "adopted_cross_role_handoff",
"lease": dict(lease2) if lease2 else dict(lease),
"assignment": dict(asn) if asn else None,
"reasons": [
"cross-role handoff: independent required-role worker consumed "
"controller allocation without abandonment"
],
"adopted_by_session_id": adopter_session_id,
"adopted_from_session_id": owner,
"required_role": required_role,
"handoff_status": "adopted",
}
# Non-active: create new lease + assignment (transfer) # Non-active: create new lease + assignment (transfer)
new_lease_id = f"lease-{uuid.uuid4().hex[:16]}" new_lease_id = f"lease-{uuid.uuid4().hex[:16]}"
new_asn_id = f"asn-{uuid.uuid4().hex[:16]}" new_asn_id = f"asn-{uuid.uuid4().hex[:16]}"
+81 -1
View File
@@ -52,7 +52,8 @@ status, onboarding checklist state, and the fail-closed error payloads (#635).
| Path | Description | | Path | Description |
|------|-------------| |------|-------------|
| `/` | Home / operator overview | | `/` | Home / operator overview |
| `/health` | JSON liveness (`status`, `service`, `mode`, `timestamp`) | | `/health` | JSON liveness (`status`, `service`, `mode`, `timestamp`, `uptime_seconds`) |
| `/api/v1/system/health` | Structured read-only system health (#634) |
| `/queue` | Live PR and issue queue dashboard (#429) | | `/queue` | Live PR and issue queue dashboard (#429) |
| `/api/queue` | JSON queue export with pagination metadata | | `/api/queue` | JSON queue export with pagination metadata |
| `/projects` | Project registry list with status and onboarding progress (#427, #635) | | `/projects` | Project registry list with status and onboarding progress (#427, #635) |
@@ -78,6 +79,85 @@ Most routes are GET-only. POST/PUT/PATCH/DELETE return `405` with
`read-only-mvp`, except `/audit` and `/api/audit` which accept POST for `read-only-mvp`, except `/audit` and `/api/audit` which accept POST for
local validator preview only (no Gitea mutations, no server-side storage). local validator preview only (no Gitea mutations, no server-side storage).
## System health API (#634)
`GET /api/v1/system/health` is the structured, read-only health surface for
automated readiness checks. It is the first console API under the `/api/v1`
prefix; the unversioned MVP exports remain as compatibility aliases.
`/health` is unchanged for existing consumers — every MVP key is still present
— and now also carries `started_at`, `uptime_seconds`, and a
`system_health_api` pointer. It stays deliberately cheap and runs no dependency
probe, because answering readiness costs real work.
**Status codes.** `200` when ready, `503` when a required dependency failed or
was never probed. Automation can branch on the code without parsing the body.
**Query flags.** The Gitea check is a network call, so it is opt-in:
`GET /api/v1/system/health?deep=1` runs it and caches the result for
`WEBUI_HEALTH_PROBE_TTL_SECONDS` (default 15s) so dashboard polling does not
amplify into remote load. Without the flag that probe reports `skipped`.
**Dependencies.** `control_plane_db` and `repository` are required and drive
readiness. `gitea` is optional: when it fails the overall `status` degrades but
`readiness.ready` stays true, because local inventory is still serveable. Each
entry carries `status`, `detail`, `required`, and `latency_ms`.
Two honesty rules are worth knowing before reading the payload:
* `stale_runtime.mutation_safe` is true only when the runtime, checkout, and
remote-tracking commits are all known and equal. An unfetched remote is
reported as indeterminate, never as safe.
* `mcp_namespaces` entries are always `unproven`. A web process runs outside
the IDE-managed MCP client and cannot invoke a namespace tool, so per #543
only a `client_namespace` probe can prove that path.
Sample response (abridged, healthy):
```json
{
"status": "ok",
"service": "mcp-control-plane-webui",
"mode": "read-only",
"api": "/api/v1/system/health",
"timestamp": "2026-07-22T11:04:18.512034+00:00",
"readiness": { "ready": true, "complete": true, "reasons": [] },
"version": {
"git_sha": "620ed6e9a9550b8da2ceb82d9ab8744e8920490f",
"git_describe": "v1.1.0-898-g620ed6e",
"control_plane_schema_version": 4,
"python_version": "3.14.5",
"known": true
},
"process": { "started_at": "2026-07-22T10:58:02.114+00:00", "uptime_seconds": 376.4 },
"deep_probes_requested": false,
"dependencies": [
{
"name": "control_plane_db",
"kind": "sqlite",
"status": "ok",
"detail": "schema v4 readable",
"required": true,
"healthy": true,
"latency_ms": 1.482,
"metadata": { "schema_version": 4, "active_leases": 3 }
},
{ "name": "repository", "kind": "git", "status": "ok", "required": true, "healthy": true },
{ "name": "gitea", "kind": "http", "status": "skipped", "required": false, "healthy": false }
],
"mcp_namespaces": [
{ "namespace": "gitea-author", "required_tool": "gitea_whoami", "status": "unproven" }
],
"stale_runtime": { "stale": false, "determinable": true, "mutation_safe": true, "reasons": [] },
"probe_errors": []
}
```
No restart, reload, or process-kill control is exposed here: those are Phase 2
at the earliest, and #630 forbids process-kill recovery outright. Every probe
opens its subject read-only — the control-plane database is opened through a
`mode=ro` URI so a health check can never create or migrate a schema.
## Report audit (#431) ## Report audit (#431)
Paste an LLM final report at `/audit` or POST JSON to `/api/audit`. The UI Paste an LLM final report at `/audit` or POST JSON to `/api/audit`. The UI
+98 -13
View File
@@ -234,12 +234,25 @@ def _effective_workspace_role() -> str:
def _profile_role_kind(profile: dict) -> str: def _profile_role_kind(profile: dict) -> str:
"""Resolve a profile's declared role before inferring from permissions.""" """Resolve a profile's declared role before inferring from permissions.
role = (profile.get("role") or profile.get("role_kind") or "").strip()
Declared ``role`` / ``role_kind`` always wins so a controller profile is
never reclassified as reconciler from permission inference (#840).
"""
role = (profile.get("role") or profile.get("role_kind") or "").strip().lower()
if role: if role:
# Normalize aliases / case.
if "control" in role:
return "controller"
return role return role
profile_name = (profile.get("profile_name") or "").strip().lower() profile_name = (profile.get("profile_name") or "").strip().lower()
for candidate in ("reconciler", "merger", "reviewer", "author"): for candidate in (
"controller",
"reconciler",
"merger",
"reviewer",
"author",
):
if candidate in profile_name: if candidate in profile_name:
return candidate return candidate
return _role_kind( return _role_kind(
@@ -15538,7 +15551,8 @@ def mcp_get_control_plane_guide(
profile = get_profile() profile = get_profile()
allowed = profile["allowed_operations"] allowed = profile["allowed_operations"]
forbidden = profile["forbidden_operations"] forbidden = profile["forbidden_operations"]
role = _role_kind(allowed, forbidden) # Prefer declared profile role so controller is not mislabeled reconciler (#840).
role = _profile_role_kind(profile)
username = _authenticated_username(h) username = _authenticated_username(h)
identity = { identity = {
@@ -15597,6 +15611,16 @@ def mcp_get_control_plane_guide(
"user, and merging requires explicit operator authorization plus the " "user, and merging requires explicit operator authorization plus the "
"'MERGE PR <n>' confirmation. " "'MERGE PR <n>' confirmation. "
"Review and merge are separate workflow roles. A reviewer approval is not merge authorization.") "Review and merge are separate workflow roles. A reviewer approval is not merge authorization.")
elif role == "controller":
guidance.append(
"Controller profile: route work via "
"gitea_route_task_session(task_type='process_work_queue') then "
"gitea_allocate_next_work (allocation_mode=cross_role by default). "
"The allocator returns exactly one authoritative selection with "
"required_role / required_profile / selected_action. Do not "
"implement, review, approve, or merge in this session — schedule "
"the matching role namespace instead. Dashboard output is "
"explanatory only and never replaces allocator selection.")
elif role == "mixed": elif role == "mixed":
guidance.append( guidance.append(
"WARNING: this profile allows both authoring and " "WARNING: this profile allows both authoring and "
@@ -15806,7 +15830,8 @@ def gitea_whoami(
"environment": profile.get("environment"), "environment": profile.get("environment"),
"service": profile.get("service"), "service": profile.get("service"),
"identity": profile.get("identity"), "identity": profile.get("identity"),
"role": profile.get("role"), "role": profile.get("role") or _profile_role_kind(profile),
"role_kind": _profile_role_kind(profile),
"profile_address": profile.get("profile_path"), "profile_address": profile.get("profile_path"),
"execution_profile": profile.get("execution_profile"), "execution_profile": profile.get("execution_profile"),
"audit_label": profile.get("audit_label"), "audit_label": profile.get("audit_label"),
@@ -20590,9 +20615,10 @@ def gitea_allocate_next_work(
candidates_json: Any = None, candidates_json: Any = None,
exclude_issue_numbers: list[int] | None = None, exclude_issue_numbers: list[int] | None = None,
expected_candidate_set_fingerprint: str | None = None, expected_candidate_set_fingerprint: str | None = None,
allocation_mode: str | None = None,
limit: int = 50, limit: int = 50,
) -> dict: ) -> dict:
"""Controller-owned next-work allocator using the #613 control-plane DB (#600). """Controller-owned next-work allocator using the #613 control-plane DB (#600/#840).
Workers must not self-select exclusive work under the standard multi-LLM Workers must not self-select exclusive work under the standard multi-LLM
workflow. Call this tool instead. workflow. Call this tool instead.
@@ -20602,6 +20628,14 @@ def gitea_allocate_next_work(
``ControlPlaneDB.assign_and_lease`` (never file locks or comment-only ``ControlPlaneDB.assign_and_lease`` (never file locks or comment-only
leases as the coordination source). leases as the coordination source).
*allocation_mode* (#840): when the active role is controller (or mode is
``cross_role``), inspect the complete queue and return exactly one
authoritative selection with selected item, action, required_role,
required_profile/namespace, pins, and allocation/lease evidence.
Role-scoped workers pass ``role=author|reviewer|merger|reconciler`` (or
omit for profile role) for single-role filtering. Controller routes only
and does not perform downstream mutations.
Outcomes include: ``assigned_work``, ``preview``, ``wait``, Outcomes include: ``assigned_work``, ``preview``, ``wait``,
``blocked_by_terminal_path``, ``no_safe_work``, ``role_ineligible``, ``blocked_by_terminal_path``, ``no_safe_work``, ``role_ineligible``,
``blocked_by_excluded_own_lease``, ``candidate_set_drift``. ``blocked_by_excluded_own_lease``, ``candidate_set_drift``.
@@ -20736,6 +20770,7 @@ def gitea_allocate_next_work(
controller_instance_id=allocator_service.resolve_controller_instance_id(), controller_instance_id=allocator_service.resolve_controller_instance_id(),
exclude_issue_numbers=exclude_issue_numbers, exclude_issue_numbers=exclude_issue_numbers,
expected_candidate_set_fingerprint=expected_candidate_set_fingerprint, expected_candidate_set_fingerprint=expected_candidate_set_fingerprint,
allocation_mode=allocation_mode,
) )
except ValueError as exc: except ValueError as exc:
return { return {
@@ -21113,10 +21148,21 @@ def gitea_adopt_workflow_lease(
remote: str = "dadeschools", remote: str = "dadeschools",
host: str | None = None, host: str | None = None,
) -> dict: ) -> dict:
"""Adopt a control-plane lease through the sanctioned path (#601). """Adopt a control-plane lease through the sanctioned path (#601 / #843).
Same-owner resume refreshes provenance. Foreign active leases are refused. Same-owner resume refreshes provenance. Foreign active leases are refused
Expired leases may be reclaimed; provenance records adopted_from/by. unless the lease is a pending controller cross-role handoff and the caller
holds the required role (independent consume without sharing the
controller session). Expired leases may be reclaimed; provenance records
adopted_from/by. Terminal (abandoned/released) leases cannot be adopted.
#843 F1: the adopter role is derived authoritatively from the active
authenticated profile never from caller input. A supplied ``role`` that
does not exactly match the profile-derived role is rejected (no silent
accept or reinterpretation), and handoff provenance ``required_profile`` /
``required_namespace`` restrictions are validated against the same
authoritative caller context. Caller-supplied role/profile/namespace can
never grant authority.
""" """
read_block = _profile_operation_gate("gitea.read") read_block = _profile_operation_gate("gitea.read")
if read_block: if read_block:
@@ -21125,25 +21171,64 @@ def gitea_adopt_workflow_lease(
"reasons": read_block, "reasons": read_block,
"permission_report": _permission_block_report("gitea.read"), "permission_report": _permission_block_report("gitea.read"),
} }
profile = get_profile()
profile_name = (profile.get("profile_name") or "").strip() or "session"
active_role = (_profile_role_kind(profile) or "").strip().lower()
if not active_role:
return {
"success": False,
"outcome": "blocked",
"mutation_performed": False,
"reasons": [
"active profile role could not be derived authoritatively; "
"refusing lease adoption (fail closed, #843)"
],
"lease_id": lease_id,
"authoritative_source": "control_plane_db",
"file_lock_only": False,
"comment_lease_only": False,
}
if role is not None and str(role).strip():
supplied_role = str(role).strip().lower()
if supplied_role != active_role:
return {
"success": False,
"outcome": "blocked",
"mutation_performed": False,
"profile_role_kind": active_role,
"supplied_role": supplied_role,
"reasons": [
f"caller-supplied role '{supplied_role}' does not match "
f"the authenticated profile-derived role '{active_role}'; "
"caller-supplied role/profile/namespace can never grant "
"authority (fail closed, #843)"
],
"lease_id": lease_id,
"authoritative_source": "control_plane_db",
"file_lock_only": False,
"comment_lease_only": False,
}
db, errs = _control_plane_db_or_error() db, errs = _control_plane_db_or_error()
if db is None: if db is None:
return {"success": False, "reasons": errs} return {"success": False, "reasons": errs}
profile = get_profile()
profile_name = (profile.get("profile_name") or "").strip() or "session"
active_role = _profile_role_kind(profile) or "author"
sid = (session_id or "").strip() or ( sid = (session_id or "").strip() or (
f"{profile_name}-{os.getpid()}-{uuid.uuid4().hex[:8]}" f"{profile_name}-{os.getpid()}-{uuid.uuid4().hex[:8]}"
) )
adopter_namespace = allocator_service.DEFAULT_ROLE_NAMESPACES.get(
active_role, f"gitea-{active_role}"
)
try: try:
return lease_lifecycle.adopt_lease( return lease_lifecycle.adopt_lease(
db, db,
lease_id=lease_id, lease_id=lease_id,
adopter_session_id=sid, adopter_session_id=sid,
role=(role or active_role).strip() or "author", role=active_role,
worktree_path=worktree_path, worktree_path=worktree_path,
expected_head_sha=expected_head_sha, expected_head_sha=expected_head_sha,
owner_pid=os.getpid(), owner_pid=os.getpid(),
operator_authorized=bool(operator_authorized), operator_authorized=bool(operator_authorized),
adopter_profile_name=profile_name,
adopter_namespace=adopter_namespace,
) )
except (lease_lifecycle.LeaseLifecycleError, control_plane_db.ControlPlaneError) as exc: except (lease_lifecycle.LeaseLifecycleError, control_plane_db.ControlPlaneError) as exc:
return { return {
+209 -13
View File
@@ -39,6 +39,7 @@ SAFE_RELEASE_OWNED = "release_owned"
SAFE_STALE_PROMPT = "stale_prompt_lease" SAFE_STALE_PROMPT = "stale_prompt_lease"
SAFE_UNKNOWN = "inspect_only" SAFE_UNKNOWN = "inspect_only"
SAFE_NO_AUTHORITY = "file_or_comment_not_authoritative" SAFE_NO_AUTHORITY = "file_or_comment_not_authoritative"
SAFE_CONSUME_CROSS_ROLE = "consume_cross_role_handoff"
LEASE_STATUS_ACTIVE = "active" LEASE_STATUS_ACTIVE = "active"
LEASE_STATUS_RELEASED = "released" LEASE_STATUS_RELEASED = "released"
@@ -250,6 +251,23 @@ def decide_safe_next_action(
"same_owner": True, "same_owner": True,
"also_allowed": [SAFE_ABANDON_ALLOWED, SAFE_RELEASE_OWNED], "also_allowed": [SAFE_ABANDON_ALLOWED, SAFE_RELEASE_OWNED],
} }
handoff = is_pending_cross_role_handoff({"lease": lease})
if handoff:
return {
"safe_next_action": SAFE_CONSUME_CROSS_ROLE,
"reasons": [
f"controller allocation pending handoff (freshness={status}); "
"required-role worker may consume without abandon/reassign; "
f"required_role={handoff['required_role']}"
],
"block": False,
"same_owner": False,
"owner_session_id": owner,
"required_role": handoff["required_role"],
"cross_role_handoff": True,
"handoff_status": "pending",
"also_allowed": [SAFE_ABANDON_ALLOWED],
}
return { return {
"safe_next_action": SAFE_ABANDON_ALLOWED, "safe_next_action": SAFE_ABANDON_ALLOWED,
"reasons": [ "reasons": [
@@ -272,6 +290,24 @@ def decide_safe_next_action(
} }
if not same_owner and status == "active": if not same_owner and status == "active":
# #843: pending cross-role handoff is consumable by required role
handoff = is_pending_cross_role_handoff({"lease": lease})
if handoff:
return {
"safe_next_action": SAFE_CONSUME_CROSS_ROLE,
"reasons": [
"controller cross-role allocation pending handoff; "
f"required_role={handoff['required_role']}; "
"consume via gitea_adopt_workflow_lease without "
"abandonment or sharing the controller session"
],
"block": False,
"same_owner": False,
"owner_session_id": owner,
"required_role": handoff["required_role"],
"cross_role_handoff": True,
"handoff_status": "pending",
}
return { return {
"safe_next_action": SAFE_WAIT_FOREIGN, "safe_next_action": SAFE_WAIT_FOREIGN,
"reasons": [ "reasons": [
@@ -440,6 +476,84 @@ def list_active_leases(
} }
def parse_lease_provenance(lease_or_state: Mapping[str, Any] | None) -> dict[str, Any]:
"""Return durable lease provenance dict (empty when absent/unparseable)."""
if not lease_or_state:
return {}
if "provenance" in lease_or_state and isinstance(lease_or_state.get("provenance"), dict):
return dict(lease_or_state["provenance"])
raw = None
if "provenance_json" in lease_or_state:
raw = lease_or_state.get("provenance_json")
elif "lease" in lease_or_state and isinstance(lease_or_state.get("lease"), Mapping):
raw = lease_or_state["lease"].get("provenance_json")
if not raw:
return {}
if isinstance(raw, dict):
return dict(raw)
try:
loaded = json.loads(raw)
except (TypeError, json.JSONDecodeError):
return {}
return dict(loaded) if isinstance(loaded, dict) else {}
def is_pending_cross_role_handoff(
state: Mapping[str, Any] | None,
) -> dict[str, Any] | None:
"""Return handoff evidence when a controller allocation awaits consume (#843).
A pending handoff is identified by durable provenance written at
cross-role apply time — not by title heuristics or session-id guessing.
"""
if not state:
return None
lease = state.get("lease") if isinstance(state.get("lease"), Mapping) else state
if not isinstance(lease, Mapping):
return None
status = str(lease.get("status") or "").strip().lower()
if status in (LEASE_STATUS_ABANDONED, LEASE_STATUS_RELEASED, LEASE_STATUS_EXPIRED):
return None
prov = parse_lease_provenance(state)
if not prov and isinstance(lease, Mapping):
prov = parse_lease_provenance(lease)
if not prov.get("cross_role_handoff"):
return None
handoff_status = str(prov.get("handoff_status") or "pending").strip().lower()
if handoff_status != "pending":
return None
adopted_by = (
lease.get("adopted_by_session_id")
or prov.get("adopted_by_session_id")
or ""
)
if str(adopted_by).strip():
return None
required_role = str(
prov.get("required_role") or lease.get("role") or ""
).strip().lower()
if not required_role:
return None
return {
"cross_role_handoff": True,
"handoff_status": "pending",
"required_role": required_role,
"allocating_session_id": str(
prov.get("allocating_session_id") or lease.get("session_id") or ""
),
"allocating_role": str(prov.get("allocating_role") or "controller"),
"lease_id": str(lease.get("lease_id") or ""),
"assignment_id": (
str(state["assignment"]["assignment_id"])
if isinstance(state.get("assignment"), Mapping)
and state["assignment"].get("assignment_id")
else None
),
"provenance": prov,
}
def adopt_lease( def adopt_lease(
db: cpd.ControlPlaneDB, db: cpd.ControlPlaneDB,
*, *,
@@ -450,8 +564,17 @@ def adopt_lease(
expected_head_sha: str | None = None, expected_head_sha: str | None = None,
owner_pid: int | None = None, owner_pid: int | None = None,
operator_authorized: bool = False, operator_authorized: bool = False,
adopter_profile_name: str | None = None,
adopter_namespace: str | None = None,
) -> dict[str, Any]: ) -> dict[str, Any]:
"""Sanctioned adopt path with provenance; never silent foreign steal.""" """Sanctioned adopt path with provenance; never silent foreign steal.
#843 F1: for a pending cross-role handoff, ``role`` must be the
authoritative profile-derived role supplied by the MCP boundary — never
caller-asserted authority. When the handoff provenance declares
``required_profile`` / ``required_namespace`` and the caller context is
provided, both are validated exactly; a mismatch fails closed.
"""
state = db.get_lease_workflow_state(lease_id) state = db.get_lease_workflow_state(lease_id)
if not state: if not state:
raise LeaseLifecycleError( raise LeaseLifecycleError(
@@ -463,11 +586,8 @@ def adopt_lease(
owner = str(lease.get("session_id") or "") owner = str(lease.get("session_id") or "")
same_owner = owner == str(adopter_session_id) same_owner = owner == str(adopter_session_id)
if freshness["freshness"] == "active" and not same_owner: handoff = is_pending_cross_role_handoff(state)
raise LeaseLifecycleError( adopter_role = (role or "").strip().lower()
f"refusing to steal active foreign lease {lease_id} owned by "
f"{owner} (fail closed)"
)
if freshness["freshness"] in ("abandoned", "released"): if freshness["freshness"] in ("abandoned", "released"):
raise LeaseLifecycleError( raise LeaseLifecycleError(
@@ -475,13 +595,64 @@ def adopt_lease(
"(fail closed)" "(fail closed)"
) )
# Expired or stale: require abandon-style safety before ownership transfer if handoff and not same_owner:
# when not same owner; same owner may reclaim. # Terminal statuses already rejected above. Freshness may be
if not same_owner and freshness["freshness"] in ( # active OR stale_dead_process (controller exited) — both are
# consumable without abandonment when handoff is still pending.
if freshness["freshness"] not in (
"active",
"stale_dead_process",
"stale_missing_worktree",
):
raise LeaseLifecycleError(
f"lease {lease_id} freshness={freshness['freshness']}; "
"terminal or non-active allocation cannot be handoff-consumed "
"(fail closed)"
)
required = handoff["required_role"]
if adopter_role != required:
raise LeaseLifecycleError(
f"wrong role for cross-role handoff consume of {lease_id}: "
f"required={required} adopter={adopter_role or 'none'} "
"(fail closed)"
)
# #843 F1: provenance profile/namespace restrictions are validated
# against the authoritative caller context when declared. Caller
# input can never widen authority; a mismatch fails closed.
handoff_prov = handoff.get("provenance") or {}
required_profile = str(
handoff_prov.get("required_profile") or ""
).strip()
if required_profile and adopter_profile_name is not None:
if str(adopter_profile_name).strip() != required_profile:
raise LeaseLifecycleError(
f"wrong profile for cross-role handoff consume of "
f"{lease_id}: required_profile={required_profile} "
f"adopter_profile={adopter_profile_name} (fail closed)"
)
required_namespace = str(
handoff_prov.get("required_namespace") or ""
).strip()
if required_namespace and adopter_namespace is not None:
if str(adopter_namespace).strip() != required_namespace:
raise LeaseLifecycleError(
f"wrong namespace for cross-role handoff consume of "
f"{lease_id}: required_namespace={required_namespace} "
f"adopter_namespace={adopter_namespace} (fail closed)"
)
reason = "cross-role-handoff-consume"
elif freshness["freshness"] == "active" and not same_owner:
raise LeaseLifecycleError(
f"refusing to steal active foreign lease {lease_id} owned by "
f"{owner} (fail closed)"
)
elif not same_owner and freshness["freshness"] in (
"expired", "expired",
"stale_dead_process", "stale_dead_process",
"stale_missing_worktree", "stale_missing_worktree",
): ):
# Expired or stale (non-handoff): require abandon-style safety before
# ownership transfer when not same owner; same owner may reclaim.
if not operator_authorized and freshness["freshness"] == "expired": if not operator_authorized and freshness["freshness"] == "expired":
# Deterministic reclaim of expired foreign lease is allowed # Deterministic reclaim of expired foreign lease is allowed
# without operator flag (sanctioned expire reclaim). # without operator flag (sanctioned expire reclaim).
@@ -492,6 +663,9 @@ def adopt_lease(
f"lease {lease_id} freshness={freshness['freshness']}; " f"lease {lease_id} freshness={freshness['freshness']}; "
"use abandon with proof before foreign adopt (fail closed)" "use abandon with proof before foreign adopt (fail closed)"
) )
reason = "sanctioned-reclaim-adopt"
else:
reason = "owner-resume-adopt" if same_owner else "sanctioned-reclaim-adopt"
provenance = build_adopt_provenance( provenance = build_adopt_provenance(
adopted_from_session_id=owner, adopted_from_session_id=owner,
@@ -504,10 +678,14 @@ def adopt_lease(
worktree_path=worktree_path, worktree_path=worktree_path,
expected_head_sha=expected_head_sha or lease.get("expected_head_sha"), expected_head_sha=expected_head_sha or lease.get("expected_head_sha"),
prior_lease_id=lease_id, prior_lease_id=lease_id,
reason=( reason=reason,
"owner-resume-adopt" if same_owner else "sanctioned-reclaim-adopt"
),
) )
if handoff and not same_owner:
provenance["cross_role_handoff"] = True
provenance["handoff_status"] = "adopted"
provenance["required_role"] = handoff["required_role"]
provenance["allocating_session_id"] = handoff["allocating_session_id"]
provenance["allocating_role"] = handoff["allocating_role"]
result = db.adopt_lease( result = db.adopt_lease(
lease_id=lease_id, lease_id=lease_id,
@@ -518,7 +696,7 @@ def adopt_lease(
owner_pid=owner_pid if owner_pid is not None else os.getpid(), owner_pid=owner_pid if owner_pid is not None else os.getpid(),
provenance=provenance, provenance=provenance,
) )
return { out = {
"success": True, "success": True,
"outcome": result.get("outcome"), "outcome": result.get("outcome"),
"same_owner": same_owner, "same_owner": same_owner,
@@ -531,6 +709,24 @@ def adopt_lease(
"comment_lease_only": False, "comment_lease_only": False,
"reasons": result.get("reasons") or [], "reasons": result.get("reasons") or [],
} }
if handoff and not same_owner:
out["cross_role_handoff"] = True
out["handoff_status"] = "adopted"
out["required_role"] = handoff["required_role"]
out["adopted_by_session_id"] = adopter_session_id
out["adopted_from_session_id"] = owner
lease_row = result.get("lease") or {}
if isinstance(lease_row, Mapping):
out["read_after_write"] = {
"lease_id": lease_row.get("lease_id"),
"session_id": lease_row.get("session_id"),
"role": lease_row.get("role"),
"status": lease_row.get("status"),
"adopted_by_session_id": lease_row.get("adopted_by_session_id"),
"adopted_from_session_id": lease_row.get("adopted_from_session_id"),
"phase": lease_row.get("phase"),
}
return out
def release_lease( def release_lease(
+19 -5
View File
@@ -24,7 +24,12 @@ ROLE_WORKTREE_ENVS: dict[str, str] = {
"reconciler": RECONCILER_WORKTREE_ENV, "reconciler": RECONCILER_WORKTREE_ENV,
} }
NON_AUTHOR_ROLES = frozenset({"reviewer", "merger", "reconciler"}) # Controller has no task worktree env — it routes only (#840).
KNOWN_ROLE_KINDS = frozenset(
{"author", "reviewer", "merger", "reconciler", "controller"}
)
NON_AUTHOR_ROLES = frozenset({"reviewer", "merger", "reconciler", "controller"})
def normalize_role_kind( def normalize_role_kind(
@@ -37,8 +42,12 @@ def normalize_role_kind(
profile = (profile_name or "").strip().lower() profile = (profile_name or "").strip().lower()
if role == "reviewer" and "merger" in profile: if role == "reviewer" and "merger" in profile:
return "merger" return "merger"
if "controller" in profile or role == "controller":
return "controller"
if role in ROLE_WORKTREE_ENVS: if role in ROLE_WORKTREE_ENVS:
return role return role
if role in KNOWN_ROLE_KINDS:
return role
return "author" return "author"
@@ -80,7 +89,7 @@ def resolve_namespace_workspace(
""" """
env_map = env if env is not None else os.environ env_map = env if env is not None else os.environ
role = normalize_role_kind(role_kind, profile_name=profile_name) role = normalize_role_kind(role_kind, profile_name=profile_name)
role_env_key = ROLE_WORKTREE_ENVS[role] role_env_key = ROLE_WORKTREE_ENVS.get(role)
# #618: durable author resolution — no silent control/master fallback. # #618: durable author resolution — no silent control/master fallback.
if role == "author" and verify_paths: if role == "author" and verify_paths:
@@ -108,13 +117,17 @@ def resolve_namespace_workspace(
) )
return workspace, source return workspace, source
role_env_candidate = (
(_env_value(env_map, role_env_key), f"{role_env_key} environment variable", True)
if role_env_key
else (None, "no role worktree env", True)
)
for candidate, source, env_sourced in ( for candidate, source, env_sourced in (
(worktree_path, "worktree_path argument", False), (worktree_path, "worktree_path argument", False),
(worktree, "worktree argument", False), (worktree, "worktree argument", False),
(_env_value(env_map, ACTIVE_WORKTREE_ENV), (_env_value(env_map, ACTIVE_WORKTREE_ENV),
f"{ACTIVE_WORKTREE_ENV} environment variable", True), f"{ACTIVE_WORKTREE_ENV} environment variable", True),
(_env_value(env_map, role_env_key), role_env_candidate,
f"{role_env_key} environment variable", True),
(session_lease_worktree if role in {"reviewer", "merger"} else None, (session_lease_worktree if role in {"reviewer", "merger"} else None,
"reviewer PR lease worktree", False), "reviewer PR lease worktree", False),
# Author lock derivation is handled by the durable path above when # Author lock derivation is handled by the durable path above when
@@ -433,7 +446,8 @@ def assess_namespace_mutation_workspace(
reasons.append( reasons.append(
f"{role} mutation blocked: workspace is the stable control checkout; " f"{role} mutation blocked: workspace is the stable control checkout; "
f"create or reconnect to a session-owned worktree under branches/ " f"create or reconnect to a session-owned worktree under branches/ "
f"or set {ROLE_WORKTREE_ENVS[role]} / {ACTIVE_WORKTREE_ENV}" f"or set {ROLE_WORKTREE_ENVS.get(role, ACTIVE_WORKTREE_ENV)} / "
f"{ACTIVE_WORKTREE_ENV}"
) )
elif ( elif (
role in {"reviewer", "merger"} role in {"reviewer", "merger"}
+35
View File
@@ -81,6 +81,12 @@ AUTHOR_TASKS = frozenset({
"reconcile_landed_pr", "reconcile_landed_pr",
}) })
CONTROLLER_TASKS = frozenset({
"process_work_queue",
"process-work-queue",
"cross_role_allocate",
})
RECONCILER_TASKS = frozenset({ RECONCILER_TASKS = frozenset({
"cleanup_merged_pr_branch", "cleanup_merged_pr_branch",
# #729: delete_branch is reconciler-owned (gitea.branch.delete is granted # #729: delete_branch is reconciler-owned (gitea.branch.delete is granted
@@ -132,6 +138,10 @@ TASK_REQUIRED_ROLE = {
"reconcile_close_superseded_pr": "reconciler", "reconcile_close_superseded_pr": "reconciler",
"reconcile_close_satisfied_issue": "reconciler", "reconcile_close_satisfied_issue": "reconciler",
"reconcile_create_followup_issue": "reconciler", "reconcile_create_followup_issue": "reconciler",
# #840: controller-owned generic queue allocation / routing.
"process_work_queue": "controller",
"process-work-queue": "controller",
"cross_role_allocate": "controller",
} }
WRONG_ROLE_REVIEWER_MSG = ( WRONG_ROLE_REVIEWER_MSG = (
@@ -147,6 +157,10 @@ WRONG_ROLE_MERGER_MSG = (
"Wrong role/session for merger task. Launch merger MCP namespace." "Wrong role/session for merger task. Launch merger MCP namespace."
) )
WRONG_ROLE_CONTROLLER_MSG = (
"Wrong role/session for controller task. Launch controller MCP namespace."
)
_session_last_route: dict | None = None _session_last_route: dict | None = None
@@ -281,6 +295,27 @@ def route_task_session(
_record_route(result) _record_route(result)
return result return result
if required_role == "controller":
result = {
"task_type": task_type,
"required_role": required_role,
"active_role": active_role_kind,
"active_profile": active_profile,
"route_result": ROUTE_WRONG_ROLE,
"downstream_allowed": False,
"reasons": [
WRONG_ROLE_CONTROLLER_MSG,
"Controller tasks (process_work_queue / cross-role allocate) "
"cannot run in author, reviewer, merger, or reconciler "
"worker sessions.",
],
"message": WRONG_ROLE_CONTROLLER_MSG,
"runtime_switching_supported": runtime_switching_supported,
"profile_switch_blocked": not runtime_switching_supported,
}
_record_route(result)
return result
if required_role == "author": if required_role == "author":
route = ROUTE_TO_AUTHOR route = ROUTE_TO_AUTHOR
message = ( message = (
+17 -2
View File
@@ -309,8 +309,10 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.pr.create", "permission": "gitea.pr.create",
"role": "author", "role": "author",
}, },
# #600: controller-owned allocator — any authenticated profile may call; # #600: workers and controller may call with gitea.read; role-scoped workers
# routing enforces role match to selected work. Uses control-plane DB (#613). # pass role=author|reviewer|merger|reconciler. Cross-role routing is the
# controller default (#840). The canonical generic queue *task type* is
# process_work_queue (controller-only below).
"allocate_next_work": { "allocate_next_work": {
"permission": "gitea.read", "permission": "gitea.read",
"role": "author", "role": "author",
@@ -319,6 +321,19 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.read", "permission": "gitea.read",
"role": "author", "role": "author",
}, },
# #840: documented generic queue task — controller routes only.
"process_work_queue": {
"permission": "gitea.read",
"role": "controller",
},
"process-work-queue": {
"permission": "gitea.read",
"role": "controller",
},
"cross_role_allocate": {
"permission": "gitea.read",
"role": "controller",
},
# #601 first-class lease lifecycle — inspect/list need read; mutations gate on # #601 first-class lease lifecycle — inspect/list need read; mutations gate on
# ownership in the control-plane DB (not a separate Gitea write permission). # ownership in the control-plane DB (not a separate Gitea write permission).
+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()
+682
View File
@@ -0,0 +1,682 @@
"""Cross-role allocation handoff consumable by independent workers (#843).
Regression coverage for the controllerrequired-role consume path:
* controller allocates author work; independent author adopts successfully
* author adoption succeeds after allocating controller process exits
* author adoption without sharing controller session identity
* wrong-role adoption rejected
* concurrent/second adoption rejected without state corruption
* terminal allocation adoption rejected
* successful adoption produces authoritative ownership evidence
* genuine abandoned-lease recovery remains valid
* process_work_queue / allocate results include consume identifiers
* same-role allocation behavior remains compatible
"""
from __future__ import annotations
import os
import tempfile
import unittest
from datetime import timedelta
from unittest.mock import patch
from allocator_service import (
ALLOCATION_MODE_CROSS_ROLE,
ALLOCATION_MODE_ROLE_SCOPED,
OUTCOME_ASSIGNED,
ROLE_AUTHOR,
ROLE_CONTROLLER,
ROLE_REVIEWER,
WorkCandidate,
allocate_next_work,
)
from control_plane_db import ControlPlaneDB, ForeignLeaseError, _ts, _utc_now
import lease_lifecycle as ll
class CrossRoleHandoffTest(unittest.TestCase):
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.db_path = os.path.join(self._tmp.name, "cp.sqlite3")
self.db = ControlPlaneDB(self.db_path)
self.db.upsert_session(
session_id="ctrl-session",
role="controller",
profile="prgs-controller",
pid=99999999, # dead-looking pid
)
self.db.upsert_session(
session_id="author-worker",
role="author",
profile="prgs-author",
pid=os.getpid(),
)
self.db.upsert_session(
session_id="author-worker-2",
role="author",
profile="prgs-author",
pid=os.getpid(),
)
self.db.upsert_session(
session_id="reviewer-worker",
role="reviewer",
profile="prgs-reviewer",
pid=os.getpid(),
)
self.wt = self._tmp.name
def tearDown(self) -> None:
self._tmp.cleanup()
def _ready_issue(self, number: int = 843, title: str = "handoff target") -> WorkCandidate:
return WorkCandidate(
kind="issue",
number=number,
labels=("status:ready", "type:bug"),
title=title,
priority=20,
)
def _controller_allocate(self, number: int = 843, **kwargs):
defaults = dict(
db=self.db,
session_id="ctrl-session",
role=ROLE_CONTROLLER,
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
candidates=[self._ready_issue(number)],
apply=True,
profile_name="prgs-controller",
username="controller-user",
allocation_mode=ALLOCATION_MODE_CROSS_ROLE,
)
defaults.update(kwargs)
return allocate_next_work(**defaults)
def test_controller_allocates_author_independent_author_adopts(self) -> None:
res = self._controller_allocate()
self.assertEqual(res["outcome"], OUTCOME_ASSIGNED)
self.assertEqual(res["required_role"], ROLE_AUTHOR)
self.assertIn("consume_allocation", res)
consume = res["consume_allocation"]
self.assertEqual(consume["tool"], "gitea_adopt_workflow_lease")
self.assertEqual(consume["required_role"], ROLE_AUTHOR)
self.assertFalse(consume["controller_session_required"])
lid = res["assignment"]["lease_id"]
self.assertEqual(consume["lease_id"], lid)
self.assertIn(lid, res["next_valid_command"])
adopted = ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="author-worker",
role=ROLE_AUTHOR,
worktree_path=self.wt,
)
self.assertTrue(adopted["success"])
self.assertEqual(adopted["outcome"], "adopted_cross_role_handoff")
self.assertEqual(adopted["adopted_by_session_id"], "author-worker")
self.assertEqual(adopted["adopted_from_session_id"], "ctrl-session")
raw = adopted["read_after_write"]
self.assertEqual(raw["session_id"], "author-worker")
self.assertEqual(raw["adopted_by_session_id"], "author-worker")
self.assertEqual(raw["status"], "active")
self.assertEqual(raw["phase"], "adopted")
# Authoritative re-read
state = self.db.get_lease_workflow_state(lid)
self.assertEqual(state["lease"]["session_id"], "author-worker")
self.assertEqual(state["lease"]["adopted_by_session_id"], "author-worker")
self.assertEqual(state["assignment"]["session_id"], "author-worker")
self.assertEqual(state["provenance"]["handoff_status"], "adopted")
def test_author_adoption_after_controller_process_exits(self) -> None:
res = self._controller_allocate(number=900)
lid = res["assignment"]["lease_id"]
# Force owner_pid dead + freshness stale_dead_process
import sqlite3
conn = sqlite3.connect(self.db_path)
try:
conn.execute(
"UPDATE leases SET owner_pid = 99999999 WHERE lease_id = ?",
(lid,),
)
conn.commit()
finally:
conn.close()
state = self.db.get_lease_workflow_state(lid)
fr = ll.classify_lease_freshness(
state["lease"], pid_checker=lambda _p: False
)
self.assertEqual(fr["freshness"], "stale_dead_process")
adopted = ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="author-worker",
role=ROLE_AUTHOR,
worktree_path=self.wt,
)
self.assertEqual(adopted["outcome"], "adopted_cross_role_handoff")
self.assertEqual(adopted["adopted_by_session_id"], "author-worker")
# No abandon required
state2 = self.db.get_lease_workflow_state(lid)
self.assertEqual(state2["lease"]["status"], "active")
self.assertNotEqual(state2["lease"]["status"], "abandoned")
def test_adoption_without_sharing_controller_session_identity(self) -> None:
res = self._controller_allocate(number=901)
lid = res["assignment"]["lease_id"]
adopted = ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="author-worker",
role=ROLE_AUTHOR,
worktree_path=self.wt,
)
self.assertNotEqual(adopted["adopted_by_session_id"], "ctrl-session")
self.assertFalse(adopted["same_owner"])
self.assertEqual(adopted["adopted_from_session_id"], "ctrl-session")
def test_wrong_role_adoption_rejected(self) -> None:
res = self._controller_allocate(number=902)
lid = res["assignment"]["lease_id"]
with self.assertRaises(ll.LeaseLifecycleError) as ctx:
ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="reviewer-worker",
role=ROLE_REVIEWER,
)
self.assertIn("wrong role", str(ctx.exception).lower())
# State unchanged
state = self.db.get_lease_workflow_state(lid)
self.assertEqual(state["lease"]["session_id"], "ctrl-session")
self.assertIsNone(state["lease"].get("adopted_by_session_id") or None)
self.assertEqual(state["provenance"]["handoff_status"], "pending")
def test_second_adoption_rejected_without_corruption(self) -> None:
res = self._controller_allocate(number=903)
lid = res["assignment"]["lease_id"]
first = ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="author-worker",
role=ROLE_AUTHOR,
worktree_path=self.wt,
)
self.assertEqual(first["outcome"], "adopted_cross_role_handoff")
with self.assertRaises(ll.LeaseLifecycleError):
ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="author-worker-2",
role=ROLE_AUTHOR,
worktree_path=self.wt,
)
state = self.db.get_lease_workflow_state(lid)
self.assertEqual(state["lease"]["session_id"], "author-worker")
self.assertEqual(state["lease"]["adopted_by_session_id"], "author-worker")
self.assertEqual(state["assignment"]["session_id"], "author-worker")
self.assertEqual(state["lease"]["status"], "active")
def test_terminal_allocation_adoption_rejected(self) -> None:
res = self._controller_allocate(number=904)
lid = res["assignment"]["lease_id"]
# Abandon as terminal
proof = ll.AbandonProof(
dead_process=True,
missing_worktree=True,
no_open_pr=True,
no_live_mutation_risk=True,
owner_pid=99999999,
worktree_path="/nonexistent/for-843",
)
# Attach dead pid / missing wt for abandon eligibility
import sqlite3
conn = sqlite3.connect(self.db_path)
try:
conn.execute(
"UPDATE leases SET owner_pid = 99999999, worktree_path = ? WHERE lease_id = ?",
("/nonexistent/for-843", lid),
)
conn.commit()
finally:
conn.close()
abandoned = ll.abandon_lease(
self.db,
lease_id=lid,
requester_session_id="author-worker",
proof=proof,
)
self.assertEqual(abandoned["outcome"], "abandoned")
with self.assertRaises(ll.LeaseLifecycleError) as ctx:
ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="author-worker",
role=ROLE_AUTHOR,
)
self.assertIn("abandoned", str(ctx.exception).lower())
def test_successful_adoption_read_after_write_ownership(self) -> None:
res = self._controller_allocate(number=905)
lid = res["assignment"]["lease_id"]
adopted = ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="author-worker",
role=ROLE_AUTHOR,
worktree_path=self.wt,
)
raw = adopted["read_after_write"]
self.assertEqual(raw["lease_id"], lid)
self.assertEqual(raw["session_id"], "author-worker")
self.assertEqual(raw["adopted_by_session_id"], "author-worker")
self.assertEqual(raw["adopted_from_session_id"], "ctrl-session")
# Re-fetch proves durable write
state = self.db.get_lease_workflow_state(lid)
self.assertEqual(state["lease"]["session_id"], raw["session_id"])
self.assertEqual(
state["lease"]["adopted_by_session_id"], raw["adopted_by_session_id"]
)
def test_genuine_abandoned_recovery_still_valid(self) -> None:
"""Same-role author lease abandoned remains reclaimable via abandon path."""
same = allocate_next_work(
self.db,
session_id="author-worker",
role=ROLE_AUTHOR,
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
candidates=[self._ready_issue(906, "same-role")],
apply=True,
profile_name="prgs-author",
username="author-user",
allocation_mode=ALLOCATION_MODE_ROLE_SCOPED,
)
self.assertEqual(same["outcome"], OUTCOME_ASSIGNED)
lid = same["assignment"]["lease_id"]
import sqlite3
conn = sqlite3.connect(self.db_path)
try:
conn.execute(
"UPDATE leases SET owner_pid = 99999999, worktree_path = ? WHERE lease_id = ?",
("/nonexistent/same-role", lid),
)
conn.commit()
finally:
conn.close()
proof = ll.AbandonProof(
dead_process=True,
missing_worktree=True,
no_open_pr=True,
no_live_mutation_risk=True,
owner_pid=99999999,
worktree_path="/nonexistent/same-role",
)
abandoned = ll.abandon_lease(
self.db,
lease_id=lid,
requester_session_id="author-worker-2",
proof=proof,
)
self.assertEqual(abandoned["outcome"], "abandoned")
# Foreign author cannot handoff-consume an abandoned non-handoff lease
with self.assertRaises(ll.LeaseLifecycleError):
ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="author-worker-2",
role=ROLE_AUTHOR,
)
# Reclaim path still works for expired/abandoned after force-expire
reclaimed = ll.reclaim_expired_lease(
self.db,
lease_id=lid,
session_id="author-worker-2",
role=ROLE_AUTHOR,
worktree_path=self.wt,
)
self.assertEqual(reclaimed["outcome"], "reclaimed")
self.assertEqual(reclaimed["assignment"]["session_id"], "author-worker-2")
def test_allocate_payload_includes_consume_identifiers(self) -> None:
res = self._controller_allocate(number=907)
self.assertIn("consume_allocation", res)
c = res["consume_allocation"]
for key in (
"tool",
"lease_id",
"assignment_id",
"required_role",
"required_profile",
"required_namespace",
"instructions",
"handoff_status",
):
self.assertIn(key, c)
self.assertEqual(c["required_namespace"], "gitea-author")
self.assertEqual(c["required_profile"], "prgs-author")
self.assertIn("gitea_adopt_workflow_lease", c["instructions"])
self.assertTrue(res["lease_proof"]["cross_role_handoff"])
self.assertEqual(res["lease_proof"]["handoff_status"], "pending")
def test_same_role_allocation_remains_compatible(self) -> None:
res = allocate_next_work(
self.db,
session_id="author-worker",
role=ROLE_AUTHOR,
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
candidates=[self._ready_issue(908)],
apply=True,
profile_name="prgs-author",
username="author-user",
)
self.assertEqual(res["outcome"], OUTCOME_ASSIGNED)
self.assertNotIn("consume_allocation", res)
lid = res["assignment"]["lease_id"]
state = self.db.get_lease_workflow_state(lid)
# No cross-role handoff provenance
prov = state.get("provenance") or {}
self.assertFalse(prov.get("cross_role_handoff"))
# Owner resume still works
resume = ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="author-worker",
role=ROLE_AUTHOR,
worktree_path=self.wt,
)
self.assertTrue(resume["same_owner"])
self.assertEqual(resume["outcome"], "adopted_owner_resume")
def test_inspect_points_required_role_at_consume(self) -> None:
res = self._controller_allocate(number=909)
lid = res["assignment"]["lease_id"]
decision = ll.inspect_lease(
self.db, lid, caller_session_id="author-worker"
)
self.assertEqual(
decision["safe_next_action"], ll.SAFE_CONSUME_CROSS_ROLE
)
self.assertFalse(decision["block"])
self.assertEqual(decision["required_role"], ROLE_AUTHOR)
def test_db_cas_rejects_concurrent_second_consume(self) -> None:
res = self._controller_allocate(number=910)
lid = res["assignment"]["lease_id"]
# First consume via DB layer directly
first = self.db.adopt_lease(
lease_id=lid,
adopter_session_id="author-worker",
role=ROLE_AUTHOR,
worktree_path=self.wt,
provenance={
"cross_role_handoff": True,
"handoff_status": "adopted",
"required_role": "author",
},
)
self.assertEqual(first["outcome"], "adopted_cross_role_handoff")
# Second CAS must fail
with self.assertRaises(ForeignLeaseError):
self.db.adopt_lease(
lease_id=lid,
adopter_session_id="author-worker-2",
role=ROLE_AUTHOR,
worktree_path=self.wt,
provenance={
"cross_role_handoff": True,
"handoff_status": "pending",
"required_role": "author",
},
)
state = self.db.get_lease_workflow_state(lid)
self.assertEqual(state["lease"]["session_id"], "author-worker")
class MCPBoundaryAdoptRoleBindingTest(unittest.TestCase):
"""#843 F1: MCP-boundary role binding for ``gitea_adopt_workflow_lease``.
The library-level wrong-role test calls ``lease_lifecycle.adopt_lease``
directly. These tests prove the MCP entry point derives the adopter role
authoritatively from the active authenticated profile and rejects any
caller-supplied role that disagrees, so a reviewer/merger profile cannot
consume an author handoff by passing ``role="author"``.
"""
AUTHOR_PROFILE = {
"profile_name": "prgs-author",
"role": "author",
"allowed_operations": [
"gitea.read",
"gitea.pr.create",
"gitea.branch.push",
],
"forbidden_operations": [],
}
REVIEWER_PROFILE = {
"profile_name": "prgs-reviewer",
"role": "reviewer",
"allowed_operations": [
"gitea.read",
"gitea.pr.review",
"gitea.pr.approve",
"gitea.pr.request_changes",
],
"forbidden_operations": ["gitea.pr.create", "gitea.branch.push"],
}
MERGER_PROFILE = {
"profile_name": "prgs-merger",
"role": "merger",
"allowed_operations": ["gitea.read", "gitea.pr.merge"],
"forbidden_operations": ["gitea.pr.create", "gitea.branch.push"],
}
FOREIGN_AUTHOR_PROFILE = {
"profile_name": "dadeschools-author",
"role": "author",
"allowed_operations": [
"gitea.read",
"gitea.pr.create",
"gitea.branch.push",
],
"forbidden_operations": [],
}
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.db_path = os.path.join(self._tmp.name, "cp.sqlite3")
self.db = ControlPlaneDB(self.db_path)
self.db.upsert_session(
session_id="ctrl-session",
role="controller",
profile="prgs-controller",
pid=99999999,
)
self.db.upsert_session(
session_id="author-worker",
role="author",
profile="prgs-author",
pid=os.getpid(),
)
self.wt = self._tmp.name
def tearDown(self) -> None:
self._tmp.cleanup()
def _ready_issue(self, number: int) -> WorkCandidate:
return WorkCandidate(
kind="issue",
number=number,
labels=("status:ready", "type:bug"),
title="handoff target",
priority=20,
)
def _handoff_lease(self, number: int = 843) -> str:
res = allocate_next_work(
db=self.db,
session_id="ctrl-session",
role=ROLE_CONTROLLER,
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
candidates=[self._ready_issue(number)],
apply=True,
profile_name="prgs-controller",
username="controller-user",
allocation_mode=ALLOCATION_MODE_CROSS_ROLE,
)
self.assertEqual(res["outcome"], OUTCOME_ASSIGNED)
self.assertEqual(res["required_role"], ROLE_AUTHOR)
return res["assignment"]["lease_id"]
def _call_adopt_tool(self, profile: dict, **kwargs):
import gitea_mcp_server as mcp_server
with (
patch.object(mcp_server, "get_profile", return_value=profile),
patch.object(
mcp_server,
"_control_plane_db_or_error",
return_value=(self.db, []),
),
):
return mcp_server.gitea_adopt_workflow_lease(
remote="prgs", **kwargs
)
def _assert_handoff_untouched(self, lease_id: str) -> None:
state = self.db.get_lease_workflow_state(lease_id)
self.assertEqual(state["lease"]["session_id"], "ctrl-session")
self.assertIsNone(state["lease"].get("adopted_by_session_id") or None)
self.assertEqual(state["lease"]["status"], "active")
self.assertEqual(state["provenance"]["handoff_status"], "pending")
def test_reviewer_profile_cannot_consume_author_handoff_via_role_author(
self,
) -> None:
lid = self._handoff_lease(920)
result = self._call_adopt_tool(
self.REVIEWER_PROFILE,
lease_id=lid,
session_id="reviewer-worker",
role="author",
worktree_path=self.wt,
)
self.assertFalse(result["success"])
self.assertEqual(result["outcome"], "blocked")
self.assertEqual(result["profile_role_kind"], "reviewer")
self.assertEqual(result["supplied_role"], "author")
self.assertIn("does not match", result["reasons"][0])
self._assert_handoff_untouched(lid)
def test_merger_profile_cannot_consume_author_handoff_via_role_author(
self,
) -> None:
lid = self._handoff_lease(921)
result = self._call_adopt_tool(
self.MERGER_PROFILE,
lease_id=lid,
session_id="merger-worker",
role="author",
worktree_path=self.wt,
)
self.assertFalse(result["success"])
self.assertEqual(result["outcome"], "blocked")
self.assertEqual(result["profile_role_kind"], "merger")
self._assert_handoff_untouched(lid)
def test_reviewer_profile_rejected_without_role_argument(self) -> None:
"""Even without a spoofed role, the profile-derived role binds."""
lid = self._handoff_lease(922)
result = self._call_adopt_tool(
self.REVIEWER_PROFILE,
lease_id=lid,
session_id="reviewer-worker",
worktree_path=self.wt,
)
self.assertFalse(result["success"])
self.assertEqual(result["outcome"], "blocked")
self.assertIn("wrong role", result["reasons"][0].lower())
self._assert_handoff_untouched(lid)
def test_author_profile_mismatching_supplied_role_rejected(self) -> None:
lid = self._handoff_lease(923)
result = self._call_adopt_tool(
self.AUTHOR_PROFILE,
lease_id=lid,
session_id="author-worker",
role="reviewer",
worktree_path=self.wt,
)
self.assertFalse(result["success"])
self.assertEqual(result["outcome"], "blocked")
self.assertEqual(result["profile_role_kind"], "author")
self.assertEqual(result["supplied_role"], "reviewer")
self._assert_handoff_untouched(lid)
def test_foreign_profile_name_rejected_for_author_handoff(self) -> None:
"""Provenance required_profile binds even when the role matches."""
lid = self._handoff_lease(924)
result = self._call_adopt_tool(
self.FOREIGN_AUTHOR_PROFILE,
lease_id=lid,
session_id="foreign-author-worker",
worktree_path=self.wt,
)
self.assertFalse(result["success"])
self.assertEqual(result["outcome"], "blocked")
self.assertIn("wrong profile", result["reasons"][0].lower())
self._assert_handoff_untouched(lid)
def test_author_profile_consumes_author_handoff(self) -> None:
lid = self._handoff_lease(925)
result = self._call_adopt_tool(
self.AUTHOR_PROFILE,
lease_id=lid,
session_id="author-worker",
role="author",
worktree_path=self.wt,
)
self.assertTrue(result["success"])
self.assertEqual(result["outcome"], "adopted_cross_role_handoff")
self.assertEqual(result["adopted_by_session_id"], "author-worker")
self.assertEqual(result["adopted_from_session_id"], "ctrl-session")
state = self.db.get_lease_workflow_state(lid)
self.assertEqual(state["lease"]["session_id"], "author-worker")
self.assertEqual(
state["lease"]["adopted_by_session_id"], "author-worker"
)
self.assertEqual(state["assignment"]["session_id"], "author-worker")
self.assertEqual(state["provenance"]["handoff_status"], "adopted")
def test_author_profile_consumes_author_handoff_without_role_argument(
self,
) -> None:
lid = self._handoff_lease(926)
result = self._call_adopt_tool(
self.AUTHOR_PROFILE,
lease_id=lid,
session_id="author-worker",
worktree_path=self.wt,
)
self.assertTrue(result["success"])
self.assertEqual(result["outcome"], "adopted_cross_role_handoff")
state = self.db.get_lease_workflow_state(lid)
self.assertEqual(state["lease"]["session_id"], "author-worker")
self.assertEqual(state["lease"]["role"], "author")
if __name__ == "__main__":
unittest.main()
+499
View File
@@ -0,0 +1,499 @@
"""Tests for the read-only system-health API (#634).
Covers the acceptance criteria directly: a structured payload with readiness
and a dependency list (AC1), version and uptime when knowable (AC2), stale
runtime reported without a false mutation-safe claim (AC3), and the healthy /
degraded-dependency / redaction cases (AC4).
"""
import json
import os
import sqlite3
import sys
import tempfile
import unittest
from pathlib import Path
from unittest import mock
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from starlette.testclient import TestClient
import control_plane_db
from webui.app import create_app
from webui.deployment_boundary import scan_text_for_client_secrets
from webui.system_health import (
API_PATH,
STATUS_DEGRADED,
STATUS_DOWN,
STATUS_OK,
STATUS_SKIPPED,
DependencyProbe,
StaleRuntime,
assess_stale_runtime,
clear_probe_cache,
load_system_health,
namespace_summaries,
probe_control_plane_db,
probe_gitea,
process_uptime,
redact,
redact_url,
snapshot_to_dict,
)
def _probe(name, status, *, required=True, detail="detail", kind="test"):
return DependencyProbe(
name=name,
kind=kind,
status=status,
detail=detail,
required=required,
latency_ms=1.5,
metadata={},
)
_ALL_HEALTHY = (
_probe("control_plane_db", STATUS_OK, kind="sqlite"),
_probe("repository", STATUS_OK, kind="git"),
_probe("gitea", STATUS_OK, required=False, kind="http"),
)
_CLEAN_PARITY = StaleRuntime(
daemon_head="abc123",
checkout_head="abc123",
remote_head="abc123",
stale=False,
determinable=True,
mutation_safe=True,
reasons=(),
)
class CleanParityMixin:
"""Pin parity for tests about aggregation rather than staleness.
Without this the assertions depend on the real checkout: a worktree whose
branch is ahead of its upstream is genuinely stale, which would degrade the
overall status and make these cases fail for an unrelated reason.
"""
def setUp(self):
super().setUp()
patcher = mock.patch(
"webui.system_health.assess_stale_runtime",
return_value=_CLEAN_PARITY,
)
patcher.start()
self.addCleanup(patcher.stop)
class TestDependencyAggregation(CleanParityMixin, unittest.TestCase):
"""AC1 — readiness and dependency list derived from probe results."""
def test_all_healthy_is_ok_and_ready(self):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="abc123")
self.assertEqual(snapshot.status, STATUS_OK)
self.assertTrue(snapshot.ready)
self.assertTrue(snapshot.readiness_complete)
self.assertEqual(snapshot.readiness_reasons, ())
self.assertEqual(len(snapshot.dependencies), 3)
def test_required_dependency_down_blocks_readiness(self):
probes = (
_probe("control_plane_db", STATUS_DOWN, detail="file missing", kind="sqlite"),
_probe("repository", STATUS_OK, kind="git"),
_probe("gitea", STATUS_OK, required=False, kind="http"),
)
snapshot = load_system_health(probes=probes, daemon_head="abc123")
self.assertEqual(snapshot.status, STATUS_DOWN)
self.assertFalse(snapshot.ready)
self.assertTrue(
any("control_plane_db" in reason for reason in snapshot.readiness_reasons)
)
def test_optional_dependency_down_degrades_but_stays_ready(self):
"""A failing optional probe must not claim the process itself is unready."""
probes = (
_probe("control_plane_db", STATUS_OK, kind="sqlite"),
_probe("repository", STATUS_OK, kind="git"),
_probe("gitea", STATUS_DOWN, required=False, detail="timeout", kind="http"),
)
snapshot = load_system_health(probes=probes, daemon_head="abc123")
self.assertEqual(snapshot.status, STATUS_DEGRADED)
self.assertTrue(snapshot.ready)
self.assertTrue(any("gitea" in reason for reason in snapshot.readiness_reasons))
def test_unrun_required_probe_leaves_readiness_incomplete(self):
"""Not probed is not the same as passing."""
probes = (
_probe("control_plane_db", STATUS_OK, kind="sqlite"),
_probe("repository", STATUS_SKIPPED, detail="offline", kind="git"),
)
snapshot = load_system_health(probes=probes, daemon_head="abc123")
self.assertFalse(snapshot.ready)
self.assertFalse(snapshot.readiness_complete)
self.assertEqual(snapshot.status, STATUS_DEGRADED)
def test_skipped_optional_probe_does_not_block_readiness(self):
probes = (
_probe("control_plane_db", STATUS_OK, kind="sqlite"),
_probe("repository", STATUS_OK, kind="git"),
_probe("gitea", STATUS_SKIPPED, required=False, kind="http"),
)
snapshot = load_system_health(probes=probes, daemon_head="abc123")
self.assertTrue(snapshot.ready)
self.assertTrue(snapshot.readiness_complete)
class TestVersionAndUptime(CleanParityMixin, unittest.TestCase):
"""AC2 — version and uptime present when knowable."""
def test_uptime_and_start_time_present(self):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="abc123")
self.assertGreaterEqual(snapshot.uptime_seconds, 0.0)
self.assertIn("T", snapshot.started_at)
def test_process_uptime_helper_matches_shape(self):
started_at, uptime = process_uptime()
self.assertIn("T", started_at)
self.assertGreaterEqual(uptime, 0.0)
def test_version_reports_python_and_schema_version(self):
probes = (
DependencyProbe(
name="control_plane_db",
kind="sqlite",
status=STATUS_OK,
detail="ok",
required=True,
latency_ms=1.0,
metadata={"schema_version": control_plane_db.SCHEMA_VERSION},
),
_probe("repository", STATUS_OK, kind="git"),
)
snapshot = load_system_health(probes=probes, daemon_head="abc123")
self.assertEqual(
snapshot.version.control_plane_schema_version,
control_plane_db.SCHEMA_VERSION,
)
self.assertTrue(snapshot.version.python_version)
def test_version_known_flag_false_when_sha_unavailable(self):
with mock.patch("webui.system_health._git", return_value=None):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="abc")
self.assertIsNone(snapshot.version.git_sha)
self.assertFalse(snapshot.version.known)
class TestStaleRuntime(unittest.TestCase):
"""AC3 — stale runtime reflected without a false mutation-safe claim."""
def test_matching_commits_are_mutation_safe(self):
assessment = assess_stale_runtime(
Path("/tmp"),
daemon_head="aaa",
git_reader=lambda *args: "aaa",
)
self.assertFalse(assessment.stale)
self.assertTrue(assessment.determinable)
self.assertTrue(assessment.mutation_safe)
def test_diverged_commits_are_stale_and_not_mutation_safe(self):
reads = {"HEAD": "aaa", "@{upstream}": "bbb"}
assessment = assess_stale_runtime(
Path("/tmp"),
daemon_head="aaa",
git_reader=lambda *args: reads.get(args[-1]),
)
self.assertTrue(assessment.stale)
self.assertFalse(assessment.mutation_safe)
self.assertTrue(assessment.reasons)
def test_unknown_remote_is_not_mutation_safe(self):
"""Indeterminate must never read as safe."""
reads = {"HEAD": "aaa", "@{upstream}": None}
assessment = assess_stale_runtime(
Path("/tmp"),
daemon_head="aaa",
git_reader=lambda *args: reads.get(args[-1]),
)
self.assertFalse(assessment.determinable)
self.assertFalse(assessment.mutation_safe)
self.assertFalse(assessment.stale)
self.assertTrue(
any("indeterminate" in reason for reason in assessment.reasons)
)
def test_unobservable_daemon_head_is_disclosed(self):
assessment = assess_stale_runtime(
Path("/tmp"),
git_reader=lambda *args: "aaa",
)
self.assertTrue(
any("not observable" in reason for reason in assessment.reasons)
)
def test_stale_runtime_degrades_overall_status(self):
reads = {"HEAD": "aaa", "@{upstream}": "bbb"}
# Pinned rather than inherited: this path uses the default git reader,
# so the assertion must hold whether or not the suite runs offline.
with mock.patch.dict(os.environ, {"WEBUI_TEST_OFFLINE": ""}), mock.patch(
"webui.system_health._git",
side_effect=lambda repo, *args: reads.get(args[-1]),
):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="aaa")
self.assertTrue(snapshot.stale_runtime.stale)
self.assertFalse(snapshot.stale_runtime.mutation_safe)
self.assertEqual(snapshot.status, STATUS_DEGRADED)
class TestControlPlaneDbProbe(unittest.TestCase):
"""The required local dependency, probed read-only."""
def setUp(self):
self.tmp = tempfile.TemporaryDirectory()
self.addCleanup(self.tmp.cleanup)
self.db_path = str(Path(self.tmp.name) / "control-plane.db")
def _build_db(self, schema_version):
conn = sqlite3.connect(self.db_path)
conn.execute("CREATE TABLE schema_meta (key TEXT PRIMARY KEY, value TEXT)")
conn.execute("CREATE TABLE leases (lease_id TEXT PRIMARY KEY, status TEXT)")
conn.execute(
"INSERT INTO schema_meta(key, value) VALUES ('schema_version', ?)",
(str(schema_version),),
)
conn.execute("INSERT INTO leases(lease_id, status) VALUES ('l1', 'active')")
conn.commit()
conn.close()
def test_missing_database_is_down(self):
probe = probe_control_plane_db(str(Path(self.tmp.name) / "absent.db"))
self.assertEqual(probe.status, STATUS_DOWN)
self.assertTrue(probe.required)
self.assertIsNotNone(probe.latency_ms)
def test_matching_schema_is_ok(self):
self._build_db(control_plane_db.SCHEMA_VERSION)
probe = probe_control_plane_db(self.db_path)
self.assertEqual(probe.status, STATUS_OK)
self.assertEqual(
probe.metadata["schema_version"], control_plane_db.SCHEMA_VERSION
)
self.assertEqual(probe.metadata["active_leases"], 1)
def test_mismatched_schema_is_degraded(self):
self._build_db(control_plane_db.SCHEMA_VERSION + 99)
probe = probe_control_plane_db(self.db_path)
self.assertEqual(probe.status, STATUS_DEGRADED)
def test_probe_does_not_create_a_database(self):
"""A health check must never initialise the substrate it inspects."""
absent = str(Path(self.tmp.name) / "never-created.db")
probe_control_plane_db(absent)
self.assertFalse(Path(absent).exists())
def test_unreadable_database_is_down_not_raised(self):
Path(self.db_path).write_text("this is not a sqlite database")
probe = probe_control_plane_db(self.db_path)
self.assertEqual(probe.status, STATUS_DOWN)
class TestRedaction(unittest.TestCase):
"""AC4 — redaction. No credential-shaped text crosses the boundary."""
def test_redacts_token_assignment(self):
cleaned = redact("failed with token=ghp_ABCDEFGHIJKLMNOPQRSTUVWXYZ012345")
self.assertNotIn("ghp_ABCDEFGHIJKLMNOPQRSTUVWXYZ012345", cleaned)
self.assertIn("[redacted]", cleaned)
def test_redacts_authorization_header_text(self):
cleaned = redact("Authorization: Bearer abcdefghijklmnopqrstuvwxyz123456")
self.assertNotIn("abcdefghijklmnopqrstuvwxyz123456", cleaned)
def test_redacts_long_opaque_strings(self):
cleaned = redact("value 0123456789abcdef0123456789abcdef here")
self.assertNotIn("0123456789abcdef0123456789abcdef", cleaned)
def test_url_userinfo_and_query_are_stripped(self):
cleaned = redact_url("https://user:[email protected]/api/v1?token=xyz")
self.assertNotIn("secretpass", cleaned)
self.assertNotIn("token=xyz", cleaned)
self.assertEqual(cleaned, "https://gitea.example.com/api/v1")
def test_url_inside_free_text_is_redacted(self):
cleaned = redact("GET https://u:[email protected]/x?token=abc failed")
self.assertNotIn("u:p@", cleaned)
self.assertNotIn("token=abc", cleaned)
def test_gitea_probe_failure_detail_is_redacted(self):
boom = RuntimeError(
"connection refused for https://user:[email protected]/api/v1/version"
)
with mock.patch("webui.system_health.get_auth_header", return_value="token x"), \
mock.patch("webui.system_health.api_request", side_effect=boom):
probe = probe_gitea("gitea.example.com")
self.assertEqual(probe.status, STATUS_DOWN)
self.assertNotIn("hunter2", probe.detail)
self.assertEqual(scan_text_for_client_secrets(probe.detail), [])
def test_credential_guard_refusal_is_a_status_not_a_crash(self):
with mock.patch(
"webui.system_health.get_auth_header",
side_effect=RuntimeError("daemon guard refused"),
):
probe = probe_gitea("gitea.example.com")
self.assertEqual(probe.status, STATUS_DEGRADED)
self.assertFalse(probe.required)
class TestNamespaceSummaries(unittest.TestCase):
"""A web process cannot prove IDE namespace health, and must not claim to."""
def test_every_namespace_reports_unproven(self):
rows = namespace_summaries()
self.assertTrue(rows)
for row in rows:
with self.subTest(namespace=row["namespace"]):
self.assertEqual(row["status"], "unproven")
self.assertFalse(row["ide_namespace_proven"])
self.assertIn("client_namespace", row["reason"])
class TestSystemHealthRoutes(CleanParityMixin, unittest.TestCase):
"""The HTTP surface: versioned path, status codes, read-only guard."""
def setUp(self):
super().setUp()
clear_probe_cache()
self.addCleanup(clear_probe_cache)
self.client = TestClient(create_app())
def _patch_snapshot(self, probes, daemon_head="abc123"):
snapshot = load_system_health(probes=probes, daemon_head=daemon_head)
patcher = mock.patch(
"webui.app.load_system_health",
return_value=snapshot,
)
patcher.start()
self.addCleanup(patcher.stop)
return snapshot
def test_versioned_route_is_registered(self):
self.assertEqual(API_PATH, "/api/v1/system/health")
self._patch_snapshot(_ALL_HEALTHY)
response = self.client.get(API_PATH)
self.assertEqual(response.status_code, 200)
def test_healthy_payload_shape(self):
self._patch_snapshot(_ALL_HEALTHY)
data = self.client.get(API_PATH).json()
self.assertEqual(data["status"], STATUS_OK)
self.assertTrue(data["readiness"]["ready"])
self.assertTrue(data["readiness"]["complete"])
self.assertEqual(data["api"], API_PATH)
self.assertEqual(len(data["dependencies"]), 3)
for key in ("version", "process", "stale_runtime", "mcp_namespaces"):
self.assertIn(key, data)
self.assertIn("uptime_seconds", data["process"])
self.assertIn("mutation_safe", data["stale_runtime"])
def test_degraded_dependency_returns_503(self):
probes = (
_probe("control_plane_db", STATUS_DOWN, detail="missing", kind="sqlite"),
_probe("repository", STATUS_OK, kind="git"),
)
self._patch_snapshot(probes)
response = self.client.get(API_PATH)
self.assertEqual(response.status_code, 503)
data = response.json()
self.assertFalse(data["readiness"]["ready"])
self.assertTrue(data["readiness"]["reasons"])
def test_dependency_entries_expose_status_and_latency(self):
self._patch_snapshot(_ALL_HEALTHY)
data = self.client.get(API_PATH).json()
names = {entry["name"] for entry in data["dependencies"]}
self.assertEqual(names, {"control_plane_db", "repository", "gitea"})
for entry in data["dependencies"]:
with self.subTest(dependency=entry["name"]):
self.assertIn("status", entry)
self.assertIn("required", entry)
self.assertIn("latency_ms", entry)
def test_response_body_carries_no_client_secrets(self):
self._patch_snapshot(_ALL_HEALTHY)
body = self.client.get(API_PATH).text
self.assertEqual(scan_text_for_client_secrets(body), [])
def test_deep_flag_is_forwarded(self):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="abc")
with mock.patch(
"webui.app.load_system_health", return_value=snapshot
) as loader:
self.client.get(f"{API_PATH}?deep=1")
loader.assert_called_once_with(deep=True)
def test_shallow_is_the_default(self):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="abc")
with mock.patch(
"webui.app.load_system_health", return_value=snapshot
) as loader:
self.client.get(API_PATH)
loader.assert_called_once_with(deep=False)
def test_route_rejects_mutation_methods(self):
for method in ("POST", "PUT", "PATCH", "DELETE"):
with self.subTest(method=method):
response = self.client.request(method, API_PATH)
self.assertEqual(response.status_code, 405)
self.assertEqual(response.json()["error"], "read-only-mvp")
def test_default_shallow_call_skips_the_network_probe(self):
"""The expensive probe must not run unless it was asked for."""
with mock.patch("webui.system_health.probe_gitea") as probe:
snapshot = load_system_health(deep=False)
probe.assert_not_called()
gitea = next(p for p in snapshot.dependencies if p.name == "gitea")
self.assertEqual(gitea.status, STATUS_SKIPPED)
class TestHealthRouteBackwardCompatibility(unittest.TestCase):
"""`/health` is expanded additively; MVP consumers must keep working."""
def setUp(self):
self.client = TestClient(create_app())
def test_mvp_keys_are_unchanged(self):
data = self.client.get("/health").json()
self.assertEqual(data["status"], "ok")
self.assertEqual(data["service"], "mcp-control-plane-webui")
self.assertEqual(data["mode"], "read-only-mvp")
self.assertIn("timestamp", data)
self.assertEqual(data["deployment"]["mode"], "internal-operator-console")
def test_health_points_at_the_versioned_api(self):
data = self.client.get("/health").json()
self.assertEqual(data["system_health_api"], API_PATH)
self.assertIn("uptime_seconds", data)
self.assertIn("started_at", data)
def test_health_runs_no_dependency_probe(self):
"""Liveness must stay cheap: no probe, no snapshot assembly."""
with mock.patch("webui.app.load_system_health") as loader:
response = self.client.get("/health")
self.assertEqual(response.status_code, 200)
loader.assert_not_called()
class TestSnapshotSerialisation(CleanParityMixin, unittest.TestCase):
def test_snapshot_dict_is_json_serialisable(self):
snapshot = load_system_health(probes=_ALL_HEALTHY, daemon_head="abc123")
encoded = json.dumps(snapshot_to_dict(snapshot))
self.assertIn("readiness", encoded)
if __name__ == "__main__":
unittest.main()
+34
View File
@@ -45,6 +45,12 @@ from webui.worktree_scanner import load_hygiene_snapshot, snapshot_to_dict as wo
from webui.worktree_views import render_worktrees_page from webui.worktree_views import render_worktrees_page
from webui.runtime_health import load_runtime_snapshot, snapshot_to_dict as runtime_snapshot_to_dict from webui.runtime_health import load_runtime_snapshot, snapshot_to_dict as runtime_snapshot_to_dict
from webui.runtime_views import render_runtime_page from webui.runtime_views import render_runtime_page
from webui.system_health import (
API_PATH as SYSTEM_HEALTH_API_PATH,
load_system_health,
process_uptime,
snapshot_to_dict as system_health_to_dict,
)
_READ_ONLY_METHODS = frozenset({"GET", "HEAD", "OPTIONS"}) _READ_ONLY_METHODS = frozenset({"GET", "HEAD", "OPTIONS"})
_AUDIT_MUTATION_PATHS = frozenset({"/audit", "/api/audit"}) _AUDIT_MUTATION_PATHS = frozenset({"/audit", "/api/audit"})
@@ -78,16 +84,43 @@ async def home(_request: Request) -> HTMLResponse:
async def health(_request: Request) -> JSONResponse: async def health(_request: Request) -> JSONResponse:
"""Liveness only — deliberately cheap, runs no dependency probe (#634).
Every MVP key is retained so existing pollers keep working; the additions
are a pointer to the structured API and the in-memory process uptime.
Readiness lives at that API because answering it costs real probes.
"""
bind_host = _request.app.state.webui_bind_host bind_host = _request.app.state.webui_bind_host
started_at, uptime_seconds = process_uptime()
return JSONResponse({ return JSONResponse({
"status": "ok", "status": "ok",
"service": "mcp-control-plane-webui", "service": "mcp-control-plane-webui",
"mode": "read-only-mvp", "mode": "read-only-mvp",
"timestamp": datetime.now(timezone.utc).isoformat(), "timestamp": datetime.now(timezone.utc).isoformat(),
"deployment": deployment_snapshot(bind_host=bind_host), "deployment": deployment_snapshot(bind_host=bind_host),
"started_at": started_at,
"uptime_seconds": uptime_seconds,
"system_health_api": SYSTEM_HEALTH_API_PATH,
}) })
def _truthy_flag(value: str | None) -> bool:
return (value or "").strip().lower() in {"1", "true", "yes", "on"}
async def api_system_health(request: Request) -> JSONResponse:
"""Structured read-only system health (#634).
`?deep=1` opts into the expensive network probe. The response status code
reflects readiness so automated checks can branch on it without parsing the
body: 200 when ready, 503 when a required dependency failed or never ran.
"""
deep = _truthy_flag(request.query_params.get("deep"))
snapshot = load_system_health(deep=deep)
payload = system_health_to_dict(snapshot)
return JSONResponse(payload, status_code=200 if snapshot.ready else 503)
async def queue(_request: Request) -> HTMLResponse: async def queue(_request: Request) -> HTMLResponse:
snapshot = load_queue_snapshot() snapshot = load_queue_snapshot()
return HTMLResponse(render_page(title="Queue", body_html=render_queue_page(snapshot))) return HTMLResponse(render_page(title="Queue", body_html=render_queue_page(snapshot)))
@@ -399,6 +432,7 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
routes=[ routes=[
Route("/", home, methods=["GET"]), Route("/", home, methods=["GET"]),
Route("/health", health, methods=["GET"]), Route("/health", health, methods=["GET"]),
Route(SYSTEM_HEALTH_API_PATH, api_system_health, methods=["GET"]),
Route("/queue", queue, methods=["GET"]), Route("/queue", queue, methods=["GET"]),
Route("/api/queue", api_queue, methods=["GET"]), Route("/api/queue", api_queue, methods=["GET"]),
Route("/projects", projects, methods=["GET"]), Route("/projects", projects, methods=["GET"]),
+682
View File
@@ -0,0 +1,682 @@
"""Read-only system-health model for the operator console API (#634).
`/health` answers liveness only. Operators automating readiness checks need a
structured view of *why* the control plane is or is not usable: which
dependencies answered, how long they took, what version of the code is running,
and whether the runtime is stale relative to its remote.
Three rules shape this module.
* **Read-only.** Every probe opens its subject read-only. The control-plane
database is opened through a ``mode=ro`` URI so a health check can never
create or migrate a schema, and no probe writes, restarts, or reloads
anything restart controls are Phase 2, and #630 forbids process-kill
recovery outright.
* **Fail-soft.** A dependency that is unreachable is a *status*, not an
exception. Probes catch their own failures and report them as a degraded or
down entry carrying a reason.
* **Never claim more than was proven.** Readiness is derived only from probes
that actually ran, ``mutation_safe`` stays false unless the parity commits are
known and equal, and an MCP namespace is reported unproven because a web
process cannot exercise the IDE-managed client path (#543).
"""
from __future__ import annotations
import os
import re
import sqlite3
import subprocess
import time
from dataclasses import dataclass
from datetime import datetime, timezone
from pathlib import Path
from typing import Any, Callable
from urllib.parse import urlsplit, urlunsplit
import control_plane_db
import mcp_namespace_health
from gitea_auth import api_request, get_auth_header, gitea_url
from webui.project_registry import load_registry
SERVICE_NAME = "mcp-control-plane-webui"
API_PATH = "/api/v1/system/health"
STATUS_OK = "ok"
STATUS_DEGRADED = "degraded"
STATUS_DOWN = "down"
STATUS_SKIPPED = "skipped"
STATUS_UNPROVEN = "unproven"
# Statuses that count as a healthy answer from a probe.
_HEALTHY_STATUSES = frozenset({STATUS_OK})
# Statuses meaning "this probe did not run", as opposed to "it ran and failed".
_NOT_RUN_STATUSES = frozenset({STATUS_SKIPPED})
_DEEP_PROBE_TTL_ENV = "WEBUI_HEALTH_PROBE_TTL_SECONDS"
_DEFAULT_DEEP_PROBE_TTL = 15.0
_GITEA_PROBE_TIMEOUT_SECONDS = 5.0
_OFFLINE_ENV = "WEBUI_TEST_OFFLINE"
# Credential-shaped material that must never reach the browser, mirroring the
# forbidden client patterns in webui/deployment_boundary.py.
_SECRET_RE = re.compile(
r"(?i)\b(token|password|passwd|secret|authorization|bearer)\b\s*[:=]?\s*\S+"
)
_LONG_OPAQUE_RE = re.compile(r"\b[A-Za-z0-9_\-]{32,}\b")
# Captured once at import so uptime measures this process, not the request.
_STARTED_AT = datetime.now(timezone.utc)
_STARTED_MONOTONIC = time.monotonic()
# TTL cache for the expensive (network) probe only.
_deep_cache: dict[str, tuple[float, "DependencyProbe"]] = {}
@dataclass(frozen=True)
class DependencyProbe:
"""One dependency check, fail-soft, with its own latency."""
name: str
kind: str
status: str
detail: str
required: bool
latency_ms: float | None = None
metadata: dict[str, Any] | None = None
@property
def healthy(self) -> bool:
return self.status in _HEALTHY_STATUSES
@property
def ran(self) -> bool:
return self.status not in _NOT_RUN_STATUSES
@dataclass(frozen=True)
class VersionInfo:
git_sha: str | None
git_describe: str | None
control_plane_schema_version: int | None
python_version: str
known: bool
@dataclass(frozen=True)
class StaleRuntime:
"""Parity between the running code, the checkout, and the remote.
``mutation_safe`` is deliberately conservative: unknown is not safe.
"""
daemon_head: str | None
checkout_head: str | None
remote_head: str | None
stale: bool
determinable: bool
mutation_safe: bool
reasons: tuple[str, ...]
@dataclass(frozen=True)
class SystemHealthSnapshot:
status: str
ready: bool
readiness_complete: bool
readiness_reasons: tuple[str, ...]
service: str
mode: str
version: VersionInfo
started_at: str
uptime_seconds: float
timestamp: str
deep_probes_requested: bool
dependencies: tuple[DependencyProbe, ...]
mcp_namespaces: tuple[dict[str, Any], ...]
stale_runtime: StaleRuntime
probe_errors: tuple[str, ...] = ()
def process_uptime() -> tuple[str, float]:
"""Process start timestamp and uptime — in-memory, safe for `/health`."""
return _STARTED_AT.isoformat(), round(time.monotonic() - _STARTED_MONOTONIC, 3)
def _offline() -> bool:
return (os.environ.get(_OFFLINE_ENV) or "").strip().lower() in {"1", "true", "yes"}
def _repo_root() -> Path:
override = (os.environ.get("WEBUI_REPO_ROOT") or "").strip()
if override:
return Path(override).resolve()
return Path(__file__).resolve().parent.parent
def _deep_probe_ttl() -> float:
raw = (os.environ.get(_DEEP_PROBE_TTL_ENV) or "").strip()
if not raw:
return _DEFAULT_DEEP_PROBE_TTL
try:
value = float(raw)
except ValueError:
return _DEFAULT_DEEP_PROBE_TTL
return value if value >= 0 else _DEFAULT_DEEP_PROBE_TTL
def redact(text: str) -> str:
"""Strip credential-shaped material from operator-visible probe text.
Probe details carry exception strings, and an exception raised by an HTTP
client can quote the request that failed. Redaction happens here, at the
boundary where those strings become part of a browser-bound payload.
"""
if not text:
return ""
cleaned = _redact_urls(text)
cleaned = _SECRET_RE.sub(lambda m: f"{m.group(1)}=[redacted]", cleaned)
return _LONG_OPAQUE_RE.sub("[redacted]", cleaned)
def _redact_urls(text: str) -> str:
return re.sub(r"https?://\S+", lambda m: redact_url(m.group(0)), text)
def redact_url(url: str) -> str:
"""Reduce a URL to scheme://host/path — no userinfo, no query, no fragment."""
try:
parts = urlsplit(url)
except ValueError:
return "[redacted-url]"
if not parts.scheme or not parts.hostname:
return "[redacted-url]"
netloc = parts.hostname
if parts.port:
netloc = f"{netloc}:{parts.port}"
return urlunsplit((parts.scheme, netloc, parts.path, "", ""))
def _git(repo: Path, *args: str) -> str | None:
try:
completed = subprocess.run(
["git", "-C", str(repo), *args],
capture_output=True,
text=True,
check=False,
timeout=10,
)
except (OSError, subprocess.SubprocessError):
return None
if completed.returncode != 0:
return None
return (completed.stdout or "").strip() or None
def _load_version(repo: Path, *, schema_version: int | None) -> VersionInfo:
import platform
git_sha = None if _offline() else _git(repo, "rev-parse", "HEAD")
describe = None if _offline() else _git(repo, "describe", "--tags", "--always")
return VersionInfo(
git_sha=git_sha,
git_describe=describe,
control_plane_schema_version=schema_version,
python_version=platform.python_version(),
known=bool(git_sha),
)
# ---------------------------------------------------------------------------
# Dependency probes
# ---------------------------------------------------------------------------
def _elapsed_ms(started: float) -> float:
return round((time.monotonic() - started) * 1000, 3)
def probe_control_plane_db(db_path: str | None = None) -> DependencyProbe:
"""Read-only reachability check for the control-plane SQLite substrate.
Opened through a ``mode=ro`` URI on purpose: ``ControlPlaneDB.__init__``
creates directories and runs schema migrations, which a health check must
never do.
"""
path = (db_path or control_plane_db.default_db_path()).strip()
started = time.monotonic()
metadata: dict[str, Any] = {"path": path}
def _result(status: str, detail: str) -> DependencyProbe:
return DependencyProbe(
name="control_plane_db",
kind="sqlite",
status=status,
detail=detail,
required=True,
latency_ms=_elapsed_ms(started),
metadata=metadata,
)
if not path or not os.path.exists(path):
return _result(STATUS_DOWN, "control-plane database file does not exist yet")
try:
conn = sqlite3.connect(f"file:{path}?mode=ro", uri=True, timeout=5)
try:
row = conn.execute(
"SELECT value FROM schema_meta WHERE key = 'schema_version'"
).fetchone()
leases = conn.execute(
"SELECT COUNT(*) FROM leases WHERE status = 'active'"
).fetchone()
finally:
conn.close()
except sqlite3.Error as exc:
return _result(STATUS_DOWN, redact(f"control-plane database unreadable: {exc}"))
schema_version = int(row[0]) if row and str(row[0]).isdigit() else None
metadata["schema_version"] = schema_version
metadata["active_leases"] = int(leases[0]) if leases else None
if schema_version is None:
return _result(
STATUS_DEGRADED, "control-plane database has no recorded schema version"
)
if schema_version != control_plane_db.SCHEMA_VERSION:
return _result(
STATUS_DEGRADED,
f"control-plane schema version {schema_version} does not match the "
f"version this code expects ({control_plane_db.SCHEMA_VERSION})",
)
return _result(STATUS_OK, f"schema v{schema_version} readable")
def probe_repository(repo: Path) -> DependencyProbe:
"""Local checkout reachability — required, cheap, no network."""
started = time.monotonic()
metadata: dict[str, Any] = {"repo_root": str(repo)}
if _offline():
return DependencyProbe(
name="repository",
kind="git",
status=STATUS_SKIPPED,
detail=f"{_OFFLINE_ENV} is set; git probe skipped",
required=True,
latency_ms=_elapsed_ms(started),
metadata=metadata,
)
head = _git(repo, "rev-parse", "HEAD")
if not head:
return DependencyProbe(
name="repository",
kind="git",
status=STATUS_DOWN,
detail=f"HEAD could not be read at {repo}",
required=True,
latency_ms=_elapsed_ms(started),
metadata=metadata,
)
branch = _git(repo, "rev-parse", "--abbrev-ref", "HEAD")
metadata["head"] = head
metadata["branch"] = branch
return DependencyProbe(
name="repository",
kind="git",
status=STATUS_OK,
detail=f"checkout readable at {branch or 'detached HEAD'}",
required=True,
latency_ms=_elapsed_ms(started),
metadata=metadata,
)
def probe_gitea(host: str) -> DependencyProbe:
"""Live Gitea reachability. Expensive (network), so opt-in via ``deep``.
Optional by design: the console stays useful for local inventory when the
remote is unreachable, so a failure here degrades status without claiming
the process itself is unready.
"""
started = time.monotonic()
metadata: dict[str, Any] = {"host": host}
def _failure(status: str, detail: str) -> DependencyProbe:
return DependencyProbe(
name="gitea",
kind="http",
status=status,
detail=detail,
required=False,
latency_ms=_elapsed_ms(started),
metadata=metadata,
)
if not host:
return _failure(STATUS_DEGRADED, "no Gitea host is configured in the registry")
try:
auth = get_auth_header(host)
except Exception as exc: # noqa: BLE001 — credential guards are a status here
return _failure(STATUS_DEGRADED, redact(f"credential lookup refused: {exc}"))
if not auth:
return _failure(STATUS_DEGRADED, f"no credentials available for {host}")
url = gitea_url(host, "/api/v1/version")
metadata["endpoint"] = redact_url(url)
try:
data = api_request("GET", url, auth, timeout=_GITEA_PROBE_TIMEOUT_SECONDS)
except Exception as exc: # noqa: BLE001 — a down dependency is a status
return _failure(STATUS_DOWN, redact(f"Gitea probe failed: {exc}"))
if isinstance(data, dict) and data.get("version"):
metadata["gitea_version"] = str(data["version"])
return DependencyProbe(
name="gitea",
kind="http",
status=STATUS_OK,
detail=f"{host} reachable",
required=False,
latency_ms=_elapsed_ms(started),
metadata=metadata,
)
def _skipped_gitea(host: str) -> DependencyProbe:
return DependencyProbe(
name="gitea",
kind="http",
status=STATUS_SKIPPED,
detail="network probe not requested; call with ?deep=1 to run it",
required=False,
latency_ms=None,
metadata={"host": host},
)
def namespace_summaries() -> tuple[dict[str, Any], ...]:
"""Declared MCP namespaces, each honestly reported as unproven.
The web process runs outside the IDE-managed MCP client, so it cannot
invoke a namespace tool. Per #543 only a ``client_namespace`` probe proves
that path, and inventing a healthy verdict here is exactly the false claim
the mutation gates exist to prevent.
"""
rows: list[dict[str, Any]] = []
for namespace, required_tool in sorted(
mcp_namespace_health.REQUIRED_NAMESPACE_TOOLS.items()
):
classification = mcp_namespace_health.classify_namespace_probe(
namespace,
required_tool=required_tool,
probe_result=None,
probe_source=mcp_namespace_health.PROBE_SOURCE_UNKNOWN,
)
rows.append(
{
"namespace": namespace,
"required_tool": required_tool,
"status": STATUS_UNPROVEN,
"ide_namespace_proven": bool(classification.get("ide_namespace_proven")),
"reason": (
"the web console cannot invoke the IDE-managed MCP client; "
"namespace health must be proven with a client_namespace "
"probe (#543)"
),
"error_type": classification.get("error_type"),
}
)
return tuple(rows)
def assess_stale_runtime(
repo: Path,
*,
daemon_head: str | None = None,
git_reader: Callable[..., str | None] | None = None,
) -> StaleRuntime:
"""Three-way parity view: running code, local checkout, remote-tracking ref.
``mutation_safe`` requires all three to be known and equal. Anything less
including "the remote ref was never fetched" is reported as not safe with
a reason, so an operator never reads an unproven green.
"""
reader = git_reader or (lambda *args: _git(repo, *args))
reasons: list[str] = []
# The offline switch suppresses real subprocess calls; an explicitly
# injected reader is already a substitute for them and is always used.
offline = _offline() and git_reader is None
checkout_head = None if offline else reader("rev-parse", "HEAD")
remote_head = None if offline else reader("rev-parse", "@{upstream}")
if offline:
reasons.append(f"{_OFFLINE_ENV} is set; parity commits were not read")
else:
if checkout_head is None:
reasons.append("local checkout HEAD could not be read")
if remote_head is None:
reasons.append(
"no remote-tracking commit is known for the current branch; "
"remote staleness is indeterminate (no fetch is performed here)"
)
effective_daemon = daemon_head if daemon_head is not None else checkout_head
if daemon_head is None:
reasons.append(
"the running MCP daemon's startup commit is not observable from the "
"web process; the checkout commit is reported in its place"
)
determinable = bool(checkout_head and remote_head and effective_daemon)
stale = bool(
determinable and len({checkout_head, remote_head, effective_daemon}) > 1
)
if stale:
reasons.append(
"runtime, checkout, and remote commits disagree; restart the MCP "
"server after updating the checkout before trusting capability gates"
)
return StaleRuntime(
daemon_head=effective_daemon,
checkout_head=checkout_head,
remote_head=remote_head,
stale=stale,
determinable=determinable,
mutation_safe=bool(determinable and not stale),
reasons=tuple(reasons),
)
# ---------------------------------------------------------------------------
# Snapshot assembly
# ---------------------------------------------------------------------------
def _default_host() -> str:
registry = load_registry()
if not registry.projects:
return ""
raw = registry.projects[0].remote_host
parts = urlsplit(raw.strip())
return parts.netloc or raw.strip().rstrip("/")
def _aggregate(
probes: tuple[DependencyProbe, ...],
) -> tuple[str, bool, bool, tuple[str, ...]]:
"""Fold probe results into overall status and readiness.
Required probes drive readiness; optional probes can only degrade status.
A probe that did not run leaves readiness incomplete rather than passing.
"""
reasons: list[str] = []
required = [probe for probe in probes if probe.required]
unrun_required = [probe for probe in required if not probe.ran]
failed_required = [probe for probe in required if probe.ran and not probe.healthy]
failed_optional = [
probe
for probe in probes
if not probe.required and probe.ran and not probe.healthy
]
for probe in unrun_required:
reasons.append(
f"required dependency '{probe.name}' was not probed: {probe.detail}"
)
for probe in failed_required:
reasons.append(
f"required dependency '{probe.name}' is {probe.status}: {probe.detail}"
)
for probe in failed_optional:
reasons.append(
f"optional dependency '{probe.name}' is {probe.status}: {probe.detail}"
)
readiness_complete = not unrun_required
ready = readiness_complete and not failed_required
if any(probe.status == STATUS_DOWN for probe in failed_required):
status = STATUS_DOWN
elif failed_required or failed_optional or unrun_required:
status = STATUS_DEGRADED
else:
status = STATUS_OK
return status, ready, readiness_complete, tuple(reasons)
def load_system_health(
*,
deep: bool = False,
host: str | None = None,
probes: tuple[DependencyProbe, ...] | None = None,
daemon_head: str | None = None,
use_cache: bool = True,
) -> SystemHealthSnapshot:
"""Assemble the read-only system-health snapshot.
``deep=True`` adds the network probe against Gitea; its result is cached for
a short TTL so repeated dashboard polls do not amplify into remote load.
"""
repo = _repo_root()
probe_errors: list[str] = []
if probes is None:
collected: list[DependencyProbe] = []
for probe_fn in (
lambda: probe_control_plane_db(),
lambda: probe_repository(repo),
):
try:
collected.append(probe_fn())
except Exception as exc: # noqa: BLE001 — a probe must not 500 the API
probe_errors.append(redact(f"probe raised: {exc}"))
resolved_host = host if host is not None else _default_host()
if deep and not _offline():
collected.append(_cached_gitea_probe(resolved_host, use_cache=use_cache))
else:
collected.append(_skipped_gitea(resolved_host))
probes = tuple(collected)
status, ready, readiness_complete, reasons = _aggregate(probes)
stale = assess_stale_runtime(repo, daemon_head=daemon_head)
if stale.stale:
if status == STATUS_OK:
status = STATUS_DEGRADED
reasons = reasons + (
"runtime is stale relative to its remote-tracking commit",
)
db_probe = next((p for p in probes if p.name == "control_plane_db"), None)
schema_version = None
if db_probe and db_probe.metadata:
schema_version = db_probe.metadata.get("schema_version")
return SystemHealthSnapshot(
status=status,
ready=ready,
readiness_complete=readiness_complete,
readiness_reasons=reasons,
service=SERVICE_NAME,
mode="read-only",
version=_load_version(repo, schema_version=schema_version),
started_at=_STARTED_AT.isoformat(),
uptime_seconds=round(time.monotonic() - _STARTED_MONOTONIC, 3),
timestamp=datetime.now(timezone.utc).isoformat(),
deep_probes_requested=deep,
dependencies=probes,
mcp_namespaces=namespace_summaries(),
stale_runtime=stale,
probe_errors=tuple(probe_errors),
)
def _cached_gitea_probe(host: str, *, use_cache: bool = True) -> DependencyProbe:
ttl = _deep_probe_ttl()
now = time.monotonic()
if use_cache and ttl > 0:
cached = _deep_cache.get(host)
if cached and (now - cached[0]) < ttl:
return cached[1]
probe = probe_gitea(host)
if use_cache and ttl > 0:
_deep_cache[host] = (now, probe)
return probe
def clear_probe_cache() -> None:
"""Drop cached deep-probe results (tests and operator-forced refresh)."""
_deep_cache.clear()
def probe_to_dict(probe: DependencyProbe) -> dict[str, Any]:
return {
"name": probe.name,
"kind": probe.kind,
"status": probe.status,
"detail": probe.detail,
"required": probe.required,
"healthy": probe.healthy,
"latency_ms": probe.latency_ms,
"metadata": dict(probe.metadata or {}),
}
def snapshot_to_dict(snapshot: SystemHealthSnapshot) -> dict[str, Any]:
return {
"status": snapshot.status,
"service": snapshot.service,
"mode": snapshot.mode,
"api": API_PATH,
"timestamp": snapshot.timestamp,
"readiness": {
"ready": snapshot.ready,
"complete": snapshot.readiness_complete,
"reasons": list(snapshot.readiness_reasons),
},
"version": {
"git_sha": snapshot.version.git_sha,
"git_describe": snapshot.version.git_describe,
"control_plane_schema_version": (
snapshot.version.control_plane_schema_version
),
"python_version": snapshot.version.python_version,
"known": snapshot.version.known,
},
"process": {
"started_at": snapshot.started_at,
"uptime_seconds": snapshot.uptime_seconds,
},
"deep_probes_requested": snapshot.deep_probes_requested,
"dependencies": [probe_to_dict(probe) for probe in snapshot.dependencies],
"mcp_namespaces": [dict(row) for row in snapshot.mcp_namespaces],
"stale_runtime": {
"daemon_head": snapshot.stale_runtime.daemon_head,
"checkout_head": snapshot.stale_runtime.checkout_head,
"remote_head": snapshot.stale_runtime.remote_head,
"stale": snapshot.stale_runtime.stale,
"determinable": snapshot.stale_runtime.determinable,
"mutation_safe": snapshot.stale_runtime.mutation_safe,
"reasons": list(snapshot.stale_runtime.reasons),
},
"probe_errors": list(snapshot.probe_errors),
}
+5 -3
View File
@@ -65,9 +65,11 @@ PROMPT_RECONCILER = (
"reconciliation (already-landed / post-merge cleanup). Do not approve or merge." "reconciliation (already-landed / post-merge cleanup). Do not approve or merge."
) )
PROMPT_CONTROLLER = ( PROMPT_CONTROLLER = (
"CONTROLLER session: inspect gitea_workflow_dashboard + control-plane leases, " "CONTROLLER session: call gitea_route_task_session(task_type='process_work_queue') "
"diagnose blocked/terminal-locked items for {remote}/{org}/{repo}, and schedule " "then gitea_allocate_next_work (cross_role default) for {remote}/{org}/{repo}; "
"exactly one fresh role-scoped cycle. Do not implement, review, or merge in-band." "use the returned required_role/profile/action to schedule exactly one downstream "
"role cycle. Dashboard is explanatory only and never replaces allocator selection. "
"Do not implement, review, approve, or merge in-band."
) )
PROMPT_IDLE = ( PROMPT_IDLE = (
"IDLE: no safe assignable work for role '{role}' on {remote}/{org}/{repo}. " "IDLE: no safe assignable work for role '{role}' on {remote}/{org}/{repo}. "