Compare commits

...
Author SHA1 Message Date
sysadminandClaude Opus 4.8 b179610e7f fix(webui): remediate PR #876 REQUEST_CHANGES for analytics (#651)
Address review blockers and medium/low findings:

- F1: HTML-escape all dynamic analytics fields (html.escape quote=True)
- F2: Gate POST /api/v1/analytics/usage through console_authz
  record_analytics_usage (fail closed for unauthenticated / Phase 1)
- F3: Enforce usage_events retention (max rows + max age)
- F4: Coerce None remote/org/repo to empty strings in load_analytics

Add tests for XSS escaping, unauthorized ingest deny, retention, and
None scope coercion. Document authz + retention in analytics guide.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-24 08:11:12 -04:00
jcwalker3 0041542fe2 Merge branch 'master' into feat/issue-651-usage-cost-analytics 2026-07-24 07:03:24 -05:00
sysadmin a87a7d1da2 Merge pull request 'Implement native author issue worktree bootstrap (#850)' (#853) from fix/issue-850-native-mcp-bootstrap into master 2026-07-24 07:02:48 -05:00
jcwalker3 9504fa8bbd Merge branch 'master' into fix/issue-850-native-mcp-bootstrap 2026-07-24 06:55:44 -05:00
jcwalker3 1f144705e2 Merge branch 'master' into feat/issue-651-usage-cost-analytics 2026-07-24 06:55:03 -05:00
sysadmin e1d844bfed fix(bootstrap): path-shaped branches ancestry without isdir (#850 review #551)
resolve_canonical_repo_root fallback now uses commonpath only — no
string split and no os.path.isdir gate — so MCP project_root =
branches/<wt> still resolves to the repo root when the path is not yet
on disk (#274 / review #551 regression).
2026-07-24 07:54:07 -04:00
sysadmin 06e95254f0 fix(bootstrap): address PR #853 REQUEST_CHANGES findings (#850)
- Remove /branches/ string-split fallback in resolve_canonical_repo_root;
  recover roots via commonpath ancestry only (review #531 F2).
- Refuse existing branches that do not contain live master; no weak
  merge-base acceptance (F3).
- Verify caller-supplied assignment_id/lease_id against the control plane
  or fail closed (F4).
- Compensating recovery releases bound workflow leases via lease_lifecycle (F5).
- Regression tests for each finding.
2026-07-24 07:46:34 -04:00
sysadmin 5e935dffb4 Merge pull request 'feat(webui): system-health dashboard (Closes #639)' (#862) from feat/issue-639-webui-system-health-dashboard into master 2026-07-24 06:28:46 -05:00
sysadmin 5c5c1fdf77 feat(webui): model usage, token cost, latency, and performance analytics (Closes #651) 2026-07-24 07:16:44 -04:00
jcwalker3 0ae05cb9bc Merge branch 'master' into fix/issue-850-native-mcp-bootstrap 2026-07-24 06:16:32 -05:00
jcwalker3 82464f4054 Merge branch 'master' into feat/issue-639-webui-system-health-dashboard 2026-07-24 06:16:12 -05:00
sysadmin c33c69b3f3 Merge pull request 'fix: conflict-fix lease lifecycle chain termination and TTL handling (#842)' (#846) from fix/issue-842-conflict-fix-lease-lifecycle into master 2026-07-24 06:11:09 -05:00
sysadminandClaude Opus 4.8 2baf726ee6 merge(master): sync PR #853 with master; keep bootstrap and #860 transitions
Resolve task_capability_map.py by unioning #850 bootstrap_author_issue_worktree
preflight transitions with master's #860 dirty-orphan recovery and work_issue
commit transitions. No feature logic discarded.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-24 07:04:20 -04:00
sysadmin 67b4889984 Merge pull request 'feat(mcp-health): MCP restart coordinator and impact analysis (Closes #658)' (#875) from feat/issue-658-mcp-restart-coordinator into master 2026-07-24 06:04:14 -05:00
jcwalker3 fd558ce5d8 Merge branch 'master' into feat/issue-658-mcp-restart-coordinator 2026-07-24 05:54:50 -05:00
sysadmin ae1161524d Merge pull request 'fix(author): bootstrap recovery for dirty orphaned issue worktrees (#860)' (#861) from fix/issue-860-dirty-orphan-worktree-recovery into master 2026-07-24 05:52:24 -05:00
jcwalker3 d456a763fa Merge branch 'master' into feat/issue-658-mcp-restart-coordinator 2026-07-24 05:48:25 -05:00
jcwalker3 cc56aeeacd Merge branch 'master' into fix/issue-860-dirty-orphan-worktree-recovery 2026-07-24 05:22:20 -05:00
sysadmin 44fe8d2eed Merge pull request 'feat(mcp-health): inventory and guard MCP restart/reload/kill paths (Closes #657)' (#870) from feat/issue-657-mcp-restart-path-inventory-guard into master 2026-07-24 04:28:10 -05:00
jcwalker3 d03d982e3b Merge branch 'master' into fix/issue-860-dirty-orphan-worktree-recovery 2026-07-24 04:11:02 -05:00
jcwalker3 657b5bc1b3 Merge branch 'master' into feat/issue-657-mcp-restart-path-inventory-guard 2026-07-24 04:10:08 -05:00
sysadmin b6ca778cef Merge pull request 'fix(reconciler): PR-scoped post-merge cleanup executor + expired reviewer-lease reclaim (Closes #855)' (#866) from fix/issue-855-pr-scoped-merged-cleanup into master 2026-07-24 03:24:53 -05:00
jcwalker3 73f82a2305 Merge branch 'master' into fix/issue-855-pr-scoped-merged-cleanup 2026-07-24 03:15:43 -05:00
sysadmin cc7dc8ac14 Merge pull request 'fix(author): refresh durable issue-lock head after branch sync; recover merge-sync-drifted dead-session locks (Closes #871)' (#874) from fix/issue-871-durable-lock-head-refresh into master 2026-07-24 03:04:24 -05:00
jcwalker3 e151759212 Merge branch 'master' into fix/issue-871-durable-lock-head-refresh 2026-07-24 02:42:52 -05:00
sysadmin 60df5087e9 Merge pull request 'docs(governance): MCP restart governance and authorization policy (Closes #656)' (#857) from docs/issue-656-mcp-restart-governance into master 2026-07-24 02:26:02 -05:00
jcwalker3 78e3befbbb Merge branch 'master' into feat/issue-658-mcp-restart-coordinator 2026-07-24 02:15:32 -05:00
jcwalker3 d7e69fbe77 Merge branch 'master' into feat/issue-639-webui-system-health-dashboard 2026-07-24 02:07:53 -05:00
jcwalker3 c9aa09f341 Merge branch 'master' into fix/issue-855-pr-scoped-merged-cleanup 2026-07-24 02:07:33 -05:00
jcwalker3 b4afc8cefd Merge branch 'master' into feat/issue-657-mcp-restart-path-inventory-guard 2026-07-24 02:07:25 -05:00
jcwalker3 0088ecaf00 Merge branch 'master' into docs/issue-656-mcp-restart-governance 2026-07-24 02:07:17 -05:00
jcwalker3 e0536d344f Merge branch 'master' into fix/issue-850-native-mcp-bootstrap 2026-07-24 02:07:03 -05:00
sysadminandClaude Opus 4.8 5d59c57c98 fix(author): refresh durable issue-lock head after branch sync; recover merge-sync-drifted dead-session locks (#871)
`gitea_update_pr_branch_by_merge` advanced a PR's remote head but never
advanced the linked durable issue lock's recorded head. After the owning
session died the drifted lock became unrecoverable and no further
synchronization was possible (PR #866 / issue #855).

Write-side (prevents future drift):
- issue_lock_store.assess/apply_durable_lock_head_refresh: on a successful
  sync, CAS-refresh the durable lock's recorded head from the exact expected
  PR head to the resulting head, re-verifying repo/issue/branch/worktree/
  identity/profile/live-session ownership, with read-after-write verification.
- gitea_update_pr_branch_by_merge now refreshes the lock after the remote
  advance and reports a PARTIAL LIFECYCLE FAILURE (success=False) when the
  refresh fails, instead of falsely reporting a full synchronization.

Read-side (recovers already-drifted locks):
- issue_lock_worktree.read_merge_sync_provenance: server-side git observation
  proving a remote head is a sanctioned base-into-branch merge that preserved
  the branch mainline back to the recorded head.
- issue_lock_recovery: new HEAD_RELATION_REMOTE_MERGE_SYNCED accepts a
  dead-session lock whose recorded head is a strict merge-sync ancestor of the
  live PR head — and only that. Rewrites, rebases, force-pushes, non-ancestor
  heads, dirty worktrees, live/competing owners, and wrong repo/issue/branch/
  identity/profile all stay protected.

No existing exact-head, branch-protection, parity, workspace, identity, role,
or mutation-safety gate is weakened. All provenance is server-derived; nothing
is reachable from an MCP caller.

Tests: tests/test_issue_871_durable_lock_head_refresh.py (32 cases) covering
first/second sync, CAS, ownership, partial-failure, merge-sync recovery
happy-path and every fail-closed branch. Full suite: 4789 passed, 13 pre-
existing baseline failures unchanged.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_01FZPyVh2DGczQrDxtqwGH5p
2026-07-24 02:56:57 -04:00
sysadminandClaude Opus 4.8 a54b16676a merge(master): sync PR #861 with master; preserve F1–F10
Resolve conflicts from master (#858/#864/#868) while keeping published F10
PID-less lock freshness and dirty-orphan recovery (Issue #860 / PR #861).

Conflict resolutions:
- issue_lock_provenance.py: keep both SOURCE_RECOVER_DIRTY_ORPHANED and SOURCE_DIRTY_SAME_CLAIMANT_REBIND
- task_capability_map.py: keep both recover and rebind task entries
- gitea_mcp_server.py: keep both MCP tools
- issue_lock_store.py: restore F1–F9 helpers + surgical F10 owner_pid cascade

Closes #860

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-24 02:31:39 -04:00
sysadminandClaude Opus 4.8 2d0d8a682b feat(mcp-health): add MCP restart coordinator and impact analysis (Closes #658)
Child of umbrella #655 (governed MCP restart coordination); builds on the
#657 restart-path inventory. Adds a central coordinator that evaluates live
control-plane state before a restart and returns a blast-radius impact
preview, so operators and the web console (#642/#652) can see what a restart
would disrupt before concurrent LLM work is destroyed.

Changes
- restart_coordinator.py (new) — pure classification: inventory -> impact
  report DTO (RestartImpactReport/SessionImpact/LeaseImpact). Verdicts:
  safe / unsafe / override. Never restarts anything; fails closed on an
  incomplete inventory.
- control_plane_db.py — additive ControlPlaneDB.list_sessions() read-only
  session inventory (the process-level unit a restart kills).
- gitea_mcp_server.py — new dry-run MCP tool gitea_request_mcp_restart:
  gathers sessions/leases/terminal-lock from the #613 DB, calls the
  coordinator, returns the report. Override authority is read from the
  environment, never self-asserted (#630/#710 F1 pattern). Apply is gated
  by a later drain proof (non-goal here).
- docs/mcp-restart-coordinator.md + docs/mcp-restart-impact-sample.json — doc
  and a real dry-run sample report.
- docs/mcp-tool-inventory.md — register the new tool (inventory sync).
- tests/test_restart_coordinator.py (new) — 15 tests: multi-session fixtures,
  deny-when-critical-section-open, fail-closed deny, override, terminal lock,
  stale heartbeat, JSON-serializable DTO, list_sessions.

Tests: pytest tests/test_restart_coordinator.py -> 15 passed. Full suite:
13 failed / 4753 passed; all 13 reproduce identically on clean master
@ef14622 (0 regressions). The residual test_issue_781 doc-registry failure
is a pre-existing baseline gap for gitea_rebind_dirty_same_claimant_author_session
(merged #864, undocumented on master) — out of scope for #658.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-24 02:14:27 -04:00
jcwalker3 eb35c75514 fix(author): remediate test suite regression F10 for Issue #860 (#861) 2026-07-24 01:09:39 -05:00
jcwalker3 d542b08ced Merge branch 'master' into feat/issue-657-mcp-restart-path-inventory-guard 2026-07-24 00:26:15 -05:00
jcwalker3 499b87c482 Merge branch 'master' into fix/issue-855-pr-scoped-merged-cleanup 2026-07-24 00:25:40 -05:00
sysadminandClaude Opus 4.8 3428fb4190 feat(mcp-health): inventory and guard MCP restart/reload/kill paths (#657)
Enumerate every code/script/host path that can restart, reload, reconnect,
kill, or force-recreate an MCP process, classify each, and link it to the
guard that constrains it.

- mcp_restart_paths.py: machine-readable registry (single source of truth)
  with classifications (sanctioned_narrow / guarded_fail_closed / forbidden /
  removed / host_residual) plus fail-closed guards:
  * assert_restart_attempt_registered() -- unknown restart attempts fail closed
  * assert_no_daemon_self_replacement() -- daemon never os.execv/os.kill/os._exit
    itself (source-tree scan; comment/docstring mentions ignored)
  * assert_auto_restart_helper_absent() -- keeps the #685-removed
    _trigger_mcp_auto_restart from returning
  * assert_registry_wellformed() -- every path classified, guarded, referenced
- docs/mcp-restart-path-inventory.md: complete inventory table linked from
  #655; documents residual host behaviors (/mcp reconnect) and rollout.
- tests/test_mcp_restart_paths.py: 17 tests -- registry well-formedness,
  unknown-attempt fail-closed, daemon-self-replacement scan (with injected
  violation + comment/docstring negative case), legacy-helper-removed
  regression, pkill-stays-contamination (#630), and doc/module lock-step.

No behavior change to existing modules; regression assertions codify invariants
that already hold (per #657 flag-free-before-hard-block rollout). Links
#652 #653 #655 #656.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-24 00:57:33 -04:00
sysadminandClaude Opus 4.8 24c52abf6b feat(reconciler): PR-scoped post-merge cleanup executor + expired reviewer-lease reclaim (Closes #855)
Adds a single-target path for post-merge cleanup so a reconciler can
complete one merged PR without a batch sweep across unrelated PRs.

## PR-scoped selector

gitea_reconcile_merged_cleanups gains an optional pr_number. When set,
only that merged PR is assessed and acted on: the PR is resolved live and
fails closed on an invalid/non-positive number, an unresolvable or
ambiguous PR, or an unmerged PR; reviewer scratch worktrees are filtered
to that PR; and the report's entry set is pinned to exactly [pr_number],
failing closed on any drift. The existing execute loop then operates on
the single pinned entry only -- worktree removal, ownership reassessment,
then remote-branch delete -- with no unrelated target. Batch behaviour is
unchanged when pr_number is omitted.

## Expired reviewer-lease reclaim (AC4)

An expired or stale reviewer lease no longer protects an already-merged
branch forever. branch_cleanup_guard.assess_expired_reviewer_lease_reclaim
makes the decision explicitly and fail-closed: reclaim only when the lease
is a reviewer lease, its status is expired/stale, the PR is proven merged,
the owner process is proven dead, and no competing active claimant uses
the branch. _collect_branch_ownership_records supplies that evidence from
authoritative state (live PR merged-state, lease owner liveness, and the
full ownership inventory for competing-claimant detection) and evaluates
it only after the complete inventory is built, so the post-worktree-removal
reassessment is what unblocks the branch delete. Any unknown fails closed.

Author/merger/controller/reconciler leases are untouched; active leases,
worktree bindings, issue locks, and live sessions still block.

## Tests

- tests/test_issue_855_expired_reviewer_reclaim.py: full fail-closed matrix
  for the reclaim decision plus collector wiring (merged+dead+uncontested
  reclaims; unmerged, live-owner, competing-worktree, and author-lease
  cases stay protective).
- tests/test_branch_cleanup_guard.py: exact-PR selector coverage (ignores
  newer PRs in the batch queue, execute mutates only the selected PR,
  unknown/not-merged/invalid fail closed, batch mode preserved).

Changed-surface suites pass; the 2 pre-existing test_branch_cleanup_guard
failures and test_reconciler_supersession_close reproduce identically on
master 6d0015ca and are unrelated to this change.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-24 00:49:52 -04:00
jcwalker3 0b60fd6557 Remediate PR #861 findings F1-F9 and TestWorktreeStart failures (#860)
- Fix F1: Prepare recovery worktree detached at remote_head without git checkout -B to avoid exit 128 when branch is held by source worktree
- Fix F2: Pass recovery_sanctioned=True in bind_session_lock and assess_same_issue_lease_conflict
- Fix F3: Add SOURCE_RECOVER_DIRTY_ORPHANED to SANCTIONED_LOCK_SOURCES
- Fix F4: Stop after Phase 4 dirty apply when conflicts exist; do not finalize session binding
- Fix F5: Dynamically query competing live locks and workflow leases in MCP server
- Fix F6: Fail closed on recovery worktree resume when HEAD does not match expected remote_head
- Fix F7: Fail closed on remote HEAD observation failure rather than copying expected_remote_head pin
- Fix F8: Enforce foreign overwrite protection requiring same claimant or sanctioned reclaim
- Fix F9: Add real multi-worktree integration tests for prepare_recovery_worktree and lock rebind
- Fix TestWorktreeStart: Bypass session lock check for dry-run and review/pr-* branches in scripts/worktree-start
2026-07-23 22:34:35 -05:00
jcwalker3 67cd2da561 fix: remediate PR #853 review #528 findings for native MCP bootstrap (#850)
- Fix module reloading bug in task capability router (F-1)
- Harden journal persistence and pending creations crash window (F-3)
- Implement dirty worktree and author commit recovery preservation (F-4)
- Fail closed on missing identity, profile, or session parameters (F-5)
- Fix branches root path traversal and symlink validation (F-6)
- Enforce O_NOFOLLOW and symlink checking on transition locks (F-7)
- Support common ancestor merge-base verification for base SHA (F-8)
- Release transition lock on compensating recovery (F-10)
- Thread journal_dir through recovery and fix guidance strings (F-11, F-12)
- Fix unittest mock import in bootstrap test suite (F-13)
2026-07-23 20:40:09 -05:00
jcwalker3andGrok 4.5 18d6583e83 fix(author): bootstrap recovery for dirty orphaned issue worktrees (#860)
Add an explicit recovery operation for same-claimant dirty registered
worktrees under malformed PID-less durable locks, with crash-safe journals,
dirty byte preservation, path-level conflict detection, and live session
binding. PID-less locks are never treated as live merely because expiry is
absent.

Closes #860

Co-Authored-By: Grok 4.5 (xAI) <[email protected]>
2026-07-23 20:38:47 -05:00
sysadmin f80e3b33b0 Merge remote-tracking branch 'prgs/master' into feat/issue-639-webui-system-health-dashboard 2026-07-23 21:14:50 -04:00
jcwalker3 dc99c15ffa Merge branch 'master' into docs/issue-656-mcp-restart-governance 2026-07-23 19:53:16 -05:00
jcwalker3 e9f6d68bd7 Merge branch 'master' into fix/issue-850-native-mcp-bootstrap 2026-07-23 19:53:00 -05:00
jcwalker3 b3859f6dad Merge branch 'master' into fix/issue-850-native-mcp-bootstrap 2026-07-23 19:13:38 -05:00
sysadminandClaude Opus 4.8 edd5f813b2 Merge master into feat/issue-639-webui-system-health-dashboard
Resolve the #638 shell landing against the #639 dashboard:

- webui/layout.py: drop the flat NAV_ITEMS tuple in favor of master's
  grouped NAV_GROUPS nav-config module.
- webui/nav.py: register /system-health as a live item in the Health
  group, satisfying issue #639 AC5 through the canonical nav source.
- docs/webui-local-dev.md: keep both additive sections (#638 shell and
  #639 dashboard).
- tests/test_webui_system_health_dashboard.py: assert the nav entry via
  iter_nav_items() instead of the removed NAV_ITEMS tuple.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-23 19:56:57 -04:00
sysadminandClaude Opus 4.8 9f759150b8 docs(governance): correct linked-issue descriptions in restart ADR (#656)
The Related section and cross-reference lines described #652, #653, #630, #642,
and #591 by roles they do not hold. Align each description with the linked
issue's actual title and scope:

- #652 is the Control Plane Web Console product vision (restart controls live in
  its capability area A), not a restart-specific vision.
- #653 is the console phased-delivery roadmap; restart controls are Phase 2.
- #630 is the manual process-kill contamination guard, not the coordinator; the
  coordinator remains an unimplemented later child of #655.
- #642 is the sanctioned restart / graceful reload console UX.
- #591 is auto-restart on master advance (closed); only #584 is transport-flap
  reconnect. They were previously collapsed into one transport-recovery label.

Documentation-only wording change. Policy IDs RG-01..RG-08, the policy version
restart-governance/v1, the authorization matrix, and every normative statement
are unchanged.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-23 19:18:47 -04:00
sysadmin 347464a057 Merge branch 'master' into docs/issue-656-mcp-restart-governance 2026-07-23 19:17:05 -04:00
sysadminandClaude Opus 4.8 1301a57de4 docs(governance): MCP restart governance and authorization policy (#656)
Adds docs/architecture/mcp-restart-governance.md, the restart-governance/v1 ADR
defining who may restart the MCP control plane and under what conditions.

- Recovery ladder (reconnect -> rebind -> scoped restart -> full restart -> host)
  with restart stated as the last resort.
- Authorization matrix across author/reviewer/merger/reconciler/controller/
  operator/admin; no LLM worker role may perform or authorize a full or host
  restart.
- v1 authority decision recorded: controller approval + automated safety gates;
  quorum deferred to a superseding ADR.
- Break-glass path with pre-declared incident and mandatory post-hoc audit.
- Ambiguous policy state denies restart.
- Stable policy IDs RG-01..RG-08 for later enforcement code to bind to.

Cross-links the ADR from docs/safety-model.md and docs/webui-deployment.md, and
adds tests/test_mcp_restart_governance_docs.py asserting acceptance criteria 1-5.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-23 19:16:28 -04:00
jcwalker3 3b2b4e1dca Remediate PR #853 in response to review #525 for Issue #850 2026-07-23 17:49:31 -05:00
sysadmin a942afe6c4 Implement native author issue worktree bootstrap (#850) 2026-07-23 17:28:23 -04:00
sysadminandClaude Opus 4.8 ecda200180 feat(webui): system-health dashboard (Closes #639)
Phase 1 child of the Web Console epic #631. Adds the operator-facing
system-health dashboard on top of the read-only system-health API landed
by #634, so runtime problems are visible on a surface instead of being
discovered late through failed LLM sessions.

- webui/system_health_views.py (new): renders the SystemHealthSnapshot as
  readiness, stale-runtime parity, version/uptime, dependency, MCP
  namespace, probe-error, and recovery cards.
- webui/app.py: GET /system-health, sharing load_system_health() with the
  JSON API so page and API cannot disagree. ?deep=1 behaves as on the API.
- webui/layout.py: nav entry and health card/badge styles.
- tests/test_webui_system_health_dashboard.py (new, 26 cases).
- docs/webui-local-dev.md: route, field authority, and redaction split.

Readiness honesty is preserved from the API: ready and readiness_complete
render separately, a probe that did not run is listed under "Not probed"
rather than counted healthy, and mutation safety is never claimed when the
runtime is stale or parity is indeterminate.

Redaction is split by field kind. Free text (probe details, reasons, probe
errors) passes through system_health.redact. Structured fields (commit
SHAs, probe names, statuses, timestamps) are HTML-escaped only: redact's
opaque-token rule matches any run of 32 or more characters, so routing a
40-character git SHA through it rendered "[redacted]" and blanked the
parity evidence the page exists to show.

Non-goals honored: no restart or reload controls (Phase 2, #642), no
manual process-kill guidance (#630). Read-only throughout.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-23 16:08:36 -04:00
49 changed files with 10746 additions and 227 deletions
File diff suppressed because it is too large Load Diff
+80 -37
View File
@@ -40,22 +40,88 @@ def _normalize_path(path: str) -> str:
return (path or "").replace("\\", "/").rstrip("/")
def get_canonical_branches_root(project_root: str | None = None) -> str:
"""Return the absolute path of the canonical branches directory for *project_root*."""
root = os.path.realpath(project_root) if project_root else os.path.realpath(os.getcwd())
canonical_repo_root = resolve_canonical_repo_root(root, root)
return os.path.realpath(os.path.join(canonical_repo_root, "branches"))
def is_path_under_branches(path: str, project_root: str | None = None) -> bool:
"""True when *path* resolves inside ``<project_root>/branches/``."""
normalized = _normalize_path(path)
if not normalized:
"""True when *path* resolves inside a canonical ``branches/`` directory."""
if not path or not str(path).strip():
return False
if "/branches/" in f"{normalized}/":
return True
if normalized.endswith("/branches"):
return True
if project_root:
root = _normalize_path(os.path.realpath(project_root))
real = _normalize_path(os.path.realpath(path))
if real.startswith(f"{root}/"):
rel = real[len(root) + 1 :]
return rel == "branches" or rel.startswith("branches/")
return False
try:
real_path = os.path.realpath(os.path.abspath(str(path).strip()))
except Exception:
return False
branches_root = get_canonical_branches_root(project_root or real_path)
try:
common = os.path.commonpath([branches_root, real_path])
except Exception:
return False
if common != branches_root:
return False
rel = os.path.relpath(real_path, branches_root)
return rel != "." and not rel.startswith("..")
def resolve_canonical_repo_root(workspace_path: str, fallback_project_root: str) -> str:
"""Return the stable repository root for *workspace_path* via git metadata (#460)."""
p = (workspace_path or "").strip()
if p:
try:
res = subprocess.run(
["git", "-C", p, "rev-parse", "--git-common-dir"],
capture_output=True,
text=True,
check=True,
)
common = _realpath_git_common_dir(p, res.stdout)
if common.endswith(f"{os.sep}.git") or os.path.basename(common) == ".git":
candidate_root = os.path.dirname(common)
real_p = os.path.realpath(p)
try:
if os.path.commonpath([candidate_root, real_p]) == candidate_root:
return candidate_root
except Exception:
pass
except Exception:
pass
# Fallback when git metadata is unavailable. Never string-split on
# "/branches/" (review #531 F2 / #551): recover the repo root only via
# resolved-path commonpath ancestry. Do **not** require on-disk isdir —
# MCP may launch with project_root = branches/<wt> before that path
# exists, and #274 path-shaped worktree-as-project-root must still resolve.
fallback = os.path.realpath(fallback_project_root or workspace_path or ".")
cur = fallback
for _ in range(64):
parent = os.path.dirname(cur)
if parent == cur:
break
branches_dir = os.path.realpath(os.path.join(parent, "branches"))
try:
# Path-shaped: fallback is under parent/branches/ (commonpath).
if os.path.commonpath([branches_dir, fallback]) == branches_dir:
return parent
except ValueError:
pass
# Fallback path itself is the branches directory.
if os.path.basename(os.path.realpath(cur)) == "branches":
try:
if os.path.commonpath([os.path.realpath(cur), fallback]) == os.path.realpath(
cur
):
return parent
except ValueError:
pass
cur = parent
return fallback
def resolve_mutation_workspace(
@@ -87,29 +153,6 @@ def _realpath_git_common_dir(workspace_path: str, common_dir: str) -> str:
return os.path.realpath(os.path.join(workspace_path, raw))
def resolve_canonical_repo_root(workspace_path: str, fallback_project_root: str) -> str:
"""Return the stable repository root for *workspace_path* via git metadata (#460)."""
path = (workspace_path or "").strip()
fallback = os.path.realpath(fallback_project_root)
if not path:
return fallback
try:
res = subprocess.run(
["git", "-C", path, "rev-parse", "--git-common-dir"],
capture_output=True,
text=True,
check=True,
)
common = _realpath_git_common_dir(path, res.stdout)
except Exception:
return fallback
if common.endswith(f"{os.sep}.git"):
return os.path.dirname(common)
if os.path.basename(common) == ".git":
return os.path.dirname(common)
return fallback
def resolve_author_mutation_context(
worktree_path: str | None,
process_project_root: str,
+84
View File
@@ -525,6 +525,90 @@ def assess_ownership_record_activity(record: dict[str, Any]) -> dict[str, Any]:
}
# Reviewer-lease reclaim is only reachable from a non-live (expired/stale) lease.
_RECLAIMABLE_REVIEWER_STATUSES = _EXPIRED_STATUSES | _STALE_STATUSES
def is_active_ownership_status(status: str | None) -> bool:
"""True when *status* denotes live/active ownership of a branch (#855).
Used to decide whether a *competing* active claimant still uses a branch
when weighing an expired reviewer lease for reclaim. Expired, stale,
released, and terminal statuses are not active.
"""
return _norm_str(status).lower() in _ACTIVE_OWNERSHIP_STATUSES
def assess_expired_reviewer_lease_reclaim(
*,
role: str,
status: str,
pr_merged: bool | None,
owner_pid_alive: bool | None,
competing_active_claimant: bool | None,
) -> dict[str, Any]:
"""Decide, explicitly and fail-closed, whether an expired reviewer lease
may stop protecting an already-merged branch (#855 AC4).
An expired reviewer lease should not protect a merged branch forever once
its work is done and no live claimant remains. Reclaim is permitted only
when **every** condition below is provably satisfied; any unknown
(``None``) or contrary value keeps the lease protective:
- the lease is a ``reviewer`` lease (author/merger/controller/reconciler
leases are out of scope and always keep protecting);
- its status is expired or stale (never an active/live lease);
- the PR is proven merged (``pr_merged is True``);
- the lease owner process is proven dead (``owner_pid_alive is False``);
- no competing active claimant uses the branch
(``competing_active_claimant is False``).
Returns a decision dict with ``reclaim_allowed`` and, when refused, the
fail-closed ``reasons``. The reasons never contain secrets — only the
role, the status, and which condition was unproven.
"""
reasons: list[str] = []
normalized_role = _norm_str(role).lower()
normalized_status = _norm_str(status).lower()
if normalized_role != "reviewer":
reasons.append(
f"lease role '{normalized_role or 'unknown'}' is not a reviewer "
"lease; expired-reviewer reclaim does not apply"
)
if normalized_status not in _RECLAIMABLE_REVIEWER_STATUSES:
reasons.append(
f"lease status '{normalized_status or 'unknown'}' is not expired "
"or stale; only a non-live reviewer lease may be reclaimed"
)
if pr_merged is not True:
reasons.append(
"PR merged state is not proven true; reclaim requires an "
"already-merged PR (fail closed)"
)
if owner_pid_alive is not False:
reasons.append(
"lease owner process liveness is not proven dead; a live owner "
"still protects the branch (fail closed)"
)
if competing_active_claimant is not False:
reasons.append(
"a competing active claimant may still use the branch; reclaim "
"requires no other active ownership (fail closed)"
)
allowed = not reasons
return {
"reclaim_allowed": allowed,
"role": normalized_role,
"status": normalized_status,
"decision": (
"reclaim_expired_reviewer_lease" if allowed else "keep_protecting"
),
"reasons": [] if allowed else reasons,
}
def assess_active_branch_ownership(
*,
remote: str,
+260 -1
View File
@@ -31,7 +31,7 @@ from typing import Any, Iterator, Sequence
import dependency_graph
SCHEMA_VERSION = 4
SCHEMA_VERSION = 5
# Assignable work kinds only — raw monitoring incidents are never work items.
WORK_KINDS = frozenset({"issue", "pr"})
@@ -186,6 +186,34 @@ CREATE INDEX IF NOT EXISTS idx_dependency_edges_target
ON dependency_edges(remote, org, repo, target_kind, target_number);
CREATE INDEX IF NOT EXISTS idx_assignments_session ON assignments(session_id, status);
CREATE INDEX IF NOT EXISTS idx_incident_gitea ON incident_links(gitea_org, gitea_repo, gitea_issue_number);
-- Model usage, token cost, latency, and performance events (#651)
CREATE TABLE IF NOT EXISTS usage_events (
usage_id INTEGER PRIMARY KEY AUTOINCREMENT,
session_id TEXT,
remote TEXT NOT NULL DEFAULT 'dadeschools',
org TEXT NOT NULL DEFAULT '',
repo TEXT NOT NULL DEFAULT '',
project_id TEXT,
role TEXT NOT NULL DEFAULT 'unknown',
model TEXT NOT NULL DEFAULT 'unknown',
issue_number INTEGER,
pr_number INTEGER,
stage TEXT NOT NULL DEFAULT 'unknown',
input_tokens INTEGER,
output_tokens INTEGER,
total_tokens INTEGER,
estimated_cost_usd REAL,
latency_ms INTEGER,
duration_ms INTEGER,
status TEXT NOT NULL DEFAULT 'success',
metadata TEXT,
created_at TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_usage_events_scope ON usage_events(remote, org, repo);
CREATE INDEX IF NOT EXISTS idx_usage_events_role_model ON usage_events(role, model);
CREATE INDEX IF NOT EXISTS idx_usage_events_stage ON usage_events(stage);
"""
@@ -339,6 +367,7 @@ class ControlPlaneDB:
self._migrate_incident_links_null_scope(conn)
self._migrate_lease_lifecycle_columns(conn)
self._migrate_session_ownership_columns(conn)
self._migrate_usage_events_table(conn)
conn.execute(
"INSERT OR REPLACE INTO schema_meta(key, value) VALUES (?, ?)",
("schema_version", str(SCHEMA_VERSION)),
@@ -518,6 +547,207 @@ class ControlPlaneDB:
f"UPDATE incident_links SET {col} = '' WHERE {col} IS NULL"
)
def _migrate_usage_events_table(self, conn: sqlite3.Connection) -> None:
"""Create usage_events table and indexes if they do not exist (#651)."""
conn.execute("""
CREATE TABLE IF NOT EXISTS usage_events (
usage_id INTEGER PRIMARY KEY AUTOINCREMENT,
session_id TEXT,
remote TEXT NOT NULL DEFAULT 'dadeschools',
org TEXT NOT NULL DEFAULT '',
repo TEXT NOT NULL DEFAULT '',
project_id TEXT,
role TEXT NOT NULL DEFAULT 'unknown',
model TEXT NOT NULL DEFAULT 'unknown',
issue_number INTEGER,
pr_number INTEGER,
stage TEXT NOT NULL DEFAULT 'unknown',
input_tokens INTEGER,
output_tokens INTEGER,
total_tokens INTEGER,
estimated_cost_usd REAL,
latency_ms INTEGER,
duration_ms INTEGER,
status TEXT NOT NULL DEFAULT 'success',
metadata TEXT,
created_at TEXT NOT NULL
);
""")
conn.execute("CREATE INDEX IF NOT EXISTS idx_usage_events_scope ON usage_events(remote, org, repo);")
conn.execute("CREATE INDEX IF NOT EXISTS idx_usage_events_role_model ON usage_events(role, model);")
conn.execute("CREATE INDEX IF NOT EXISTS idx_usage_events_stage ON usage_events(stage);")
# #651 retention: cap growth so unauthenticated or high-volume ingest
# cannot DoS the control-plane DB (PR #876 F3). Applied after every write.
USAGE_EVENTS_MAX_ROWS = 10_000
USAGE_EVENTS_RETENTION_DAYS = 90
def record_usage_event(
self,
*,
session_id: str | None = None,
remote: str = "dadeschools",
org: str = "",
repo: str = "",
project_id: str | None = None,
role: str = "unknown",
model: str = "unknown",
issue_number: int | None = None,
pr_number: int | None = None,
stage: str = "unknown",
input_tokens: int | None = None,
output_tokens: int | None = None,
total_tokens: int | None = None,
estimated_cost_usd: float | None = None,
latency_ms: int | None = None,
duration_ms: int | None = None,
status: str = "success",
metadata: str | dict[str, Any] | None = None,
created_at: str | None = None,
) -> int:
"""Record a model usage, token cost, latency, or stage performance event (#651)."""
ts = created_at or _ts()
meta_str: str | None = None
if metadata is not None:
from webui import console_redaction
redacted_meta = console_redaction.redact_payload(metadata)
if isinstance(redacted_meta, str):
meta_str = redacted_meta
else:
try:
meta_str = json.dumps(redacted_meta, default=str)
except Exception:
meta_str = str(redacted_meta)
if total_tokens is None and (input_tokens is not None or output_tokens is not None):
total_tokens = (input_tokens or 0) + (output_tokens or 0)
with self._tx(immediate=True) as conn:
cursor = conn.execute(
"""
INSERT INTO usage_events (
session_id, remote, org, repo, project_id, role, model,
issue_number, pr_number, stage, input_tokens, output_tokens,
total_tokens, estimated_cost_usd, latency_ms, duration_ms,
status, metadata, created_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
session_id,
remote,
org,
repo,
project_id,
role,
model,
issue_number,
pr_number,
stage,
input_tokens,
output_tokens,
total_tokens,
estimated_cost_usd,
latency_ms,
duration_ms,
status,
meta_str,
ts,
),
)
usage_id = cursor.lastrowid
self._enforce_usage_events_retention(conn)
return usage_id
def _enforce_usage_events_retention(self, conn: sqlite3.Connection) -> None:
"""Drop aged and excess usage_events rows (PR #876 F3)."""
# Age-based: ISO-8601 UTC timestamps compare lexicographically.
cutoff = (
datetime.now(timezone.utc)
- timedelta(days=int(self.USAGE_EVENTS_RETENTION_DAYS))
).strftime("%Y-%m-%dT%H:%M:%SZ")
conn.execute(
"DELETE FROM usage_events WHERE created_at < ?",
(cutoff,),
)
# Count-based: keep the newest USAGE_EVENTS_MAX_ROWS by usage_id.
max_rows = int(self.USAGE_EVENTS_MAX_ROWS)
if max_rows > 0:
conn.execute(
"""
DELETE FROM usage_events
WHERE usage_id NOT IN (
SELECT usage_id FROM usage_events
ORDER BY usage_id DESC
LIMIT ?
)
""",
(max_rows,),
)
def query_usage_events(
self,
*,
remote: str | None = None,
org: str | None = None,
repo: str | None = None,
project_id: str | None = None,
role: str | None = None,
model: str | None = None,
issue_number: int | None = None,
pr_number: int | None = None,
stage: str | None = None,
session_id: str | None = None,
limit: int = 500,
offset: int = 0,
) -> list[dict[str, Any]]:
"""Query stored usage events matching filters (#651)."""
conditions = []
params = []
if remote:
conditions.append("remote = ?")
params.append(remote)
if org:
conditions.append("org = ?")
params.append(org)
if repo:
conditions.append("repo = ?")
params.append(repo)
if project_id:
conditions.append("project_id = ?")
params.append(project_id)
if role:
conditions.append("role = ?")
params.append(role)
if model:
conditions.append("model = ?")
params.append(model)
if issue_number is not None:
conditions.append("issue_number = ?")
params.append(issue_number)
if pr_number is not None:
conditions.append("pr_number = ?")
params.append(pr_number)
if stage:
conditions.append("stage = ?")
params.append(stage)
if session_id:
conditions.append("session_id = ?")
params.append(session_id)
where_clause = f"WHERE {' AND '.join(conditions)}" if conditions else ""
sql = f"""
SELECT * FROM usage_events
{where_clause}
ORDER BY usage_id ASC
LIMIT ? OFFSET ?
"""
params.extend([limit, offset])
with self._tx(immediate=False) as conn:
cursor = conn.execute(sql, params)
rows = cursor.fetchall()
return [dict(row) for row in rows]
# ── sessions ──────────────────────────────────────────────────────────
def upsert_session(
@@ -599,6 +829,35 @@ class ControlPlaneDB:
(_ts(), session_id),
)
def list_sessions(
self,
*,
statuses: Sequence[str] | None = None,
limit: int = 500,
) -> list[dict[str, Any]]:
"""List session rows for restart / impact analysis (#658).
Read-only. Sessions are the process-level unit an MCP restart
disrupts, so the restart coordinator inventories them to compute blast
radius. Optional ``statuses`` filter (e.g. ``('active',)``) narrows to
live rows. Never returns secrets — only operational metadata.
"""
clauses: list[str] = []
params: list[Any] = []
if statuses:
placeholders = ", ".join("?" for _ in statuses)
clauses.append(f"status IN ({placeholders})")
params.extend(statuses)
where = ("WHERE " + " AND ".join(clauses)) if clauses else ""
sql = (
f"SELECT * FROM sessions {where} "
"ORDER BY last_heartbeat_at DESC LIMIT ?"
)
params.append(max(1, int(limit)))
with self._tx(immediate=False) as conn:
rows = conn.execute(sql, params).fetchall()
return [dict(r) for r in rows]
# ── work items ────────────────────────────────────────────────────────
def upsert_work_item(
+2 -1
View File
@@ -247,7 +247,8 @@ def bootstrap_permits_control_checkout(
"""
if not isinstance(assessment, dict):
return False
if not is_create_issue_task(task):
import author_issue_bootstrap
if not is_create_issue_task(task) and not author_issue_bootstrap.is_author_issue_bootstrap_task(task):
return False
# Positive proof: the assessment must affirmatively allow, with no
File diff suppressed because it is too large Load Diff
+223
View File
@@ -0,0 +1,223 @@
# ADR: MCP restart governance and authorization policy
- **Status:** Accepted (policy effective immediately for LLM and operator sessions; enforcement tooling may lag)
- **Date:** 2026-07-23
- **Tracking issue:** [#656](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/656)
- **Policy version:** `restart-governance/v1`
- **Related:**
- Umbrella: [#655](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/655) — governed MCP restart coordination and zero-disruption recovery
- Vision: [#652](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/652) — MCP Control Plane Web Console product vision (§A system health and process control)
- Roadmap: [#653](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/653) — Control Plane Web Console phased delivery (Phase 2 restart controls)
- Contamination guard: [#630](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/630) — blocks manual process-kill recovery
- Console restart UX: [#642](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/642) — sanctioned restart and graceful reload
- Existing restart / reconnect paths to inventory: [#591](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/591) — auto-restart on master advance (closed); [#584](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/584) — host auto-reconnect on transport flap
- Stable-control runtime split: `docs/architecture/mcp-stable-control-runtime-policy-adr.md` (#615)
- Client-namespace health: `docs/mcp-namespace-health.md` (#543)
- Reconnect-only EOF recovery: `docs/mcp-namespace-eof-recovery.md`
## 1. Context
The Gitea MCP server is the **control plane** for real issue and PR mutations
(create, comment, lock, review, merge, reconcile). The same process serves every
role namespace (`gitea-author`, `gitea-reviewer`, `gitea-merger`,
`gitea-reconciler`, `gitea-controller`) and holds the in-memory capability-gate
code loaded at startup.
Restarting that process is destructive to concurrent work:
- It resets every session's identity, preflight, and capability-lease binding.
- It can interrupt a mutation mid-critical-section (a lock acquire, a review
submit, a merge), leaving durable state half-written.
- Relaunching from the wrong checkout or worktree silently changes which code
the control plane runs, defeating master-parity gates (#420 / #615).
Today there is **no durable written policy** stating who may restart MCP, under
what conditions, that restart is a last resort, and how controller approval,
automated safety gates, and break-glass interact. Operators and LLM sessions
therefore invent restart behavior ad hoc, which makes concurrent multi-role work
unsafe. #630 and #642 need this policy as their backbone.
This ADR defines that policy. It does **not** implement coordinator code or HA
multi-instance restart (those are later children of #655).
## 2. Decision
### 2.1 v1 decision (recorded)
**Restart authority in v1 is `controller approval + automated safety gates`.**
A restart of the stable control runtime is authorized only when **both** hold:
1. A **controller** role explicitly approves the restart, recording an audit
entry (who, why, scope, affected sessions), **and**
2. The **automated safety gates** pass: a completed drain acknowledgement (no
affected session is mid-critical-section) or a declared break-glass incident
(§2.5).
Quorum among multiple controllers is **not** required day-one. It is deferred
unless a later investigation (tracked under #653) proves single-controller
approval is insufficient. This ADR records the v1 decision so enforcement code
(#630) has a fixed target; changing it requires a superseding ADR.
### 2.2 Restart is a last resort — the recovery ladder
Restart is the **last** rung. Before any restart, exhaust the narrower
recoveries, in order:
1. **Reconnect** the IDE/client MCP namespace (transport EOF, `client is
closing: EOF`, transient `#584` flap). No process change. See
`docs/mcp-namespace-eof-recovery.md`.
2. **Refresh / rebind** the session workspace: re-run `gitea_whoami`,
`gitea_resolve_task_capability`, and pass an explicit validated
`worktree_path`. Fixes stale session context without touching the process.
3. **Scoped restart** of a single misbehaving namespace/service (where the
deployment supports per-service restart) rather than the whole control plane.
4. **Full restart** of the stable control runtime process — operator-owned,
controller-approved, drained.
5. **Host / infrastructure restart** — the broadest action; same authorization
as a full restart plus infrastructure ownership.
A session **must** try rungs 12 and record why they were insufficient before
requesting a restart at rung 3 or above. Skipping straight to restart is a
policy violation.
### 2.3 Authorization matrix
| Role | Reconnect (1) | Refresh/rebind (2) | Scoped restart (3) | Full restart (4) | Host restart (5) |
|---|---|---|---|---|---|
| **author** | self | self | request only | **forbidden** | forbidden |
| **reviewer** | self | self | request only | **forbidden** | forbidden |
| **merger** | self | self | request only | **forbidden** | forbidden |
| **reconciler** | self | self | request only | **forbidden** | forbidden |
| **controller** | self | self | **approve** (+gates) | **approve** (+gates) | request to operator |
| **operator** | self | self | execute (controller-approved) | execute (controller-approved) | execute (controller-approved) |
| **admin** | self | self | execute | execute | execute (break-glass) |
Legend: *self* = may perform for its own client session; *request only* = may
raise a restart request but not authorize or execute it; *approve* = may
authorize under §2.1 gates; *execute* = may perform the process action after the
authorization is recorded.
Key invariants:
- **No LLM worker role (author/reviewer/merger/reconciler) may perform or
authorize a full or host restart.** They may only reconnect/rebind their own
client and file a restart request.
- **Controller approval authorizes; operator/admin executes.** The approving
controller and the executing operator may be the same human, but both the
approval and the execution are audited.
- Privileged process actions (full restart, host restart) are reserved to
**operator/admin**, never to an automated worker.
### 2.4 Approved conditions
A restart at rung 3+ is approved only under one of these recorded conditions:
- **No affected sessions:** the control plane has no live session that would be
interrupted (verified, not assumed).
- **Full drain acknowledged:** every affected session has drained
(no open critical section — no held mutation lease mid-write) and the drain is
acknowledged in the audit record.
- **Controller + gates:** controller approval plus passing automated safety
gates (§2.1), the standard v1 path.
- **Quorum:** not required in v1; reserved for a future superseding ADR.
- **Break-glass:** an incident-backed emergency exception (§2.5).
Restart **never** bypasses mutation gates mid-critical-section. Drain before
restart is mandatory except under break-glass with a declared incident.
### 2.5 Break-glass
Break-glass is a **separate, narrower** authorization path for emergencies where
the normal drain-and-approve path cannot complete (e.g. the control plane is
wedged and cannot drain).
Break-glass conditions:
- A declared incident record exists (id, timestamp, declarer) **before** the
action.
- The action is taken by **operator or admin** authority only — never by an LLM
worker role, and never unilaterally by an operator with active peers when a
controller is reachable.
- The scope is the minimum necessary rung of the ladder.
- A **mandatory post-hoc audit** entry is filed: what was restarted, why the
normal path was impossible, which sessions were affected, and the incident id.
Break-glass suspends the drain requirement, not the audit requirement.
### 2.6 Explicit prohibitions
- **A unilateral LLM or operator full restart while active peer sessions
exist is forbidden.** An LLM worker role must not kill, restart, or relaunch
the MCP process; a lone operator must not full-restart over live peer work
without controller approval or a break-glass incident.
- Process-kill recovery is forbidden as a routine tool (#630). This ADR does not
introduce a kill path.
- Ambiguous policy state **denies** restart (§4).
## 3. Security requirements
- Full restart and host restart are **privileged**; only operator/admin execute
them, only after a controller approval or break-glass incident is recorded.
- Break-glass is a distinct authorization path with its own audit mandate; it is
never the default and never silent.
- **Every approval and every restart action is audited** (who approved, who
executed, scope, affected sessions, condition, policy version). No restart is
authorized without a durable audit entry.
## 4. Failure behavior
**Ambiguous policy → deny restart.** If it cannot be established that a
restart is authorized under §2 — unknown affected-session state, missing
controller approval, absent break-glass incident, or an unclassifiable request —
the safe action is to **refuse** the restart and stop with a recovery report,
never to restart on assumption.
## 5. Policy IDs (for enforcement code)
Enforcement code — the restart coordinator (a later child of #655), the #630
contamination guard, and the #642 console restart UX — binds to these stable
policy identifiers rather than to prose:
| Policy ID | Statement |
|---|---|
| `RG-01` | Restart is last resort; rungs 12 must be tried and recorded first (§2.2). |
| `RG-02` | v1 authority = controller approval + automated safety gates (§2.1). |
| `RG-03` | No LLM worker role performs or authorizes full/host restart (§2.3). |
| `RG-04` | Full/host restart executed by operator/admin only, post approval (§2.3). |
| `RG-05` | Drain before restart is mandatory except break-glass with incident (§2.4). |
| `RG-06` | Break-glass requires a pre-declared incident and post-hoc audit (§2.5). |
| `RG-07` | Unilateral LLM/operator full restart with active peers is forbidden (§2.6). |
| `RG-08` | Ambiguous policy state denies restart (§4). |
The `restart-governance/v1` **policy version** field is emitted on future
restart audit events so approvals can be reconciled against the policy revision
in force.
## 6. Dogfooding
Gitea-Tools governs its own MCP control plane by this policy. Author, reviewer,
merger, and reconciler sessions operating on this repository use the recovery
ladder (§2.2) — reconnect and rebind, never self-restart — and any real restart
of the Gitea-Tools stable control runtime follows the controller-approval +
drain path defined here.
## 7. Acceptance and cross-links
This ADR is the authoritative restart-governance policy. It **must** stay
cross-linked from the safety model and the web-console deployment boundary:
- `docs/safety-model.md` § Process restart governance references this ADR.
- `docs/webui-deployment.md` references this ADR for restart/reload disposition.
It is linked to its issue lineage — umbrella **#655**, vision **#652**, roadmap
**#653**, contamination guard **#630**, and console restart UX **#642** — in
§ Related above.
## 8. Non-goals
- Implementing the restart coordinator or approval state machine (#630, later
children of #655).
- Implementing HA multi-instance restart or quorum machinery.
- Introducing any process-kill or auto-restart tool; existing auto-restart
behavior must be inventoried before any new restart tool is enabled.
+95
View File
@@ -0,0 +1,95 @@
# MCP restart coordinator and impact analysis (#658)
Before any sanctioned MCP restart, a central coordinator evaluates the live
control-plane state and produces an **impact preview** so operators and the web
console (#642 / #652) can see the blast radius *before* concurrent LLM work is
disrupted. Uncoordinated restarts destroy in-flight author/reviewer/merger work
and give operators no way to see what they are about to break.
This lands the coordinator + impact DTO + a dry-run MCP tool. It is the single
sanctioned entry point for restart evaluation post-#657 (which inventoried the
restart/reload/kill paths). The **mutative apply** path — actually performing a
restart — is a later child gated by a drain proof and is explicitly out of
scope here.
## Components
| Piece | Where | Responsibility |
|-------|-------|----------------|
| `restart_coordinator.evaluate_restart_impact` | `restart_coordinator.py` | Pure classification: inventory → impact report DTO. No I/O, no restart. |
| `RestartImpactReport` / `SessionImpact` / `LeaseImpact` | `restart_coordinator.py` | Console-facing DTO (`.as_dict()` is JSON-serializable). |
| `ControlPlaneDB.list_sessions` | `control_plane_db.py` | Read-only session inventory (the process-level unit a restart kills). |
| `gitea_request_mcp_restart` | `gitea_mcp_server.py` | MCP tool: gathers inventory from the #613 DB, calls the coordinator, returns the report. Dry-run only. |
## Dimensions evaluated
The coordinator classifies the inventory across the dimensions #658 requires:
- **Sessions** — every active MCP session; a restart terminates all of them.
Liveness = `status == active` **and** the owner pid is alive **and** the
heartbeat is fresh (default window 15 min). Dead/stale sessions do not count
toward blast radius.
- **Leases / locks** — control-plane leases joined with work items and their
freshness (`lease_lifecycle.classify_lease_freshness`). Only `active` (live
owner) leases are *disruptive*; expired / released / dead-process leases never
withhold a restart.
- **Issue / PR work** — the issues and PRs behind disruptive leases.
- **Mutations / critical sections** — a live lease carrying an author worktree
or a mutating phase (`implementing`, `publishing`, `merging`, …) is a
critical section a restart must not sever.
- **Terminal (merge) lock** — an active terminal lock always makes a restart
unsafe.
- **Prior recovery attempts** — narrower recovery already tried (e.g. sanctioned
client reconnects) is echoed so the operator sees the escalation history.
## Verdict
Exactly three verdicts, matching the acceptance criteria:
| Verdict | `allow_restart` | Meaning |
|---------|-----------------|---------|
| `safe` | `true` | No other live sessions, no live leases, no terminal lock. |
| `unsafe` | `false` | Live work would be disrupted and no operator override is present — **or** the inventory could not be completed (fail closed). |
| `override` | `true` | Live work present, but an operator override accepts the blast radius. |
`override_would_allow` tells the console whether an override path exists for the
current state. `blast_radius` is a `none` / `low` / `medium` / `high` severity
band derived from the affected session and work counts.
### Fail closed
If the control-plane inventory cannot be completed (DB unavailable, a listing
failed), `inventory_complete` is `false` and the verdict is `unsafe` / deny. An
incomplete evaluation must never green-light a restart.
### Operator override authority
Override authority is read from the environment variable
`GITEA_OPERATOR_RESTART_OVERRIDE_AUTHORIZATION` and **never** from a tool
argument. A worker session cannot set an environment variable on an
already-running daemon, so override cannot be self-asserted (same pattern as the
#630 daemon-maintenance authorization). The `request_override` tool argument only
expresses caller intent; it takes effect solely when the environment
authorization is present.
## The tool
```text
gitea_request_mcp_restart(remote, host, org, repo,
dry_run=True, request_override=False,
session_id=None, limit=200)
```
Read-only, dry-run, and it **never restarts anything**. `apply_supported` is
always `false`; passing `dry_run=False` performs no restart and reports that
apply is gated by a drain proof (a separate child).
## Audit
Every evaluation carries an `audit_record` (event, coordinator version, verdict,
allow decision, blast radius, counts, timestamp) so restart decisions are
auditable. No secrets flow through the coordinator — session ids, pids, and
profiles are operational metadata only.
A representative dry-run report is in
[`mcp-restart-impact-sample.json`](./mcp-restart-impact-sample.json).
+148
View File
@@ -0,0 +1,148 @@
{
"coordinator_version": "1.0.0-issue-658",
"evaluated_at": "2026-07-24T06:00:00+00:00",
"dry_run": true,
"restart_performed": false,
"inventory_complete": true,
"incomplete_reasons": [],
"verdict": "unsafe",
"allow_restart": false,
"override_would_allow": true,
"operator_override": false,
"blast_radius": "high",
"reasons": [
"live work would be disrupted; restart denied without operator override",
"1 critical section(s) in flight (active lease with a live owner)"
],
"affected_sessions": [
{
"session_id": "prgs-author-30988-d6f43c25",
"role": "author",
"profile": "prgs-author",
"pid": 1,
"status": "active",
"alive": true,
"heartbeat_stale": false,
"is_requester": false,
"live": true
},
{
"session_id": "prgs-reviewer-4157-0ce9",
"role": "reviewer",
"profile": "prgs-reviewer",
"pid": 1,
"status": "active",
"alive": true,
"heartbeat_stale": false,
"is_requester": true,
"live": true
}
],
"affected_leases": [
{
"lease_id": "lease-abc",
"session_id": "prgs-author-30988-d6f43c25",
"role": "author",
"phase": "implementing",
"freshness": "active",
"work_kind": "issue",
"work_number": 658,
"worktree_path": "/repo/branches/feat-issue-658",
"disruptive": true,
"is_mutation": true,
"is_critical_section": true
},
{
"lease_id": "lease-dead",
"session_id": "prgs-author-91485",
"role": "author",
"phase": "allocated",
"freshness": "stale_dead_process",
"work_kind": "issue",
"work_number": 651,
"worktree_path": null,
"disruptive": false,
"is_mutation": false,
"is_critical_section": false
}
],
"critical_sections": [
{
"lease_id": "lease-abc",
"session_id": "prgs-author-30988-d6f43c25",
"role": "author",
"phase": "implementing",
"freshness": "active",
"work_kind": "issue",
"work_number": 658,
"worktree_path": "/repo/branches/feat-issue-658",
"disruptive": true,
"is_mutation": true,
"is_critical_section": true
}
],
"affected_issues": [
658
],
"affected_prs": [],
"mutations": [
{
"lease_id": "lease-abc",
"session_id": "prgs-author-30988-d6f43c25",
"role": "author",
"phase": "implementing",
"freshness": "active",
"work_kind": "issue",
"work_number": 658,
"worktree_path": "/repo/branches/feat-issue-658",
"disruptive": true,
"is_mutation": true,
"is_critical_section": true
}
],
"terminal_lock": null,
"ack_state": {
"prgs-author-30988-d6f43c25": "pending"
},
"prior_recovery_attempts": [
{
"kind": "client_reconnect",
"at": "2026-07-24T06:00:00+00:00",
"outcome": "insufficient"
}
],
"counts": {
"sessions_total": 2,
"sessions_live_other": 1,
"leases_total": 2,
"leases_disruptive": 1,
"critical_sections": 1,
"mutations": 1,
"affected_issues": 1,
"affected_prs": 0,
"prior_recovery_attempts": 1
},
"audit_record": {
"event": "restart_impact_evaluated",
"coordinator_version": "1.0.0-issue-658",
"evaluated_at": "2026-07-24T06:00:00+00:00",
"dry_run": true,
"operator_override": false,
"requesting_session_id": "prgs-reviewer-4157-0ce9",
"inventory_complete": true,
"verdict": "unsafe",
"allow_restart": false,
"blast_radius": "high",
"counts": {
"sessions_total": 2,
"sessions_live_other": 1,
"leases_total": 2,
"leases_disruptive": 1,
"critical_sections": 1,
"mutations": 1,
"affected_issues": 1,
"affected_prs": 0,
"prior_recovery_attempts": 1
}
}
}
+90
View File
@@ -0,0 +1,90 @@
# MCP restart / reload / kill path inventory (#657)
Complete inventory of every code, script, and host path that can **restart,
reload, reconnect, kill, or force-recreate** an MCP process in this project,
with each path classified and linked to the guard that constrains it.
This document is the human-readable companion to the machine-readable registry
in [`mcp_restart_paths.py`](../mcp_restart_paths.py). The two are kept in
lock-step by [`tests/test_mcp_restart_paths.py`](../tests/test_mcp_restart_paths.py):
every `path_id` below must appear in this file, and the source guards are run
against the live tree.
Roadmap linkage: this inventory is the enumeration step of the restart
governance work — parent **#655**, restart-governance ADR **#656**, vision
**#652**, roadmap **#653**. Related detection/guard work: master-advance
staleness **#591**/**#420**, side-effect-free resolver **#685**, transport flap
**#584**, manual-kill contamination **#630**.
## Classifications
| Classification | Meaning |
|---|---|
| `sanctioned_narrow_recovery` | One-shot, safe-by-construction recovery that never targets the running daemon. |
| `guarded_fail_closed` | Detects a restart-requiring condition, then fails mutations closed and emits reconnect guidance. Never self-restarts. |
| `forbidden` | A workflow-safety violation; where an LLM tool could invoke it, it is marked contamination. |
| `removed` | A previously-existing unguarded restart primitive that has been deleted; a regression guard keeps it absent. |
| `host_residual` | Behavior owned by the host/IDE, outside this process's control. Documented, not code-guarded here. |
## The rule
**No component may perform an unguarded full restart of the MCP daemon.** The
in-process daemon (`gitea_mcp_server.py`, `mcp_server.py`,
`role_session_router.py`) must never replace or terminate its own process:
replacing the process after the host has wired up the stdio pipes desyncs the
JSON-RPC transport (observed with Antigravity/Cascade hosts). Recovery is owned
by the host/operator via a client reconnect — the daemon only ever *detects*
and *fails closed*.
## Inventory
| path_id | Classification | Mechanism | Guard | Refs |
|---|---|---|---|---|
| `cli_venv_bootstrap_execv` | sanctioned_narrow_recovery | CLI wrapper scripts re-exec into `venv/bin/python3` via `os.execv`, guarded by `sys.executable != venv_python`. | One-shot pre-import bootstrap; runs before any MCP transport exists and only when not already on the venv interpreter; idempotent guard prevents a re-exec loop. | #657 |
| `daemon_self_replacement` | forbidden | The daemon replacing/terminating its own process (`os.execv`/`os.kill`/`os._exit`) to reload code. | Forbidden by design; enforced against the source tree by `assert_no_daemon_self_replacement()`. | #657, #584 |
| `legacy_auto_restart_helper` | removed | A helper (`_trigger_mcp_auto_restart`) that actively restarted the server from the read-only resolver path. | Removed in #685; kept absent by `assert_auto_restart_helper_absent()`. | #685, #657 |
| `config_touch_reload` | removed | Touching (utime) the MCP client config to make the host reload the server. | Removed from the resolver in #685: stale detection is report-only, never mutating config, spawning threads, or calling `os._exit`. | #685, #657 |
| `master_advance_auto_restart` | guarded_fail_closed | On-disk master advancing past the running code. | `master_parity_gate` captures startup parity and blocks mutations while stale, emitting restart guidance; the process never self-restarts. | #420, #591, #657 |
| `stale_runtime_resolver_reconnect` | guarded_fail_closed | The capability resolver detecting a stale serving process. | Report-only (#685): returns `restart_required`/`stop_required` and an exact reconnect action; no restart, thread, config touch, or `os._exit`. | #685, #657 |
| `manual_daemon_kill` | forbidden | Shell kills of the daemon: `pkill -f mcp_server.py`, `killall`, broad `pkill -f python` sweeps, or `kill <pid>` of a daemon pid. | Forbidden (#630): `runtime_recovery_guard` classifies these as contamination and `gitea_record_daemon_process_kill_attempt` writes a durable marker that fails later mutations closed. Operator maintenance authorization is read only from the environment. | #630, #657 |
| `conflict_marker_infra_stop` | guarded_fail_closed | The daemon entrypoint scans for unresolved merge-conflict markers at startup and stops (`sys.exit(1)`). | Fail-closed startup stop, not a restart: the process exits and waits for the operator to resolve conflicts and relaunch; never loops. | #657 |
| `ide_client_reconnect` | host_residual | A manual `/mcp reconnect` (or equivalent host action) that recreates the MCP client connection. | Outside this process's control; the sanctioned recovery the gates point operators toward. No in-process code initiates it. | #584, #656, #657 |
| `profile_switch_runtime` | sanctioned_narrow_recovery | Switching the active execution profile at runtime (dynamic-profile mode). | In-process and restart-free: `runtime_switching_supported` is true, so a switch rebinds capability without recreating the process. | #656, #657 |
## Guards enforced in CI
`tests/test_mcp_restart_paths.py` asserts, against the live source tree:
1. **Registry well-formedness** — every path has a valid classification, a
non-empty guard description, references, and locations; ids are unique; all
five classifications are represented.
2. **Unknown restart attempts fail closed**
`assert_restart_attempt_registered()` raises `UnknownRestartPathError` for
any path id not in this inventory, so a novel/unnamed restart primitive
cannot slip through silently.
3. **Daemon never self-replaces**`assert_no_daemon_self_replacement()` scans
the daemon modules for `os.execv`/`os.kill`/`os._exit`/`os.abort` calls
(comment/docstring mentions are ignored) and finds none.
4. **Legacy helper stays removed**`assert_auto_restart_helper_absent()`
confirms `_trigger_mcp_auto_restart` has not returned.
5. **pkill stays forbidden** — a daemon `pkill` command still classifies as
contamination via `runtime_recovery_guard`.
## Residual host behaviors (outside process control)
* `/mcp reconnect` in the IDE/host — the sanctioned recovery for stale-runtime,
transport-flap (#584), and worktree-binding conditions. The daemon can only
emit guidance toward it.
* Host-level process management (the operator relaunching the daemon after a
fail-closed stop, or after resolving merge conflicts).
These are documented rather than code-guarded because the process cannot
observe or gate them from inside itself.
## Rollout
Per #657, guards are introduced flag-free as **regression assertions** (they
codify invariants that already hold) before any hard runtime block is layered
on. When the restart coordinator (#655/#656) lands, registered paths gain a
coordinator token/capability check; unregistered attempts already fail closed
today via `assert_restart_attempt_registered()`.
+2
View File
@@ -69,6 +69,7 @@ that gates each call, not which tools exist.
- `gitea_audit_worktree_cleanup`
- `gitea_authorize_reconciliation_cleanup_phase`
- `gitea_authorize_review_correction`
- `gitea_bootstrap_author_issue_worktree`
- `gitea_capability_stop_terminal_report`
- `gitea_capture_branches_worktree_snapshot`
- `gitea_check_pr_eligibility`
@@ -135,6 +136,7 @@ that gates each call, not which tools exist.
- `gitea_release_merger_pr_lease`
- `gitea_release_reviewer_pr_lease`
- `gitea_release_workflow_lease`
- `gitea_request_mcp_restart`
- `gitea_resolve_task_capability`
- `gitea_resume_review_draft`
- `gitea_review_pr`
@@ -0,0 +1,143 @@
# Model Usage, Token Cost, Latency, and Workflow Analytics (Phase 4)
- **Tracking Issue:** [#651](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/651)
- **Parent Epic:** [#631](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/651)
- **Console Surface:** `/analytics`, `/api/v1/analytics`, `/api/v1/analytics/usage`
## 1. Overview
The Web Console Analytics module provides durable, aggregate visibility into **model usage, token cost, latency percentiles, and workflow-stage performance** across projects, worker roles, AI models, issues, and PRs.
### Non-Goals
- No mandatory client-side telemetry that leaks prompts or secret keys.
- No third-party payment provider or billing integration.
- No automatic model routing changes without controller policy (#647).
---
## 2. Event Schema (`usage_events`)
Usage metrics are stored in the control-plane database under table `usage_events`.
| Column | Type | Description |
|---|---|---|
| `usage_id` | `INTEGER` | Primary key (autoincrement) |
| `session_id` | `TEXT` | Optional active session identifier |
| `remote` | `TEXT` | Known Gitea instance (`dadeschools` or `prgs`) |
| `org` | `TEXT` | Repository owner / organization |
| `repo` | `TEXT` | Repository name |
| `project_id` | `TEXT` | Optional project identifier |
| `role` | `TEXT` | Active worker role (`author`, `reviewer`, `merger`, `reconciler`, `controller`) |
| `model` | `TEXT` | LLM model identifier (e.g. `gemini-3.6-flash`, `claude-3-5-sonnet`) |
| `issue_number` | `INTEGER` | Correlated Gitea issue number (optional) |
| `pr_number` | `INTEGER` | Correlated Gitea PR number (optional) |
| `stage` | `TEXT` | Workflow stage (`preflight`, `implementation`, `review`, `merge`, `reconciliation`) |
| `input_tokens` | `INTEGER` | Input token count (optional / nullable) |
| `output_tokens` | `INTEGER` | Output token count (optional / nullable) |
| `total_tokens` | `INTEGER` | Total token count (optional / nullable) |
| `estimated_cost_usd` | `REAL` | Estimated USD cost (optional / nullable) |
| `latency_ms` | `INTEGER` | Request latency in milliseconds (optional / nullable) |
| `duration_ms` | `INTEGER` | Stage execution duration in milliseconds (optional / nullable) |
| `status` | `TEXT` | Outcome status (`success`, `failure`, `timeout`) |
| `metadata` | `TEXT` | Redacted metadata or summary string |
| `created_at` | `TEXT` | ISO 8601 UTC timestamp |
---
## 3. Handling of Missing Data ("Unknown" vs. Zero Fabrication)
To ensure operational metrics accurately reflect evidence:
- **Untracked or missing metrics are displayed as `Unknown`**, never zero-fabricated.
- If an event omits `estimated_cost_usd`, `latency_ms`, or token counts, the aggregator marks those fields as missing (`None`) rather than defaulting to `0` or `$0.00`.
- Summary tables and KPI cards explicitly indicate when data is unmeasured or partially reported.
---
## 4. Redaction & Security Rules
Per `#633` security policy:
- Free-text fields (`metadata`, `prompt_summary`, `session_id`) are run through `console_redaction.redact_text` before persistence and output serialization.
- Secret tokens, keychain commands, authorization headers, passwords, and JWTs are stripped automatically.
---
## 5. Opt-in Instrumentation Guide
Applications, MCP servers, and background sessions can report usage metrics through either Python API or HTTP ingestion.
### Python Ingestion
```python
from webui.analytics_loader import record_usage
record_usage(
remote="dadeschools",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
role="author",
model="gemini-3.6-flash",
issue_number=651,
stage="implementation",
input_tokens=1420,
output_tokens=380,
total_tokens=1800,
estimated_cost_usd=0.00045,
latency_ms=320,
duration_ms=4500,
status="success",
metadata={"note": "Implementation of analytics module"},
)
```
### HTTP Ingestion API (authorized write)
`POST /api/v1/analytics/usage` is a **gated write**. It runs through
`console_authz` action `record_analytics_usage` (operator+, Phase 2 execution).
Unauthenticated or phase-inactive requests receive **403** and do not write.
Prefer in-process `record_usage` for MCP / session instrumentation.
```http
POST /api/v1/analytics/usage HTTP/1.1
Content-Type: application/json
# Requires authenticated principal with record_analytics_usage execution enabled
{
"remote": "dadeschools",
"org": "Scaled-Tech-Consulting",
"repo": "Gitea-Tools",
"role": "author",
"model": "gemini-3.6-flash",
"issue_number": 651,
"stage": "implementation",
"input_tokens": 1420,
"output_tokens": 380,
"total_tokens": 1800,
"estimated_cost_usd": 0.00045,
"latency_ms": 320,
"duration_ms": 4500,
"status": "success",
"metadata": "Analytics schema landed"
}
```
### Retention
`usage_events` is retained with hard caps applied on every write:
| Limit | Default |
|---|---|
| Max rows | 10,000 (`ControlPlaneDB.USAGE_EVENTS_MAX_ROWS`) |
| Max age | 90 days (`ControlPlaneDB.USAGE_EVENTS_RETENTION_DAYS`) |
Older rows (by `created_at`) and excess oldest rows (by `usage_id`) are deleted
after each insert so unbounded growth / DoS-by-volume cannot fill the DB.
---
## 6. Querying Analytics API
```http
GET /api/v1/analytics?role=author&stage=implementation HTTP/1.1
```
Returns `AnalyticsSnapshot` JSON containing aggregations (`by_model`, `by_stage`, `by_role`, `by_work_item`, `by_project`) and latency percentiles (`p50`, `p90`, `p95`, `p99`).
+14
View File
@@ -46,3 +46,17 @@ If shell helpers are unavailable and MCP commit cannot run, stop with a recovery
report (restart session, clear hung terminals, use MCP-native commit). See
[`llm-workflow-runbooks.md`](llm-workflow-runbooks.md) § MCP-native commit path
(#260) and agent temp artifact cleanup (#261).
## 7. Process restart governance
Restarting the MCP control-plane process is destructive to concurrent multi-role
work and is governed by a dedicated policy. Restart is a **last resort** behind
narrower recoveries (reconnect, rebind), full/host restart is reserved to
operator/admin under **controller approval + automated safety gates**, a
unilateral LLM or operator full restart with active peers is **forbidden**, and
ambiguous policy state **denies** restart. Break-glass is a separate,
incident-backed path with a mandatory audit.
See [`architecture/mcp-restart-governance.md`](architecture/mcp-restart-governance.md)
(#656) for the authorization matrix, the recovery ladder, break-glass
conditions, and the `RG-01``RG-08` policy IDs.
+1
View File
@@ -91,6 +91,7 @@ already define, and a regression test asserts each mapping matches.
| `close_pr` | controller | privileged | `gitea.pr.close` | Yes | No | No | 3 |
| `merge_pr` | controller | privileged | `gitea.pr.merge` | Yes | **Yes** | **Yes** | 3 |
| `delete_branch` | admin | destructive | `gitea.branch.delete` | Yes | **Yes** | **Yes** | 3 |
| `record_analytics_usage` | operator | gated_write | `runtime.record_analytics_usage` | Yes | No | No | 2 |
**Dual control** means the acting principal may not be the sole authority: a
second distinct principal must confirm. **Break-glass** means the action is
+9
View File
@@ -55,6 +55,15 @@ shipped to the browser.
assumption paths, and the client-secret policy. Use it to verify an instance is
configured for internal-only operation.
## Process restart / reload disposition
The console never exposes a restart or reload control; process restart of the
MCP control-plane runtime is governed separately. Restart is a last resort behind
reconnect/rebind, full restart is operator/admin-only under controller approval
plus safety gates, and break-glass is an incident-backed path. See
[`architecture/mcp-restart-governance.md`](architecture/mcp-restart-governance.md)
(#656).
## Non-goals (MVP)
- Full SSO or session login in the UI
+32
View File
@@ -54,6 +54,7 @@ status, onboarding checklist state, and the fail-closed error payloads (#635).
| `/` | Home / operator overview |
| `/health` | JSON liveness (`status`, `service`, `mode`, `timestamp`, `uptime_seconds`) |
| `/api/v1/system/health` | Structured read-only system health (#634) |
| `/system-health` | System-health dashboard — readiness, version/uptime, dependencies, MCP namespaces, stale-runtime parity (#639) |
| `/queue` | Live PR and issue queue dashboard (#429) |
| `/api/queue` | JSON queue export with pagination metadata |
| `/projects` | Project registry list with status and onboarding progress (#427, #635) |
@@ -258,6 +259,37 @@ Not-yet-implemented surfaces (`/sessions`, `/inventory`, `/timeline`,
surfaces are backed by #636). Mutating methods on stub routes still fail closed
with `read-only-mvp`.
## System-health dashboard (#639)
`/system-health` renders the same snapshot the `/api/v1/system/health` API
returns, so the page and the API can never disagree. Cards: overall readiness,
stale-runtime parity, version and uptime, dependency probes, MCP namespaces,
probe errors (only when present), and recovery pointers. `?deep=1` opts into
the network probe exactly as the API does; the plain page load stays cheap.
Field authority and honesty rules:
* `ready` and `readiness_complete` are shown separately. A snapshot whose
required probes never ran is not the same as one that ran them and passed,
and the page never collapses the two into an unproven green.
* A probe that did not run appears under **Not probed**, never as healthy.
* `stale_runtime.mutation_safe` is displayed verbatim from the API. When the
runtime is stale, or when parity is indeterminate, the page warns and does
not claim mutation safety.
* MCP namespaces are reported `unproven`: the web process runs outside the
IDE-managed MCP client and cannot prove that path (#543).
Redaction is split by field kind. Free text — probe details, readiness and
parity reasons, probe errors — passes through `system_health.redact`.
Structured fields — commit SHAs, probe names, statuses, timestamps — are
HTML-escaped only, because `redact`'s opaque-token rule matches any run of 32
or more characters and would otherwise blank every 40-character git SHA, which
is precisely the evidence the parity view exists to show.
The dashboard is read-only: no restart, reload, or process-kill control. Those
arrive in Phase 2 (#642). Recovery guidance points at the sanctioned client
reconnect / operator restart path — never a manual daemon kill (#630).
## Deployment boundary (#435)
MVP serves on loopback by default. Binding `0.0.0.0` or `::` is **refused**
+1001 -162
View File
File diff suppressed because it is too large Load Diff
+2
View File
@@ -16,6 +16,7 @@ ISSUE_LOCK_FILE = os.environ.get("GITEA_ISSUE_LOCK_FILE", "/tmp/gitea_issue_lock
SOURCE_LOCK_ISSUE = "gitea_lock_issue"
SOURCE_LOCK_ADOPTION = "gitea_lock_issue_adoption"
SOURCE_OPERATOR_OVERRIDE = "operator_override"
SOURCE_RECOVER_DIRTY_ORPHANED = "gitea_recover_dirty_orphaned_issue_worktree"
# #864: dirty-preserving same-claimant author-session rebind (dead owner PID).
SOURCE_DIRTY_SAME_CLAIMANT_REBIND = (
"gitea_rebind_dirty_same_claimant_author_session"
@@ -25,6 +26,7 @@ SANCTIONED_LOCK_SOURCES = frozenset({
SOURCE_LOCK_ISSUE,
SOURCE_LOCK_ADOPTION,
SOURCE_OPERATOR_OVERRIDE,
SOURCE_RECOVER_DIRTY_ORPHANED,
SOURCE_DIRTY_SAME_CLAIMANT_REBIND,
})
+130 -8
View File
@@ -85,6 +85,12 @@ HEAD_RELATION_STRICT_DESCENDANT = "strict_descendant"
# #772: an unpublished claim has no recorded head to compare against at all, so
# its head is measured against the base the branch was cut from instead.
HEAD_RELATION_DESCENDS_FROM_BASE = "descends_from_recorded_base"
# #871: the remote/PR head advanced *past* the recorded head via a sanctioned
# merge-based branch synchronization (``gitea_update_pr_branch_by_merge``) while
# the local worktree stayed at the recorded head. This is the inverse of the
# #768 descendant relation — here the *remote* strictly descends the local head,
# and only because a base was merged into the branch, proven server-side.
HEAD_RELATION_REMOTE_MERGE_SYNCED = "remote_merge_synced"
# Which body of evidence a recovery was decided on (#772 AC10). These are not
# interchangeable: a published claim proves ownership against a remote/PR head,
@@ -266,6 +272,70 @@ def _assess_base_descendancy(
]
def _assess_remote_merge_synced(
sync_provenance: Mapping[str, Any] | None,
*,
recorded_head: str,
remote_head: str,
) -> tuple[bool, list[str]]:
"""Did ``remote_head`` advance past ``recorded_head`` via a sanctioned
merge-based branch sync (#871)?
``sync_provenance`` is the server-side git observation from
``issue_lock_worktree.read_merge_sync_provenance``. Its own
``prior_head_sha`` / ``synced_head_sha`` are re-checked against the heads
this assessment is actually reasoning about, so an observation taken for some
other pair of commits stale, mismatched, or hand-built can never
authorize recovery. This is the inverse of ``_assess_strict_descendant``: the
recorded head is the ancestor and the *remote* head is the descendant, and it
is accepted only because the remote head is a base-into-branch merge that
preserved the branch mainline back to the recorded head.
Returns ``(proven, notes)``. Notes name the exact missing element so a
refused caller sees why, never a bare "unproven".
"""
if not isinstance(sync_provenance, Mapping):
return False, [
"no server-derived merge-sync provenance observation was available; a "
"remote head ahead of the recorded head cannot be accepted"
]
probe_prior = _text(sync_provenance.get("prior_head_sha"))
probe_synced = _text(sync_provenance.get("synced_head_sha"))
if probe_prior != recorded_head or probe_synced != remote_head:
return False, [
f"merge-sync observation covers {probe_prior or 'unknown'} -> "
f"{probe_synced or 'unknown'}, not the heads under assessment "
f"({recorded_head} -> {remote_head})"
]
if not sync_provenance.get("probe_ok"):
return False, (
list(sync_provenance.get("reasons") or [])
or ["merge-sync provenance probe did not complete; provenance unproven"]
)
if not sync_provenance.get("prior_is_ancestor"):
return False, [
f"recorded head {recorded_head} is not an ancestor of remote head "
f"{remote_head}; a rewritten or force-moved head cannot be recovered"
]
if not sync_provenance.get("is_merge_sync"):
return False, (
list(sync_provenance.get("reasons") or [])
or [
f"remote head {remote_head} is not a sanctioned merge-based sync "
f"of the base into the branch above {recorded_head}"
]
)
proof = _text(sync_provenance.get("proof")) or (
f"{remote_head} merged the base into the branch above {recorded_head}"
)
return True, [
f"remote head {remote_head} advanced past recorded head {recorded_head} "
f"via a sanctioned merge-based branch sync ({proof})"
]
def assess_dead_session_lock_recovery(
existing_lock: Mapping[str, Any] | None,
*,
@@ -290,6 +360,7 @@ def assess_dead_session_lock_recovery(
remote_branch_exists: bool | None = None,
recorded_base_sha: str | None = None,
base_ancestry: Mapping[str, Any] | None = None,
sync_provenance: Mapping[str, Any] | None = None,
) -> dict[str, Any]:
"""Decide whether a dead-session author lock may be natively recovered.
@@ -469,19 +540,39 @@ def assess_dead_session_lock_recovery(
head_relation = HEAD_RELATION_STRICT_DESCENDANT
ancestry_proof = notes[0] if notes else None
else:
reasons.append(
f"local head {local_head} does not match remote branch head "
f"{remote_head}"
# #871: the reverse relation — the remote head advanced past
# the recorded/local head via a sanctioned merge-based branch
# sync while the local worktree stayed put. Accepted only on
# server-proven merge-sync provenance, never a caller claim.
synced, sync_notes = _assess_remote_merge_synced(
sync_provenance,
recorded_head=local_head,
remote_head=remote_head,
)
reasons.extend(notes)
if synced:
head_relation = HEAD_RELATION_REMOTE_MERGE_SYNCED
ancestry_proof = sync_notes[0] if sync_notes else None
else:
reasons.append(
f"local head {local_head} does not match remote branch "
f"head {remote_head}"
)
reasons.extend(notes)
reasons.extend(sync_notes)
evidence["recorded_base"] = recorded_base or None
evidence["local_head"] = local_head or None
evidence["remote_head"] = remote_head or None
# ``recorded_head`` is the head recovery is being measured against;
# ``accepted_head`` is the head this recovery actually adopts. They differ
# only in the descendant case, and downstream gates need both (#768 AC2/AC7).
# #871: in the merge-sync case the branch/PR already carries the synced
# remote head, so that is the head recovery adopts; the local worktree stays
# at the ancestor recorded head.
evidence["recorded_head"] = remote_head or None
evidence["accepted_head"] = local_head or None
if head_relation == HEAD_RELATION_REMOTE_MERGE_SYNCED:
evidence["accepted_head"] = remote_head or None
else:
evidence["accepted_head"] = local_head or None
evidence["head_relation"] = head_relation
evidence["ancestry_proof"] = ancestry_proof
@@ -493,12 +584,17 @@ def assess_dead_session_lock_recovery(
# contradictory; re-stating it as a head mismatch would only obscure why.
if not unpublished and local_head and pr_head != local_head:
# A descendant recovery has not been published yet, so the open PR
# legitimately still points at the recorded head. Any other
# disagreement is a real mismatch.
# legitimately still points at the recorded head. A merge-sync
# recovery's PR legitimately sits at the advanced remote head. Any
# other disagreement is a real mismatch.
if not (
head_relation == HEAD_RELATION_STRICT_DESCENDANT
and remote_head
and pr_head == remote_head
) and not (
head_relation == HEAD_RELATION_REMOTE_MERGE_SYNCED
and remote_head
and pr_head == remote_head
):
reasons.append(
f"open PR #{pr_number} head {pr_head} does not match local head "
@@ -619,7 +715,11 @@ def assess_dead_session_lock_recovery(
)
if (
head_relation
in (HEAD_RELATION_STRICT_DESCENDANT, HEAD_RELATION_DESCENDS_FROM_BASE)
in (
HEAD_RELATION_STRICT_DESCENDANT,
HEAD_RELATION_DESCENDS_FROM_BASE,
HEAD_RELATION_REMOTE_MERGE_SYNCED,
)
and ancestry_proof
):
proof.append(ancestry_proof)
@@ -697,6 +797,16 @@ def owning_pr_recovery_evidence(
return None
if accepted_head != local_head:
return None
elif relation == HEAD_RELATION_REMOTE_MERGE_SYNCED:
# #871: the PR already sits at the advanced remote head; the local
# worktree is the ancestor the merge preserved. The head the open PR
# shows and the head recovery adopts are both the synced remote head.
if not remote_head or pr_head != remote_head:
return None
if accepted_head != remote_head:
return None
if not local_head or local_head == remote_head:
return None
else:
return None
try:
@@ -765,6 +875,18 @@ def recovered_owning_pr_from_lock(
return None
if not accepted_head or accepted_head == recorded_head:
return None
elif relation == HEAD_RELATION_REMOTE_MERGE_SYNCED:
# #871: PR sits at the advanced remote head, which is both the recorded
# measured-against head and the adopted head; the local worktree is the
# ancestor the merge preserved.
remote_head = _text(record.get("remote_head"))
local_head = _text(record.get("local_head"))
if not remote_head or pr_head != remote_head:
return None
if accepted_head and accepted_head != remote_head:
return None
if not local_head or local_head == remote_head:
return None
else:
return None
try:
+389 -6
View File
@@ -169,6 +169,7 @@ def bind_session_lock(
*,
expected_generation: int | None = None,
renewal_sanctioned: bool = False,
recovery_sanctioned: bool = False,
) -> str:
"""Persist a keyed lock and bind it to the current process session.
@@ -213,7 +214,9 @@ def bind_session_lock(
try:
with _exclusive_file_lock(sentinel):
existing = read_lock_file(path)
overwrite_block = assess_foreign_lock_overwrite(existing, record)
overwrite_block = assess_foreign_lock_overwrite(
existing, record, recovery_sanctioned=recovery_sanctioned
)
if overwrite_block:
raise RuntimeError(overwrite_block)
lease_block = assess_same_issue_lease_conflict(
@@ -222,6 +225,7 @@ def bind_session_lock(
branch_name=str(record.get("branch_name") or ""),
worktree_path=str(record.get("worktree_path") or ""),
renewal_sanctioned=renewal_sanctioned,
recovery_sanctioned=recovery_sanctioned,
)
if lease_block:
raise RuntimeError(lease_block)
@@ -380,7 +384,18 @@ def assess_lock_freshness(
pid = lock_data.get("session_pid")
if pid is None:
pid = lock_data.get("pid")
pid_alive = is_process_alive(pid) if pid is not None else False
if pid is None:
pid = lock_data.get("owner_pid")
pid_missing = pid is None or str(pid).strip() == ""
try:
pid_int = int(pid) if not pid_missing else None
if pid_int is not None and pid_int <= 0:
pid_missing = True
pid_int = None
except (TypeError, ValueError):
pid_missing = True
pid_int = None
pid_alive = is_process_alive(pid_int) if pid_int is not None else False
if expires_at and expires_at <= current:
return {
@@ -389,15 +404,36 @@ def assess_lock_freshness(
"stale": True,
"reason": f"lease expired at {expires_at.isoformat()}",
"pid_alive": pid_alive,
"pid_missing": pid_missing,
}
if pid is not None and not pid_alive:
# #860: a PID-less lock must never be considered live merely because
# expiration / heartbeat fields are absent. Missing PID is insufficient
# evidence of a live owner; treat as malformed/stale so recovery routes
# can evaluate corroborating pins instead of blocking on a false live flag.
if pid_missing:
return {
"status": "malformed",
"live": False,
"stale": True,
"reason": (
"lock has no usable session pid; cannot prove live ownership "
"(PID-less locks are never live by missing expiry alone)"
),
"pid_alive": False,
"pid_missing": True,
"heartbeat_at": heartbeat_at.isoformat() if heartbeat_at else None,
"expires_at": expires_at.isoformat() if expires_at else None,
}
if pid_int is not None and not pid_alive:
return {
"status": "stale",
"live": False,
"stale": True,
"reason": f"owner pid {pid} is not alive",
"reason": f"owner pid {pid_int} is not alive",
"pid_alive": False,
"pid_missing": False,
}
return {
@@ -406,6 +442,7 @@ def assess_lock_freshness(
"stale": False,
"reason": "lock heartbeat and lease are fresh",
"pid_alive": pid_alive,
"pid_missing": False,
"heartbeat_at": heartbeat_at.isoformat() if heartbeat_at else None,
"expires_at": expires_at.isoformat() if expires_at else None,
}
@@ -486,6 +523,7 @@ def assess_same_issue_lease_conflict(
worktree_path: str,
operation_type: str = AUTHOR_ISSUE_WORK_LEASE,
renewal_sanctioned: bool = False,
recovery_sanctioned: bool = False,
now: datetime | None = None,
) -> str | None:
"""Return a fail-closed error when a competing live lease blocks acquisition.
@@ -517,6 +555,8 @@ def assess_same_issue_lease_conflict(
existing_branch == branch_name
and _same_realpath(str(existing_worktree or ""), worktree_path)
)
if recovery_sanctioned and existing_issue == issue_number and existing_branch == branch_name:
return None
if is_lease_expired(existing_lock, now=now):
# #760 AC1/AC2: exact-owner renewal is a different disposition from
# foreign takeover and is evaluated first. Before this, both branches
@@ -547,10 +587,26 @@ def assess_same_issue_lease_conflict(
)
def _lock_claimant(lock: dict[str, Any] | None) -> dict[str, str]:
if not isinstance(lock, dict):
return {}
claimant = lock.get("claimant")
if not isinstance(claimant, dict):
lease = lock.get("work_lease")
claimant = lease.get("claimant") if isinstance(lease, dict) else None
if not isinstance(claimant, dict):
return {}
return {
"username": str(claimant.get("username") or ""),
"profile": str(claimant.get("profile") or ""),
}
def assess_foreign_lock_overwrite(
existing_lock: dict[str, Any] | None,
incoming_lock: dict[str, Any],
*,
recovery_sanctioned: bool = False,
now: datetime | None = None,
) -> str | None:
"""Block writes that would clobber an unrelated live lease on the same key."""
@@ -565,8 +621,31 @@ def assess_foreign_lock_overwrite(
)
if same_issue and same_branch and same_worktree:
return None
if not is_lease_live(existing_lock, now=now):
existing_claimant = _lock_claimant(existing_lock)
incoming_claimant = _lock_claimant(incoming_lock)
same_claimant = (
bool(existing_claimant.get("username"))
and existing_claimant.get("username") == incoming_claimant.get("username")
and existing_claimant.get("profile") == incoming_claimant.get("profile")
)
if recovery_sanctioned and same_issue and same_branch and same_claimant:
return None
if not is_lease_live(existing_lock, now=now):
# #860 F8: A non-live or PID-less lock still blocks foreign overwrite
# unless same claimant or sanctioned reclaim is proven.
if not same_claimant and same_issue:
reclaim = assess_expired_lock_reclaim(existing_lock, now=now)
if not reclaim.get("reclaim_allowed"):
return (
"Refusing foreign overwrite of non-live issue lock "
f"(issue #{existing_lock.get('issue_number')}, owner '{existing_claimant.get('username')}') "
"without sanctioned reclaim proof (fail closed)"
)
return None
return (
"Refusing to overwrite a live foreign issue lock "
f"(issue #{existing_lock.get('issue_number')}, "
@@ -737,4 +816,308 @@ def format_lock_proof(
parts.append("lock released")
elif released is False:
parts.append("lock retained")
return "; ".join(parts)
return "; ".join(parts)
# ── #871: durable linked-issue lock head refresh after branch synchronization ──
_FULL_SHA_RE = re.compile(r"^[0-9a-f]{40}$", re.IGNORECASE)
# Provenance recorded on the lock when the head is refreshed by a sanctioned
# merge-based branch synchronization (``gitea_update_pr_branch_by_merge``).
LOCK_HEAD_REFRESH_PROVENANCE_MERGE_SYNC = "gitea_update_pr_branch_by_merge"
def _norm_sha(value: Any) -> str | None:
text = str(value or "").strip().lower()
return text if _FULL_SHA_RE.match(text) else None
def _lock_claimant_view(lock: dict[str, Any] | None) -> dict[str, Any]:
if not isinstance(lock, dict):
return {}
claimant = lock.get("claimant")
if not isinstance(claimant, dict):
lease = lock.get("work_lease")
claimant = lease.get("claimant") if isinstance(lease, dict) else None
return dict(claimant) if isinstance(claimant, dict) else {}
def assess_durable_lock_head_refresh(
existing_lock: dict[str, Any] | None,
*,
remote: str,
org: str,
repo: str,
issue_number: int,
branch_name: str,
worktree_path: str,
pr_number: int | None,
identity: str | None,
profile: str | None,
current_pid: int | None,
expected_old_head: str | None,
new_head: str | None,
base_head: str | None = None,
) -> dict[str, Any]:
"""Fail-closed assessment for refreshing a durable lock's recorded head (#871).
A successful ``gitea_update_pr_branch_by_merge`` advances the *remote* PR head
but must also advance the durable linked-issue lock so a later dead-session
recovery can prove ownership. This decides whether that refresh is permitted;
it mutates nothing.
Every element of durable ownership is re-verified against the persisted lock
repository, issue, branch, worktree, claimant identity/profile, and the live
owning session and the recorded head is compare-and-swapped: the lock's
currently recorded synced head (if any) must equal ``expected_old_head``, so a
lock whose head or provenance changed concurrently is never overwritten.
"""
reasons: list[str] = []
old = _norm_sha(expected_old_head)
new = _norm_sha(new_head)
evidence: dict[str, Any] = {
"issue_number": issue_number,
"branch_name": branch_name,
"worktree_path": worktree_path,
"pr_number": pr_number,
"expected_old_head": old,
"new_head": new,
"base_head": _norm_sha(base_head),
}
if not isinstance(existing_lock, dict) or not existing_lock:
reasons.append("no durable lock exists for this issue; nothing to refresh")
return {"allowed": False, "reasons": reasons, "evidence": evidence,
"expected_generation": 0}
lock = dict(existing_lock)
evidence["current_generation"] = lock_generation(lock)
if lock.get("issue_number") != issue_number:
reasons.append(
f"durable lock targets issue #{lock.get('issue_number')}, not "
f"#{issue_number}; refusing head refresh"
)
for field, expected in (("remote", remote), ("org", org), ("repo", repo)):
actual = str(lock.get(field) or "").strip()
if actual != str(expected or "").strip():
reasons.append(
f"lock {field} '{actual}' does not match requested "
f"'{str(expected or '').strip()}'"
)
locked_branch = str(lock.get("branch_name") or "").strip()
if locked_branch != str(branch_name or "").strip():
reasons.append(
f"lock branch '{locked_branch}' does not match requested "
f"'{str(branch_name or '').strip()}'"
)
locked_worktree = str(lock.get("worktree_path") or "").strip()
try:
same_wt = bool(locked_worktree) and bool(worktree_path) and (
os.path.realpath(locked_worktree) == os.path.realpath(worktree_path)
)
except OSError:
same_wt = locked_worktree == (worktree_path or "")
if not same_wt:
reasons.append(
f"lock worktree '{locked_worktree}' does not match declared "
f"'{str(worktree_path or '').strip()}'"
)
claimant = _lock_claimant_view(lock)
locked_identity = str(claimant.get("username") or "").strip()
locked_profile = str(claimant.get("profile") or "").strip()
if not locked_identity or not locked_profile:
reasons.append(
"durable lock does not record a claimant identity/profile; "
"ownership could not be proven for head refresh"
)
if not str(identity or "").strip() or not str(profile or "").strip():
reasons.append(
"active session identity/profile is unknown; ownership could not be "
"proven for head refresh"
)
if locked_identity and str(identity or "").strip() and locked_identity != str(identity).strip():
reasons.append(
f"lock claimant '{locked_identity}' does not match active identity "
f"'{str(identity).strip()}'"
)
if locked_profile and str(profile or "").strip() and locked_profile != str(profile).strip():
reasons.append(
f"lock profile '{locked_profile}' does not match active profile "
f"'{str(profile).strip()}'"
)
# The refresh is written by the LIVE owning author session. A refresh is not
# a recovery: the current process must be the recorded owner.
recorded_pid = lock.get("session_pid")
if recorded_pid is None:
recorded_pid = lock.get("pid")
evidence["recorded_pid"] = recorded_pid
evidence["current_pid"] = current_pid
if current_pid is None:
reasons.append("current session pid is unknown; cannot prove live ownership")
else:
try:
if recorded_pid is None or int(recorded_pid) != int(current_pid):
reasons.append(
f"durable lock is owned by pid {recorded_pid}, not the current "
f"session pid {current_pid}; head refresh requires the live owner"
)
except (TypeError, ValueError):
reasons.append(
"durable lock owner pid is malformed; cannot prove live ownership"
)
if not old:
reasons.append("expected_old_head is not a full 40-char hex SHA (fail closed)")
if not new:
reasons.append("new_head is not a full 40-char hex SHA (fail closed)")
if old and new and old == new:
reasons.append(
"new head equals the expected old head; a sync must advance the head"
)
# Compare-and-swap on the recorded head: if the lock already records a synced
# head it must be exactly the expected old head, else another sync moved it.
recorded_synced = _norm_sha(lock.get("synced_pr_head"))
evidence["recorded_synced_pr_head"] = recorded_synced
if recorded_synced is not None and old is not None and recorded_synced != old:
reasons.append(
f"durable lock already records synced head {recorded_synced}, not the "
f"expected old head {old}; a concurrent sync changed it (CAS fail closed)"
)
if reasons:
return {"allowed": False, "reasons": reasons, "evidence": evidence,
"expected_generation": lock_generation(lock)}
return {
"allowed": True,
"reasons": [
f"durable lock for issue #{issue_number} branch '{locked_branch}' is "
f"owned by the live session; refresh recorded head {old} -> {new}"
],
"evidence": evidence,
"expected_generation": lock_generation(lock),
}
def apply_durable_lock_head_refresh(
*,
remote: str,
org: str,
repo: str,
issue_number: int,
branch_name: str,
worktree_path: str,
pr_number: int | None,
identity: str | None,
profile: str | None,
current_pid: int | None,
expected_old_head: str | None,
new_head: str | None,
synced_at: str,
base_head: str | None = None,
provenance: str = LOCK_HEAD_REFRESH_PROVENANCE_MERGE_SYNC,
lock_dir: str | None = None,
) -> dict[str, Any]:
"""CAS-refresh the durable lock's recorded head after a branch sync (#871).
Reads the durable lock from disk, re-asserts ownership via
``assess_durable_lock_head_refresh``, and only when permitted writes the
new synced head through ``bind_session_lock`` with a generation compare-and-
swap. Then re-reads the lock and proves it records the complete new head
(read-after-write). Any failure at any step returns ``refreshed=False`` with
reasons; the caller must treat that as a partial lifecycle failure and never
report a fully successful synchronization.
"""
existing = load_issue_lock(
remote=remote, org=org, repo=repo, issue_number=issue_number, lock_dir=lock_dir
)
assessment = assess_durable_lock_head_refresh(
existing,
remote=remote,
org=org,
repo=repo,
issue_number=issue_number,
branch_name=branch_name,
worktree_path=worktree_path,
pr_number=pr_number,
identity=identity,
profile=profile,
current_pid=current_pid,
expected_old_head=expected_old_head,
new_head=new_head,
base_head=base_head,
)
result: dict[str, Any] = {
"refreshed": False,
"read_after_write_ok": False,
"prior_head": _norm_sha(expected_old_head),
"new_head": _norm_sha(new_head),
"reasons": list(assessment.get("reasons") or []),
"evidence": assessment.get("evidence"),
}
if not assessment.get("allowed"):
return result
new = _norm_sha(new_head)
old = _norm_sha(expected_old_head)
record = dict(existing or {})
sync_block = {
"last_synced_pr_head": new,
"prior_pr_head": old,
"base_head": _norm_sha(base_head),
"pr_number": pr_number,
"provenance": provenance,
"synced_at": synced_at,
"synced_by_pid": current_pid,
"synced_by": {
"username": str(identity or "").strip() or None,
"profile": str(profile or "").strip() or None,
},
}
record["synced_pr_head"] = new
record["branch_sync"] = sync_block
history = record.get("branch_sync_history")
if not isinstance(history, list):
history = []
history = list(history)
history.append(sync_block)
record["branch_sync_history"] = history
try:
bind_session_lock(
record,
lock_dir=lock_dir,
expected_generation=assessment.get("expected_generation"),
renewal_sanctioned=True,
)
except Exception as exc: # CAS miss or write failure — partial lifecycle failure
result["reasons"].append(
f"durable lock head refresh write failed (fail closed): {exc}"
)
return result
after = load_issue_lock(
remote=remote, org=org, repo=repo, issue_number=issue_number, lock_dir=lock_dir
)
after_head = _norm_sha((after or {}).get("synced_pr_head"))
result["lock_generation_after"] = lock_generation(after)
if after_head == new and new is not None:
result["refreshed"] = True
result["read_after_write_ok"] = True
result["reasons"].append(
f"durable lock recorded head refreshed to {new} and verified by "
"read-after-write"
)
else:
result["reasons"].append(
"read-after-write verification failed: durable lock does not record "
f"the new head {new} (found {after_head}); partial lifecycle failure"
)
return result
+169
View File
@@ -145,6 +145,175 @@ def read_head_ancestry(
return result
def read_merge_sync_provenance(
worktree_path: str,
*,
prior_head_sha: str | None,
synced_head_sha: str | None,
) -> dict:
"""Observe whether ``synced_head_sha`` is a sanctioned merge-based branch sync
that advanced the PR branch past ``prior_head_sha`` (#871).
``gitea_update_pr_branch_by_merge`` advances a PR branch by merging the base
branch *into* the branch (``POST /pulls/{n}/update?style=merge``). The result
is a merge commit ``M`` on the branch whose **first** parent is the prior
branch head and whose second parent is the base tip. When the owning session
then dies without the durable lock's recorded head being refreshed, the local
worktree still sits at ``prior_head_sha`` while the live PR head is ``M``.
Recovering that drift safely requires proving the remote head is *exactly*
such a merge-sync not a rewrite, rebase, force-push, or an unrelated
commit. This is that server-side observation. It reports facts only; the
disposition lives in ``issue_lock_recovery``. Every field is read from git in
the declared worktree nothing is supplied by, or reachable from, an MCP
caller (#871).
Provenance is proven only when ALL hold:
* both commits are present (a rewritten/force-moved prior head leaves the
object graph and fails closed);
* ``prior_head_sha`` is a strict ancestor of ``synced_head_sha`` (the branch
history is preserved, never replaced);
* ``synced_head_sha`` is a merge commit (two or more parents), i.e. a base
merged in a plain fast-forward of new direct commits is not a sync;
* ``prior_head_sha`` is an ancestor of the merge's **first** parent, so the
branch mainline (first-parent lineage) still reaches the prior head a
rebase/force-push that re-authored the branch side fails this.
"""
path = (worktree_path or "").strip()
prior = (prior_head_sha or "").strip()
synced = (synced_head_sha or "").strip()
result: dict = {
"prior_head_sha": prior or None,
"synced_head_sha": synced or None,
"probe_ok": False,
"prior_present": False,
"synced_present": False,
"prior_is_ancestor": False,
"synced_is_merge": False,
"first_parent_reaches_prior": False,
"is_merge_sync": False,
"first_parent_sha": None,
"parent_count": None,
"proof": None,
"reasons": [],
}
if not path or not prior or not synced:
result["reasons"].append(
"merge-sync provenance probe requires a worktree path and both "
"commit SHAs"
)
return result
if prior == synced:
result["reasons"].append(
"prior and synced heads are identical; no branch sync occurred"
)
return result
def _present(sha: str) -> bool:
res = subprocess.run(
["git", "-C", path, "rev-parse", "--verify", "--quiet", f"{sha}^{{commit}}"],
capture_output=True,
text=True,
check=False,
)
return res.returncode == 0
def _is_ancestor(ancestor: str, descendant: str) -> bool | None:
res = subprocess.run(
["git", "-C", path, "merge-base", "--is-ancestor", ancestor, descendant],
capture_output=True,
text=True,
check=False,
)
if res.returncode == 0:
return True
if res.returncode == 1:
return False
return None # failed probe — never a silent "no"
try:
result["prior_present"] = _present(prior)
result["synced_present"] = _present(synced)
except OSError as exc: # git unavailable — fail closed, never assume
result["reasons"].append(f"merge-sync provenance probe could not run: {exc}")
return result
if not result["prior_present"]:
result["reasons"].append(
f"prior head {prior} is not reachable in '{path}'; history may have "
"been rewritten or force-moved"
)
if not result["synced_present"]:
result["reasons"].append(
f"synced head {synced} is not reachable in '{path}'"
)
if not (result["prior_present"] and result["synced_present"]):
return result
ancestor = _is_ancestor(prior, synced)
if ancestor is None:
result["reasons"].append(
"ancestry probe failed; merge-sync provenance unproven"
)
return result
result["prior_is_ancestor"] = bool(ancestor)
if not ancestor:
result["reasons"].append(
f"prior head {prior} is not an ancestor of synced head {synced}; "
"the branch history was not preserved (not a merge-based sync)"
)
return result
parents_res = subprocess.run(
["git", "-C", path, "rev-list", "--parents", "-n", "1", synced],
capture_output=True,
text=True,
check=False,
)
if parents_res.returncode != 0:
result["reasons"].append(
f"could not read parents of {synced}; merge-sync provenance unproven"
)
return result
tokens = (parents_res.stdout or "").split()
# tokens[0] is the commit itself; the rest are its parents.
parents = tokens[1:]
result["parent_count"] = len(parents)
result["synced_is_merge"] = len(parents) >= 2
if not result["synced_is_merge"]:
result["probe_ok"] = True
result["reasons"].append(
f"synced head {synced} has {len(parents)} parent(s); a merge-based "
"branch sync produces a merge commit (two or more parents)"
)
return result
first_parent = parents[0]
result["first_parent_sha"] = first_parent
fp_reaches = _is_ancestor(prior, first_parent) if prior != first_parent else True
if fp_reaches is None:
result["reasons"].append(
"first-parent ancestry probe failed; merge-sync provenance unproven"
)
return result
result["first_parent_reaches_prior"] = bool(fp_reaches)
result["probe_ok"] = True
if not fp_reaches:
result["reasons"].append(
f"merge first parent {first_parent} does not reach prior head "
f"{prior}; the branch mainline was re-authored (not a sanctioned sync)"
)
return result
result["is_merge_sync"] = True
result["proof"] = (
f"synced head {synced} is a merge commit (parents={len(parents)}) whose "
f"first-parent lineage reaches prior head {prior}; base merged into branch"
)
return result
def read_recorded_base(
worktree_path: str,
*,
+475
View File
@@ -0,0 +1,475 @@
"""Inventory and fail-closed guards for MCP restart/reload/kill paths (#657).
Single source of truth enumerating every code/script/doc path that can
restart, reload, reconnect, kill, or force-recreate an MCP process. Each path
is classified and linked to the guard that constrains it. The companion
human-readable inventory lives in ``docs/mcp-restart-path-inventory.md`` and is
kept in lock-step with this module by ``tests/test_mcp_restart_paths.py``.
Design intent (aligns with #655 restart-coordinator roadmap):
* **No unguarded full restart.** The in-process MCP daemon
(``gitea_mcp_server.py`` / ``mcp_server.py`` / ``role_session_router.py``)
must never replace or kill its own process replacing the process after the
host wired up the stdio pipes desyncs the JSON-RPC transport (observed with
Antigravity/Cascade hosts). ``assert_no_daemon_self_replacement`` enforces
this against the live source tree.
* **No legacy auto-restart helper.** ``_trigger_mcp_auto_restart`` was removed
when the stale-runtime resolver became side-effect free (#685);
``assert_auto_restart_helper_absent`` keeps it removed.
* **Unknown restart attempts fail closed.** LLM tools must route any restart
intent through a *registered* path. ``assert_restart_attempt_registered``
raises ``UnknownRestartPathError`` for anything not in this inventory.
* **pkill stays forbidden (#630).** Manual daemon kills are classified as
contamination by :mod:`runtime_recovery_guard`; this module records that path
and the test asserts the classification still holds.
This module performs no restarts, spawns no threads, and touches no config or
process state. It is pure inventory + read-only source assertions.
"""
from __future__ import annotations
import os
from dataclasses import dataclass
from pathlib import Path
from typing import Iterable
# --- Classifications -------------------------------------------------------
#: A narrow, one-shot recovery that is safe by construction (e.g. a CLI wrapper
#: re-execing into the venv interpreter before importing anything, or an
#: in-process profile switch). Never targets the running MCP daemon process.
CLASS_SANCTIONED_NARROW = "sanctioned_narrow_recovery"
#: The path detects a condition that would require a restart, then *fails
#: closed* on mutations and emits restart/reconnect guidance. It never restarts
#: the process itself (recovery is owned by the host/operator).
CLASS_GUARDED_FAIL_CLOSED = "guarded_fail_closed"
#: The path is forbidden. Attempting it is a workflow-safety violation and,
#: where an LLM tool could invoke it, is marked as contamination.
CLASS_FORBIDDEN = "forbidden"
#: A previously-existing unguarded restart primitive that has been deleted. A
#: regression guard keeps it absent.
CLASS_REMOVED = "removed"
#: Behavior that lives in the host/IDE and is outside this process's control
#: (e.g. a manual ``/mcp reconnect``). Documented, not code-guarded here.
CLASS_HOST_RESIDUAL = "host_residual"
VALID_CLASSIFICATIONS = frozenset(
{
CLASS_SANCTIONED_NARROW,
CLASS_GUARDED_FAIL_CLOSED,
CLASS_FORBIDDEN,
CLASS_REMOVED,
CLASS_HOST_RESIDUAL,
}
)
#: The in-process MCP daemon modules. These must never self-replace/self-kill.
DAEMON_MODULES = (
"gitea_mcp_server.py",
"mcp_server.py",
"role_session_router.py",
)
#: The legacy auto-restart helper removed in #685. Must stay removed.
LEGACY_AUTO_RESTART_HELPER = "_trigger_mcp_auto_restart"
#: Call patterns that would let the daemon replace or terminate its own
#: process. Matched as calls (trailing ``(``) so prose/docstring mentions such
#: as "we do NOT os.execv() here" or "never calls ``os._exit``" do not trip the
#: scanner (comment lines are stripped first regardless).
DAEMON_SELF_REPLACEMENT_PRIMITIVES = (
"os.execv(",
"os.execve(",
"os.execvp(",
"os.execvpe(",
"os.kill(",
"os.killpg(",
"os._exit(",
"os.abort(",
)
@dataclass(frozen=True)
class RestartPath:
"""One classified restart/reload/kill path in the inventory."""
path_id: str
title: str
mechanism: str
classification: str
guard: str
locations: tuple[str, ...]
references: tuple[str, ...]
residual_host: bool = False
notes: str = ""
class UnknownRestartPathError(RuntimeError):
"""Raised when a restart attempt is not a registered, classified path."""
# --- The inventory ---------------------------------------------------------
_RESTART_PATHS: tuple[RestartPath, ...] = (
RestartPath(
path_id="cli_venv_bootstrap_execv",
title="CLI wrapper venv re-exec",
mechanism=(
"Standalone CLI scripts re-exec into venv/bin/python3 via os.execv "
"at import top, guarded by `sys.executable != venv_python`."
),
classification=CLASS_SANCTIONED_NARROW,
guard=(
"One-shot, pre-import bootstrap; runs before any MCP transport "
"exists and only when not already on the venv interpreter, so it "
"cannot desync a live daemon. Idempotent guard condition prevents "
"a re-exec loop."
),
locations=(
"create_pr.py",
"create_issue.py",
"close_issue.py",
"merge_pr.py",
"review_pr.py",
"edit_pr.py",
"delete_branch.py",
"mark_issue.py",
"manage_labels.py",
"list_issues.py",
"list_prs.py",
),
references=("#657",),
),
RestartPath(
path_id="daemon_self_replacement",
title="MCP daemon self-replacement",
mechanism=(
"The in-process MCP daemon replacing/terminating its own process "
"(os.execv/os.kill/os._exit) to reload code."
),
classification=CLASS_FORBIDDEN,
guard=(
"Forbidden by design: replacing the process after the host wired "
"up stdio desyncs JSON-RPC (Antigravity/Cascade). Enforced against "
"the source tree by assert_no_daemon_self_replacement()."
),
locations=("gitea_mcp_server.py:~155 (decision comment)",) + DAEMON_MODULES,
references=("#657", "#584"),
),
RestartPath(
path_id="legacy_auto_restart_helper",
title="Legacy _trigger_mcp_auto_restart helper",
mechanism=(
"A helper that actively restarted the MCP server from the "
"read-only resolver path."
),
classification=CLASS_REMOVED,
guard=(
"Removed in #685 when the resolver became side-effect free. Kept "
"absent by assert_auto_restart_helper_absent()."
),
locations=("gitea_mcp_server.py", "mcp_server.py"),
references=("#685", "#657"),
),
RestartPath(
path_id="config_touch_reload",
title="MCP client config-touch reload",
mechanism=(
"Touching (utime) the MCP client config file to make the host "
"reload/recreate the server process."
),
classification=CLASS_REMOVED,
guard=(
"Removed from the resolver in #685: stale-runtime detection is "
"report-only and never mutates client config, spawns threads, or "
"calls os._exit."
),
locations=("gitea_mcp_server.py (resolve_task_capability)",),
references=("#685", "#657"),
),
RestartPath(
path_id="master_advance_auto_restart",
title="Master-advance staleness gate",
mechanism=(
"On-disk master advancing past the running code. The master-parity "
"gate detects it and fails mutations closed with restart guidance."
),
classification=CLASS_GUARDED_FAIL_CLOSED,
guard=(
"Detect + fail closed only; the process never self-restarts. "
"master_parity_gate captures startup parity and blocks mutations "
"while stale, emitting restart/reconnect guidance."
),
locations=(
"master_parity_gate.py",
"gitea_mcp_server.py (gitea_assess_master_parity)",
),
references=("#420", "#591", "#657"),
),
RestartPath(
path_id="stale_runtime_resolver_reconnect",
title="Stale-runtime resolver reconnect guidance",
mechanism=(
"The capability resolver detecting a stale serving process and "
"reporting restart_required/stop_required for a client reconnect."
),
classification=CLASS_GUARDED_FAIL_CLOSED,
guard=(
"Report-only (#685): returns restart_required/stop_required and an "
"exact_safe_next_action pointing at IDE/client reconnect; performs "
"no restart, thread spawn, config touch, or os._exit."
),
locations=("gitea_mcp_server.py (gitea_resolve_task_capability)",),
references=("#685", "#657"),
),
RestartPath(
path_id="manual_daemon_kill",
title="Manual daemon kill (pkill/killall/kill)",
mechanism=(
"Shell kills of the MCP daemon: `pkill -f mcp_server.py`, "
"`killall`, broad `pkill -f python` sweeps, or `kill <pid>` of a "
"daemon pid."
),
classification=CLASS_FORBIDDEN,
guard=(
"Forbidden (#630): runtime_recovery_guard classifies these as "
"contamination and gitea_record_daemon_process_kill_attempt writes "
"a durable marker that fails subsequent mutations closed. Operator "
"maintenance authorization is read only from the environment, not "
"from a tool argument."
),
locations=(
"runtime_recovery_guard.py",
"gitea_mcp_server.py (gitea_record_daemon_process_kill_attempt)",
),
references=("#630", "#657"),
),
RestartPath(
path_id="conflict_marker_infra_stop",
title="Startup conflict-marker infra stop",
mechanism=(
"The daemon entrypoint scans for unresolved merge-conflict markers "
"at startup and stops (sys.exit(1)) if found."
),
classification=CLASS_GUARDED_FAIL_CLOSED,
guard=(
"Fail-closed startup stop, not a restart: the process exits and "
"waits for the operator to resolve conflicts and relaunch. Never "
"self-restarts or loops."
),
locations=("mcp_server.py (check_conflict_markers)",),
references=("#657",),
),
RestartPath(
path_id="ide_client_reconnect",
title="Host/IDE MCP reconnect",
mechanism=(
"A manual `/mcp reconnect` (or equivalent host action) that the "
"IDE performs to recreate the MCP client connection."
),
classification=CLASS_HOST_RESIDUAL,
guard=(
"Outside this process's control. It is the sanctioned recovery the "
"gates point operators toward; documented as residual host "
"behavior. No in-process code initiates it."
),
locations=("host/IDE",),
references=("#584", "#656", "#657"),
residual_host=True,
),
RestartPath(
path_id="profile_switch_runtime",
title="Runtime profile switch",
mechanism=(
"Switching the active execution profile at runtime "
"(dynamic-profile mode)."
),
classification=CLASS_SANCTIONED_NARROW,
guard=(
"In-process and restart-free: runtime_switching_supported is true, "
"so a profile switch rebinds capability without recreating the "
"process. No restart primitive is invoked."
),
locations=("gitea_mcp_server.py (gitea_activate_profile)",),
references=("#656", "#657"),
),
)
_BY_ID: dict[str, RestartPath] = {p.path_id: p for p in _RESTART_PATHS}
# --- Read-only accessors ---------------------------------------------------
def iter_restart_paths() -> tuple[RestartPath, ...]:
"""Return the full inventory as an immutable tuple."""
return _RESTART_PATHS
def restart_path_ids() -> frozenset[str]:
"""Return the set of registered path ids."""
return frozenset(_BY_ID)
def get_restart_path(path_id: str) -> RestartPath:
"""Return the registered path, or raise :class:`UnknownRestartPathError`."""
try:
return _BY_ID[path_id]
except KeyError as exc:
raise UnknownRestartPathError(
f"unknown restart path id {path_id!r}; not in the #657 inventory"
) from exc
def paths_by_classification(classification: str) -> tuple[RestartPath, ...]:
"""Return all registered paths with the given classification."""
if classification not in VALID_CLASSIFICATIONS:
raise ValueError(f"unknown classification {classification!r}")
return tuple(p for p in _RESTART_PATHS if p.classification == classification)
def assert_restart_attempt_registered(path_id: str) -> RestartPath:
"""Fail closed unless ``path_id`` is a registered, classified restart path.
LLM tools that intend to trigger any restart/reload/reconnect must name a
registered path so an unknown/novel restart primitive cannot slip through
silently. Forbidden and removed paths are registered too this only
asserts the attempt is *known*, not that it is *permitted*; callers must
still honor the classification.
"""
return get_restart_path(path_id)
def assert_registry_wellformed() -> None:
"""Validate the inventory's own invariants (fail closed on drift)."""
seen: set[str] = set()
for path in _RESTART_PATHS:
if path.path_id in seen:
raise ValueError(f"duplicate restart path id {path.path_id!r}")
seen.add(path.path_id)
if path.classification not in VALID_CLASSIFICATIONS:
raise ValueError(
f"{path.path_id!r} has invalid classification "
f"{path.classification!r}"
)
if not path.guard.strip():
raise ValueError(f"{path.path_id!r} is missing a guard description")
if not path.references:
raise ValueError(f"{path.path_id!r} is missing references")
if not path.locations:
raise ValueError(f"{path.path_id!r} is missing locations")
if path.classification == CLASS_HOST_RESIDUAL and not path.residual_host:
raise ValueError(
f"{path.path_id!r} is host_residual but residual_host is False"
)
# --- Source-tree guards ----------------------------------------------------
def _repo_root(root: str | os.PathLike[str] | None = None) -> Path:
if root is not None:
return Path(root)
return Path(__file__).resolve().parent
def _iter_code_lines(text: str) -> Iterable[tuple[int, str]]:
"""Yield (1-based lineno, line) for lines that are not full-line comments."""
for lineno, line in enumerate(text.splitlines(), start=1):
if line.lstrip().startswith("#"):
continue
yield lineno, line
def scan_daemon_self_replacement(
root: str | os.PathLike[str] | None = None,
) -> list[dict[str, object]]:
"""Return violations where a daemon module could self-replace/self-kill.
Scans :data:`DAEMON_MODULES` for calls in
:data:`DAEMON_SELF_REPLACEMENT_PRIMITIVES`. Full-line comments are ignored,
and only call forms (with a trailing ``(``) match, so decision comments and
docstrings that merely mention the primitives do not produce false hits.
"""
repo = _repo_root(root)
violations: list[dict[str, object]] = []
for module in DAEMON_MODULES:
path = repo / module
if not path.exists():
continue
text = path.read_text(encoding="utf-8", errors="replace")
for lineno, line in _iter_code_lines(text):
for primitive in DAEMON_SELF_REPLACEMENT_PRIMITIVES:
if primitive in line:
violations.append(
{
"module": module,
"line": lineno,
"primitive": primitive,
"text": line.strip(),
}
)
return violations
def assert_no_daemon_self_replacement(
root: str | os.PathLike[str] | None = None,
) -> None:
"""Fail closed if any daemon module can restart/kill its own process."""
violations = scan_daemon_self_replacement(root)
if violations:
rendered = "; ".join(
f"{v['module']}:{v['line']} {v['primitive']}" for v in violations
)
raise AssertionError(
"MCP daemon must never self-replace/self-kill (#657); found: "
f"{rendered}"
)
def scan_auto_restart_helper(
root: str | os.PathLike[str] | None = None,
) -> list[dict[str, object]]:
"""Return occurrences of a *definition* of the legacy auto-restart helper."""
repo = _repo_root(root)
needle = f"def {LEGACY_AUTO_RESTART_HELPER}"
hits: list[dict[str, object]] = []
for module in DAEMON_MODULES:
path = repo / module
if not path.exists():
continue
text = path.read_text(encoding="utf-8", errors="replace")
for lineno, line in _iter_code_lines(text):
if needle in line:
hits.append({"module": module, "line": lineno})
return hits
def assert_auto_restart_helper_absent(
root: str | os.PathLike[str] | None = None,
) -> None:
"""Fail closed if the removed ``_trigger_mcp_auto_restart`` reappears."""
hits = scan_auto_restart_helper(root)
if hits:
rendered = "; ".join(f"{h['module']}:{h['line']}" for h in hits)
raise AssertionError(
f"{LEGACY_AUTO_RESTART_HELPER} was removed in #685 and must not "
f"return (#657); found definition at: {rendered}"
)
+451
View File
@@ -0,0 +1,451 @@
"""MCP restart coordinator and impact analysis (#658).
Before any sanctioned MCP restart, a central coordinator must evaluate the
live control-plane state active sessions, leases/locks, in-flight issue/PR
work, mutations, worktrees, and recovery history and produce an *impact
preview* so operators (and the web console, #642/#652) can see the blast
radius **before** concurrent LLM work is disrupted.
Design rules (mirrors the read-only posture of ``workflow_dashboard`` /
``lease_lifecycle``):
* **Pure classification.** :func:`evaluate_restart_impact` takes an already
gathered inventory and returns a structured report. It never touches the
network, the filesystem, or a live process, so multi-session fixtures can
drive every branch in unit tests. The coordinator *never restarts anything*;
a mutative apply path is a later child gated by a drain proof (non-goal here).
* **Fail closed.** If the inventory is not explicitly complete, the verdict is
``unsafe`` / deny an incomplete evaluation must never green-light a restart.
* **No secrets.** Session ids, pids, and profiles are operational metadata, not
credentials; nothing secret flows through this module.
The single sanctioned entry point post-#657 is the MCP tool
``gitea_request_mcp_restart`` (dry-run by default), which gathers the inventory
from the #613 control-plane DB and calls :func:`evaluate_restart_impact`.
"""
from __future__ import annotations
from dataclasses import dataclass, field
from datetime import datetime, timezone
from typing import Any, Mapping, Sequence
import lease_lifecycle
COORDINATOR_VERSION = "1.0.0-issue-658"
# Restart verdicts. Exactly the three the acceptance criteria name.
VERDICT_SAFE = "safe"
VERDICT_UNSAFE = "unsafe"
VERDICT_OVERRIDE = "override"
# Blast-radius severity bands.
BLAST_NONE = "none"
BLAST_LOW = "low"
BLAST_MEDIUM = "medium"
BLAST_HIGH = "high"
# A live lease with a live owner process is treated as active in-flight work.
LEASE_FRESHNESS_LIVE = "active"
# Default staleness window for a session heartbeat (seconds). A session whose
# last heartbeat is older than this is not counted as live even if its row is
# still marked ``active`` — it is assumed dead/detached.
DEFAULT_SESSION_HEARTBEAT_STALE_SECONDS = 900
def _utc_now() -> datetime:
return datetime.now(timezone.utc)
def _parse_ts(value: str | None) -> datetime | None:
return lease_lifecycle._parse_ts(value)
@dataclass(frozen=True)
class SessionImpact:
"""One MCP session a restart would terminate."""
session_id: str
role: str | None
profile: str | None
pid: int | None
status: str | None
alive: bool | None
heartbeat_stale: bool
is_requester: bool
live: bool
def as_dict(self) -> dict[str, Any]:
return {
"session_id": self.session_id,
"role": self.role,
"profile": self.profile,
"pid": self.pid,
"status": self.status,
"alive": self.alive,
"heartbeat_stale": self.heartbeat_stale,
"is_requester": self.is_requester,
"live": self.live,
}
@dataclass(frozen=True)
class LeaseImpact:
"""One control-plane lease a restart would disrupt."""
lease_id: str | None
session_id: str | None
role: str | None
phase: str | None
freshness: str | None
work_kind: str | None
work_number: int | None
worktree_path: str | None
disruptive: bool
is_mutation: bool
is_critical_section: bool
def as_dict(self) -> dict[str, Any]:
return {
"lease_id": self.lease_id,
"session_id": self.session_id,
"role": self.role,
"phase": self.phase,
"freshness": self.freshness,
"work_kind": self.work_kind,
"work_number": self.work_number,
"worktree_path": self.worktree_path,
"disruptive": self.disruptive,
"is_mutation": self.is_mutation,
"is_critical_section": self.is_critical_section,
}
@dataclass(frozen=True)
class RestartImpactReport:
"""Impact preview DTO returned to the console / operator (#642/#652)."""
coordinator_version: str
evaluated_at: str
dry_run: bool
restart_performed: bool
inventory_complete: bool
verdict: str
allow_restart: bool
override_would_allow: bool
operator_override: bool
blast_radius: str
reasons: list[str]
affected_sessions: list[SessionImpact]
affected_leases: list[LeaseImpact]
critical_sections: list[LeaseImpact]
affected_issues: list[int]
affected_prs: list[int]
mutations: list[LeaseImpact]
terminal_lock: dict[str, Any] | None
ack_state: dict[str, str]
prior_recovery_attempts: list[dict[str, Any]]
counts: dict[str, int]
audit_record: dict[str, Any]
incomplete_reasons: list[str] = field(default_factory=list)
def as_dict(self) -> dict[str, Any]:
return {
"coordinator_version": self.coordinator_version,
"evaluated_at": self.evaluated_at,
"dry_run": self.dry_run,
"restart_performed": self.restart_performed,
"inventory_complete": self.inventory_complete,
"incomplete_reasons": list(self.incomplete_reasons),
"verdict": self.verdict,
"allow_restart": self.allow_restart,
"override_would_allow": self.override_would_allow,
"operator_override": self.operator_override,
"blast_radius": self.blast_radius,
"reasons": list(self.reasons),
"affected_sessions": [s.as_dict() for s in self.affected_sessions],
"affected_leases": [l.as_dict() for l in self.affected_leases],
"critical_sections": [l.as_dict() for l in self.critical_sections],
"affected_issues": list(self.affected_issues),
"affected_prs": list(self.affected_prs),
"mutations": [l.as_dict() for l in self.mutations],
"terminal_lock": self.terminal_lock,
"ack_state": dict(self.ack_state),
"prior_recovery_attempts": list(self.prior_recovery_attempts),
"counts": dict(self.counts),
"audit_record": dict(self.audit_record),
}
def _classify_session(
row: Mapping[str, Any],
*,
now: datetime,
requesting_session_id: str | None,
heartbeat_stale_seconds: int,
) -> SessionImpact:
session_id = str(row.get("session_id") or "")
pid = row.get("pid")
status = (row.get("status") or "").strip().lower() or None
alive = lease_lifecycle.is_process_alive(pid) if pid is not None else None
hb = _parse_ts(row.get("last_heartbeat_at"))
heartbeat_stale = bool(
hb is not None and (now - hb).total_seconds() > heartbeat_stale_seconds
)
live = bool(status == "active" and alive is not False and not heartbeat_stale)
return SessionImpact(
session_id=session_id,
role=row.get("role"),
profile=row.get("profile"),
pid=pid,
status=status,
alive=alive,
heartbeat_stale=heartbeat_stale,
is_requester=bool(
requesting_session_id and session_id == requesting_session_id
),
live=live,
)
# Lease phases that represent an active mutation in flight (as opposed to a
# mere allocation/claim with no work committed yet). An active lease in any of
# these phases is a critical section a restart must not sever.
_MUTATING_PHASES = frozenset(
{
"implementing",
"publishing",
"pushing",
"committing",
"reviewing",
"merging",
"reconciling",
"conflict_fix",
}
)
def _classify_lease(row: Mapping[str, Any]) -> LeaseImpact:
freshness_obj = row.get("freshness")
if isinstance(freshness_obj, Mapping):
freshness = str(freshness_obj.get("freshness") or "").strip().lower() or None
else:
freshness = str(freshness_obj or "").strip().lower() or None
phase = (row.get("phase") or "").strip().lower() or None
worktree = row.get("worktree_path")
disruptive = freshness == LEASE_FRESHNESS_LIVE
# A live lease is a mutation-in-flight if it carries an author worktree or
# its phase names a mutating step. All disruptive leases are critical
# sections a restart would sever regardless.
is_mutation = bool(
disruptive and (bool(worktree) or (phase in _MUTATING_PHASES))
)
number = row.get("work_number")
try:
number = int(number) if number is not None else None
except (TypeError, ValueError):
number = None
return LeaseImpact(
lease_id=row.get("lease_id"),
session_id=row.get("session_id"),
role=row.get("role"),
phase=phase,
freshness=freshness,
work_kind=(str(row.get("work_kind") or "").strip().lower() or None),
work_number=number,
worktree_path=worktree,
disruptive=disruptive,
is_mutation=is_mutation,
is_critical_section=disruptive,
)
def _blast_radius(*, session_count: int, work_count: int, mutation_count: int) -> str:
if mutation_count > 0 or work_count >= 3 or session_count >= 3:
return BLAST_HIGH
if work_count > 0 or session_count == 2:
return BLAST_MEDIUM
if session_count == 1:
return BLAST_LOW
return BLAST_NONE
def evaluate_restart_impact(
inventory: Mapping[str, Any],
*,
now: datetime | None = None,
operator_override: bool = False,
requesting_session_id: str | None = None,
dry_run: bool = True,
session_heartbeat_stale_seconds: int = DEFAULT_SESSION_HEARTBEAT_STALE_SECONDS,
) -> RestartImpactReport:
"""Evaluate a proposed MCP restart and return an impact preview.
``inventory`` is a mapping with:
* ``sessions`` session rows (session_id, role, profile, pid, status,
last_heartbeat_at).
* ``leases`` control-plane lease rows, each ideally carrying an enriched
``freshness`` dict (as :func:`lease_lifecycle.list_active_leases` returns);
a bare string freshness is also accepted.
* ``terminal_lock`` the active terminal (merge) lock, if any.
* ``prior_recovery_attempts`` narrower recovery attempts already tried
(e.g. sanctioned reconnects) so the operator sees escalation history.
* ``inventory_complete`` bool. **Must** be explicitly True; a missing or
falsy value forces a deny (fail closed).
* ``incomplete_reasons`` optional reasons the inventory is incomplete.
The coordinator never restarts anything: ``restart_performed`` is always
False and the mutative apply path is a later drain-gated child.
"""
moment = now or _utc_now()
reasons: list[str] = []
inventory_complete = bool(inventory.get("inventory_complete", False))
incomplete_reasons = [str(r) for r in (inventory.get("incomplete_reasons") or [])]
sessions_raw: Sequence[Mapping[str, Any]] = inventory.get("sessions") or []
leases_raw: Sequence[Mapping[str, Any]] = inventory.get("leases") or []
terminal_lock = inventory.get("terminal_lock") or None
prior_recovery_attempts = [
dict(a) for a in (inventory.get("prior_recovery_attempts") or [])
]
session_impacts = [
_classify_session(
s,
now=moment,
requesting_session_id=requesting_session_id,
heartbeat_stale_seconds=session_heartbeat_stale_seconds,
)
for s in sessions_raw
]
lease_impacts = [_classify_lease(l) for l in leases_raw]
# Only *other* live sessions and live leases constitute blast radius: a
# restart that would kill only the requesting session with no other work in
# flight is safe.
other_live_sessions = [
s for s in session_impacts if s.live and not s.is_requester
]
disruptive_leases = [l for l in lease_impacts if l.disruptive]
critical_sections = [l for l in lease_impacts if l.is_critical_section]
mutations = [l for l in lease_impacts if l.is_mutation]
affected_issues = sorted(
{
l.work_number
for l in disruptive_leases
if l.work_kind == "issue" and l.work_number is not None
}
)
affected_prs = sorted(
{
l.work_number
for l in disruptive_leases
if l.work_kind == "pr" and l.work_number is not None
}
)
disruptive = bool(disruptive_leases or other_live_sessions or terminal_lock)
if not inventory_complete:
verdict = VERDICT_UNSAFE
allow_restart = False
reasons.append(
"inventory incomplete: restart evaluation cannot confirm blast "
"radius — deny (fail closed, #658)"
)
reasons.extend(incomplete_reasons)
elif not disruptive:
verdict = VERDICT_SAFE
allow_restart = True
reasons.append("no other live sessions, live leases, or terminal lock")
elif operator_override:
verdict = VERDICT_OVERRIDE
allow_restart = True
reasons.append(
"live work present; operator override accepts the blast radius"
)
else:
verdict = VERDICT_UNSAFE
allow_restart = False
reasons.append(
"live work would be disrupted; restart denied without operator "
"override"
)
if critical_sections and inventory_complete:
reasons.append(
f"{len(critical_sections)} critical section(s) in flight "
"(active lease with a live owner)"
)
if terminal_lock:
reasons.append("active terminal (merge) lock present")
override_would_allow = bool(inventory_complete and disruptive)
blast_radius = _blast_radius(
session_count=len(other_live_sessions),
work_count=len(affected_issues) + len(affected_prs),
mutation_count=len(mutations),
)
# Acknowledgement is a later child (drain protocol); expose per-session
# placeholders so the console can render the ack column now.
ack_state = {s.session_id: "pending" for s in other_live_sessions}
counts = {
"sessions_total": len(session_impacts),
"sessions_live_other": len(other_live_sessions),
"leases_total": len(lease_impacts),
"leases_disruptive": len(disruptive_leases),
"critical_sections": len(critical_sections),
"mutations": len(mutations),
"affected_issues": len(affected_issues),
"affected_prs": len(affected_prs),
"prior_recovery_attempts": len(prior_recovery_attempts),
}
audit_record = {
"event": "restart_impact_evaluated",
"coordinator_version": COORDINATOR_VERSION,
"evaluated_at": moment.isoformat(),
"dry_run": dry_run,
"operator_override": bool(operator_override),
"requesting_session_id": requesting_session_id,
"inventory_complete": inventory_complete,
"verdict": verdict,
"allow_restart": allow_restart,
"blast_radius": blast_radius,
"counts": counts,
}
return RestartImpactReport(
coordinator_version=COORDINATOR_VERSION,
evaluated_at=moment.isoformat(),
dry_run=dry_run,
restart_performed=False,
inventory_complete=inventory_complete,
verdict=verdict,
allow_restart=allow_restart,
override_would_allow=override_would_allow,
operator_override=bool(operator_override),
blast_radius=blast_radius,
reasons=reasons,
affected_sessions=session_impacts,
affected_leases=lease_impacts,
critical_sections=critical_sections,
affected_issues=affected_issues,
affected_prs=affected_prs,
mutations=mutations,
terminal_lock=dict(terminal_lock)
if isinstance(terminal_lock, Mapping)
else terminal_lock,
ack_state=ack_state,
prior_recovery_attempts=prior_recovery_attempts,
counts=counts,
audit_record=audit_record,
incomplete_reasons=incomplete_reasons,
)
+2
View File
@@ -73,6 +73,8 @@ AUTHOR_TASKS = frozenset({
"claim_issue",
"create_branch",
"push_branch",
"bootstrap_author_issue_worktree",
"gitea_bootstrap_author_issue_worktree",
"create_pr",
"comment_pr",
"address_pr_change_requests",
+10 -8
View File
@@ -43,19 +43,21 @@ repo_root="$(cd "$script_dir/.." && pwd)"
# Enforce issue-linked, traceable branch names (issue → branch → worktree → PR).
if [[ "$allow_unlinked" -eq 0 ]]; then
locked_branch=$(python3 -c "
if [[ "$dry_run" -eq 0 ]] && [[ ! "$branch" =~ ^review/pr-[0-9]+-.+ ]]; then
locked_branch=$(python3 -c "
import sys
sys.path.insert(0, '$repo_root')
import issue_lock_store
print(issue_lock_store.resolve_locked_branch_for_session('$branch'))
")
if [[ -z "$locked_branch" ]]; then
echo "Error: No session issue lock is bound. Call gitea_lock_issue before branch creation (fail closed)." >&2
exit 2
fi
if [[ "$branch" != "$locked_branch" ]]; then
echo "Error: Requested branch '$branch' does not match locked branch '$locked_branch' (fail closed)." >&2
exit 2
if [[ -z "$locked_branch" ]]; then
echo "Error: No session issue lock is bound. Call gitea_lock_issue before branch creation (fail closed)." >&2
exit 2
fi
if [[ "$branch" != "$locked_branch" ]]; then
echo "Error: Requested branch '$branch' does not match locked branch '$locked_branch' (fail closed)." >&2
exit 2
fi
fi
if [[ "$branch" =~ ^(fix|feat|docs|chore)/issue-[0-9]+-.+ ]] \
+38
View File
@@ -32,6 +32,15 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.issue.comment",
"role": "author",
},
# #860: dirty orphaned same-claimant worktree recovery (explicit operation).
"recover_dirty_orphaned_issue_worktree": {
"permission": "gitea.issue.comment",
"role": "author",
},
"gitea_recover_dirty_orphaned_issue_worktree": {
"permission": "gitea.issue.comment",
"role": "author",
},
# #864: dirty-preserving same-claimant author-session rebind (dead owner PID).
# Author MCP tool path. Reconciler execute is gated inside the tool via
# authorize_reconciler_execute + role_kind checks (not this map entry).
@@ -69,6 +78,14 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.branch.create",
"role": "author",
},
"bootstrap_author_issue_worktree": {
"permission": "gitea.branch.create",
"role": "author",
},
"gitea_bootstrap_author_issue_worktree": {
"permission": "gitea.branch.create",
"role": "author",
},
"push_branch": {
"permission": "gitea.branch.push",
"role": "author",
@@ -478,6 +495,15 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.issue.comment",
"role": "author",
},
# #651 console analytics ingest — control-plane DB write, not a Gitea API
# call. Authority comes from console RBAC (operator+) plus phase gating;
# permission string is a non-Gitea runtime capability so no Gitea profile
# can satisfy it by accident.
"record_analytics_usage": {
"permission": "runtime.record_analytics_usage",
"role": "author",
},
}
@@ -488,6 +514,15 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
# merger lease (#763).
_PREFLIGHT_TASK_TRANSITIONS = frozenset({
("review_pr", "acquire_reviewer_pr_lease"),
# #850: native author issue worktree bootstrap
("work_issue", "bootstrap_author_issue_worktree"),
("bootstrap_author_issue_worktree", "lock_issue"),
# #860: dirty-orphan recovery and related work_issue transitions (master)
("work_issue", "lock_issue"),
("work_issue", "recover_dirty_orphaned_issue_worktree"),
("work_issue", "gitea_recover_dirty_orphaned_issue_worktree"),
("work_issue", "commit_files"),
("work_issue", "gitea_commit_files"),
})
@@ -534,6 +569,8 @@ ROLE_EXCLUSIVE_TASKS: frozenset[str] = frozenset(
"gitea_release_merger_pr_lease",
"create_branch",
"push_branch",
"bootstrap_author_issue_worktree",
"gitea_bootstrap_author_issue_worktree",
"publish_unpublished_branch",
"create_pr",
"commit_files",
@@ -559,6 +596,7 @@ ISSUE_MUTATION_TOOL_TASKS: dict[str, str] = {
"gitea_set_issue_labels": "set_issue_labels",
"gitea_cleanup_terminal_pr_labels": "cleanup_terminal_pr_labels",
"gitea_create_label": "create_label",
"gitea_bootstrap_author_issue_worktree": "bootstrap_author_issue_worktree",
"gitea_commit_files": "commit_files",
}
+601
View File
@@ -0,0 +1,601 @@
"""Regression test suite for native author issue worktree bootstrap (#850)."""
from __future__ import annotations
import json
import os
import shutil
import subprocess
import tempfile
import unittest
from unittest import mock
import author_issue_bootstrap
import task_capability_map
def _concurrent_bootstrap_worker(args: tuple[str, int, str, str, str, str]) -> dict:
repo_dir, issue_num, key, lock_dir, journal_dir, master_sha = args
os.environ["GITEA_BOOTSTRAP_JOURNAL_DIR"] = journal_dir
return author_issue_bootstrap.bootstrap_author_issue_worktree(
issue_number=issue_num,
canonical_repo_root=repo_dir,
expected_base_sha=master_sha,
idempotency_key=key,
lock_dir=lock_dir,
owner_session="session-concurrent-test",
active_identity="jcwalker3",
active_profile="prgs-author",
)
class TestAuthorIssueBootstrap(unittest.TestCase):
"""Test suite covering AC1-AC10 and comment #14959 specification."""
def setUp(self):
self.tmp_dir = tempfile.mkdtemp(prefix="test_bootstrap_")
self.repo_dir = os.path.join(self.tmp_dir, "repo")
os.makedirs(self.repo_dir)
# Initialize synthetic git repo
subprocess.run(["git", "init", "-b", "master"], cwd=self.repo_dir, check=True, capture_output=True)
subprocess.run(["git", "config", "user.name", "Test User"], cwd=self.repo_dir, check=True)
subprocess.run(["git", "config", "user.email", "[email protected]"], cwd=self.repo_dir, check=True)
readme = os.path.join(self.repo_dir, "README.md")
with open(readme, "w", encoding="utf-8") as f:
f.write("# Test Repo\n")
subprocess.run(["git", "add", "README.md"], cwd=self.repo_dir, check=True, capture_output=True)
subprocess.run(["git", "commit", "-m", "initial commit"], cwd=self.repo_dir, check=True, capture_output=True)
rev_res = subprocess.run(["git", "rev-parse", "HEAD"], cwd=self.repo_dir, capture_output=True, text=True, check=True)
self.master_sha = rev_res.stdout.strip()
self.branches_dir = os.path.join(self.repo_dir, "branches")
os.makedirs(self.branches_dir, exist_ok=True)
self.lock_dir = os.path.join(self.tmp_dir, "locks")
os.makedirs(self.lock_dir, exist_ok=True)
self.journal_dir = os.path.join(self.tmp_dir, "journals")
os.makedirs(self.journal_dir, exist_ok=True)
os.environ["GITEA_BOOTSTRAP_JOURNAL_DIR"] = self.journal_dir
def tearDown(self):
os.environ.pop("GITEA_BOOTSTRAP_JOURNAL_DIR", None)
shutil.rmtree(self.tmp_dir, ignore_errors=True)
def test_bootstrap_success_path(self):
"""AC1/AC3/AC8: Successful bootstrap creates branch, worktree, registration, and lock proof."""
key = "test_key_success_1"
res = author_issue_bootstrap.bootstrap_author_issue_worktree(
issue_number=850,
canonical_repo_root=self.repo_dir,
# assignment_id/lease_id omitted: optional unless verified live.
expected_base_sha=self.master_sha,
idempotency_key=key,
remote="prgs",
lock_dir=self.lock_dir,
owner_session="session-test-1234",
)
self.assertTrue(res.get("success"), f"Bootstrap failed: {res}")
self.assertFalse(res.get("replayed"))
self.assertEqual(res.get("issue_number"), 850)
self.assertEqual(res.get("base_sha"), self.master_sha)
self.assertIn("branches/fix-issue-850-native-mcp-bootstrap", res.get("worktree_path"))
# Verify worktree directory exists and is registered
worktree_path = res["worktree_path"]
self.assertTrue(os.path.isdir(worktree_path))
wt_list = subprocess.run(["git", "-C", self.repo_dir, "worktree", "list"], capture_output=True, text=True, check=True)
self.assertIn(worktree_path, wt_list.stdout)
# Verify phase journal written
journal = author_issue_bootstrap.load_phase_journal(key, journal_dir=self.lock_dir)
self.assertIsNotNone(journal)
self.assertTrue(journal.get("completed"))
self.assertEqual(journal.get("current_phase"), author_issue_bootstrap.PHASE_7_TRANSITION_COMPLETED)
def test_idempotent_replay(self):
"""Item 2: Replaying with identical key returns cached transition without duplicate creation."""
key = "test_key_idempotent_1"
res1 = author_issue_bootstrap.bootstrap_author_issue_worktree(
issue_number=850,
canonical_repo_root=self.repo_dir,
idempotency_key=key,
lock_dir=self.lock_dir,
owner_session="session-test-1234",
)
self.assertTrue(res1["success"], f"res1 failed: {res1}")
self.assertFalse(res1.get("replayed"))
# Second call
res2 = author_issue_bootstrap.bootstrap_author_issue_worktree(
issue_number=850,
canonical_repo_root=self.repo_dir,
idempotency_key=key,
lock_dir=self.lock_dir,
owner_session="session-test-1234",
)
self.assertTrue(res2["success"], f"res2 failed: {res2}")
self.assertTrue(res2.get("replayed"))
self.assertEqual(res1["worktree_path"], res2["worktree_path"])
def test_stale_concurrency_pin_refusal(self):
"""Item 3: Mismatched expected base SHA fails closed without silent rebasing."""
stale_sha = "0000000000000000000000000000000000000000"
res = author_issue_bootstrap.bootstrap_author_issue_worktree(
issue_number=850,
canonical_repo_root=self.repo_dir,
expected_base_sha=stale_sha,
lock_dir=self.lock_dir,
owner_session="session-test-1234",
)
self.assertFalse(res["success"])
self.assertEqual(res.get("reason_code"), "stale_concurrency_pin")
self.assertIn("exact_next_action", res)
def test_path_outside_branches_root_refusal(self):
"""Item 6: Worktree path outside branches/ root is refused."""
outside_path = os.path.join(self.tmp_dir, "outside_worktree")
res = author_issue_bootstrap.bootstrap_author_issue_worktree(
issue_number=850,
canonical_repo_root=self.repo_dir,
worktree_path=outside_path,
lock_dir=self.lock_dir,
owner_session="session-test-1234",
)
self.assertFalse(res["success"])
self.assertEqual(res.get("reason_code"), "path_outside_canonical_branches_root")
def test_preexisting_dirty_worktree_preservation(self):
"""Item 6: Preexisting dirty worktree fails closed and is NOT modified or cleaned."""
branch = "fix/issue-850-dirty-test"
wt_path = os.path.join(self.branches_dir, "fix-issue-850-dirty-test")
subprocess.run(["git", "-C", self.repo_dir, "worktree", "add", "-b", branch, wt_path], check=True, capture_output=True)
# Create dirty untracked file
dirty_file = os.path.join(wt_path, "dirty.txt")
with open(dirty_file, "w") as f:
f.write("dirty edits\n")
res = author_issue_bootstrap.bootstrap_author_issue_worktree(
issue_number=850,
canonical_repo_root=self.repo_dir,
branch_name=branch,
worktree_path=wt_path,
lock_dir=self.lock_dir,
owner_session="session-test-1234",
)
self.assertFalse(res["success"])
self.assertEqual(res.get("reason_code"), "preexisting_dirty_worktree")
# Prove dirty file is preserved byte-for-byte
self.assertTrue(os.path.exists(dirty_file))
with open(dirty_file, "r") as f:
self.assertEqual(f.read(), "dirty edits\n")
def test_compensating_recovery_on_failed_phase(self):
"""AC4/Item 4: Failure during transition rolls back ONLY newly created artifacts."""
key = "test_key_recovery_1"
# Simulate partial progress in journal
journal = {
"idempotency_key": key,
"issue_number": 850,
"branch_name": "fix/issue-850-recovery-test",
"worktree_path": os.path.join(self.branches_dir, "fix-issue-850-recovery-test"),
"artifacts_created": {
"branch_created": True,
"worktree_dir_created": True,
"worktree_registered": True,
"lock_created": False,
},
"failure_reason": "simulated lock failure",
"current_phase": author_issue_bootstrap.PHASE_5_REGISTRATION_VERIFIED,
"completed": False,
}
# Create the branch and worktree manually to simulate partial state
subprocess.run(["git", "-C", self.repo_dir, "branch", journal["branch_name"]], check=True, capture_output=True)
subprocess.run(["git", "-C", self.repo_dir, "worktree", "add", journal["worktree_path"], journal["branch_name"]], check=True, capture_output=True)
# Run compensating recovery
rec = author_issue_bootstrap.run_compensating_recovery(journal, self.repo_dir)
self.assertTrue(rec["executed"])
self.assertIn(f"worktree_path:{journal['worktree_path']}", rec["rolled_back"])
self.assertIn(f"branch:{journal['branch_name']}", rec["rolled_back"])
# Prove worktree directory and branch were rolled back
self.assertFalse(os.path.exists(journal["worktree_path"]))
branch_check = subprocess.run(["git", "-C", self.repo_dir, "rev-parse", "--verify", journal["branch_name"]], capture_output=True, text=True, check=False)
self.assertNotEqual(branch_check.returncode, 0)
def test_cross_process_concurrency(self):
"""Review #525 Finding 1: Genuine cross-process concurrency locking prevents corruption."""
import concurrent.futures
key = "test_concurrent_key_850"
args = (self.repo_dir, 850, key, self.lock_dir, self.journal_dir, self.master_sha)
with concurrent.futures.ProcessPoolExecutor(max_workers=2) as executor:
fut1 = executor.submit(_concurrent_bootstrap_worker, args)
fut2 = executor.submit(_concurrent_bootstrap_worker, args)
res1 = fut1.result(timeout=10)
res2 = fut2.result(timeout=10)
self.assertTrue(res1["success"], f"res1 failed: {res1}")
self.assertTrue(res2["success"], f"res2 failed: {res2}")
# One process performs creation, the other process receives idempotent replay
replayed_count = sum(1 for r in (res1, res2) if r.get("replayed"))
created_count = sum(1 for r in (res1, res2) if not r.get("replayed"))
self.assertEqual(replayed_count, 1)
self.assertEqual(created_count, 1)
self.assertEqual(res1["worktree_path"], res2["worktree_path"])
def test_interrupted_replay_preserves_artifacts_created_provenance(self):
"""Review #525 Finding 2: Replaying incomplete journal preserves creation provenance monotonically."""
key = "test_key_interrupted_replay_1"
branch = "fix/issue-850-interrupted-replay"
wt_path = os.path.join(self.branches_dir, "fix-issue-850-interrupted-replay")
# Simulate Phase 2/3 completion where branch and worktree directory were created by this transition
journal = {
"idempotency_key": key,
"issue_number": 850,
"branch_name": branch,
"worktree_path": wt_path,
"active_identity": "jcwalker3",
"active_profile": "prgs-author",
"remote": "prgs",
"org": "Scaled-Tech-Consulting",
"repo": "Gitea-Tools",
"phases": {
author_issue_bootstrap.PHASE_1_REQUEST_ACCEPTED: {"status": "completed"},
author_issue_bootstrap.PHASE_2_BRANCH_CONFIRMED: {"status": "completed", "created": True},
},
"artifacts_created": {
"branch_created": True,
"worktree_dir_created": True,
"worktree_registered": True,
"lock_created": False,
},
"current_phase": author_issue_bootstrap.PHASE_3_PATH_RESERVED,
"completed": False,
}
# Pre-create the branch and worktree on disk to simulate partial state after crash
subprocess.run(["git", "-C", self.repo_dir, "branch", branch, self.master_sha], check=True, capture_output=True)
subprocess.run(["git", "-C", self.repo_dir, "worktree", "add", wt_path, branch], check=True, capture_output=True)
author_issue_bootstrap.save_phase_journal(journal, journal_dir=self.lock_dir)
# Now resume/replay the transition but simulate lock binding failure during Phase 6
with mock.patch("issue_lock_store.bind_session_lock", side_effect=RuntimeError("Lock failure test")):
res = author_issue_bootstrap.bootstrap_author_issue_worktree(
issue_number=850,
canonical_repo_root=self.repo_dir,
branch_name=branch,
worktree_path=wt_path,
idempotency_key=key,
lock_dir=self.lock_dir,
owner_session="session-test-1234",
)
self.assertFalse(res["success"])
self.assertEqual(res.get("reason_code"), "issue_lock_acquisition_failed")
# Verify that compensating recovery correctly deleted transition-created branch & worktree
# because creation provenance was preserved across replay (NOT downgraded to False!)
self.assertFalse(os.path.exists(wt_path))
branch_check = subprocess.run(["git", "-C", self.repo_dir, "rev-parse", "--verify", branch], capture_output=True, text=True, check=False)
self.assertNotEqual(branch_check.returncode, 0)
def test_transition_created_only_compensation(self):
"""Review #525 Finding 4: Preexisting branch is NOT deleted by compensation when only worktree was transition-created."""
key = "test_key_preexisting_branch_compensation"
preexisting_branch = "fix/issue-850-preexisting"
wt_path = os.path.join(self.branches_dir, "fix-issue-850-preexisting")
# Create branch BEFORE bootstrap (preexisting branch)
subprocess.run(["git", "-C", self.repo_dir, "branch", preexisting_branch, self.master_sha], check=True, capture_output=True)
# Call bootstrap with simulated failure during Phase 6 (lock binding)
with mock.patch("issue_lock_store.bind_session_lock", side_effect=RuntimeError("Simulated lock failure")):
res = author_issue_bootstrap.bootstrap_author_issue_worktree(
issue_number=850,
canonical_repo_root=self.repo_dir,
branch_name=preexisting_branch,
worktree_path=wt_path,
idempotency_key=key,
lock_dir=self.lock_dir,
owner_session="session-test-1234",
)
self.assertFalse(res["success"])
# Worktree dir was created by transition -> removed by compensation
self.assertFalse(os.path.exists(wt_path))
# Preexisting branch was NOT created by transition -> MUST BE PRESERVED!
branch_check = subprocess.run(["git", "-C", self.repo_dir, "rev-parse", "--verify", preexisting_branch], capture_output=True, text=True, check=False)
self.assertEqual(branch_check.returncode, 0, "Preexisting branch was deleted by mistake!")
def test_incompatible_idempotency_replay_refusal(self):
"""Review #525 Finding 4: Replaying key with incompatible parameters returns refusal."""
key = "test_key_incompatible_replay"
res1 = author_issue_bootstrap.bootstrap_author_issue_worktree(
issue_number=850,
canonical_repo_root=self.repo_dir,
branch_name="fix/issue-850-param-a",
idempotency_key=key,
lock_dir=self.lock_dir,
owner_session="session-test-1234",
)
self.assertTrue(res1["success"])
# Second call with different branch_name
res2 = author_issue_bootstrap.bootstrap_author_issue_worktree(
issue_number=850,
canonical_repo_root=self.repo_dir,
branch_name="fix/issue-850-param-b",
idempotency_key=key,
lock_dir=self.lock_dir,
owner_session="session-test-1234",
)
self.assertFalse(res2["success"])
self.assertEqual(res2.get("reason_code"), "incompatible_idempotency_replay")
def test_exact_next_action_satisfiable_via_mcp(self):
"""Review #525 Finding 4: exact_next_action provides satisfiable MCP actions, not shell commands."""
key = "test_key_next_action_mcp"
stale_sha = "0000000000000000000000000000000000000000"
res = author_issue_bootstrap.bootstrap_author_issue_worktree(
issue_number=850,
canonical_repo_root=self.repo_dir,
expected_base_sha=stale_sha,
lock_dir=self.lock_dir,
owner_session="session-test-1234",
)
next_action = res.get("exact_next_action", "")
self.assertNotIn("scripts/worktree-start", next_action)
self.assertNotIn("git worktree add", next_action)
self.assertNotIn("bash", next_action.lower())
def test_missing_owner_session_refusal(self):
"""Finding D: Missing owner_session context fails closed with typed refusal and zero mutation."""
res = author_issue_bootstrap.bootstrap_author_issue_worktree(
issue_number=850,
canonical_repo_root=self.repo_dir,
owner_session=None,
lock_dir=self.lock_dir,
)
self.assertFalse(res["success"])
self.assertEqual(res.get("reason_code"), "missing_owner_session")
self.assertIn("exact_next_action", res)
def test_symlink_lock_file_refusal(self):
"""Finding C: BootstrapTransitionLock refuses to follow symlinks."""
key = "test_symlink_lock_key"
safe_key = "".join(c if c.isalnum() or c in ("-", "_", ".") else "_" for c in key)
lock_path = os.path.join(self.lock_dir, f"{safe_key}.lock")
target_file = os.path.join(self.tmp_dir, "fake_target")
with open(target_file, "w") as f:
f.write("target")
os.symlink(target_file, lock_path)
with self.assertRaises(RuntimeError) as ctx:
with author_issue_bootstrap.BootstrapTransitionLock(key, journal_dir=self.lock_dir):
pass
self.assertIn("symlink", str(ctx.exception).lower())
def test_lock_directory_escape_refusal(self):
"""Finding C: BootstrapTransitionLock refuses keys that escape lock directory."""
with mock.patch("os.path.abspath", return_value="/tmp/outside/evil_key.lock"):
with self.assertRaises(RuntimeError) as ctx:
author_issue_bootstrap.BootstrapTransitionLock("key", journal_dir=self.lock_dir)
self.assertIn("escapes", str(ctx.exception).lower())
def test_missing_active_identity_refusal(self):
"""F-5: Missing active_identity parameter fails closed."""
res = author_issue_bootstrap.bootstrap_author_issue_worktree(
issue_number=850,
canonical_repo_root=self.repo_dir,
owner_session="session-test-1234",
active_identity=None,
active_profile="prgs-author",
lock_dir=self.lock_dir,
)
self.assertFalse(res["success"])
self.assertEqual(res.get("reason_code"), "missing_active_identity")
def test_missing_active_profile_refusal(self):
"""F-5: Missing active_profile parameter fails closed."""
res = author_issue_bootstrap.bootstrap_author_issue_worktree(
issue_number=850,
canonical_repo_root=self.repo_dir,
owner_session="session-test-1234",
active_identity="jcwalker3",
active_profile=None,
lock_dir=self.lock_dir,
)
self.assertFalse(res["success"])
self.assertEqual(res.get("reason_code"), "missing_active_profile")
def test_dirty_worktree_preserved_during_recovery(self):
"""F-4: Compensating recovery does not delete dirty worktree."""
branch = "fix/issue-850-rec-dirty"
wt_path = os.path.join(self.branches_dir, "fix-issue-850-rec-dirty")
subprocess.run(["git", "-C", self.repo_dir, "worktree", "add", "-b", branch, wt_path], check=True, capture_output=True)
dirty_file = os.path.join(wt_path, "dirty.txt")
with open(dirty_file, "w") as f:
f.write("uncommitted work")
journal = {
"idempotency_key": "test_dirty_rec",
"issue_number": 850,
"branch_name": branch,
"worktree_path": wt_path,
"artifacts_created": {
"worktree_dir_created": True,
"worktree_registered": True,
},
"failure_reason": "test dirty recovery",
}
rec = author_issue_bootstrap.run_compensating_recovery(journal, self.repo_dir, journal_dir=self.lock_dir)
self.assertTrue(os.path.exists(wt_path))
self.assertIn(f"worktree_path_preserved_dirty:{wt_path}", rec["rolled_back"])
def test_branch_with_commits_preserved_during_recovery(self):
"""F-4: Compensating recovery does not delete branch with author commits."""
branch = "fix/issue-850-rec-commits"
subprocess.run(["git", "-C", self.repo_dir, "branch", branch, self.master_sha], check=True, capture_output=True)
# Add a commit on the branch
wt_path = os.path.join(self.branches_dir, "fix-issue-850-rec-commits")
subprocess.run(["git", "-C", self.repo_dir, "worktree", "add", wt_path, branch], check=True, capture_output=True)
cfile = os.path.join(wt_path, "commit.txt")
with open(cfile, "w") as f:
f.write("author commit")
subprocess.run(["git", "-C", wt_path, "add", "commit.txt"], check=True, capture_output=True)
subprocess.run(["git", "-C", wt_path, "commit", "-m", "author commit"], check=True, capture_output=True)
subprocess.run(["git", "-C", self.repo_dir, "worktree", "remove", "--force", wt_path], check=True, capture_output=True)
journal = {
"idempotency_key": "test_commits_rec",
"issue_number": 850,
"branch_name": branch,
"resolved_base_sha": self.master_sha,
"artifacts_created": {
"branch_created": True,
},
"failure_reason": "test commit branch recovery",
}
rec = author_issue_bootstrap.run_compensating_recovery(journal, self.repo_dir, journal_dir=self.lock_dir)
branch_check = subprocess.run(["git", "-C", self.repo_dir, "rev-parse", "--verify", branch], capture_output=True, text=True, check=False)
self.assertEqual(branch_check.returncode, 0, "Branch with commits was deleted!")
self.assertIn(f"branch_preserved_commits:{branch}", rec["rolled_back"])
def test_task_capability_map_integration(self):
"""Verify task_capability_map has bootstrap_author_issue_worktree configured correctly."""
self.assertEqual(task_capability_map.required_role("bootstrap_author_issue_worktree"), "author")
self.assertEqual(task_capability_map.required_permission("bootstrap_author_issue_worktree"), "gitea.branch.create")
self.assertTrue(task_capability_map.preflight_task_matches("work_issue", "bootstrap_author_issue_worktree"))
self.assertTrue(task_capability_map.preflight_task_matches("bootstrap_author_issue_worktree", "lock_issue"))
def test_unverified_assignment_lease_ids_fail_closed(self):
"""Review #531 Finding 4: fabricated assignment/lease IDs are refused."""
res = author_issue_bootstrap.bootstrap_author_issue_worktree(
issue_number=850,
canonical_repo_root=self.repo_dir,
assignment_id="asn-fabricated",
lease_id="lease-fabricated",
expected_base_sha=self.master_sha,
lock_dir=self.lock_dir,
owner_session="session-test-1234",
)
self.assertFalse(res["success"])
self.assertIn(
res.get("reason_code"),
{
"unknown_lease_id",
"assignment_lease_lookup_failed",
"incomplete_assignment_lease_ids",
},
)
def test_partial_assignment_lease_ids_fail_closed(self):
"""Review #531 Finding 4: one of assignment_id/lease_id alone is incomplete."""
res = author_issue_bootstrap.bootstrap_author_issue_worktree(
issue_number=850,
canonical_repo_root=self.repo_dir,
assignment_id="asn-only",
lock_dir=self.lock_dir,
owner_session="session-test-1234",
)
self.assertFalse(res["success"])
self.assertEqual(res.get("reason_code"), "incomplete_assignment_lease_ids")
def test_stale_diverged_branch_is_not_accepted_via_merge_base(self):
"""Review #531 Finding 3: any common ancestor is not enough; require master ⊆ branch."""
branch = "fix/issue-850-stale-divergent"
# Create branch from current master, then advance master so branch lacks tip.
subprocess.run(
["git", "-C", self.repo_dir, "branch", branch, self.master_sha],
check=True,
capture_output=True,
)
# Make a new commit on master (orphan path so branch does not contain it).
marker = os.path.join(self.repo_dir, "master-advance.txt")
with open(marker, "w") as f:
f.write("advance master\n")
subprocess.run(["git", "-C", self.repo_dir, "add", "master-advance.txt"], check=True, capture_output=True)
subprocess.run(
["git", "-C", self.repo_dir, "commit", "-m", "advance master past branch"],
check=True,
capture_output=True,
)
new_master = subprocess.run(
["git", "-C", self.repo_dir, "rev-parse", "HEAD"],
capture_output=True,
text=True,
check=True,
).stdout.strip()
res = author_issue_bootstrap.bootstrap_author_issue_worktree(
issue_number=850,
canonical_repo_root=self.repo_dir,
branch_name=branch,
expected_base_sha=new_master,
lock_dir=self.lock_dir,
owner_session="session-test-1234",
)
self.assertFalse(res["success"])
self.assertEqual(res.get("reason_code"), "incompatible_existing_branch")
def test_compensating_recovery_attempts_lease_release(self):
"""Review #531 Finding 5: recovery invokes lease release when lease_id is present."""
from unittest import mock
journal = {
"idempotency_key": "test_lease_rec",
"issue_number": 850,
"owner_session": "session-test-1234",
"lease_id": "lease-abc",
"branch_name": "fix/issue-850-lease-rec",
"artifacts_created": {"lock_created": True},
"failure_reason": "simulated",
"completed": False,
}
with mock.patch.object(
author_issue_bootstrap.lease_lifecycle,
"release_lease",
return_value={"success": True},
) as rel, mock.patch.object(
author_issue_bootstrap.control_plane_db,
"ControlPlaneDB",
return_value=mock.Mock(),
):
rec = author_issue_bootstrap.run_compensating_recovery(
journal, self.repo_dir, journal_dir=self.lock_dir
)
self.assertTrue(rec["executed"])
rel.assert_called_once()
self.assertIn("lease:lease-abc", rec["rolled_back"])
class TestCanonicalRootNoStringSplit(unittest.TestCase):
def test_fallback_uses_commonpath_not_substring_split(self):
"""Review #531 Finding 2: no norm.split('/branches/') fallback."""
import inspect
import author_mutation_worktree as amw
src = inspect.getsource(amw.resolve_canonical_repo_root)
self.assertNotIn('split("/branches/")', src)
self.assertNotIn("split('/branches/')", src)
# Fallback recovers repo root from a nested branches worktree path.
with tempfile.TemporaryDirectory() as tmp:
repo = os.path.join(tmp, "repo")
wt = os.path.join(repo, "branches", "fix-issue-850-x")
os.makedirs(wt)
# git unavailable path: pass missing workspace so fallback is used.
root = amw.resolve_canonical_repo_root("/missing/path", wt)
self.assertEqual(root, os.path.realpath(repo))
if __name__ == "__main__":
unittest.main()
+22
View File
@@ -28,6 +28,21 @@ class TestPathUnderBranches(unittest.TestCase):
amw.is_path_under_branches("/repo/other-checkout", self.ROOT)
)
def test_unrelated_branches_dir_fails(self):
self.assertFalse(
amw.is_path_under_branches("/tmp/branches/evil", self.ROOT)
)
def test_prefix_confusion_fails(self):
self.assertFalse(
amw.is_path_under_branches(f"{self.ROOT}/branches-other/foo", self.ROOT)
)
def test_traversal_fails(self):
self.assertFalse(
amw.is_path_under_branches(f"{self.ROOT}/branches/../evil", self.ROOT)
)
class TestAssessAuthorMutationWorktree(unittest.TestCase):
ROOT = "/repo/Gitea-Tools"
@@ -71,6 +86,13 @@ class TestAssessAuthorMutationWorktree(unittest.TestCase):
self.assertTrue(result["proven"])
self.assertFalse(result["block"])
def test_path_shaped_branches_ancestry_without_isdir(self):
"""Review #551: commonpath recovery must not require on-disk isdir."""
fake_wt = "/repo/Gitea-Tools/branches/issue-274"
root = amw.resolve_canonical_repo_root(fake_wt, fake_wt)
self.assertEqual(root, "/repo/Gitea-Tools")
self.assertNotIn('split("/branches/")', open(amw.__file__).read())
class TestPreflightIntegration(unittest.TestCase):
def test_verify_preflight_blocks_control_checkout_with_test_porcelain(self):
+352
View File
@@ -1639,6 +1639,358 @@ class TestSecondRemediationIntegration(unittest.TestCase):
self.assertTrue(ownership_calls)
class TestIssue855ExactPrSelector(unittest.TestCase):
"""#855: exact pr_number pin for reconcile_merged_cleanups (#851 lifecycle)."""
def setUp(self):
self._remotes = patch.dict(
mcp_server.REMOTES,
{
"prgs": {
"host": "gitea.example.com",
"org": "Scaled-Tech-Consulting",
"repo": "Gitea-Tools",
}
},
)
self._remotes.start()
patch("gitea_audit.audit_enabled", return_value=False).start()
self.mock_api = patch("mcp_server.api_request").start()
self.mock_all = patch("mcp_server.api_get_all", return_value=[]).start()
patch("mcp_server.get_auth_header", return_value=FAKE_AUTH).start()
patch(
"mcp_server.merged_cleanup_reconcile.is_head_ancestor_of_ref",
return_value=True,
).start()
patch(
"mcp_server.get_profile",
return_value=dict(RECONCILER_WITH_DELETE),
).start()
patch(
"mcp_server._profile_operation_gate",
return_value=[],
).start()
patch(
"mcp_server._collect_branch_ownership_records",
return_value={"records": [], "inventory_error": False},
).start()
patch(
"mcp_server.merged_cleanup_reconcile.discover_reviewer_scratch_worktrees",
return_value=[],
).start()
patch("mcp_server.verify_preflight_purity", return_value=None).start()
patch(
"mcp_server.audit_reconciliation_mode.check_cleanup_execution_allowed",
return_value=(True, []),
).start()
def tearDown(self):
patch.stopall()
def _merged_pr(self, number, branch, sha="c" * 40):
return {
"number": number,
"title": f"PR {number}",
"body": f"Closes #{number - 4}",
"merged": True,
"merged_at": "2026-07-23T12:00:00Z",
"merge_commit_sha": "f" * 40,
"state": "closed",
"head": {"ref": branch, "sha": sha},
"base": {"ref": "master"},
}
def test_exact_pr_848_ignores_newer_852_in_batch_queue(self):
"""pr_number=848 selects only #848 even when #852 is newer/first."""
from mcp_server import gitea_reconcile_merged_cleanups
pr_848 = self._merged_pr(
848, "fix/issue-844-exclude-epic-containers", sha="c3f282ba" + "0" * 32
)
# Closed list would rank #852 first in batch mode; exact pin must ignore it.
closed_batch = [
self._merged_pr(852, "fix/issue-851-cleanup-worktree-before-remote-delete"),
pr_848,
self._merged_pr(849, "fix/issue-849-other"),
self._merged_pr(846, "fix/issue-846-other"),
self._merged_pr(845, "fix/issue-845-other"),
]
batch_fetch_calls = []
def fake_api(method, url, *args, **kwargs):
if method == "GET" and url.rstrip("/").endswith("/pulls/848"):
return dict(pr_848)
if method == "GET" and "/pulls/" in url:
raise AssertionError(f"unexpected PR fetch: {url}")
if method == "GET" and "/branches/" in url:
return {"name": "present"}
return {}
def fake_all(url, auth, limit=None):
batch_fetch_calls.append((url, limit))
if "state=open" in url:
return []
if "state=closed" in url:
# Exact mode must not use the closed batch list.
raise AssertionError(
"exact pr_number mode must not page closed PRs: " + url
)
return []
self.mock_api.side_effect = fake_api
self.mock_all.side_effect = fake_all
patch(
"mcp_server._remote_branch_exists",
return_value=True,
).start()
patch(
"mcp_server.merged_cleanup_reconcile.build_reconciliation_report",
side_effect=lambda **kwargs: {
"entries": [
{
"pr_number": int(pr["number"]),
"head_branch": (pr.get("head") or {}).get("ref"),
"issue_number": 844,
"remote_branch": {
"safe_to_delete_remote": True,
"head_branch": (pr.get("head") or {}).get("ref"),
},
"local_worktree": {
"safe_to_remove_worktree": True,
"worktree_path": (
"/tmp/branches/fix-issue-844-exclude-epic-containers"
),
},
"planned_execution_order": (
mcp_server.merged_cleanup_reconcile.plan_cleanup_execution_order(
remote_assessment={"safe_to_delete_remote": True},
local_assessment={"safe_to_remove_worktree": True},
)
),
}
for pr in kwargs.get("closed_prs") or []
if pr.get("merged_at") or pr.get("merged")
],
"reviewer_scratch_entries": [],
"merged_pr_count": len(kwargs.get("closed_prs") or []),
},
).start()
res = gitea_reconcile_merged_cleanups(
dry_run=True,
pr_number=848,
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
)
self.assertTrue(res.get("success"))
self.assertFalse(res.get("performed"))
self.assertEqual(res.get("selection_mode"), "exact_pr")
self.assertEqual(res.get("selected_pr_number"), 848)
entries = res.get("entries") or []
self.assertEqual(len(entries), 1, entries)
self.assertEqual(entries[0].get("pr_number"), 848)
self.assertEqual(
entries[0].get("head_branch"),
"fix/issue-844-exclude-epic-containers",
)
# No other PR appears in plan.
self.assertEqual(list((res.get("planned_execution_orders") or {}).keys()), ["848"])
plan = (res.get("planned_execution_orders") or {}).get("848") or []
actions = [s.get("action") for s in plan]
self.assertEqual(
actions,
[
"remove_local_worktree",
"reassess_branch_ownership",
"delete_remote_branch",
],
)
# Prove we never scanned the multi-PR closed batch.
self.assertFalse(any("state=closed" in (u or "") for u, _ in batch_fetch_calls))
# closed_batch fixture must remain unused (sanity).
self.assertEqual(closed_batch[0]["number"], 852)
def test_exact_pr_execute_only_mutates_selected_pr(self):
"""Execute with pr_number must never touch #845/#846/#849/#852."""
from mcp_server import gitea_reconcile_merged_cleanups
pr_848 = self._merged_pr(848, "fix/issue-844-exclude-epic-containers")
worktree_path = "/tmp/branches/fix-issue-844-exclude-epic-containers"
remove_calls = []
delete_api_calls = []
ownership_branches = []
def fake_api(method, url, *args, **kwargs):
if method == "GET" and url.rstrip("/").endswith("/pulls/848"):
return dict(pr_848)
if method == "DELETE":
delete_api_calls.append(url)
# Forbid foreign PR branch deletion by URL content.
for forbidden in ("845", "846", "849", "852"):
self.assertNotIn(forbidden, url)
return {}
def fake_remove(project_root, branch, worktree_path=None):
remove_calls.append({"branch": branch, "worktree_path": worktree_path})
return {
"success": True,
"performed": True,
"message": f"removed {worktree_path}",
"worktree_path": worktree_path,
}
def fake_collect(**kwargs):
ownership_branches.append(kwargs.get("branch"))
return {"records": [], "inventory_error": False}
def fake_probe(h, o, r, auth, br):
return guard.classify_branch_readback_http_status(
404, not_found_scope=guard.NOT_FOUND_SCOPE_BRANCH
)
self.mock_api.side_effect = fake_api
self.mock_all.side_effect = lambda url, auth, limit=None: []
patch("mcp_server._remote_branch_exists", return_value=True).start()
patch(
"mcp_server.merged_cleanup_reconcile.build_reconciliation_report",
return_value={
"entries": [
{
"pr_number": 848,
"head_branch": "fix/issue-844-exclude-epic-containers",
"remote_branch": {"safe_to_delete_remote": True},
"local_worktree": {
"safe_to_remove_worktree": True,
"worktree_path": worktree_path,
},
"planned_execution_order": [
{"action": "remove_local_worktree", "phase": 1},
{"action": "reassess_branch_ownership", "phase": 2},
{"action": "delete_remote_branch", "phase": 3},
],
}
],
"reviewer_scratch_entries": [
# Foreign scratch must be filtered before report execute loop;
# if present here it would still be a test failure if acted on.
],
"merged_pr_count": 1,
},
).start()
patch(
"mcp_server.merged_cleanup_reconcile.remove_local_worktree",
side_effect=fake_remove,
).start()
patch(
"mcp_server._collect_branch_ownership_records",
side_effect=fake_collect,
).start()
patch("mcp_server._probe_remote_branch", side_effect=fake_probe).start()
res = gitea_reconcile_merged_cleanups(
dry_run=False,
execute_confirmed=True,
pr_number=848,
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
)
self.assertTrue(res.get("performed") or res.get("executed"))
self.assertEqual(res.get("selection_mode"), "exact_pr")
self.assertEqual(res.get("selected_pr_number"), 848)
actions = res.get("actions") or []
pr_numbers_touched = {
a.get("pr_number") for a in actions if a.get("pr_number") is not None
}
self.assertTrue(pr_numbers_touched.issubset({None, 848}) or not pr_numbers_touched)
removes = [a for a in actions if a.get("action") == "remove_local_worktree"]
deletes = [a for a in actions if a.get("action") == "delete_remote_branch"]
self.assertEqual(len(removes), 1)
self.assertEqual(remove_calls[0]["branch"], "fix/issue-844-exclude-epic-containers")
self.assertEqual(len(deletes), 1)
self.assertTrue(deletes[0].get("success"))
self.assertTrue(deletes[0].get("after_worktree_removal"))
self.assertEqual(len(delete_api_calls), 1)
self.assertEqual(
ownership_branches, ["fix/issue-844-exclude-epic-containers"]
)
def test_exact_pr_unknown_fails_closed_without_mutation(self):
from mcp_server import gitea_reconcile_merged_cleanups
def fake_api(method, url, *args, **kwargs):
if method == "GET" and "/pulls/99999" in url:
raise RuntimeError("HTTP 404 Not Found")
raise AssertionError(f"unexpected API call {method} {url}")
self.mock_api.side_effect = fake_api
res = gitea_reconcile_merged_cleanups(
dry_run=True,
pr_number=99999,
remote="prgs",
)
self.assertFalse(res.get("success"))
self.assertFalse(res.get("performed"))
self.assertEqual(res.get("blocker_kind"), "pr_unresolvable")
self.assertIn("99999", " ".join(res.get("reasons") or []))
def test_exact_pr_not_merged_fails_closed(self):
from mcp_server import gitea_reconcile_merged_cleanups
def fake_api(method, url, *args, **kwargs):
if method == "GET" and url.rstrip("/").endswith("/pulls/900"):
return {
"number": 900,
"merged": False,
"merged_at": None,
"state": "open",
"head": {"ref": "feat/x", "sha": "a" * 40},
}
raise AssertionError(f"unexpected {method} {url}")
self.mock_api.side_effect = fake_api
res = gitea_reconcile_merged_cleanups(
dry_run=False,
execute_confirmed=True,
pr_number=900,
remote="prgs",
)
self.assertFalse(res.get("success"))
self.assertFalse(res.get("performed"))
self.assertEqual(res.get("blocker_kind"), "pr_not_merged")
def test_exact_pr_invalid_number_fails_closed(self):
from mcp_server import gitea_reconcile_merged_cleanups
res = gitea_reconcile_merged_cleanups(
dry_run=True,
pr_number=0,
remote="prgs",
)
self.assertFalse(res.get("success"))
self.assertEqual(res.get("blocker_kind"), "invalid_pr_number")
self.mock_api.assert_not_called()
def test_batch_mode_still_works_without_pr_number(self):
"""Unfiltered batch path remains backward compatible."""
from mcp_server import gitea_reconcile_merged_cleanups
self.mock_all.side_effect = lambda url, auth, limit=None: []
self.mock_api.side_effect = lambda *a, **k: {}
patch(
"mcp_server.merged_cleanup_reconcile.build_reconciliation_report",
return_value={
"entries": [],
"reviewer_scratch_entries": [],
"merged_pr_count": 0,
},
).start()
res = gitea_reconcile_merged_cleanups(dry_run=True, remote="prgs", limit=10)
self.assertTrue(res.get("success"))
self.assertEqual(res.get("selection_mode"), "batch")
self.assertIsNone(res.get("selected_pr_number"))
if __name__ == "__main__":
unittest.main()
+1 -1
View File
@@ -36,7 +36,7 @@ class ControlPlaneDBTest(unittest.TestCase):
rows = dict(conn.execute("SELECT key, value FROM schema_meta").fetchall())
finally:
conn.close()
self.assertEqual(rows["schema_version"], "4")
self.assertEqual(rows["schema_version"], "5")
self.assertIn("DB coordinates", rows["architecture"])
self.assertIn("bridge", rows["architecture"].lower())
@@ -0,0 +1,483 @@
"""Synthetic regression coverage for dirty orphaned worktree recovery (#860).
Modeled on the #850 / #855 shape without mutating their real state.
"""
from __future__ import annotations
import json
import os
import shutil
import tempfile
import unittest
from unittest import mock
import dirty_orphan_worktree_recovery as dorec
import issue_lock_store
DEAD_PID = 999_999_999
LIVE_PID = os.getpid()
BRANCH = "fix/issue-901-dirty-orphan"
SOURCE_WT = "/repo/branches/issue-901-dirty-orphan"
RECOVERY_WT_NAME = "recovery-issue-901-dirty-orphan"
LOCAL_HEAD = "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa"
REMOTE_HEAD = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"
OTHER_HEAD = "cccccccccccccccccccccccccccccccccccccccc"
FP_A = dorec.sha256_bytes(b"dirty-a")
FP_B = dorec.sha256_bytes(b"dirty-b")
FP_C = dorec.sha256_bytes(b"dirty-c-conflict")
def durable_lock(**overrides):
"""#850-shaped PID-less malformed same-claimant lock."""
lock = {
"issue_number": 901,
"branch_name": BRANCH,
"worktree_path": SOURCE_WT,
"remote": "prgs",
"org": "Example-Org",
"repo": "Example-Repo",
# intentionally no pid / session_pid / work_lease expiry
"claimant": {"username": "author-user", "profile": "prgs-author"},
}
lock.update(overrides)
return lock
def base_kwargs(**overrides):
kwargs = {
"issue_number": 901,
"branch_name": BRANCH,
"source_worktree_path": SOURCE_WT,
"remote": "prgs",
"org": "Example-Org",
"repo": "Example-Repo",
"identity": "author-user",
"profile": "prgs-author",
"expected_local_head": LOCAL_HEAD,
"expected_remote_head": REMOTE_HEAD,
"expected_dirty_fingerprints": {"a.py": FP_A, "b.py": FP_B},
"current_branch": BRANCH,
"porcelain_status": " M a.py\n M b.py\n",
"observed_local_head": LOCAL_HEAD,
"observed_remote_head": REMOTE_HEAD,
"observed_dirty_fingerprints": {"a.py": FP_A, "b.py": FP_B},
"competing_live_locks": [],
"competing_live_sessions": [],
"workflow_lease_active": False,
"workflow_lease_expired": True,
"canonical_repo_root": "/repo",
"worktree_registered": True,
"current_pid": LIVE_PID,
}
kwargs.update(overrides)
return kwargs
def assess(lock=None, **overrides):
return dorec.assess_dirty_orphan_recovery(
durable_lock() if lock is None else lock, **base_kwargs(**overrides)
)
class FreshnessPidLess(unittest.TestCase):
def test_pid_less_lock_is_not_live(self):
freshness = issue_lock_store.assess_lock_freshness(durable_lock())
self.assertFalse(freshness["live"])
self.assertTrue(freshness.get("pid_missing"))
self.assertEqual(freshness["status"], "malformed")
def test_pid_less_with_far_future_expiry_still_not_live(self):
lock = durable_lock(
work_lease={
"operation_type": "author_issue_work",
"expires_at": "2999-01-01T00:00:00Z",
"last_heartbeat_at": "2999-01-01T00:00:00Z",
}
)
freshness = issue_lock_store.assess_lock_freshness(lock)
self.assertFalse(freshness["live"])
self.assertTrue(freshness.get("pid_missing"))
class EligibilityGranted(unittest.TestCase):
def test_dead_same_claimant_pid_less_dirty(self):
result = assess()
self.assertEqual(result["outcome"], dorec.ELIGIBLE)
self.assertTrue(result["eligible"])
def test_expired_workflow_lease_corroboration(self):
result = assess(workflow_lease_active=False, workflow_lease_expired=True)
self.assertTrue(result["eligible"])
def test_older_local_newer_remote_heads(self):
result = assess()
self.assertTrue(result["evidence"].get("heads_diverged"))
self.assertTrue(result["eligible"])
class EligibilityRefused(unittest.TestCase):
def test_active_owner_with_pid(self):
lock = durable_lock(pid=LIVE_PID, session_pid=LIVE_PID)
result = assess(lock=lock, owner_process_alive_override=True)
self.assertEqual(result["outcome"], dorec.REFUSED)
self.assertFalse(result["eligible"])
self.assertTrue(any("alive" in r for r in result["reasons"]))
def test_foreign_claimant(self):
result = assess(identity="other-user")
self.assertEqual(result["outcome"], dorec.REFUSED)
self.assertTrue(any("foreign claimant identity" in r for r in result["reasons"]))
def test_foreign_profile(self):
result = assess(profile="prgs-reviewer")
self.assertEqual(result["outcome"], dorec.REFUSED)
def test_fingerprint_mismatch(self):
result = assess(observed_dirty_fingerprints={"a.py": "0" * 64, "b.py": FP_B})
self.assertEqual(result["outcome"], dorec.REFUSED)
self.assertTrue(any("fingerprint mismatch" in r for r in result["reasons"]))
def test_head_mismatch(self):
result = assess(observed_local_head=OTHER_HEAD)
self.assertEqual(result["outcome"], dorec.REFUSED)
def test_remote_head_mismatch(self):
result = assess(observed_remote_head=OTHER_HEAD)
self.assertEqual(result["outcome"], dorec.REFUSED)
def test_path_not_under_branches(self):
result = assess(
source_worktree_path="/tmp/branches/evil",
# lock path also changed so worktree agreement holds
lock=durable_lock(worktree_path="/tmp/branches/evil"),
)
self.assertEqual(result["outcome"], dorec.REFUSED)
self.assertTrue(any("canonical branches" in r for r in result["reasons"]))
def test_unregistered_worktree(self):
result = assess(worktree_registered=False)
self.assertEqual(result["outcome"], dorec.REFUSED)
def test_active_workflow_lease(self):
result = assess(workflow_lease_active=True, workflow_lease_expired=False)
self.assertEqual(result["outcome"], dorec.REFUSED)
def test_unsafe_dirty_path_pin(self):
result = assess(
expected_dirty_fingerprints={"../etc/passwd": FP_A},
observed_dirty_fingerprints={"../etc/passwd": FP_A},
)
self.assertEqual(result["outcome"], dorec.REFUSED)
def test_symlink_escape_rejected_by_ancestry(self):
ok, reasons = dorec.is_path_under_canonical_branches(
"/tmp/branches/evil", canonical_repo_root="/repo"
)
self.assertFalse(ok)
self.assertTrue(reasons)
class ConflictDetection(unittest.TestCase):
def test_overlapping_upstream_change(self):
conflicts = dorec.detect_path_conflicts(
dirty_paths=["c.py"],
local_head_contents={"c.py": b"local-base"},
remote_head_contents={"c.py": b"remote-changed"},
dirty_contents={"c.py": b"dirty-c-conflict"},
)
self.assertEqual(len(conflicts), 1)
self.assertEqual(conflicts[0]["path"], "c.py")
def test_unchanged_upstream_no_conflict(self):
conflicts = dorec.detect_path_conflicts(
dirty_paths=["a.py"],
local_head_contents={"a.py": b"same"},
remote_head_contents={"a.py": b"same"},
dirty_contents={"a.py": b"dirty-a"},
)
self.assertEqual(conflicts, [])
class CrashSafeRecovery(unittest.TestCase):
def setUp(self):
self.tmp = tempfile.mkdtemp(prefix="dirty-orphan-")
self.repo = os.path.join(self.tmp, "repo")
self.branches = os.path.join(self.repo, "branches")
self.source = os.path.join(self.branches, "issue-901-dirty-orphan")
self.recovery = os.path.join(self.branches, RECOVERY_WT_NAME)
os.makedirs(self.source, exist_ok=True)
os.makedirs(self.branches, exist_ok=True)
# seed dirty files in source
with open(os.path.join(self.source, "a.py"), "wb") as fh:
fh.write(b"dirty-a")
with open(os.path.join(self.source, "b.py"), "wb") as fh:
fh.write(b"dirty-b")
self.journal_dir = os.path.join(self.tmp, "journals")
self.lock = durable_lock(worktree_path=self.source)
self.assessment = dorec.assess_dirty_orphan_recovery(
self.lock,
**base_kwargs(
source_worktree_path=self.source,
canonical_repo_root=self.repo,
),
)
class FakeGit(dorec.GitOps):
def __init__(self, recovery_path, head):
self.recovery_path = recovery_path
self.head = head
self.calls = []
def run(self, args, *, cwd):
self.calls.append((args, cwd))
if args[:3] == ["git", "worktree", "add"]:
os.makedirs(self.recovery_path, exist_ok=True)
return mock.Mock(returncode=0, stdout="", stderr="")
if args[:2] == ["git", "checkout"]:
return mock.Mock(returncode=0, stdout="", stderr="")
if args[:2] == ["git", "rev-parse"]:
return mock.Mock(returncode=0, stdout=self.head + "\n", stderr="")
return mock.Mock(returncode=0, stdout="", stderr="")
self.git = FakeGit(self.recovery, REMOTE_HEAD)
self.written_locks = []
def lock_writer(record):
self.written_locks.append(record)
self.lock_writer = lock_writer
def tearDown(self):
shutil.rmtree(self.tmp, ignore_errors=True)
def _run(self, **overrides):
kwargs = {
"assessment": self.assessment,
"existing_lock": self.lock,
"issue_number": 901,
"branch_name": BRANCH,
"source_worktree_path": self.source,
"recovery_worktree_path": self.recovery,
"remote": "prgs",
"org": "Example-Org",
"repo": "Example-Repo",
"identity": "author-user",
"profile": "prgs-author",
"expected_local_head": LOCAL_HEAD,
"expected_remote_head": REMOTE_HEAD,
"expected_dirty_fingerprints": {"a.py": FP_A, "b.py": FP_B},
"dirty_contents": {"a.py": b"dirty-a", "b.py": b"dirty-b"},
"local_head_contents": {"a.py": b"base-a", "b.py": b"base-b"},
"remote_head_contents": {"a.py": b"base-a", "b.py": b"base-b"},
"canonical_repo_root": self.repo,
"bind_lock": True,
"lock_writer": self.lock_writer,
"git_ops": self.git,
"journal_dir": self.journal_dir,
"session_pid": LIVE_PID,
}
kwargs.update(overrides)
return dorec.run_dirty_orphan_recovery(**kwargs)
def test_success_preserves_dirty_bytes_and_source(self):
result = self._run()
self.assertTrue(result["success"])
self.assertEqual(result["outcome"], dorec.RECOVERY_COMPLETED)
self.assertTrue(os.path.isdir(self.source))
with open(os.path.join(self.source, "a.py"), "rb") as fh:
self.assertEqual(fh.read(), b"dirty-a")
with open(os.path.join(self.recovery, "a.py"), "rb") as fh:
self.assertEqual(fh.read(), b"dirty-a")
with open(os.path.join(self.recovery, "b.py"), "rb") as fh:
self.assertEqual(fh.read(), b"dirty-b")
self.assertEqual(len(self.written_locks), 1)
rec = self.written_locks[0]
self.assertEqual(rec["session_pid"], LIVE_PID)
self.assertTrue(rec["dirty_orphan_recovery"]["recovered"])
self.assertTrue(rec["dirty_orphan_recovery"]["source_frozen"])
def test_conflict_leaves_governed_state(self):
result = self._run(
expected_dirty_fingerprints={"c.py": FP_C},
dirty_contents={"c.py": b"dirty-c-conflict"},
local_head_contents={"c.py": b"local-base"},
remote_head_contents={"c.py": b"remote-changed"},
)
# #860 F4: session binding is NOT finalized while conflicts remain
self.assertFalse(result["success"])
self.assertEqual(result["outcome"], dorec.CONFLICTS_PRESENT)
sidecar = os.path.join(self.recovery, "c.py.recovered-dirty")
self.assertTrue(os.path.isfile(sidecar))
state = os.path.join(
self.recovery, dorec.CONFLICT_STATE_DIR, dorec.CONFLICT_STATE_FILE
)
self.assertTrue(os.path.isfile(state))
with open(state, "r", encoding="utf-8") as fh:
payload = json.load(fh)
self.assertEqual(payload["resolution"], "author_edit_required")
def test_interrupt_before_journal_no_artifacts(self):
result = self._run(interrupt_after_phase=dorec.PHASE_ELIGIBILITY)
self.assertFalse(result["success"])
self.assertEqual(result["outcome"], "INTERRUPTED")
self.assertFalse(os.path.isdir(self.recovery))
def test_interrupt_after_journal_then_retry_idempotent(self):
first = self._run(interrupt_after_phase=dorec.PHASE_JOURNAL_PERSISTED)
self.assertEqual(first["outcome"], "INTERRUPTED")
self.assertTrue(first["journal"]["artifacts_created"]["journal"])
second = self._run()
self.assertTrue(second["success"])
# source still recoverable
with open(os.path.join(self.source, "a.py"), "rb") as fh:
self.assertEqual(fh.read(), b"dirty-a")
def test_interrupt_after_worktree_then_retry(self):
first = self._run(interrupt_after_phase=dorec.PHASE_RECOVERY_WORKTREE)
self.assertEqual(first["outcome"], "INTERRUPTED")
self.assertTrue(os.path.isdir(self.recovery))
second = self._run()
self.assertTrue(second["success"])
def test_interrupt_after_binding_then_retry_complete(self):
first = self._run(interrupt_after_phase=dorec.PHASE_BINDING)
self.assertEqual(first["outcome"], "INTERRUPTED")
second = self._run()
self.assertTrue(second["success"])
# completed journal makes further retries no-ops
third = self._run()
self.assertEqual(third["outcome"], dorec.RECOVERY_RESUMED)
def test_source_worktree_never_deleted(self):
self._run()
self.assertTrue(os.path.isdir(self.source))
self.assertTrue(os.path.isfile(os.path.join(self.source, "a.py")))
def test_fingerprint_drift_refuses_without_mutation(self):
result = self._run(dirty_contents={"a.py": b"CHANGED", "b.py": b"dirty-b"})
self.assertFalse(result["success"])
self.assertFalse(os.path.isdir(self.recovery))
class SessionBindingPreflight(unittest.TestCase):
def test_canonical_session_binding_recognized(self):
lock = {
"worktree_path": "/repo/branches/recovery",
"session_pid": LIVE_PID,
"dirty_orphan_recovery": {
"recovered": True,
"conflicts": [],
"recovery_worktree_path": "/repo/branches/recovery",
"source_worktree_path": SOURCE_WT,
"accepted_head": REMOTE_HEAD,
},
}
result = dorec.preflight_recognizes_recovered_provenance(lock)
self.assertTrue(result["recognized"])
def test_conflicts_block_commit_preflight(self):
lock = {
"worktree_path": "/repo/branches/recovery",
"session_pid": LIVE_PID,
"dirty_orphan_recovery": {
"recovered": True,
"conflicts": [{"path": "c.py"}],
},
}
result = dorec.preflight_recognizes_recovered_provenance(lock)
self.assertFalse(result["recognized"])
def test_active_foreign_does_not_mutate(self):
# assess-only path: foreign refused before run
result = assess(identity="intruder")
self.assertFalse(result["eligible"])
class JournalSymlinkRefusal(unittest.TestCase):
def test_symlink_journal_path_refused_on_load(self):
tmp = tempfile.mkdtemp()
try:
real = os.path.join(tmp, "real.json")
with open(real, "w", encoding="utf-8") as fh:
fh.write("{}")
link = os.path.join(tmp, "link.json")
os.symlink(real, link)
key = "symlink-test"
jdir = tmp
path = dorec._journal_path(key, journal_dir=jdir)
with open(path, "w", encoding="utf-8") as fh:
json.dump({"idempotency_key": key}, fh)
os.remove(path)
os.symlink(real, path)
with self.assertRaises(ValueError):
dorec.load_journal(key, journal_dir=jdir)
finally:
shutil.rmtree(tmp, ignore_errors=True)
class RealGitMultiWorktreeIntegration(unittest.TestCase):
def setUp(self):
import subprocess
self.tmp = tempfile.mkdtemp(prefix="git-integration-")
self.repo = os.path.join(self.tmp, "repo")
os.makedirs(self.repo, exist_ok=True)
subprocess.run(["git", "init"], cwd=self.repo, check=True, capture_output=True)
subprocess.run(["git", "config", "user.name", "Test User"], cwd=self.repo, check=True)
subprocess.run(["git", "config", "user.email", "[email protected]"], cwd=self.repo, check=True)
with open(os.path.join(self.repo, "init.txt"), "w") as fh:
fh.write("init")
subprocess.run(["git", "add", "."], cwd=self.repo, check=True)
subprocess.run(["git", "commit", "-m", "init"], cwd=self.repo, check=True)
branch = "fix/issue-999-test"
subprocess.run(["git", "branch", branch], cwd=self.repo, check=True)
self.branches = os.path.join(self.repo, "branches")
self.source = os.path.join(self.branches, "issue-999-test")
subprocess.run(["git", "worktree", "add", self.source, branch], cwd=self.repo, check=True)
self.dirty_path = os.path.join(self.source, "dirty.txt")
with open(self.dirty_path, "w") as fh:
fh.write("dirty-data")
def tearDown(self):
shutil.rmtree(self.tmp, ignore_errors=True)
def test_prepare_recovery_worktree_detached_no_exit_128(self):
import subprocess
head_sha = subprocess.check_output(["git", "rev-parse", "HEAD"], cwd=self.repo, text=True).strip()
rec_wt = os.path.join(self.branches, "recovery-issue-999-test")
res = dorec.prepare_recovery_worktree(
canonical_repo_root=self.repo,
recovery_worktree_path=rec_wt,
branch_name="fix/issue-999-test",
remote_head=head_sha,
)
self.assertTrue(res["success"], res.get("reasons"))
self.assertTrue(os.path.isdir(rec_wt))
def test_real_lock_rebind_recovery_sanctioned(self):
lock_dir = os.path.join(self.tmp, "locks")
rec_wt = os.path.join(self.branches, "recovery-issue-999-test")
os.makedirs(rec_wt, exist_ok=True)
record = {
"remote": "prgs",
"org": "Example-Org",
"repo": "Example-Repo",
"issue_number": 999,
"branch_name": "fix/issue-999-test",
"worktree_path": rec_wt,
"claimant": {"username": "author-user", "profile": "prgs-author"},
}
record_src = dict(record)
record_src["worktree_path"] = self.source
issue_lock_store.bind_session_lock(record_src, lock_dir=lock_dir)
path = issue_lock_store.bind_session_lock(
record,
lock_dir=lock_dir,
recovery_sanctioned=True,
)
self.assertTrue(os.path.isfile(path))
if __name__ == "__main__":
unittest.main()
@@ -0,0 +1,244 @@
"""#855 AC4: an expired reviewer lease must not indefinitely protect an
already-merged branch when no live claimant exists.
Two layers are covered:
* ``branch_cleanup_guard.assess_expired_reviewer_lease_reclaim`` the pure,
fail-closed reclaim decision. Every condition must be provably satisfied or
the lease keeps protecting the branch.
* ``gitea_mcp_server._collect_branch_ownership_records`` the wiring that
supplies authoritative evidence (PR merged state, owner-process liveness,
competing ownership) to that decision, and flips an expired reviewer lease
to reclaimable only under the full policy.
All inputs are fabricated; no real repository, lease, or credential is used.
"""
import importlib
import unittest
from unittest.mock import patch
import branch_cleanup_guard
mcp_server = importlib.import_module("gitea_mcp_server")
FAKE_AUTH = "token fake"
REMOTE = "prgs"
ORG = "Scaled-Tech-Consulting"
REPO = "Gitea-Tools"
HOST = "gitea.prgs.cc"
BRANCH = "feat/issue-638-webui-app-shell-phase1"
PR_NUMBER = 818
class TestAssessExpiredReviewerLeaseReclaim(unittest.TestCase):
"""Pure fail-closed reclaim decision (#855 AC4)."""
def _call(self, **overrides):
base = dict(
role="reviewer",
status="expired",
pr_merged=True,
owner_pid_alive=False,
competing_active_claimant=False,
)
base.update(overrides)
return branch_cleanup_guard.assess_expired_reviewer_lease_reclaim(**base)
def test_full_policy_satisfied_allows_reclaim(self):
out = self._call()
self.assertTrue(out["reclaim_allowed"])
self.assertEqual(out["reasons"], [])
self.assertEqual(out["decision"], "reclaim_expired_reviewer_lease")
def test_stale_dead_process_reviewer_also_reclaimable(self):
out = self._call(status="stale_dead_process")
self.assertTrue(out["reclaim_allowed"])
def test_non_reviewer_role_never_reclaims(self):
for role in ("author", "merger", "controller", "reconciler", "unknown"):
with self.subTest(role=role):
out = self._call(role=role)
self.assertFalse(out["reclaim_allowed"])
self.assertTrue(out["reasons"])
self.assertEqual(out["decision"], "keep_protecting")
def test_active_status_never_reclaims(self):
out = self._call(status="active")
self.assertFalse(out["reclaim_allowed"])
def test_pr_not_merged_blocks_reclaim(self):
out = self._call(pr_merged=False)
self.assertFalse(out["reclaim_allowed"])
def test_pr_merged_unknown_fails_closed(self):
out = self._call(pr_merged=None)
self.assertFalse(out["reclaim_allowed"])
def test_owner_process_alive_blocks_reclaim(self):
out = self._call(owner_pid_alive=True)
self.assertFalse(out["reclaim_allowed"])
def test_owner_liveness_unknown_fails_closed(self):
out = self._call(owner_pid_alive=None)
self.assertFalse(out["reclaim_allowed"])
def test_competing_active_claimant_blocks_reclaim(self):
out = self._call(competing_active_claimant=True)
self.assertFalse(out["reclaim_allowed"])
def test_competing_claimant_unknown_fails_closed(self):
out = self._call(competing_active_claimant=None)
self.assertFalse(out["reclaim_allowed"])
def test_reasons_never_leak_secrets(self):
out = self._call(role="author")
blob = " ".join(out["reasons"]).lower()
self.assertNotIn("token", blob)
self.assertNotIn("password", blob)
class _FakeLease(dict):
pass
class TestCollectorExpiredReviewerReclaimWiring(unittest.TestCase):
"""`_collect_branch_ownership_records` supplies authoritative evidence and
flips an expired reviewer lease to reclaimable only under the full policy."""
def _run(
self,
*,
lease_role="reviewer",
lease_freshness="stale_dead_process",
owner_pid_alive=False,
pr_merged=True,
extra_leases=None,
worktree_on_branch=False,
):
lease = _FakeLease(
role=lease_role,
work_kind="pr",
work_number=PR_NUMBER,
branch=BRANCH,
status="active",
owner_pid=999999,
remote=REMOTE,
org=ORG,
repo=REPO,
host=HOST,
freshness={
"freshness": lease_freshness,
"owner_pid": 999999,
"owner_pid_alive": owner_pid_alive,
"expired_by_time": lease_freshness == "expired",
},
)
leases = [lease] + list(extra_leases or [])
pr_payload = {
"number": PR_NUMBER,
"merged": pr_merged,
"merged_at": "2026-07-23T00:00:00Z" if pr_merged else None,
"head": {"ref": BRANCH},
}
def fake_api_request(method, url, *a, **k):
if method == "GET" and f"/pulls/{PR_NUMBER}" in url:
return pr_payload
raise AssertionError(f"unexpected api_request {method} {url}")
wt_entries = []
if worktree_on_branch:
wt_entries = [{"branch": BRANCH, "path": f"/x/branches/{BRANCH}"}]
with patch.object(
mcp_server.lease_lifecycle,
"list_active_leases",
return_value={"leases": leases},
), patch.object(
mcp_server.control_plane_db, "get_db", return_value=object(), create=True
), patch.object(
mcp_server.issue_lock_store, "iter_lock_files", return_value=[]
), patch.object(
mcp_server.worktree_cleanup_audit,
"list_worktrees",
return_value=wt_entries,
), patch.object(
mcp_server, "api_get_all", return_value=[]
), patch.object(
mcp_server, "api_request", side_effect=fake_api_request
):
return mcp_server._collect_branch_ownership_records(
remote=REMOTE,
host=HOST,
org=ORG,
repo=REPO,
branch=BRANCH,
pr_number=PR_NUMBER,
project_root="/x",
auth=FAKE_AUTH,
base_api="https://gitea.prgs.cc/api/v1/repos/x/y",
)
def _reviewer_records(self, bundle):
return [
rec
for rec in bundle["records"]
if rec.get("category")
== branch_cleanup_guard.OWNERSHIP_CATEGORY_REVIEWER_LEASE
]
def test_merged_dead_uncontested_reviewer_lease_is_reclaimable(self):
bundle = self._run()
self.assertFalse(bundle["inventory_error"])
recs = self._reviewer_records(bundle)
self.assertEqual(len(recs), 1)
self.assertTrue(recs[0]["reclaim_allowed"])
# And the guard consequently does not block deletion on it.
ownership = branch_cleanup_guard.assess_active_branch_ownership(
remote=REMOTE, org=ORG, repo=REPO, branch=BRANCH, host=HOST,
records=bundle["records"],
)
self.assertFalse(ownership["block"])
def test_unmerged_pr_keeps_reviewer_lease_protective(self):
bundle = self._run(pr_merged=False)
recs = self._reviewer_records(bundle)
self.assertEqual(len(recs), 1)
self.assertFalse(recs[0]["reclaim_allowed"])
ownership = branch_cleanup_guard.assess_active_branch_ownership(
remote=REMOTE, org=ORG, repo=REPO, branch=BRANCH, host=HOST,
records=bundle["records"],
)
self.assertTrue(ownership["block"])
def test_owner_process_alive_keeps_reviewer_lease_protective(self):
bundle = self._run(owner_pid_alive=True, lease_freshness="expired")
recs = self._reviewer_records(bundle)
self.assertFalse(recs[0]["reclaim_allowed"])
def test_competing_worktree_binding_keeps_reviewer_lease_protective(self):
bundle = self._run(worktree_on_branch=True)
recs = self._reviewer_records(bundle)
self.assertFalse(recs[0]["reclaim_allowed"])
ownership = branch_cleanup_guard.assess_active_branch_ownership(
remote=REMOTE, org=ORG, repo=REPO, branch=BRANCH, host=HOST,
records=bundle["records"],
)
self.assertTrue(ownership["block"])
def test_expired_author_lease_never_reclaimed_by_reviewer_policy(self):
bundle = self._run(lease_role="author")
author_recs = [
rec
for rec in bundle["records"]
if rec.get("category")
== branch_cleanup_guard.OWNERSHIP_CATEGORY_AUTHOR_LEASE
]
self.assertEqual(len(author_recs), 1)
self.assertFalse(author_recs[0]["reclaim_allowed"])
if __name__ == "__main__":
unittest.main()
@@ -0,0 +1,630 @@
"""Durable linked-issue lock head refresh + merge-sync dead-session recovery (#871).
``gitea_update_pr_branch_by_merge`` advances a PR's *remote* head but historically
never advanced the linked durable issue lock's recorded head. After the owning
session died the drifted lock became unrecoverable and no further synchronization
was possible (PR #866 / issue #855).
Two halves are covered:
* the write-side refresh (``issue_lock_store.assess/apply_durable_lock_head_refresh``)
that records the new synced head under compare-and-swap with read-after-write; and
* the read-side recovery relation (``issue_lock_recovery`` +
``issue_lock_worktree.read_merge_sync_provenance``) that lets a dead-session lock
whose recorded head is a merge-sync *ancestor* of the live PR head be recovered
and nothing else.
"""
from __future__ import annotations
import os
import subprocess
import sys
import tempfile
import unittest
from datetime import datetime, timedelta, timezone
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
import issue_lock_recovery # noqa: E402
import issue_lock_store # noqa: E402
import issue_lock_worktree # noqa: E402
ISSUE = 8710
PR_NUMBER = 8711
BRANCH = f"fix/issue-{ISSUE}-durable-lock-head-refresh"
IDENTITY = "example-user"
PROFILE = "example-author"
OLD = "a" * 40
NEW1 = "b" * 40
NEW2 = "c" * 40
BASE = "d" * 40
REMOTE = "prgs"
ORG = "ExampleOrg"
REPO = "ExampleRepo"
def dead_pid() -> int:
proc = subprocess.Popen([sys.executable, "-c", "pass"])
proc.wait()
return proc.pid
def future_ts(hours: int = 4) -> str:
return (
(datetime.now(timezone.utc) + timedelta(hours=hours))
.isoformat()
.replace("+00:00", "Z")
)
def _git(cwd, *args):
return subprocess.run(
["git", "-C", cwd, *args],
capture_output=True,
text=True,
check=True,
)
def _rev(cwd, ref="HEAD") -> str:
return _git(cwd, "rev-parse", ref).stdout.strip()
def build_merge_sync_repo(tmp: str) -> dict:
"""Build a repo where a feature branch was synced by merging master in.
Returns a dict with the prior (branch) head, the synced merge-commit head,
the master tip, plus a rebase-style linear descendant and an unrelated head.
"""
_git(tmp, "init", "-q", "-b", "master")
_git(tmp, "config", "user.email", "[email protected]")
_git(tmp, "config", "user.name", "T")
Path(tmp, "base.txt").write_text("base\n")
_git(tmp, "add", "-A")
_git(tmp, "commit", "-q", "-m", "root")
# Feature branch cut from root, one commit — this is the PRIOR/recorded head.
_git(tmp, "checkout", "-q", "-b", BRANCH)
Path(tmp, "feature.txt").write_text("feature\n")
_git(tmp, "add", "-A")
_git(tmp, "commit", "-q", "-m", "feature work")
prior = _rev(tmp)
# Master advances (the base the sync will merge in).
_git(tmp, "checkout", "-q", "master")
Path(tmp, "base.txt").write_text("base\nmore\n")
_git(tmp, "add", "-A")
_git(tmp, "commit", "-q", "-m", "master advance")
master_tip = _rev(tmp)
# Sync: merge master INTO the feature branch → merge commit, first parent = prior.
_git(tmp, "checkout", "-q", BRANCH)
_git(tmp, "merge", "-q", "--no-ff", "-m", "Merge master into feature", "master")
synced = _rev(tmp)
# A plain linear descendant of prior (NOT a merge) — a rebase/extra-commit shape.
_git(tmp, "checkout", "-q", "-b", "linear-branch", prior)
Path(tmp, "extra.txt").write_text("extra\n")
_git(tmp, "add", "-A")
_git(tmp, "commit", "-q", "-m", "extra linear commit")
linear = _rev(tmp)
# An unrelated root (force-push / rewritten history shape).
unrelated_dir = tempfile.mkdtemp()
_git(unrelated_dir, "init", "-q", "-b", "x")
_git(unrelated_dir, "config", "user.email", "[email protected]")
_git(unrelated_dir, "config", "user.name", "T")
Path(unrelated_dir, "z.txt").write_text("z\n")
_git(unrelated_dir, "add", "-A")
_git(unrelated_dir, "commit", "-q", "-m", "unrelated")
unrelated = _rev(unrelated_dir)
# Leave the worktree checked out on the feature branch at the PRIOR head, as
# a dead author session that never advanced would have left it.
_git(tmp, "checkout", "-q", BRANCH)
_git(tmp, "reset", "-q", "--hard", prior)
return {
"prior": prior,
"master_tip": master_tip,
"synced": synced,
"linear": linear,
"unrelated": unrelated,
}
# ─────────────────────────── write-side refresh ───────────────────────────
class TestDurableLockHeadRefresh(unittest.TestCase):
def setUp(self):
self.lock_dir = tempfile.mkdtemp()
self.wt = tempfile.mkdtemp()
lock_data = {
"issue_number": ISSUE,
"branch_name": BRANCH,
"worktree_path": self.wt,
"remote": REMOTE,
"org": ORG,
"repo": REPO,
"claimant": {"username": IDENTITY, "profile": PROFILE},
"work_lease": {
"operation_type": issue_lock_store.AUTHOR_ISSUE_WORK_LEASE,
"issue_number": ISSUE,
"branch": BRANCH,
"worktree_path": self.wt,
"claimant": {"username": IDENTITY, "profile": PROFILE},
"expires_at": future_ts(),
},
}
issue_lock_store.bind_session_lock(lock_data, lock_dir=self.lock_dir)
def _apply(self, **over):
kw = dict(
remote=REMOTE, org=ORG, repo=REPO, issue_number=ISSUE,
branch_name=BRANCH, worktree_path=self.wt, pr_number=PR_NUMBER,
identity=IDENTITY, profile=PROFILE, current_pid=os.getpid(),
expected_old_head=OLD, new_head=NEW1, synced_at=future_ts(0),
base_head=BASE, lock_dir=self.lock_dir,
)
kw.update(over)
return issue_lock_store.apply_durable_lock_head_refresh(**kw)
def _load(self):
return issue_lock_store.load_issue_lock(
remote=REMOTE, org=ORG, repo=REPO, issue_number=ISSUE,
lock_dir=self.lock_dir,
)
def test_first_sync_updates_recorded_head(self):
"""AC1: first base sync writes the resulting head to the durable lock."""
res = self._apply()
self.assertTrue(res["refreshed"], res["reasons"])
self.assertTrue(res["read_after_write_ok"])
self.assertEqual(self._load().get("synced_pr_head"), NEW1)
def test_second_sync_after_master_advance(self):
"""AC2: a later master advance permits a second sanctioned sync."""
self.assertTrue(self._apply()["refreshed"])
res2 = self._apply(expected_old_head=NEW1, new_head=NEW2)
self.assertTrue(res2["refreshed"], res2["reasons"])
self.assertEqual(self._load().get("synced_pr_head"), NEW2)
history = self._load().get("branch_sync_history")
self.assertEqual(len(history), 2)
self.assertEqual(history[0]["last_synced_pr_head"], NEW1)
self.assertEqual(history[1]["prior_pr_head"], NEW1)
def test_cas_detects_concurrent_head_change(self):
"""AC6: CAS refuses when the recorded synced head is not the old head."""
self.assertTrue(self._apply()["refreshed"]) # recorded head now NEW1
# A second sync claiming the old head is still OLD must fail closed.
res = self._apply(expected_old_head=OLD, new_head=NEW2)
self.assertFalse(res["refreshed"])
self.assertTrue(any("CAS" in r or "concurrent" in r for r in res["reasons"]))
self.assertEqual(self._load().get("synced_pr_head"), NEW1)
def test_wrong_issue_fails_closed(self):
res = self._apply(issue_number=999999)
self.assertFalse(res["refreshed"])
def test_wrong_branch_fails_closed(self):
res = self._apply(branch_name="fix/issue-8710-wrong")
self.assertFalse(res["refreshed"])
def test_wrong_repo_fails_closed(self):
res = self._apply(repo="OtherRepo")
self.assertFalse(res["refreshed"])
def test_wrong_identity_fails_closed(self):
res = self._apply(identity="intruder")
self.assertFalse(res["refreshed"])
def test_wrong_profile_fails_closed(self):
res = self._apply(profile="prgs-reviewer")
self.assertFalse(res["refreshed"])
def test_foreign_session_fails_closed(self):
"""A refresh is not a recovery: the current process must own the lock."""
path = issue_lock_store.lock_file_path(
remote=REMOTE, org=ORG, repo=REPO, issue_number=ISSUE,
lock_dir=self.lock_dir,
)
rec = issue_lock_store.read_lock_file(path)
rec["session_pid"] = dead_pid()
rec["pid"] = rec["session_pid"]
issue_lock_store.save_lock_file(path, rec)
res = self._apply()
self.assertFalse(res["refreshed"])
self.assertTrue(any("current session" in r or "live owner" in r for r in res["reasons"]))
def test_new_equals_old_fails_closed(self):
res = self._apply(expected_old_head=OLD, new_head=OLD)
self.assertFalse(res["refreshed"])
def test_non_full_sha_fails_closed(self):
self.assertFalse(self._apply(new_head="deadbeef")["refreshed"])
self.assertFalse(self._apply(expected_old_head="xyz")["refreshed"])
def test_no_lock_fails_closed(self):
assessment = issue_lock_store.assess_durable_lock_head_refresh(
None, remote=REMOTE, org=ORG, repo=REPO, issue_number=ISSUE,
branch_name=BRANCH, worktree_path=self.wt, pr_number=PR_NUMBER,
identity=IDENTITY, profile=PROFILE, current_pid=os.getpid(),
expected_old_head=OLD, new_head=NEW1,
)
self.assertFalse(assessment["allowed"])
# ─────────────────────── merge-sync provenance (real git) ───────────────────
class TestMergeSyncProvenanceObservation(unittest.TestCase):
def setUp(self):
self.tmp = tempfile.mkdtemp()
self.shas = build_merge_sync_repo(self.tmp)
def test_merge_sync_is_recognized(self):
obs = issue_lock_worktree.read_merge_sync_provenance(
self.tmp, prior_head_sha=self.shas["prior"],
synced_head_sha=self.shas["synced"],
)
self.assertTrue(obs["is_merge_sync"], obs["reasons"])
self.assertTrue(obs["prior_is_ancestor"])
self.assertTrue(obs["synced_is_merge"])
self.assertTrue(obs["first_parent_reaches_prior"])
def test_linear_descendant_is_not_a_merge_sync(self):
"""A plain non-merge descendant (rebase/extra commit) is not a sync."""
obs = issue_lock_worktree.read_merge_sync_provenance(
self.tmp, prior_head_sha=self.shas["prior"],
synced_head_sha=self.shas["linear"],
)
self.assertTrue(obs["probe_ok"])
self.assertFalse(obs["is_merge_sync"])
self.assertFalse(obs["synced_is_merge"])
def test_unrelated_history_fails_closed(self):
"""A rewritten/force-pushed head where prior is unreachable fails closed."""
obs = issue_lock_worktree.read_merge_sync_provenance(
self.tmp, prior_head_sha=self.shas["prior"],
synced_head_sha=self.shas["unrelated"],
)
self.assertFalse(obs["is_merge_sync"])
def test_missing_args_fail_closed(self):
obs = issue_lock_worktree.read_merge_sync_provenance(
self.tmp, prior_head_sha=None, synced_head_sha=self.shas["synced"],
)
self.assertFalse(obs["is_merge_sync"])
# ──────────────────── merge-sync dead-session recovery ──────────────────────
def make_dead_lock(worktree, **over):
pid = dead_pid()
lock = {
"issue_number": ISSUE,
"branch_name": BRANCH,
"worktree_path": worktree,
"remote": REMOTE,
"org": ORG,
"repo": REPO,
"session_pid": pid,
"pid": pid,
"claimant": {"username": IDENTITY, "profile": PROFILE},
"work_lease": {
"operation_type": issue_lock_store.AUTHOR_ISSUE_WORK_LEASE,
"issue_number": ISSUE,
"branch": BRANCH,
"worktree_path": worktree,
"claimant": {"username": IDENTITY, "profile": PROFILE},
"expires_at": future_ts(),
},
}
lock.update(over)
return lock
def sync_prov(prior, synced, **over):
d = {
"prior_head_sha": prior,
"synced_head_sha": synced,
"probe_ok": True,
"prior_present": True,
"synced_present": True,
"prior_is_ancestor": True,
"synced_is_merge": True,
"first_parent_reaches_prior": True,
"is_merge_sync": True,
"first_parent_sha": prior,
"parent_count": 2,
"proof": f"{synced} merged base into branch above {prior}",
"reasons": [],
}
d.update(over)
return d
class TestMergeSyncRecovery(unittest.TestCase):
def setUp(self):
self.tmp = tempfile.mkdtemp()
self.shas = build_merge_sync_repo(self.tmp)
self.prior = self.shas["prior"]
self.synced = self.shas["synced"]
def _assess(self, **over):
lock = over.pop("_lock", None) or make_dead_lock(self.tmp)
kw = dict(
issue_number=ISSUE, branch_name=BRANCH, worktree_path=self.tmp,
remote=REMOTE, org=ORG, repo=REPO, identity=IDENTITY, profile=PROFILE,
current_branch=BRANCH, porcelain_status="",
head_sha=self.prior, remote_head_sha=self.synced,
pr_head_sha=self.synced, pr_number=PR_NUMBER,
competing_live_locks=[], candidate_branches=[BRANCH],
current_pid=os.getpid(),
remote_branch_exists=True,
sync_provenance=sync_prov(self.prior, self.synced),
)
kw.update(over)
return issue_lock_recovery.assess_dead_session_lock_recovery(lock, **kw)
def test_merge_sync_drift_is_recoverable(self):
"""AC3/AC4: dead session, recorded head is a merge-sync ancestor of PR head."""
res = self._assess()
self.assertEqual(res["outcome"], issue_lock_recovery.RECOVERY_SANCTIONED, res["reasons"])
self.assertEqual(
res["evidence"]["head_relation"],
issue_lock_recovery.HEAD_RELATION_REMOTE_MERGE_SYNCED,
)
self.assertEqual(res["evidence"]["accepted_head"], self.synced)
def test_missing_provenance_fails_closed(self):
"""No server-derived provenance → cannot accept a remote ahead of local."""
res = self._assess(sync_provenance=None)
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_non_ancestor_recorded_head_fails_closed(self):
"""AC7: provenance that does not prove ancestry is rejected."""
res = self._assess(
sync_provenance=sync_prov(
self.prior, self.synced, prior_is_ancestor=False, is_merge_sync=False,
reasons=["prior head is not an ancestor"],
)
)
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_force_pushed_history_fails_closed(self):
"""AC8: a rewritten head (not a merge sync) stays protected."""
res = self._assess(
sync_provenance=sync_prov(
self.prior, self.synced, is_merge_sync=False, synced_is_merge=False,
reasons=["not a merge-based sync"],
)
)
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_provenance_for_other_commits_fails_closed(self):
"""Provenance whose endpoints differ from the heads under assessment is rejected."""
res = self._assess(
sync_provenance=sync_prov("f" * 40, self.synced),
)
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_dirty_worktree_fails_closed(self):
"""AC11: dirty worktrees remain protected."""
res = self._assess(porcelain_status=" M feature.txt\n")
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_live_owner_fails_closed(self):
"""AC10: a live recorded owner is not a dead-session recovery."""
lock = make_dead_lock(self.tmp, session_pid=os.getpid(), pid=os.getpid())
res = self._assess(_lock=lock)
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_competing_claimant_fails_closed(self):
"""AC13: a competing live lock blocks recovery."""
res = self._assess(
competing_live_locks=[{
"issue_number": ISSUE, "branch_name": BRANCH,
"worktree_path": "/some/other/wt", "pid": os.getpid(),
}]
)
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_wrong_branch_fails_closed(self):
"""AC9: worktree on a different branch fails closed."""
res = self._assess(current_branch="fix/issue-8710-other")
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_wrong_identity_fails_closed(self):
res = self._assess(identity="intruder")
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_pr_head_mismatch_fails_closed(self):
"""The open PR must sit at the synced remote head."""
res = self._assess(pr_head_sha="e" * 40)
self.assertEqual(res["outcome"], issue_lock_recovery.REFUSED)
def test_owning_pr_evidence_for_merge_sync(self):
res = self._assess()
ev = issue_lock_recovery.owning_pr_recovery_evidence(res)
self.assertIsNotNone(ev)
self.assertEqual(ev["pr_number"], PR_NUMBER)
self.assertEqual(ev["head_sha"], self.synced)
self.assertEqual(
ev["head_relation"],
issue_lock_recovery.HEAD_RELATION_REMOTE_MERGE_SYNCED,
)
def test_recovered_owning_pr_from_persisted_record(self):
res = self._assess()
record = issue_lock_recovery.build_recovery_record(res, recovered_at=future_ts(0))
lock = {"issue_number": ISSUE, "branch_name": BRANCH,
"dead_session_recovery": record}
rebuilt = issue_lock_recovery.recovered_owning_pr_from_lock(lock)
self.assertIsNotNone(rebuilt)
self.assertEqual(rebuilt["head_sha"], self.synced)
self.assertEqual(
rebuilt["head_relation"],
issue_lock_recovery.HEAD_RELATION_REMOTE_MERGE_SYNCED,
)
class TestExistingRelationsUnchanged(unittest.TestCase):
"""AC14/AC15: equal-head recovery still works; merge-sync did not weaken it."""
def setUp(self):
self.tmp = tempfile.mkdtemp()
self.shas = build_merge_sync_repo(self.tmp)
def test_equal_head_recovery_still_sanctioned(self):
# Worktree at prior head; remote also at prior head → the #753 equal case.
prior = self.shas["prior"]
lock = make_dead_lock(self.tmp)
res = issue_lock_recovery.assess_dead_session_lock_recovery(
lock, issue_number=ISSUE, branch_name=BRANCH, worktree_path=self.tmp,
remote=REMOTE, org=ORG, repo=REPO, identity=IDENTITY, profile=PROFILE,
current_branch=BRANCH, porcelain_status="",
head_sha=prior, remote_head_sha=prior,
pr_head_sha=prior, pr_number=PR_NUMBER,
competing_live_locks=[], candidate_branches=[BRANCH],
current_pid=os.getpid(), remote_branch_exists=True,
)
self.assertEqual(res["outcome"], issue_lock_recovery.RECOVERY_SANCTIONED, res["reasons"])
self.assertEqual(
res["evidence"]["head_relation"], issue_lock_recovery.HEAD_RELATION_EQUAL,
)
class TestUpdatePrWrapperPartialFailure(unittest.TestCase):
"""AC5/AC16: the tool advances the remote head then refreshes the durable lock.
When the durable refresh fails after the remote advance, the tool must report a
partial lifecycle failure and NOT a fully successful synchronization. Exact PR-
head / base-head pinning is preserved (delegated to the real preflight, stubbed
here only to isolate the post-update lifecycle branch).
"""
def setUp(self):
import gitea_mcp_server as gms # noqa: E402
self.gms = gms
self._orig = {}
def _patch(name, value):
self._orig[name] = getattr(gms, name)
setattr(gms, name, value)
_patch("get_profile", lambda *a, **k: {
"allowed_operations": ["gitea.branch.push"],
"forbidden_operations": [],
"profile_name": "prgs-author",
})
_patch("_role_kind", lambda *a, **k: "author")
_patch("_profile_operation_gate", lambda *a, **k: None)
_patch("_permission_block_report", lambda *a, **k: {})
_patch("_resolve", lambda *a, **k: ("gitea.prgs.cc", ORG, REPO))
_patch("_verify_role_mutation_workspace", lambda *a, **k: None)
_patch("_get_workspace_porcelain", lambda *a, **k: "")
_patch("_canonical_local_git_root", lambda *a, **k: "/x")
_patch("_master_parity_block", lambda *a, **k: None)
_patch("_auth", lambda *a, **k: {"token": "x"})
_patch("repo_api_url", lambda *a, **k: "http://api")
_patch("_redact", lambda s: s)
_patch("_work_lease_claimant", lambda *a, **k: {
"username": IDENTITY, "profile": PROFILE,
})
_patch("_prove_author_ownership_for_pr", lambda *a, **k: {
"has_author_lock": True, "matched_issue": ISSUE,
"matched_via": "branch", "linked_issues": [ISSUE],
"recovered_owning_pr": None, "reasons": [],
})
# Real preflight is unit-tested elsewhere; stub it to isolate the
# post-update durable-lock lifecycle branch under test.
orig_pf = gms.pr_sync_status.assess_update_pr_branch_preflight
self._orig_pf = orig_pf
gms.pr_sync_status.assess_update_pr_branch_preflight = (
lambda *a, **k: {"mutation_allowed": True, "reasons": [], "performed": False}
)
# Sequence the two GET /pulls calls: OLD before update, NEW after.
self._pull_calls = {"n": 0}
def fake_api_request(method, url, auth, *a, **k):
m = method.upper()
if m == "GET" and url.endswith(f"/pulls/{PR_NUMBER}"):
self._pull_calls["n"] += 1
head = OLD if self._pull_calls["n"] == 1 else NEW1
return {
"state": "open",
"head": {"sha": head, "ref": BRANCH},
"base": {"sha": BASE, "ref": "master"},
"mergeable": True, "title": "t", "body": "b",
}
if m == "GET" and "/branches/" in url:
return {"commit": {"id": BASE}}
if m == "POST" and "/update" in url:
return {}
return {}
_patch("api_request", fake_api_request)
def tearDown(self):
for name, value in self._orig.items():
setattr(self.gms, name, value)
self.gms.pr_sync_status.assess_update_pr_branch_preflight = self._orig_pf
def _run(self):
return self.gms.gitea_update_pr_branch_by_merge(
pr_number=PR_NUMBER,
expected_pr_head_sha=OLD,
expected_base_head_sha=BASE,
remote=REMOTE,
worktree_path="/tmp/branches/wt-871",
)
def test_partial_failure_when_refresh_fails(self):
self._orig["apply_durable_lock_head_refresh"] = (
self.gms.issue_lock_store.apply_durable_lock_head_refresh
)
self.gms.issue_lock_store.apply_durable_lock_head_refresh = (
lambda **k: {"refreshed": False, "reasons": ["forced refresh failure"]}
)
try:
res = self._run()
finally:
self.gms.issue_lock_store.apply_durable_lock_head_refresh = (
self._orig["apply_durable_lock_head_refresh"]
)
self.assertTrue(res["performed"])
self.assertEqual(res["new_pr_head_sha"], NEW1)
self.assertFalse(res["success"])
self.assertTrue(res["partial_lifecycle_failure"])
self.assertFalse(res["durable_lock_refreshed"])
def test_full_success_when_refresh_succeeds(self):
self._orig["apply_durable_lock_head_refresh"] = (
self.gms.issue_lock_store.apply_durable_lock_head_refresh
)
self.gms.issue_lock_store.apply_durable_lock_head_refresh = (
lambda **k: {"refreshed": True, "read_after_write_ok": True,
"new_head": NEW1, "reasons": ["ok"]}
)
try:
res = self._run()
finally:
self.gms.issue_lock_store.apply_durable_lock_head_refresh = (
self._orig["apply_durable_lock_head_refresh"]
)
self.assertTrue(res["success"])
self.assertTrue(res["performed"])
self.assertTrue(res["durable_lock_refreshed"])
self.assertTrue(res["fully_synchronized"])
self.assertEqual(res["new_pr_head_sha"], NEW1)
if __name__ == "__main__":
unittest.main()
+6
View File
@@ -24,6 +24,8 @@ def _lease(expires_at: str) -> dict:
def _lock_record(**overrides) -> dict:
# #860: live locks require a usable session pid; PID-less records are never
# classified live merely because expiry/heartbeat fields are present.
record = {
"issue_number": 420,
"branch_name": "feat/issue-420-server-code-parity",
@@ -31,6 +33,8 @@ def _lock_record(**overrides) -> dict:
"org": "Scaled-Tech-Consulting",
"repo": "Gitea-Tools",
"worktree_path": "/tmp/wt-420",
"session_pid": os.getpid(),
"pid": os.getpid(),
"work_lease": _lease("2999-01-01T00:00:00Z"),
}
record.update(overrides)
@@ -88,6 +92,8 @@ class TestIssueLockStore(unittest.TestCase):
existing = _lock_record(
branch_name="feat/issue-420-other",
worktree_path="/tmp/other",
session_pid=os.getpid(),
pid=os.getpid(),
work_lease=_lease("2999-01-01T00:00:00Z"),
)
path = ils.lock_file_path(
+107
View File
@@ -0,0 +1,107 @@
"""Documentation acceptance for the MCP restart governance ADR (#656).
Enforces issue #656 acceptance criteria:
* AC1 policy document exists with an authorization matrix and the recorded
v1 decision (controller approval + automated safety gates).
* AC2 restart is stated as a last resort with enumerated narrower recoveries.
* AC3 a unilateral LLM full restart with affected sessions is forbidden.
* AC4 break-glass conditions are listed.
* AC5 the ADR is linked to #655, #652, #653, #630, #642, and is cross-linked
from the safety model and the web-console deployment boundary docs.
"""
from pathlib import Path
REPO_ROOT = Path(__file__).resolve().parent.parent
ADR = REPO_ROOT / "docs" / "architecture" / "mcp-restart-governance.md"
ADR_BASENAME = "mcp-restart-governance.md"
CROSS_LINK_DOCS = (
REPO_ROOT / "docs" / "safety-model.md",
REPO_ROOT / "docs" / "webui-deployment.md",
)
LINKED_ISSUES = ("#655", "#652", "#653", "#630", "#642")
POLICY_IDS = ("RG-01", "RG-02", "RG-03", "RG-04", "RG-05", "RG-06", "RG-07", "RG-08")
def _read(path: Path) -> str:
assert path.is_file(), f"missing {path.relative_to(REPO_ROOT)}"
return path.read_text(encoding="utf-8")
def test_ac1_adr_exists_with_matrix_and_v1_decision():
text = _read(ADR)
lower = text.lower()
assert text.lstrip().startswith("#"), "ADR lacks a title"
assert "#656" in text
assert "authorization matrix" in lower
# The matrix is a real table with the worker and privileged roles.
for role in ("author", "reviewer", "merger", "reconciler", "controller",
"operator", "admin"):
assert role in lower, f"authorization matrix missing role {role!r}"
# Recorded v1 decision.
assert "restart-governance/v1" in text
assert "controller approval" in lower and "automated safety gates" in lower
def test_ac2_restart_is_last_resort_with_narrower_recoveries():
text = _read(ADR)
lower = text.lower()
assert "last resort" in lower
# Enumerated narrower recoveries precede full restart on the ladder.
for rung in ("reconnect", "rebind", "scoped restart", "full restart",
"host"):
assert rung in lower, f"recovery ladder missing rung {rung!r}"
def test_ac3_forbids_unilateral_llm_full_restart_with_affected_sessions():
text = _read(ADR)
lower = text.lower()
assert "forbidden" in lower
assert "llm" in lower and "restart" in lower
assert "unilateral" in lower
# A worker role must not perform or authorize full/host restart.
assert "must not" in lower
def test_ac4_break_glass_conditions_listed():
text = _read(ADR)
lower = text.lower()
assert "break-glass" in lower
assert "incident" in lower
assert "audit" in lower
def test_ac5_adr_links_issue_lineage():
text = _read(ADR)
for issue in LINKED_ISSUES:
assert issue in text, f"ADR must link issue {issue}"
def test_ac5_safety_model_and_deployment_cross_link_adr():
for path in CROSS_LINK_DOCS:
text = _read(path)
assert ADR_BASENAME in text, (
f"{path.relative_to(REPO_ROOT)} must cross-link {ADR_BASENAME} "
f"(issue #656 acceptance criterion 5)"
)
def test_policy_ids_present_for_enforcement_code():
text = _read(ADR)
for pid in POLICY_IDS:
assert pid in text, f"policy id {pid} missing from ADR"
def test_failure_behavior_denies_on_ambiguity():
text = _read(ADR)
lower = text.lower()
assert "ambiguous" in lower and "deny" in lower
def test_cross_links_do_not_embed_secrets():
for path in (ADR,) + CROSS_LINK_DOCS:
text = _read(path)
for marker in ("ghp_", "BEGIN PRIVATE KEY", "Authorization: Bearer"):
assert marker not in text, f"{path} contains {marker!r}"
+146
View File
@@ -0,0 +1,146 @@
"""Tests for the MCP restart-path inventory and guards (#657).
Covers:
* the registry is well-formed and every path is classified;
* unknown restart attempts fail closed (AC "fail closed on unknown restart");
* the previously-unguarded full-restart primitives stay guarded/absent
against the real source tree (AC "tests for at least one previously
unguarded path");
* pkill of the daemon is still classified as contamination (#630, AC3);
* the inventory doc and module stay in lock-step.
"""
import os
import tempfile
import unittest
from pathlib import Path
import mcp_restart_paths as rp
import runtime_recovery_guard
REPO_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
DOC_PATH = os.path.join(REPO_ROOT, "docs", "mcp-restart-path-inventory.md")
class TestRegistryWellformed(unittest.TestCase):
def test_registry_is_wellformed(self):
# Must not raise.
rp.assert_registry_wellformed()
def test_every_path_has_valid_classification(self):
for path in rp.iter_restart_paths():
self.assertIn(path.classification, rp.VALID_CLASSIFICATIONS)
self.assertTrue(path.guard.strip(), path.path_id)
self.assertTrue(path.references, path.path_id)
self.assertTrue(path.locations, path.path_id)
def test_ids_are_unique(self):
ids = [p.path_id for p in rp.iter_restart_paths()]
self.assertEqual(len(ids), len(set(ids)))
def test_covers_every_classification(self):
present = {p.classification for p in rp.iter_restart_paths()}
self.assertEqual(present, set(rp.VALID_CLASSIFICATIONS))
class TestUnknownAttemptFailsClosed(unittest.TestCase):
def test_unknown_path_raises(self):
with self.assertRaises(rp.UnknownRestartPathError):
rp.assert_restart_attempt_registered("totally_novel_restart_hack")
def test_get_unknown_raises(self):
with self.assertRaises(rp.UnknownRestartPathError):
rp.get_restart_path("nope")
def test_registered_attempt_returns_path(self):
path = rp.assert_restart_attempt_registered("manual_daemon_kill")
self.assertEqual(path.classification, rp.CLASS_FORBIDDEN)
class TestDaemonNeverSelfReplaces(unittest.TestCase):
"""Previously-unguarded full-restart primitive: daemon self-replacement."""
def test_no_self_replacement_in_source(self):
# The live daemon modules must contain no os.execv/os.kill/os._exit
# self-restart call. Must not raise.
rp.assert_no_daemon_self_replacement(REPO_ROOT)
def test_scanner_flags_injected_violation(self):
# Guard the guard: prove the scanner catches a real self-replace call.
with tempfile.TemporaryDirectory() as tmp:
bad = Path(tmp) / "gitea_mcp_server.py"
bad.write_text(
"import os\n"
"def restart():\n"
" os.execv('/usr/bin/python', ['python'])\n",
encoding="utf-8",
)
found = rp.scan_daemon_self_replacement(tmp)
self.assertTrue(found)
with self.assertRaises(AssertionError):
rp.assert_no_daemon_self_replacement(tmp)
def test_scanner_ignores_comment_and_docstring_mentions(self):
with tempfile.TemporaryDirectory() as tmp:
ok = Path(tmp) / "gitea_mcp_server.py"
ok.write_text(
"import os\n"
"# NOT os.execv() to re-point the interpreter here.\n"
'"""Never calls os._exit to restart."""\n'
"value = 1\n",
encoding="utf-8",
)
self.assertEqual(rp.scan_daemon_self_replacement(tmp), [])
class TestLegacyAutoRestartHelperRemoved(unittest.TestCase):
"""Previously-unguarded full-restart path: _trigger_mcp_auto_restart."""
def test_helper_absent_in_source(self):
# Must not raise: helper was removed in #685.
rp.assert_auto_restart_helper_absent(REPO_ROOT)
def test_scanner_flags_reintroduced_helper(self):
with tempfile.TemporaryDirectory() as tmp:
bad = Path(tmp) / "mcp_server.py"
bad.write_text(
"def _trigger_mcp_auto_restart():\n return True\n",
encoding="utf-8",
)
with self.assertRaises(AssertionError):
rp.assert_auto_restart_helper_absent(tmp)
class TestPkillStaysForbidden(unittest.TestCase):
"""AC3: pkill of the daemon remains forbidden/contaminating (#630)."""
def test_manual_daemon_kill_registered_as_forbidden(self):
path = rp.get_restart_path("manual_daemon_kill")
self.assertEqual(path.classification, rp.CLASS_FORBIDDEN)
def test_pkill_classified_as_contamination(self):
assessment = runtime_recovery_guard.assess_recovery_command(
"pkill -f mcp_server.py"
)
self.assertTrue(assessment["contaminated"])
def test_read_only_probe_not_contamination(self):
assessment = runtime_recovery_guard.assess_recovery_command(
"ps aux | grep mcp_server"
)
self.assertFalse(assessment["contaminated"])
class TestInventoryDocInSync(unittest.TestCase):
def test_doc_exists(self):
self.assertTrue(os.path.exists(DOC_PATH), DOC_PATH)
def test_doc_mentions_every_path_id(self):
with open(DOC_PATH, encoding="utf-8") as handle:
doc = handle.read()
for path in rp.iter_restart_paths():
self.assertIn(path.path_id, doc, f"doc missing {path.path_id}")
if __name__ == "__main__":
unittest.main()
+16 -2
View File
@@ -37,6 +37,7 @@ def _live_lock(
"operation_type": issue_lock_store.AUTHOR_ISSUE_WORK_LEASE,
"acquired_at": now.isoformat(),
"expires_at": (now + timedelta(hours=2)).isoformat(),
"session_pid": os.getpid(),
"owner_pid": os.getpid(),
"status": "active",
}
@@ -177,11 +178,24 @@ class TestAuthorOwnershipIssuePrMismatch(unittest.TestCase):
self.assertFalse(result["proven"], result)
self.assertTrue(any("branch" in r for r in result["reasons"]))
def test_no_lock_fail_closed(self):
def test_pidless_durable_lock_rejected(self):
"""A lock without any PID identity must be classified as malformed/non-live and fail closed."""
lock = _live_lock(issue_number=727)
lock.pop("session_pid", None)
lock.pop("owner_pid", None)
lock.pop("pid", None)
path = issue_lock_store.lock_file_path(
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
issue_number=727,
lock_dir=self.lock_dir,
)
issue_lock_store.save_lock_file(path, lock)
result = mcp._prove_author_ownership_for_pr(
pr_number=728,
pr_title="feat: pr sync",
pr_body="Closes #727",
pr_body="Fixes #727",
source_branch="feat/issue-727-pr-sync-status",
remote="prgs",
host=None,
+340
View File
@@ -0,0 +1,340 @@
"""Tests for the MCP restart coordinator and impact analysis (#658).
Multi-session fixtures exercise every verdict branch: safe, unsafe (live work),
override, and the fail-closed deny on incomplete inventory. Also covers the
critical-section deny path and the new ``ControlPlaneDB.list_sessions``.
"""
from __future__ import annotations
import os
import tempfile
import unittest
from datetime import datetime, timedelta, timezone
import restart_coordinator as rc
from control_plane_db import ControlPlaneDB
NOW = datetime(2026, 7, 24, 6, 0, 0, tzinfo=timezone.utc)
def _ts(dt: datetime) -> str:
return dt.isoformat()
def _live_pid() -> int:
return os.getpid()
def _dead_pid() -> int:
# A pid that is essentially never alive. os.kill(0) on it raises
# ProcessLookupError → is_process_alive False.
return 2_000_000_000
def _session(session_id, *, pid, status="active", heartbeat=None, role="author"):
return {
"session_id": session_id,
"role": role,
"profile": "prgs-author",
"pid": pid,
"status": status,
"last_heartbeat_at": _ts(heartbeat or NOW),
}
def _lease(
lease_id,
*,
session_id,
freshness,
kind="issue",
number=658,
phase="allocated",
worktree=None,
role="author",
):
return {
"lease_id": lease_id,
"session_id": session_id,
"role": role,
"phase": phase,
"work_kind": kind,
"work_number": number,
"worktree_path": worktree,
"freshness": {"freshness": freshness},
}
class EvaluateRestartImpactTest(unittest.TestCase):
def test_incomplete_inventory_denies_fail_closed(self) -> None:
report = rc.evaluate_restart_impact(
{"inventory_complete": False, "incomplete_reasons": ["db down"]},
now=NOW,
)
self.assertEqual(report.verdict, rc.VERDICT_UNSAFE)
self.assertFalse(report.allow_restart)
self.assertFalse(report.restart_performed)
self.assertIn("db down", report.incomplete_reasons)
self.assertTrue(
any("fail closed" in reasoning for reasoning in report.reasons)
)
def test_missing_completeness_flag_denies(self) -> None:
# No inventory_complete key at all → treated as incomplete.
report = rc.evaluate_restart_impact({}, now=NOW)
self.assertEqual(report.verdict, rc.VERDICT_UNSAFE)
self.assertFalse(report.allow_restart)
def test_no_other_work_is_safe(self) -> None:
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [_session("requester", pid=_live_pid())],
"leases": [],
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(report.verdict, rc.VERDICT_SAFE)
self.assertTrue(report.allow_restart)
self.assertEqual(report.blast_radius, rc.BLAST_NONE)
self.assertEqual(report.affected_issues, [])
def test_dead_foreign_session_and_lease_are_not_disruptive(self) -> None:
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [
_session("requester", pid=_live_pid()),
_session("dead", pid=_dead_pid()),
],
"leases": [
_lease("l-dead", session_id="dead", freshness="stale_dead_process")
],
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(report.verdict, rc.VERDICT_SAFE)
self.assertTrue(report.allow_restart)
self.assertEqual(report.counts["leases_disruptive"], 0)
self.assertEqual(report.counts["sessions_live_other"], 0)
def test_live_foreign_lease_denies_without_override(self) -> None:
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [
_session("requester", pid=_live_pid()),
_session("worker", pid=_live_pid()),
],
"leases": [
_lease(
"l1",
session_id="worker",
freshness="active",
worktree="/tmp/wt-658",
phase="implementing",
)
],
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(report.verdict, rc.VERDICT_UNSAFE)
self.assertFalse(report.allow_restart)
# Critical section detected: active lease with a live owner.
self.assertEqual(len(report.critical_sections), 1)
self.assertEqual(report.affected_issues, [658])
self.assertEqual(report.counts["mutations"], 1)
self.assertTrue(report.override_would_allow)
self.assertEqual(report.blast_radius, rc.BLAST_HIGH)
# Placeholder ack state for the affected session.
self.assertEqual(report.ack_state.get("worker"), "pending")
def test_operator_override_allows_despite_live_work(self) -> None:
inv = {
"inventory_complete": True,
"sessions": [
_session("requester", pid=_live_pid()),
_session("worker", pid=_live_pid()),
],
"leases": [_lease("l1", session_id="worker", freshness="active")],
}
report = rc.evaluate_restart_impact(
inv,
now=NOW,
requesting_session_id="requester",
operator_override=True,
)
self.assertEqual(report.verdict, rc.VERDICT_OVERRIDE)
self.assertTrue(report.allow_restart)
self.assertFalse(report.restart_performed)
def test_deny_when_critical_section_open(self) -> None:
# A single live author lease in a mutating phase is a critical section
# that must deny an un-overridden restart.
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [_session("worker", pid=_live_pid())],
"leases": [
_lease(
"l1",
session_id="worker",
freshness="active",
phase="merging",
kind="pr",
number=900,
)
],
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(report.verdict, rc.VERDICT_UNSAFE)
self.assertFalse(report.allow_restart)
self.assertEqual(report.affected_prs, [900])
self.assertEqual(len(report.critical_sections), 1)
def test_terminal_lock_makes_restart_unsafe(self) -> None:
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [_session("requester", pid=_live_pid())],
"leases": [],
"terminal_lock": {"terminal_pr": 812},
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(report.verdict, rc.VERDICT_UNSAFE)
self.assertFalse(report.allow_restart)
self.assertIsNotNone(report.terminal_lock)
self.assertTrue(
any("terminal" in reasoning for reasoning in report.reasons)
)
def test_other_live_session_without_lease_is_disruptive(self) -> None:
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [
_session("requester", pid=_live_pid()),
_session("idle-but-live", pid=_live_pid()),
],
"leases": [],
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(report.verdict, rc.VERDICT_UNSAFE)
self.assertEqual(report.counts["sessions_live_other"], 1)
def test_stale_heartbeat_session_not_counted_live(self) -> None:
stale = NOW - timedelta(hours=2)
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [
_session("requester", pid=_live_pid()),
_session("stale", pid=_live_pid(), heartbeat=stale),
],
"leases": [],
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(report.verdict, rc.VERDICT_SAFE)
self.assertEqual(report.counts["sessions_live_other"], 0)
def test_prior_recovery_attempts_echoed(self) -> None:
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [_session("requester", pid=_live_pid())],
"leases": [],
"prior_recovery_attempts": [
{"kind": "client_reconnect", "at": _ts(NOW)}
],
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(len(report.prior_recovery_attempts), 1)
self.assertEqual(report.counts["prior_recovery_attempts"], 1)
def test_bare_string_freshness_accepted(self) -> None:
lease = _lease("l1", session_id="worker", freshness="active")
lease["freshness"] = "active" # bare string, not a dict
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [_session("worker", pid=_live_pid())],
"leases": [lease],
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(report.counts["leases_disruptive"], 1)
def test_as_dict_is_serializable_dto(self) -> None:
import json
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": [_session("requester", pid=_live_pid())],
"leases": [],
},
now=NOW,
requesting_session_id="requester",
)
payload = report.as_dict()
# Round-trips through JSON — safe for the console DTO.
encoded = json.dumps(payload)
decoded = json.loads(encoded)
self.assertEqual(decoded["verdict"], rc.VERDICT_SAFE)
self.assertIn("audit_record", decoded)
self.assertEqual(decoded["audit_record"]["event"], "restart_impact_evaluated")
self.assertFalse(decoded["restart_performed"])
self.assertIn("coordinator_version", decoded)
class ListSessionsTest(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 test_list_sessions_filters_by_status(self) -> None:
self.db.upsert_session(session_id="a", role="author", pid=1, status="active")
self.db.upsert_session(session_id="b", role="author", pid=2, status="ended")
active = self.db.list_sessions(statuses=("active",))
ids = {row["session_id"] for row in active}
self.assertEqual(ids, {"a"})
every = self.db.list_sessions()
self.assertEqual({row["session_id"] for row in every}, {"a", "b"})
def test_list_sessions_feeds_coordinator(self) -> None:
self.db.upsert_session(
session_id="requester", role="author", pid=os.getpid(), status="active"
)
report = rc.evaluate_restart_impact(
{
"inventory_complete": True,
"sessions": self.db.list_sessions(statuses=("active",)),
"leases": [],
},
now=NOW,
requesting_session_id="requester",
)
self.assertEqual(report.counts["sessions_total"], 1)
if __name__ == "__main__": # pragma: no cover
unittest.main()
@@ -139,6 +139,8 @@ EXPECTED_ROLE_EXCLUSIVE_TASKS = frozenset(
"gitea_release_merger_pr_lease",
"create_branch",
"push_branch",
"bootstrap_author_issue_worktree",
"gitea_bootstrap_author_issue_worktree",
# #812 AC20: publishing an unpublished local head is author-only for the
# same reason every other push is — it writes a branch to the remote.
"publish_unpublished_branch",
+252
View File
@@ -0,0 +1,252 @@
"""Unit and integration tests for Model Usage & Performance Analytics (#651)."""
from __future__ import annotations
import os
import tempfile
import unittest
from starlette.testclient import TestClient
import control_plane_db
from webui.analytics_loader import (
ANALYTICS_SCHEMA_VERSION,
compute_percentile,
load_analytics,
record_usage,
)
from webui.app import create_app
from webui import console_redaction
class AnalyticsLoaderTest(unittest.TestCase):
def setUp(self) -> None:
self.temp_dir = tempfile.TemporaryDirectory()
self.db_path = os.path.join(self.temp_dir.name, "test_control_plane.sqlite3")
os.environ["GITEA_CONTROL_PLANE_DB"] = self.db_path
self.db = control_plane_db.ControlPlaneDB(db_path=self.db_path)
def tearDown(self) -> None:
self.temp_dir.cleanup()
def test_compute_percentile(self) -> None:
self.assertIsNone(compute_percentile([], 50.0))
self.assertEqual(compute_percentile([100], 50.0), 100.0)
# 2 elements: [100, 200]
self.assertEqual(compute_percentile([100, 200], 50.0), 150.0)
# 100 elements: 1..100
vals = list(range(1, 101))
self.assertEqual(compute_percentile(vals, 50.0), 50.5)
self.assertAlmostEqual(compute_percentile(vals, 90.0), 90.1)
def test_record_and_aggregate_usage(self) -> None:
# Record event 1 (complete data)
u1 = record_usage(
db_path=self.db_path,
remote="dadeschools",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
role="author",
model="gemini-3.6-flash",
issue_number=651,
stage="implementation",
input_tokens=1000,
output_tokens=500,
estimated_cost_usd=0.0015,
latency_ms=200,
duration_ms=3000,
metadata={"secret_key": "secret123", "note": "token=secret123"},
)
self.assertGreater(u1, 0)
# Record event 2 (missing tokens and cost -> unknown)
u2 = record_usage(
db_path=self.db_path,
remote="dadeschools",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
role="reviewer",
model="claude-3-5-sonnet",
pr_number=846,
stage="review",
latency_ms=500,
duration_ms=6000,
)
self.assertGreater(u2, u1)
snapshot = load_analytics(
db_path=self.db_path,
remote="dadeschools",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
)
self.assertTrue(snapshot.ok)
self.assertEqual(snapshot.schema_version, ANALYTICS_SCHEMA_VERSION)
self.assertEqual(snapshot.total_events, 2)
# Verify overall summary
summary = snapshot.overall_summary
self.assertEqual(summary.total_events, 2)
self.assertEqual(summary.events_with_tokens, 1)
self.assertEqual(summary.total_tokens, 1500)
self.assertEqual(summary.events_with_cost, 1)
self.assertEqual(summary.estimated_cost_usd, 0.0015)
self.assertEqual(summary.events_with_latency, 2)
self.assertEqual(summary.latency_p50_ms, 350.0)
# Verify missing data handling (AC 3: not zero-fabricated)
reviewer_model = snapshot.by_model.get("claude-3-5-sonnet")
self.assertIsNotNone(reviewer_model)
self.assertEqual(reviewer_model.total_events, 1)
self.assertEqual(reviewer_model.events_with_tokens, 0)
self.assertIsNone(reviewer_model.total_tokens)
self.assertEqual(reviewer_model.display_tokens, "Unknown")
self.assertEqual(reviewer_model.events_with_cost, 0)
self.assertIsNone(reviewer_model.estimated_cost_usd)
self.assertEqual(reviewer_model.display_cost, "Unknown")
# Verify redaction (AC 4)
e1 = [e for e in snapshot.events if e.usage_id == u1][0]
self.assertIsNotNone(e1.metadata)
self.assertNotIn("secret123", e1.metadata)
self.assertIn("[REDACTED]", e1.metadata)
def test_missing_db_fail_soft(self) -> None:
invalid_path = "/nonexistent_path_dir/db.sqlite3"
snapshot = load_analytics(db_path=invalid_path)
self.assertFalse(snapshot.ok)
self.assertIn("control_plane_db_unavailable", snapshot.reason)
self.assertEqual(snapshot.overall_summary.display_tokens, "Unknown")
class AnalyticsWebUITest(unittest.TestCase):
def setUp(self) -> None:
self.temp_dir = tempfile.TemporaryDirectory()
self.db_path = os.path.join(self.temp_dir.name, "test_webui.sqlite3")
os.environ["GITEA_CONTROL_PLANE_DB"] = self.db_path
self.app = create_app()
self.client = TestClient(self.app)
record_usage(
db_path=self.db_path,
remote="dadeschools",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
role="author",
model="gemini-3.6-flash",
issue_number=651,
stage="implementation",
input_tokens=2000,
output_tokens=1000,
estimated_cost_usd=0.003,
latency_ms=150,
duration_ms=2500,
)
def tearDown(self) -> None:
self.temp_dir.cleanup()
def test_analytics_html_route(self) -> None:
response = self.client.get("/analytics")
self.assertEqual(response.status_code, 200)
self.assertIn("Model Usage & Performance Analytics", response.text)
self.assertIn("gemini-3.6-flash", response.text)
self.assertIn("3,000", response.text)
def test_analytics_api_route(self) -> None:
response = self.client.get("/api/v1/analytics")
self.assertEqual(response.status_code, 200)
data = response.json()
self.assertTrue(data["ok"])
self.assertEqual(data["total_events"], 1)
self.assertIn("gemini-3.6-flash", data["by_model"])
def test_analytics_ingest_unauthorized_denied(self) -> None:
"""F2: unauthenticated POST must not write the control-plane DB."""
payload = {
"remote": "dadeschools",
"org": "Scaled-Tech-Consulting",
"repo": "Gitea-Tools",
"role": "reviewer",
"model": "claude-3-5-sonnet",
"pr_number": 846,
"stage": "review",
"input_tokens": 500,
"output_tokens": 100,
"latency_ms": 400,
"metadata": "Review note token=secret456",
}
response = self.client.post("/api/v1/analytics/usage", json=payload)
self.assertEqual(response.status_code, 403)
res_json = response.json()
self.assertFalse(res_json.get("ok", True))
self.assertEqual(res_json.get("error"), "unauthorized")
authorization = res_json.get("authorization") or {}
self.assertFalse(authorization.get("allowed"))
self.assertFalse(authorization.get("execution_enabled"))
# No new row written
res2 = self.client.get("/api/v1/analytics")
self.assertEqual(res2.status_code, 200)
self.assertEqual(res2.json()["total_events"], 1)
def test_html_escapes_script_bearing_model_role_stage(self) -> None:
"""F1: stored XSS — dynamic model/role/stage must render escaped."""
xss = '<script>alert(1)</script>'
record_usage(
db_path=self.db_path,
remote="dadeschools",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
role=xss,
model=xss,
stage=xss,
issue_number=999,
status="success",
)
response = self.client.get("/analytics")
self.assertEqual(response.status_code, 200)
# Raw tag must not appear; escaped form must.
self.assertNotIn("<script>alert(1)</script>", response.text)
self.assertIn("&lt;script&gt;alert(1)&lt;/script&gt;", response.text)
def test_load_analytics_coerces_none_scope(self) -> None:
"""F4: None remote/org/repo become empty strings, never None."""
snapshot = load_analytics(db_path=self.db_path, remote=None, org=None, repo=None)
self.assertIsInstance(snapshot.remote, str)
self.assertIsInstance(snapshot.org, str)
self.assertIsInstance(snapshot.repo, str)
self.assertEqual(snapshot.remote, "")
self.assertEqual(snapshot.org, "")
self.assertEqual(snapshot.repo, "")
def test_usage_events_retention_max_rows(self) -> None:
"""F3: record_usage_event enforces USAGE_EVENTS_MAX_ROWS."""
db = control_plane_db.ControlPlaneDB(db_path=self.db_path)
original_max = db.USAGE_EVENTS_MAX_ROWS
try:
db.USAGE_EVENTS_MAX_ROWS = 3
for i in range(5):
db.record_usage_event(
remote="dadeschools",
org="org",
repo="repo",
role="author",
model=f"model-{i}",
stage="test",
)
rows = db.query_usage_events(limit=100)
self.assertLessEqual(len(rows), 3)
# Newest three retained
models = {r["model"] for r in rows}
self.assertEqual(models, {"model-2", "model-3", "model-4"})
finally:
db.USAGE_EVENTS_MAX_ROWS = original_max
if __name__ == "__main__":
unittest.main()
+345
View File
@@ -0,0 +1,345 @@
"""Tests for the system-health dashboard view (#639).
Covers the acceptance criteria directly: the page renders the health DTO
fields (AC1), degraded dependencies are visible (AC2), stale runtime is warned
prominently and never rendered as mutation-safe (AC3), healthy and degraded
fixtures both render (AC4), and the shell carries a nav entry (AC5).
"""
import sys
import unittest
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from starlette.testclient import TestClient
from webui.app import create_app
from webui.deployment_boundary import scan_text_for_client_secrets
from webui.layout import render_page
from webui.nav import iter_nav_items
from webui.system_health import (
STATUS_DEGRADED,
STATUS_DOWN,
STATUS_OK,
STATUS_SKIPPED,
STATUS_UNPROVEN,
DependencyProbe,
StaleRuntime,
SystemHealthSnapshot,
VersionInfo,
)
from webui.system_health_views import render_system_health_page
DASHBOARD_PATH = "/system-health"
def _version(*, known: bool = True) -> VersionInfo:
return VersionInfo(
git_sha="1c455b6ec0f9cb761fe6248de68c17e061fb5ecd" if known else None,
git_describe="v0.4.1-12-g1c455b6" if known else None,
control_plane_schema_version=4 if known else None,
python_version="3.13.1",
known=known,
)
def _parity(*, stale: bool = False, determinable: bool = True) -> StaleRuntime:
if stale:
return StaleRuntime(
daemon_head="aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa",
checkout_head="bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb",
remote_head="cccccccccccccccccccccccccccccccccccccccc",
stale=True,
determinable=True,
mutation_safe=False,
reasons=("runtime, checkout, and remote commits disagree",),
)
if not determinable:
return StaleRuntime(
daemon_head=None,
checkout_head=None,
remote_head=None,
stale=False,
determinable=False,
mutation_safe=False,
reasons=("local checkout HEAD could not be read",),
)
return StaleRuntime(
daemon_head="1c455b6ec0f9cb761fe6248de68c17e061fb5ecd",
checkout_head="1c455b6ec0f9cb761fe6248de68c17e061fb5ecd",
remote_head="1c455b6ec0f9cb761fe6248de68c17e061fb5ecd",
stale=False,
determinable=True,
mutation_safe=True,
reasons=(),
)
def _snapshot(
*,
status: str = STATUS_OK,
ready: bool = True,
readiness_complete: bool = True,
readiness_reasons: tuple[str, ...] = (),
dependencies: tuple[DependencyProbe, ...] | None = None,
parity: StaleRuntime | None = None,
namespaces: tuple[dict, ...] = (),
probe_errors: tuple[str, ...] = (),
version_known: bool = True,
) -> SystemHealthSnapshot:
if dependencies is None:
dependencies = (
DependencyProbe(
name="control_plane_db",
kind="sqlite",
status=STATUS_OK,
detail="schema version 4",
required=True,
latency_ms=1.25,
metadata={"schema_version": 4},
),
)
return SystemHealthSnapshot(
status=status,
ready=ready,
readiness_complete=readiness_complete,
readiness_reasons=readiness_reasons,
service="mcp-control-plane-webui",
mode="read-only",
version=_version(known=version_known),
started_at="2026-07-23T19:50:47+00:00",
uptime_seconds=3661.5,
timestamp="2026-07-23T20:51:48+00:00",
deep_probes_requested=False,
dependencies=dependencies,
mcp_namespaces=namespaces,
stale_runtime=parity if parity is not None else _parity(),
probe_errors=probe_errors,
)
class TestHealthyRender(unittest.TestCase):
"""AC1 / AC4 — every health DTO field reaches the page."""
def setUp(self):
self.html = render_system_health_page(_snapshot())
def test_readiness_fields_render(self):
self.assertIn("System health", self.html)
self.assertIn("Ready", self.html)
self.assertIn("mcp-control-plane-webui", self.html)
self.assertIn("read-only", self.html)
self.assertIn("2026-07-23T20:51:48+00:00", self.html)
def test_version_and_uptime_render(self):
self.assertIn("1c455b6ec0f9cb761fe6248de68c17e061fb5ecd", self.html)
self.assertIn("v0.4.1-12-g1c455b6", self.html)
self.assertIn("3.13.1", self.html)
self.assertIn("3661.500s", self.html)
self.assertIn("1.02h", self.html)
def test_dependency_row_renders_with_latency(self):
self.assertIn("control_plane_db", self.html)
self.assertIn("sqlite", self.html)
self.assertIn("schema version 4", self.html)
self.assertIn("1.2 ms", self.html)
def test_healthy_page_shows_no_stale_warning(self):
self.assertNotIn("Stale runtime:", self.html)
self.assertNotIn("Staleness", self.html)
def test_unknown_version_is_labelled_not_faked(self):
html = render_system_health_page(_snapshot(version_known=False))
self.assertIn("unknown", html)
self.assertIn("unresolved", html)
class TestDegradedRender(unittest.TestCase):
"""AC2 — a degraded or unrun dependency is visible, not swallowed."""
def setUp(self):
self.deps = (
DependencyProbe(
name="control_plane_db",
kind="sqlite",
status=STATUS_OK,
detail="schema version 4",
required=True,
latency_ms=0.9,
),
DependencyProbe(
name="repository",
kind="git",
status=STATUS_DOWN,
detail="repository root is not a git checkout",
required=True,
latency_ms=4.0,
),
DependencyProbe(
name="gitea",
kind="http",
status=STATUS_SKIPPED,
detail="deep probe not requested",
required=False,
),
)
self.html = render_system_health_page(
_snapshot(
status=STATUS_DEGRADED,
ready=False,
readiness_complete=False,
readiness_reasons=("required dependency 'repository' is down",),
dependencies=self.deps,
)
)
def test_degraded_banner_names_the_dependency(self):
self.assertIn("Degraded dependencies:", self.html)
self.assertIn("repository", self.html)
def test_not_run_probe_is_reported_separately(self):
self.assertIn("Not probed:", self.html)
self.assertIn("gitea", self.html)
self.assertIn("not counted", self.html)
def test_not_ready_headline_and_reason(self):
self.assertIn("Not ready", self.html)
self.assertIn("required dependency &#x27;repository&#x27; is down", self.html)
def test_degraded_status_badge_present(self):
self.assertIn("badge-health-degraded", self.html)
self.assertIn("badge-health-down", self.html)
def test_ready_but_incomplete_is_not_shown_as_plain_ready(self):
html = render_system_health_page(
_snapshot(ready=True, readiness_complete=False)
)
self.assertIn("Ready (incomplete evidence)", html)
class TestStaleRuntimeWarning(unittest.TestCase):
"""AC3 — staleness is prominent and never claims mutation safety."""
def test_stale_runtime_warns_and_denies_mutation_safety(self):
html = render_system_health_page(_snapshot(parity=_parity(stale=True)))
self.assertIn("Stale runtime:", html)
self.assertIn("do not treat this runtime as mutation-safe", html)
self.assertIn("<tr><th>Mutation safe</th><td>False</td></tr>", html)
def test_indeterminate_parity_is_not_reported_safe(self):
html = render_system_health_page(
_snapshot(parity=_parity(determinable=False))
)
self.assertIn("Staleness", html)
self.assertIn("<tr><th>Mutation safe</th><td>False</td></tr>", html)
self.assertIn("<tr><th>Determinable</th><td>False</td></tr>", html)
def test_healthy_parity_reports_mutation_safe_true(self):
html = render_system_health_page(_snapshot())
self.assertIn("<tr><th>Mutation safe</th><td>True</td></tr>", html)
class TestNamespacesAndErrors(unittest.TestCase):
def test_unproven_namespace_rows_render(self):
html = render_system_health_page(
_snapshot(
namespaces=(
{
"namespace": "gitea-author",
"required_tool": "gitea_lock_issue",
"status": STATUS_UNPROVEN,
"ide_namespace_proven": False,
"reason": "the web console cannot invoke the IDE-managed MCP client",
},
)
)
)
self.assertIn("gitea-author", html)
self.assertIn("gitea_lock_issue", html)
self.assertIn("badge-health-unproven", html)
def test_no_namespaces_degrades_gracefully(self):
html = render_system_health_page(_snapshot(namespaces=()))
self.assertIn("No MCP namespaces are declared.", html)
def test_probe_errors_render_when_present(self):
html = render_system_health_page(
_snapshot(probe_errors=("probe raised: disk offline",))
)
self.assertIn("Probe errors", html)
self.assertIn("disk offline", html)
def test_probe_error_card_absent_when_clean(self):
self.assertNotIn("Probe errors", render_system_health_page(_snapshot()))
class TestReadOnlyAndRedaction(unittest.TestCase):
def test_no_restart_or_kill_controls(self):
html = render_system_health_page(_snapshot())
self.assertNotIn("<button", html)
self.assertNotIn("<form", html)
self.assertNotIn("pkill", html)
self.assertIn("read-only", html)
def test_recovery_points_at_sanctioned_path(self):
html = render_system_health_page(_snapshot())
self.assertIn("Reconnect the MCP client", html)
self.assertIn("Never kill the daemon process manually", html)
def test_secret_shaped_detail_is_redacted(self):
leaky = DependencyProbe(
name="gitea",
kind="http",
status=STATUS_DOWN,
detail="auth failed for token=ghp_ABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789",
required=False,
latency_ms=12.0,
)
html = render_system_health_page(_snapshot(dependencies=(leaky,)))
self.assertNotIn("ghp_ABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789", html)
def test_html_in_detail_is_escaped(self):
hostile = DependencyProbe(
name="repository",
kind="git",
status=STATUS_DOWN,
detail="<script>alert(1)</script>",
required=True,
)
html = render_system_health_page(_snapshot(dependencies=(hostile,)))
self.assertNotIn("<script>", html)
self.assertIn("&lt;script&gt;", html)
class TestNavAndRoute(unittest.TestCase):
"""AC5 — the shell links the dashboard, and the route serves it."""
def setUp(self):
self.client = TestClient(create_app())
def test_nav_contains_system_health(self):
self.assertIn(
(DASHBOARD_PATH, "System health"),
[(item.href, item.label) for item in iter_nav_items()],
)
def test_rendered_shell_links_dashboard(self):
page = render_page(title="Home", body_html="<p>x</p>")
self.assertIn(f'href="{DASHBOARD_PATH}"', page)
def test_route_renders_dashboard(self):
response = self.client.get(DASHBOARD_PATH)
self.assertEqual(response.status_code, 200)
self.assertIn("System health", response.text)
self.assertIn("Stale-runtime parity", response.text)
def test_route_is_read_only(self):
self.assertEqual(self.client.post(DASHBOARD_PATH).status_code, 405)
def test_live_page_leaks_no_client_secret(self):
findings = scan_text_for_client_secrets(self.client.get(DASHBOARD_PATH).text)
self.assertEqual(findings, [])
if __name__ == "__main__": # pragma: no cover
unittest.main()
+3 -1
View File
@@ -35,8 +35,10 @@ class TestCanonicalRepoRoot(unittest.TestCase):
self.assertEqual(root, CONTROL_ROOT)
def test_falls_back_when_git_unavailable(self):
# When fallback is a path under <repo>/branches/<worktree>, recover
# <repo> via commonpath ancestry (never string-split on "/branches/").
root = amw.resolve_canonical_repo_root("/missing/path", MCP_PROCESS_ROOT)
self.assertEqual(root, os.path.realpath(MCP_PROCESS_ROOT))
self.assertEqual(root, os.path.realpath(CONTROL_ROOT))
class TestWorkspaceRepoMembership(unittest.TestCase):
+433
View File
@@ -0,0 +1,433 @@
"""Model usage, token cost, latency, and workflow-performance analytics (#651, Phase 4).
Ingests session instrumentation metrics, aggregates usage/cost/latency percentiles
by project, role, model, issue/PR, and stage, enforcing secret redaction and
explicitly rendering missing metrics as "Unknown" without zero-fabrication.
"""
from __future__ import annotations
import math
from dataclasses import asdict, dataclass
from typing import Any, Sequence
import control_plane_db
from webui import console_redaction
ANALYTICS_SCHEMA_VERSION = 1
@dataclass(frozen=True)
class UsageEvent:
usage_id: int
session_id: str | None
remote: str
org: str
repo: str
project_id: str | None
role: str
model: str
issue_number: int | None
pr_number: int | None
stage: str
input_tokens: int | None
output_tokens: int | None
total_tokens: int | None
estimated_cost_usd: float | None
latency_ms: int | None
duration_ms: int | None
status: str
metadata: str | None
created_at: str
def to_dict(self) -> dict[str, Any]:
d = asdict(self)
if d["metadata"]:
d["metadata"] = console_redaction.redact_text(d["metadata"])
return d
@dataclass(frozen=True)
class GroupMetrics:
name: str
total_events: int
events_with_tokens: int
input_tokens: int | None
output_tokens: int | None
total_tokens: int | None
events_with_cost: int
estimated_cost_usd: float | None
events_with_latency: int
latency_p50_ms: float | None
latency_p90_ms: float | None
latency_p95_ms: float | None
latency_p99_ms: float | None
latency_avg_ms: float | None
events_with_duration: int
duration_avg_ms: float | None
display_tokens: str
display_cost: str
display_latency_p50: str
display_latency_p90: str
display_duration_avg: str
def to_dict(self) -> dict[str, Any]:
return asdict(self)
@dataclass(frozen=True)
class AnalyticsSnapshot:
ok: bool
reason: str
schema_version: int
remote: str
org: str
repo: str
total_events: int
overall_summary: GroupMetrics
by_project: dict[str, GroupMetrics]
by_role: dict[str, GroupMetrics]
by_model: dict[str, GroupMetrics]
by_work_item: dict[str, GroupMetrics]
by_stage: dict[str, GroupMetrics]
events: tuple[UsageEvent, ...]
def to_dict(self) -> dict[str, Any]:
return {
"ok": self.ok,
"reason": self.reason,
"schema_version": self.schema_version,
"remote": self.remote,
"org": self.org,
"repo": self.repo,
"total_events": self.total_events,
"overall_summary": self.overall_summary.to_dict(),
"by_project": {k: v.to_dict() for k, v in self.by_project.items()},
"by_role": {k: v.to_dict() for k, v in self.by_role.items()},
"by_model": {k: v.to_dict() for k, v in self.by_model.items()},
"by_work_item": {k: v.to_dict() for k, v in self.by_work_item.items()},
"by_stage": {k: v.to_dict() for k, v in self.by_stage.items()},
"events": [e.to_dict() for e in self.events],
}
def compute_percentile(values: Sequence[float | int], percentile: float) -> float | None:
if not values:
return None
sorted_vals = sorted(values)
n = len(sorted_vals)
if n == 1:
return float(sorted_vals[0])
k = (n - 1) * (percentile / 100.0)
f = math.floor(k)
c = math.ceil(k)
if f == c:
return float(sorted_vals[int(f)])
d0 = sorted_vals[int(f)] * (c - k)
d1 = sorted_vals[int(c)] * (k - f)
return float(d0 + d1)
def aggregate_events(group_name: str, events: Sequence[UsageEvent]) -> GroupMetrics:
total_events = len(events)
if total_events == 0:
return GroupMetrics(
name=group_name,
total_events=0,
events_with_tokens=0,
input_tokens=None,
output_tokens=None,
total_tokens=None,
events_with_cost=0,
estimated_cost_usd=None,
events_with_latency=0,
latency_p50_ms=None,
latency_p90_ms=None,
latency_p95_ms=None,
latency_p99_ms=None,
latency_avg_ms=None,
events_with_duration=0,
duration_avg_ms=None,
display_tokens="Unknown",
display_cost="Unknown",
display_latency_p50="Unknown",
display_latency_p90="Unknown",
display_duration_avg="Unknown",
)
token_events = [
e for e in events
if e.total_tokens is not None or e.input_tokens is not None or e.output_tokens is not None
]
events_with_tokens = len(token_events)
if events_with_tokens > 0:
input_tokens = sum(e.input_tokens or 0 for e in token_events)
output_tokens = sum(e.output_tokens or 0 for e in token_events)
total_tokens = sum(
e.total_tokens if e.total_tokens is not None else ((e.input_tokens or 0) + (e.output_tokens or 0))
for e in token_events
)
display_tokens = f"{total_tokens:,}"
else:
input_tokens = None
output_tokens = None
total_tokens = None
display_tokens = "Unknown"
cost_events = [e for e in events if e.estimated_cost_usd is not None]
events_with_cost = len(cost_events)
if events_with_cost > 0:
estimated_cost_usd = round(sum(e.estimated_cost_usd for e in cost_events), 6)
display_cost = f"${estimated_cost_usd:.4f}"
else:
estimated_cost_usd = None
display_cost = "Unknown"
latency_vals = [e.latency_ms for e in events if e.latency_ms is not None]
events_with_latency = len(latency_vals)
if events_with_latency > 0:
latency_p50_ms = compute_percentile(latency_vals, 50.0)
latency_p90_ms = compute_percentile(latency_vals, 90.0)
latency_p95_ms = compute_percentile(latency_vals, 95.0)
latency_p99_ms = compute_percentile(latency_vals, 99.0)
latency_avg_ms = round(sum(latency_vals) / events_with_latency, 2)
display_latency_p50 = f"{round(latency_p50_ms, 1)} ms" if latency_p50_ms is not None else "Unknown"
display_latency_p90 = f"{round(latency_p90_ms, 1)} ms" if latency_p90_ms is not None else "Unknown"
else:
latency_p50_ms = None
latency_p90_ms = None
latency_p95_ms = None
latency_p99_ms = None
latency_avg_ms = None
display_latency_p50 = "Unknown"
display_latency_p90 = "Unknown"
duration_vals = [e.duration_ms for e in events if e.duration_ms is not None]
events_with_duration = len(duration_vals)
if events_with_duration > 0:
duration_avg_ms = round(sum(duration_vals) / events_with_duration, 2)
display_duration_avg = f"{round(duration_avg_ms / 1000.0, 2)} s" if duration_avg_ms >= 1000 else f"{round(duration_avg_ms, 1)} ms"
else:
duration_avg_ms = None
display_duration_avg = "Unknown"
return GroupMetrics(
name=group_name,
total_events=total_events,
events_with_tokens=events_with_tokens,
input_tokens=input_tokens,
output_tokens=output_tokens,
total_tokens=total_tokens,
events_with_cost=events_with_cost,
estimated_cost_usd=estimated_cost_usd,
events_with_latency=events_with_latency,
latency_p50_ms=latency_p50_ms,
latency_p90_ms=latency_p90_ms,
latency_p95_ms=latency_p95_ms,
latency_p99_ms=latency_p99_ms,
latency_avg_ms=latency_avg_ms,
events_with_duration=events_with_duration,
duration_avg_ms=duration_avg_ms,
display_tokens=display_tokens,
display_cost=display_cost,
display_latency_p50=display_latency_p50,
display_latency_p90=display_latency_p90,
display_duration_avg=display_duration_avg,
)
def record_usage(
*,
db_path: str | None = None,
session_id: str | None = None,
remote: str = "dadeschools",
org: str = "",
repo: str = "",
project_id: str | None = None,
role: str = "unknown",
model: str = "unknown",
issue_number: int | None = None,
pr_number: int | None = None,
stage: str = "unknown",
input_tokens: int | None = None,
output_tokens: int | None = None,
total_tokens: int | None = None,
estimated_cost_usd: float | None = None,
latency_ms: int | None = None,
duration_ms: int | None = None,
status: str = "success",
metadata: str | dict[str, Any] | None = None,
created_at: str | None = None,
) -> int:
"""Ingest/record a single usage event with optional metrics."""
db = control_plane_db.ControlPlaneDB(db_path=db_path)
return db.record_usage_event(
session_id=session_id,
remote=remote,
org=org,
repo=repo,
project_id=project_id,
role=role,
model=model,
issue_number=issue_number,
pr_number=pr_number,
stage=stage,
input_tokens=input_tokens,
output_tokens=output_tokens,
total_tokens=total_tokens,
estimated_cost_usd=estimated_cost_usd,
latency_ms=latency_ms,
duration_ms=duration_ms,
status=status,
metadata=metadata,
created_at=created_at,
)
def load_analytics(
*,
db_path: str | None = None,
remote: str | None = None,
org: str | None = None,
repo: str | None = None,
project_id: str | None = None,
role: str | None = None,
model: str | None = None,
stage: str | None = None,
issue_number: int | None = None,
pr_number: int | None = None,
limit: int = 500,
) -> AnalyticsSnapshot:
"""Load analytics snapshot aggregated by project, role, model, issue/PR, and stage."""
remote_filter = (remote or "").strip() or None
org_filter = (org or "").strip() or None
repo_filter = (repo or "").strip() or None
role_filter = (role or "").strip() or None
model_filter = (model or "").strip() or None
stage_filter = (stage or "").strip() or None
# F4: coerce optional scope filters to str so AnalyticsSnapshot never holds None.
scope_remote = (remote or "").strip()
scope_org = (org or "").strip()
scope_repo = (repo or "").strip()
try:
db = control_plane_db.ControlPlaneDB(db_path=db_path)
rows = db.query_usage_events(
remote=remote_filter,
org=org_filter,
repo=repo_filter,
project_id=project_id,
role=role_filter,
model=model_filter,
stage=stage_filter,
issue_number=issue_number,
pr_number=pr_number,
limit=limit,
)
except Exception as exc:
empty_summary = aggregate_events("Overall", [])
return AnalyticsSnapshot(
ok=False,
reason=f"control_plane_db_unavailable: {exc}",
schema_version=ANALYTICS_SCHEMA_VERSION,
remote=scope_remote,
org=scope_org,
repo=scope_repo,
total_events=0,
overall_summary=empty_summary,
by_project={},
by_role={},
by_model={},
by_work_item={},
by_stage={},
events=(),
)
parsed_events: list[UsageEvent] = []
for r in rows:
meta = console_redaction.redact_text(r.get("metadata")) if r.get("metadata") else None
parsed_events.append(
UsageEvent(
usage_id=r["usage_id"],
session_id=r.get("session_id"),
remote=r.get("remote") or scope_remote,
org=r.get("org") or scope_org,
repo=r.get("repo") or scope_repo,
project_id=r.get("project_id"),
role=r.get("role") or "unknown",
model=r.get("model") or "unknown",
issue_number=r.get("issue_number"),
pr_number=r.get("pr_number"),
stage=r.get("stage") or "unknown",
input_tokens=r.get("input_tokens"),
output_tokens=r.get("output_tokens"),
total_tokens=r.get("total_tokens"),
estimated_cost_usd=r.get("estimated_cost_usd"),
latency_ms=r.get("latency_ms"),
duration_ms=r.get("duration_ms"),
status=r.get("status") or "success",
metadata=meta,
created_at=r.get("created_at") or "",
)
)
overall_summary = aggregate_events("Overall", parsed_events)
# Group by project
groups_by_project: dict[str, list[UsageEvent]] = {}
for e in parsed_events:
key = e.project_id or (f"{e.org}/{e.repo}" if e.org and e.repo else "default")
groups_by_project.setdefault(key, []).append(e)
by_project = {k: aggregate_events(k, v) for k, v in groups_by_project.items()}
# Group by role
groups_by_role: dict[str, list[UsageEvent]] = {}
for e in parsed_events:
groups_by_role.setdefault(e.role, []).append(e)
by_role = {k: aggregate_events(k, v) for k, v in groups_by_role.items()}
# Group by model
groups_by_model: dict[str, list[UsageEvent]] = {}
for e in parsed_events:
groups_by_model.setdefault(e.model, []).append(e)
by_model = {k: aggregate_events(k, v) for k, v in groups_by_model.items()}
# Group by work item
groups_by_work_item: dict[str, list[UsageEvent]] = {}
for e in parsed_events:
if e.issue_number:
key = f"issue #{e.issue_number}"
elif e.pr_number:
key = f"pr #{e.pr_number}"
else:
key = "unlinked"
groups_by_work_item.setdefault(key, []).append(e)
by_work_item = {k: aggregate_events(k, v) for k, v in groups_by_work_item.items()}
# Group by stage
groups_by_stage: dict[str, list[UsageEvent]] = {}
for e in parsed_events:
groups_by_stage.setdefault(e.stage, []).append(e)
by_stage = {k: aggregate_events(k, v) for k, v in groups_by_stage.items()}
return AnalyticsSnapshot(
ok=True,
reason="ok",
schema_version=ANALYTICS_SCHEMA_VERSION,
remote=scope_remote,
org=scope_org,
repo=scope_repo,
total_events=len(parsed_events),
overall_summary=overall_summary,
by_project=by_project,
by_role=by_role,
by_model=by_model,
by_work_item=by_work_item,
by_stage=by_stage,
events=tuple(parsed_events),
)
def snapshot_to_dict(snapshot: AnalyticsSnapshot) -> dict[str, Any]:
return snapshot.to_dict()
+248
View File
@@ -0,0 +1,248 @@
"""HTML views for the Model Usage & Performance Analytics console (#651)."""
from __future__ import annotations
import html
from webui.analytics_loader import AnalyticsSnapshot, GroupMetrics, UsageEvent
from webui.layout import render_page
def _escape(text: object) -> str:
"""HTML-escape dynamic analytics fields (mirrors audit_views / project_views)."""
return html.escape(str(text), quote=True)
def _render_badge(text: str, badge_type: str = "muted") -> str:
return f'<span class="badge badge-{_escape(badge_type)}">{_escape(text)}</span>'
def _render_group_table(title: str, groups: dict[str, GroupMetrics], key_header: str = "Group") -> str:
if not groups:
return (
f"<h3>{_escape(title)}</h3>"
'<div class="card"><p class="muted">No telemetry events recorded for this dimension.</p></div>'
)
rows = []
for key, g in sorted(groups.items(), key=lambda x: x[1].total_events, reverse=True):
cost_cell = (
f'<span class="accent">{_escape(g.display_cost)}</span>'
if g.events_with_cost > 0
else _render_badge("Unknown")
)
tokens_cell = (
_escape(g.display_tokens)
if g.events_with_tokens > 0
else _render_badge("Unknown")
)
lat_p50 = (
_escape(g.display_latency_p50)
if g.events_with_latency > 0
else _render_badge("Unknown")
)
lat_p90 = (
_escape(g.display_latency_p90)
if g.events_with_latency > 0
else _render_badge("Unknown")
)
dur_avg = (
_escape(g.display_duration_avg)
if g.events_with_duration > 0
else _render_badge("Unknown")
)
rows.append(
"<tr>"
f"<td><strong>{_escape(key)}</strong></td>"
f"<td>{g.total_events}</td>"
f"<td>{tokens_cell}</td>"
f"<td>{cost_cell}</td>"
f"<td>{lat_p50}</td>"
f"<td>{lat_p90}</td>"
f"<td>{dur_avg}</td>"
"</tr>"
)
rows_html = "".join(rows)
return f"""
<h3>{_escape(title)}</h3>
<div class="card" style="overflow-x: auto;">
<table class="data-table">
<thead>
<tr>
<th>{_escape(key_header)}</th>
<th>Events</th>
<th>Total Tokens</th>
<th>Est. Cost</th>
<th>Latency (p50)</th>
<th>Latency (p90)</th>
<th>Avg Stage Duration</th>
</tr>
</thead>
<tbody>
{rows_html}
</tbody>
</table>
</div>
"""
def _render_events_table(events: tuple[UsageEvent, ...]) -> str:
if not events:
return (
"<h3>Recent Usage & Instrumentation Events</h3>"
'<div class="card"><p class="muted">No individual telemetry events recorded yet. Opt-in instrumentation via session logging or authorized POST /api/v1/analytics/usage.</p></div>'
)
rows = []
for e in list(events)[-50:]: # Display latest 50
if e.issue_number is not None:
work_item = f"issue #{e.issue_number}"
elif e.pr_number is not None:
work_item = f"pr #{e.pr_number}"
else:
work_item = "unlinked"
tokens = (
_escape(f"{e.total_tokens:,}")
if e.total_tokens is not None
else _render_badge("Unknown")
)
cost = (
_escape(f"${e.estimated_cost_usd:.4f}")
if e.estimated_cost_usd is not None
else _render_badge("Unknown")
)
latency = (
_escape(f"{e.latency_ms} ms")
if e.latency_ms is not None
else _render_badge("Unknown")
)
duration = (
_escape(f"{e.duration_ms} ms")
if e.duration_ms is not None
else _render_badge("Unknown")
)
status_badge = _render_badge(
e.status, "success" if e.status == "success" else "danger"
)
rows.append(
"<tr>"
f"<td>#{e.usage_id}</td>"
f"<td><small>{_escape(e.created_at)}</small></td>"
f"<td><span class=\"badge\">{_escape(e.role)}</span></td>"
f"<td><strong>{_escape(e.model)}</strong></td>"
f"<td>{_escape(e.stage)}</td>"
f"<td>{_escape(work_item)}</td>"
f"<td>{tokens}</td>"
f"<td>{cost}</td>"
f"<td>{latency}</td>"
f"<td>{duration}</td>"
f"<td>{status_badge}</td>"
"</tr>"
)
rows_html = "".join(rows)
return f"""
<h3>Recent Telemetry Events</h3>
<div class="card" style="overflow-x: auto;">
<table class="data-table">
<thead>
<tr>
<th>ID</th>
<th>Timestamp</th>
<th>Role</th>
<th>Model</th>
<th>Stage</th>
<th>Work Item</th>
<th>Tokens</th>
<th>Cost</th>
<th>Latency</th>
<th>Duration</th>
<th>Status</th>
</tr>
</thead>
<tbody>
{rows_html}
</tbody>
</table>
</div>
"""
def render_analytics_page(snapshot: AnalyticsSnapshot) -> str:
"""Render the main Model Usage & Performance Analytics console page."""
summary = snapshot.overall_summary
kpi_tokens = (
_escape(summary.display_tokens)
if summary.events_with_tokens > 0
else _render_badge("Unknown")
)
kpi_cost = (
_escape(summary.display_cost)
if summary.events_with_cost > 0
else _render_badge("Unknown")
)
kpi_lat_p50 = (
_escape(summary.display_latency_p50)
if summary.events_with_latency > 0
else _render_badge("Unknown")
)
kpi_dur_avg = (
_escape(summary.display_duration_avg)
if summary.events_with_duration > 0
else _render_badge("Unknown")
)
status_notice = ""
if not snapshot.ok:
status_notice = (
f'<div class="card warning-card"><strong>Degraded Data Source:</strong> '
f'{_escape(snapshot.reason)}</div>'
)
body_html = f"""
<h2>Model Usage & Performance Analytics (Phase 4)</h2>
<p class="muted">
Durable console analytics for model usage, token cost, latency percentiles, and workflow-stage performance correlated to issues, PRs, and worker roles.
</p>
{status_notice}
<div class="notice-card" style="background: rgba(91, 159, 212, 0.1); border: 1px solid var(--border); padding: 0.75rem 1rem; border-radius: 6px; margin-bottom: 1.5rem;">
<small><strong>Note on telemetry fidelity:</strong> Missing data or untracked metrics are explicitly labeled as <em>Unknown</em>. No token costs or latency metrics are zero-fabricated.</small>
</div>
<div class="card-grid" style="display: grid; grid-template-columns: repeat(auto-fit, minmax(180px, 1fr)); gap: 1rem; margin-bottom: 1.5rem;">
<div class="card">
<span class="muted" style="font-size: 0.85rem;">Total Events</span>
<h3 style="margin: 0.25rem 0 0 0;">{summary.total_events}</h3>
</div>
<div class="card">
<span class="muted" style="font-size: 0.85rem;">Total Tokens</span>
<h3 style="margin: 0.25rem 0 0 0;">{kpi_tokens}</h3>
</div>
<div class="card">
<span class="muted" style="font-size: 0.85rem;">Est. Token Cost</span>
<h3 style="margin: 0.25rem 0 0 0;">{kpi_cost}</h3>
</div>
<div class="card">
<span class="muted" style="font-size: 0.85rem;">Latency (p50)</span>
<h3 style="margin: 0.25rem 0 0 0;">{kpi_lat_p50}</h3>
</div>
<div class="card">
<span class="muted" style="font-size: 0.85rem;">Avg Stage Duration</span>
<h3 style="margin: 0.25rem 0 0 0;">{kpi_dur_avg}</h3>
</div>
</div>
{_render_group_table("Usage & Cost by Model", snapshot.by_model, "Model")}
{_render_group_table("Performance by Workflow Stage", snapshot.by_stage, "Stage")}
{_render_group_table("Usage & Cost by Role", snapshot.by_role, "Role")}
{_render_group_table("Work Item Analytics", snapshot.by_work_item, "Work Item")}
{_render_events_table(snapshot.events)}
"""
return render_page(title="Model Usage & Performance Analytics", body_html=body_html)
+138
View File
@@ -47,12 +47,19 @@ 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_views import render_runtime_page
from webui.timeline import load_timeline, snapshot_to_dict as timeline_snapshot_to_dict
from webui.analytics_loader import (
load_analytics,
record_usage,
snapshot_to_dict as analytics_snapshot_to_dict,
)
from webui.analytics_views import render_analytics_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,
)
from webui.system_health_views import render_system_health_page
_READ_ONLY_METHODS = frozenset({"GET", "HEAD", "OPTIONS"})
_AUDIT_MUTATION_PATHS = frozenset({"/audit", "/api/audit"})
@@ -161,6 +168,24 @@ async def api_system_health(request: Request) -> JSONResponse:
return JSONResponse(payload, status_code=200 if snapshot.ready else 503)
async def system_health(request: Request) -> HTMLResponse:
"""Read-only system-health dashboard (#639).
Shares the #634 snapshot loader with the JSON API so the page can never
disagree with it. `?deep=1` opts into the network probe exactly as the API
does; the default page load stays cheap. The response is always 200: this
is an operator view that must render the degraded state, not withhold it.
"""
deep = _truthy_flag(request.query_params.get("deep"))
snapshot = load_system_health(deep=deep)
return HTMLResponse(
render_page(
title="System health",
body_html=render_system_health_page(snapshot),
)
)
async def queue(_request: Request) -> HTMLResponse:
snapshot = load_queue_snapshot()
return HTMLResponse(render_page(title="Queue", body_html=render_queue_page(snapshot)))
@@ -548,6 +573,114 @@ async def api_v1_timeline(request: Request) -> JSONResponse:
return JSONResponse(timeline_snapshot_to_dict(snapshot), status_code=status_code)
async def analytics(request: Request) -> HTMLResponse:
"""Read-only model usage, token cost, latency, and performance analytics HTML view (#651)."""
snapshot = load_analytics(
remote=request.query_params.get("remote"),
org=request.query_params.get("org"),
repo=request.query_params.get("repo"),
role=request.query_params.get("role"),
model=request.query_params.get("model"),
stage=request.query_params.get("stage"),
issue_number=_query_int(request, "issue"),
pr_number=_query_int(request, "pr"),
limit=_query_int(request, "limit") or 200,
)
return HTMLResponse(render_analytics_page(snapshot))
async def api_v1_analytics(request: Request) -> JSONResponse:
"""Read-only model usage, token cost, latency, and performance analytics API (#651)."""
snapshot = load_analytics(
remote=request.query_params.get("remote"),
org=request.query_params.get("org"),
repo=request.query_params.get("repo"),
role=request.query_params.get("role"),
model=request.query_params.get("model"),
stage=request.query_params.get("stage"),
issue_number=_query_int(request, "issue"),
pr_number=_query_int(request, "pr"),
limit=_query_int(request, "limit") or 500,
)
status_code = 200 if snapshot.ok else 500
return JSONResponse(analytics_snapshot_to_dict(snapshot), status_code=status_code)
async def api_v1_analytics_ingest(request: Request) -> JSONResponse:
"""Optional session instrumentation ingestion endpoint (#651).
Fail-closed write: every request is authorized through console_authz
(``record_analytics_usage``) before any control-plane DB mutation. Phase 1
keeps ``execution_enabled=False`` and denies unauthenticated callers, so
this route cannot be used as an unauthenticated write or XSS injection
vector (PR #876 F2).
"""
try:
body = await request.json()
except Exception:
body = {}
if not isinstance(body, dict):
body = {}
principal = resolve_principal(headers=dict(request.headers))
decision = authorize(
"record_analytics_usage", principal, for_execution=True
)
allowed = bool(decision.allowed and decision.execution_enabled)
console_audit.record_event(
action_id="record_analytics_usage",
result=(
console_audit.RESULT_ALLOWED
if allowed
else console_audit.RESULT_DENIED
),
decision=decision,
principal=principal,
target=_audit_target("record_analytics_usage", body),
request_id=_request_id(),
detail=decision.detail,
)
authorization = decision.to_dict()
if not allowed:
return JSONResponse(
{
"ok": False,
"error": "unauthorized",
"detail": (
"POST /api/v1/analytics/usage requires an authenticated "
"principal with record_analytics_usage execution enabled"
),
"authorization": authorization,
},
status_code=403,
)
usage_id = record_usage(
session_id=body.get("session_id"),
remote=body.get("remote", "dadeschools"),
org=body.get("org", ""),
repo=body.get("repo", ""),
project_id=body.get("project_id"),
role=body.get("role", "unknown"),
model=body.get("model", "unknown"),
issue_number=body.get("issue_number") or body.get("issue"),
pr_number=body.get("pr_number") or body.get("pr"),
stage=body.get("stage", "unknown"),
input_tokens=body.get("input_tokens"),
output_tokens=body.get("output_tokens"),
total_tokens=body.get("total_tokens"),
estimated_cost_usd=body.get("estimated_cost_usd"),
latency_ms=body.get("latency_ms"),
duration_ms=body.get("duration_ms"),
status=body.get("status", "success"),
metadata=body.get("metadata"),
)
return JSONResponse(
{"ok": True, "usage_id": usage_id, "authorization": authorization},
status_code=201,
)
async def method_not_allowed(request: Request, _exc: Exception) -> Response:
path = request.url.path
if path in _AUDIT_MUTATION_PATHS and request.method == "POST":
@@ -571,6 +704,7 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
Route("/", home, methods=["GET"]),
Route("/health", health, methods=["GET"]),
Route(SYSTEM_HEALTH_API_PATH, api_system_health, methods=["GET"]),
Route("/system-health", system_health, methods=["GET"]),
Route("/queue", queue, methods=["GET"]),
Route("/api/queue", api_queue, methods=["GET"]),
Route("/projects", projects, methods=["GET"]),
@@ -588,6 +722,10 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
Route("/runtime", runtime, methods=["GET"]),
Route("/api/runtime", api_runtime, methods=["GET"]),
Route("/api/v1/timeline", api_v1_timeline, methods=["GET"]),
Route("/analytics", analytics, methods=["GET"]),
Route("/api/analytics", api_v1_analytics, methods=["GET"]),
Route("/api/v1/analytics", api_v1_analytics, methods=["GET"]),
Route("/api/v1/analytics/usage", api_v1_analytics_ingest, methods=["POST"]),
Route("/audit", audit, methods=["GET", "POST"]),
Route("/api/audit", api_audit, methods=["GET", "POST"]),
Route("/worktrees", worktrees, methods=["GET"]),
+13
View File
@@ -236,6 +236,19 @@ _ACTION_SPECS: tuple[ConsoleAction, ...] = (
phase=3,
summary="Remove a remote feature branch.",
),
# #651 analytics ingest: local control-plane write, not a Gitea mutation.
# Phase 2 gated write so Phase 1 (ACTIVE_PHASE=1) fails closed on execution.
ConsoleAction(
action_id="record_analytics_usage",
task_key="record_analytics_usage",
action_class=CLASS_WRITE,
minimum_role=OPERATOR,
requires_confirmation=True,
dual_control=False,
break_glass=False,
phase=2,
summary="Ingest a model-usage / latency analytics event into the control-plane DB.",
),
)
ACTIONS: dict[str, ConsoleAction] = {a.action_id: a for a in _ACTION_SPECS}
+19
View File
@@ -236,6 +236,25 @@ def render_page(*, title: str, body_html: str, extra_head: str = "") -> str:
.badge-in-review {{ color: #9ec8f0; border-color: #3d5f7a; }}
.badge-duplicate {{ color: #e0c27a; border-color: #6b5730; }}
.badge-stale {{ color: #c9b8e8; border-color: #5a4a78; }}
.badge-health-ok {{ color: #8fd19e; border-color: #3d6b4a; }}
.badge-health-degraded {{ color: #e0c27a; border-color: #6b5730; }}
.badge-health-down {{ color: #f0a8a8; border-color: #7a3b3b; }}
.badge-health-skipped {{ color: var(--muted); }}
.badge-health-unproven {{ color: #c9b8e8; border-color: #5a4a78; }}
.health-card {{
margin: 1.25rem 0;
padding: 0.85rem 1rem 1rem;
border: 1px solid var(--border);
border-radius: 8px;
background: var(--surface);
}}
.health-card h3 {{ margin: 0 0 0.5rem; font-size: 1.05rem; }}
.health-card h4 {{ margin: 1rem 0 0.35rem; font-size: 0.92rem; color: var(--muted); }}
.health-headline {{ color: var(--text); font-size: 1rem; margin: 0 0 0.5rem; }}
.health-degraded {{ border-left-color: #e0c27a; }}
.health-stale {{ border-left-color: #f0a8a8; }}
ul.reasons {{ margin: 0.35rem 0; padding-left: 1.15rem; color: var(--muted); font-size: 0.9rem; }}
ul.reasons li {{ margin-bottom: 0.3rem; }}
</style>
{extra_head}
</head>
+2
View File
@@ -38,6 +38,7 @@ class NavGroup:
NAV_GROUPS: tuple[NavGroup, ...] = (
NavGroup("Health", (
NavItem("/health", "Liveness"),
NavItem("/system-health", "System health"),
)),
NavGroup("Traffic", (
NavItem("/queue", "Queue"),
@@ -64,6 +65,7 @@ NAV_GROUPS: tuple[NavGroup, ...] = (
)),
NavGroup("Insights", (
NavItem("/insights", "Insights", "stub"),
NavItem("/analytics", "Analytics"),
NavItem("/audit", "Audit"),
)),
)
+307
View File
@@ -0,0 +1,307 @@
"""HTML views for the system-health dashboard (#639).
Renders the read-only :class:`~webui.system_health.SystemHealthSnapshot`
produced by the Phase 1 system-health API (#634). The page offers no restart,
reload, or process-kill control: those are Phase 2 work, and manual process
kills are the contamination path #630 exists to prevent.
Every free-text field passes through :func:`webui.system_health.redact` before
it reaches HTML, so a probe detail that captured a token or a credentialed URL
cannot leak through the dashboard even though the API redacts it already.
"""
from __future__ import annotations
import html
from webui.system_health import (
STATUS_DEGRADED,
STATUS_DOWN,
STATUS_OK,
STATUS_SKIPPED,
STATUS_UNPROVEN,
DependencyProbe,
SystemHealthSnapshot,
redact,
)
_STATUS_BADGE_CLASS = {
STATUS_OK: "badge-health-ok",
STATUS_DEGRADED: "badge-health-degraded",
STATUS_DOWN: "badge-health-down",
STATUS_SKIPPED: "badge-health-skipped",
STATUS_UNPROVEN: "badge-health-unproven",
}
def _safe(value: object) -> str:
"""Escape free text for HTML after redacting anything secret-shaped.
Use this for every value that can carry arbitrary text probe details,
reasons, probe errors because those are where a credential could ride
along.
"""
return html.escape(redact(str(value)))
def _esc(value: object) -> str:
"""Escape a structured field for HTML without redacting it.
Commit SHAs, probe names, statuses, and timestamps are enumerated or
machine-generated, never credential-bearing. They must not go through
:func:`redact`: its opaque-token rule matches any 32-plus-character run,
so a 40-character git SHA would render as ``[redacted]`` and the parity
view the one thing an operator reads this page for would be blank.
"""
return html.escape(str(value))
def _status_badge(status: str) -> str:
css = _STATUS_BADGE_CLASS.get(status, "badge-health-unproven")
return f'<span class="badge {css}">{_esc(status)}</span>'
def _reason_list(reasons: tuple[str, ...], *, empty: str) -> str:
if not reasons:
return f"<p class='muted'>{html.escape(empty)}</p>"
items = "".join(f"<li>{_safe(reason)}</li>" for reason in reasons)
return f"<ul class='reasons'>{items}</ul>"
def _readiness_card(snapshot: SystemHealthSnapshot) -> str:
"""Overall readiness.
``ready`` and ``readiness_complete`` are shown separately on purpose: a
snapshot whose required probes never ran is not the same as one that ran
them and passed, and collapsing the two would render an unproven green.
"""
if snapshot.ready and snapshot.readiness_complete:
headline = "Ready"
elif snapshot.ready:
headline = "Ready (incomplete evidence)"
else:
headline = "Not ready"
return (
"<section class='health-card'>"
f"<h3>Readiness {_status_badge(snapshot.status)}</h3>"
f"<p class='health-headline'>{html.escape(headline)}</p>"
"<table class='detail'>"
f"<tr><th>Service</th><td><code>{_esc(snapshot.service)}</code></td></tr>"
f"<tr><th>Mode</th><td>{_esc(snapshot.mode)}</td></tr>"
f"<tr><th>Ready</th><td>{_esc(snapshot.ready)}</td></tr>"
"<tr><th>Readiness evidence complete</th>"
f"<td>{_esc(snapshot.readiness_complete)}</td></tr>"
"<tr><th>Deep probes requested</th>"
f"<td>{_esc(snapshot.deep_probes_requested)}</td></tr>"
f"<tr><th>Observed at</th><td><code>{_esc(snapshot.timestamp)}</code></td></tr>"
"</table>"
"<h4>Readiness reasons</h4>"
f"{_reason_list(snapshot.readiness_reasons, empty='No readiness objections recorded.')}"
"</section>"
)
def _version_card(snapshot: SystemHealthSnapshot) -> str:
version = snapshot.version
uptime_hours = snapshot.uptime_seconds / 3600.0
known = (
"resolved"
if version.known
else "unresolved — version fields could not be read from the checkout"
)
schema = version.control_plane_schema_version
return (
"<section class='health-card'>"
"<h3>Version and uptime</h3>"
"<table class='detail'>"
f"<tr><th>Git SHA</th><td><code>{_esc(version.git_sha or 'unknown')}</code></td></tr>"
"<tr><th>Git describe</th>"
f"<td><code>{_esc(version.git_describe or 'unknown')}</code></td></tr>"
"<tr><th>Control-plane schema</th>"
f"<td>{_esc(schema if schema is not None else 'unknown')}</td></tr>"
f"<tr><th>Python</th><td><code>{_esc(version.python_version)}</code></td></tr>"
f"<tr><th>Version status</th><td>{html.escape(known)}</td></tr>"
f"<tr><th>Started at</th><td><code>{_esc(snapshot.started_at)}</code></td></tr>"
"<tr><th>Uptime</th>"
f"<td>{snapshot.uptime_seconds:.3f}s ({uptime_hours:.2f}h)</td></tr>"
"</table>"
"</section>"
)
def _dependency_rows(probes: tuple[DependencyProbe, ...]) -> str:
if not probes:
return "<p class='muted'>No dependency probes were reported.</p>"
rows = []
for probe in probes:
latency = (
f"{probe.latency_ms:.1f} ms" if probe.latency_ms is not None else "n/a"
)
rows.append(
"<tr>"
f"<td><code>{_esc(probe.name)}</code></td>"
f"<td>{_esc(probe.kind)}</td>"
f"<td>{_status_badge(probe.status)}</td>"
f"<td>{_esc('required' if probe.required else 'optional')}</td>"
f"<td>{html.escape(latency)}</td>"
f"<td>{_safe(probe.detail)}</td>"
"</tr>"
)
return (
"<table class='registry'><thead><tr>"
"<th>Dependency</th><th>Kind</th><th>Status</th><th>Requirement</th>"
"<th>Latency</th><th>Detail</th>"
"</tr></thead><tbody>"
f"{''.join(rows)}</tbody></table>"
)
def _dependency_card(snapshot: SystemHealthSnapshot) -> str:
degraded = [probe for probe in snapshot.dependencies if probe.ran and not probe.healthy]
not_run = [probe for probe in snapshot.dependencies if not probe.ran]
banner = ""
if degraded:
names = ", ".join(sorted(probe.name for probe in degraded))
banner += (
"<div class='stub health-degraded'><p><strong>Degraded dependencies:</strong> "
f"{_esc(names)}</p></div>"
)
if not_run:
names = ", ".join(sorted(probe.name for probe in not_run))
banner += (
"<div class='stub'><p><strong>Not probed:</strong> "
f"{_esc(names)} — these contribute no evidence and are not counted "
"as healthy.</p></div>"
)
return (
"<section class='health-card'>"
"<h3>Dependencies</h3>"
f"{banner}"
f"{_dependency_rows(snapshot.dependencies)}"
"<p class='muted'>Details are redacted at the API boundary and again "
"before rendering; credentials are never displayed.</p>"
"</section>"
)
def _namespace_card(snapshot: SystemHealthSnapshot) -> str:
if not snapshot.mcp_namespaces:
body = "<p class='muted'>No MCP namespaces are declared.</p>"
else:
rows = []
for entry in snapshot.mcp_namespaces:
rows.append(
"<tr>"
f"<td><code>{_esc(entry.get('namespace'))}</code></td>"
f"<td><code>{_esc(entry.get('required_tool'))}</code></td>"
f"<td>{_status_badge(str(entry.get('status') or STATUS_UNPROVEN))}</td>"
f"<td>{_esc(entry.get('ide_namespace_proven'))}</td>"
f"<td>{_safe(entry.get('reason'))}</td>"
"</tr>"
)
body = (
"<table class='registry'><thead><tr>"
"<th>Namespace</th><th>Required tool</th><th>Status</th>"
"<th>IDE-proven</th><th>Reason</th>"
"</tr></thead><tbody>"
f"{''.join(rows)}</tbody></table>"
)
return (
"<section class='health-card'>"
"<h3>MCP namespaces</h3>"
f"{body}"
"<p class='muted'>The web process runs outside the IDE-managed MCP "
"client, so namespace health is reported as unproven rather than "
"guessed (#543).</p>"
"</section>"
)
def _stale_runtime_card(snapshot: SystemHealthSnapshot) -> str:
stale = snapshot.stale_runtime
if stale.stale:
warning = (
"<div class='stub health-stale'><p><strong>Stale runtime:</strong> "
"the running code, the checkout, and the remote-tracking commit "
"disagree. Capability gates may be evaluating obsolete code — "
"do not treat this runtime as mutation-safe.</p></div>"
)
elif not stale.determinable:
warning = (
"<div class='stub health-stale'><p><strong>Staleness "
"indeterminate:</strong> parity could not be proven, so this "
"runtime is not reported as mutation-safe.</p></div>"
)
else:
warning = ""
return (
"<section class='health-card'>"
"<h3>Stale-runtime parity</h3>"
f"{warning}"
"<table class='detail'>"
"<tr><th>Daemon head</th>"
f"<td><code>{_esc(stale.daemon_head or 'unknown')}</code></td></tr>"
"<tr><th>Checkout head</th>"
f"<td><code>{_esc(stale.checkout_head or 'unknown')}</code></td></tr>"
"<tr><th>Remote head</th>"
f"<td><code>{_esc(stale.remote_head or 'unknown')}</code></td></tr>"
f"<tr><th>Stale</th><td>{_esc(stale.stale)}</td></tr>"
f"<tr><th>Determinable</th><td>{_esc(stale.determinable)}</td></tr>"
f"<tr><th>Mutation safe</th><td>{_esc(stale.mutation_safe)}</td></tr>"
"</table>"
f"{_reason_list(stale.reasons, empty='Runtime, checkout, and remote agree.')}"
"</section>"
)
def _probe_error_card(snapshot: SystemHealthSnapshot) -> str:
if not snapshot.probe_errors:
return ""
return (
"<section class='health-card'>"
"<h3>Probe errors</h3>"
f"{_reason_list(snapshot.probe_errors, empty='')}"
"</section>"
)
def _recovery_card() -> str:
"""Sanctioned recovery pointers only — never a manual process kill (#630)."""
return (
"<section class='health-card'>"
"<h3>Recovery</h3>"
"<p class='muted'>This dashboard is read-only. Restart and reload "
"controls arrive in Phase 2 (#642); until then recovery runs through "
"the sanctioned client reconnect / operator restart path.</p>"
"<ul class='reasons'>"
"<li><a href='/runtime'>Runtime and session view</a> — active profile, "
"workflow hashes, and shell health.</li>"
"<li>Reconnect the MCP client from the IDE, then re-run the blocked "
"cycle. Never kill the daemon process manually: unmanaged kills are "
"recorded as runtime contamination (#630).</li>"
"<li>See <code>docs/webui-local-dev.md</code> for the documented "
"recovery sequence.</li>"
"</ul>"
"</section>"
)
def render_system_health_page(snapshot: SystemHealthSnapshot) -> str:
"""Render the full system-health dashboard body."""
return (
"<h2>System health</h2>"
"<p class='meta'>Read-only view of the Phase 1 system-health API "
"(<code>/api/v1/system/health</code>). Reload this page to refresh; "
"nothing here polls or mutates on your behalf.</p>"
f"{_readiness_card(snapshot)}"
f"{_stale_runtime_card(snapshot)}"
f"{_version_card(snapshot)}"
f"{_dependency_card(snapshot)}"
f"{_namespace_card(snapshot)}"
f"{_probe_error_card(snapshot)}"
f"{_recovery_card()}"
)