Compare commits

..
Author SHA1 Message Date
sysadmin 17ba1ff035 Merge pull request 'feat(transport): add transport-neutral MCP bind seam (Closes #931)' (#966) from feat/issue-931-transport-neutral-bind-seam into master 2026-07-28 13:36:49 -05:00
jcwalker3andClaude Opus 4.8 0104a76eea docs(remote-mcp): restamp the commit the #956 anchors resolve at
The execution-authorization gate shifted two cited lines in
mcp_daemon_guard.py, so the fixture and the document must name the commit
where the re-anchored citations actually resolve. Both still named
c1626081, where the two moved anchors no longer point at the claimed code.

Point both at a143cd06. This commit changes no .py file, so every anchor
that resolves at a143cd06 resolves here too. The historical note about
#930's drift between 7bf4f125 and aad5c8b4 stays as written.

Refs #931, #956. Addresses review 635.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-28 11:17:47 -05:00
jcwalker3andClaude Opus 4.8 a143cd065b fix(transport): gate transport execution behind commissioning (review 635)
Addresses B1 and B2 from formal review 635 on PR #966. Both are one defect:
the seam conflated two questions that must stay separate.

  recognition:  "is this an identifier this system knows, and may it be
                 bound, pinned and recorded?"
  execution:    "is this entrypoint commissioned to actually serve on it?"

SUPPORTED_TRANSPORTS answered the first and was then used to answer the
second, so registering streamable-http was the same act as executing it.
With mcp.run fed the bound value, GITEA_MCP_TRANSPORT=streamable-http
reached FastMCP.run -> run_streamable_http_async and started a uvicorn
HTTP server over the whole tool surface with auth=None, no TLS and no
per-request principal. That endpoint is owned by #938 and is an explicit
non-goal of #931. B2 is why B1 was reachable: bound_transport was read
only by reporting sites, so no decision anywhere depended on it.

The fix adds an execution-authorization layer rather than removing the
identifier or mapping it to stdio, both of which would have destroyed the
seam #931 exists to create.

mcp_transport_config gains EXECUTABLE_TRANSPORTS, a strict subset of
SUPPORTED_TRANSPORTS containing only the local transport, plus
TRANSPORT_EXECUTION_OWNER recording that #938 commissions the remote
listener. assess_transport_execution returns a structured verdict:
transport, recognized, executable, allowed, blocker_kind, owner_issue,
reasons and exact_next_action.

mcp_daemon_guard gains assess_serve_authorization, which is the decision
that consumes bound_transport and is what stops it being reporting-only
metadata, and authorize_transport_execution, which enforces two ordered
boundaries: assert_transport_bound first, so an unbound runtime keeps its
pre-existing #695 failure and reason code; then execution authorization,
so a registered but uncommissioned transport is refused before any
listener exists and before any tool can dispatch. TransportExecutionError
subclasses UnsanctionedRuntimeError, so every existing fail-closed handler
still catches it while a caller that cares can distinguish the two cases.
native_runtime_status now reports executable_transports, serve_authorized
and the full serve_authorization verdict.

The entrypoint serves through authorize_transport_execution.

Net effect. stdio: unchanged, still binds and still reaches the runner.
streamable-http: still recognized, still validated, still pinned, still
recorded in decision-lock and audit provenance -- and refused at the serve
boundary with blocker transport_listener_not_commissioned naming #938.
Unregistered identifiers still fail earlier, at bind validation. Unbound
execution keeps the #695 contract. Rebind idempotence and conflict
rejection, session isolation, provenance and secret handling are untouched.

Tests: 63 in tests/test_issue_931_transport_bind_seam.py, up from 42. New
coverage proves stdio reaches the real runner; the remote identifier binds
and is recorded yet never reaches run(); no socket is bound during the
refusal, asserted with a socket.bind tripwire; the refusal names the
transport, the unmet requirement and #938; the decision genuinely consumes
bound_transport, proven by flipping only that value; no mutation is
authorized after denial; verdicts do not leak across runtimes; and the
refusal carries no credential material. Per review 635 the two file-list
scans are now globs over the production surface, and a new test asserts
there is exactly one mcp.run site and that it routes through the guard.

#956 anchors re-anchored for the two lines this change shifted,
mcp_daemon_guard.py 178->195 and 519->583. No other anchor moved.

Full suite from the branches/ worktree: 28 failed, 5864 passed, 6 skipped,
1047 subtests. Base 9b80e75c: 28 failed, 5801 passed, 6 skipped, 1047
subtests. Failing test IDs are identical sets in both directions; the +63
passed are this module.

No listener, socket, authentication, TLS or principal code is added. #938,
#957, #958 and #959 remain unimplemented.

Refs #931, #938. Addresses review 635.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-28 11:15:00 -05:00
jcwalker3andClaude Opus 4.8 2fb835a1aa docs(remote-mcp): stamp the commit the #956 anchors were re-derived at
The anchor fixture records the commit its anchors were taken at, and the
threat model must state the same commit, so a reviewer can resolve every
cited line at a named revision. Both still named aad5c8b4, where the
anchors no longer resolve after the transport seam shifted the lines they
point at.

Point both at c1626081, the seam commit whose tree the anchors were
re-derived from. This commit changes no .py file, so every anchor that
resolves at c1626081 resolves here too. The historical note about #930's
drift between 7bf4f125 and aad5c8b4 stays as written -- it is the reason
the guard exists.

Refs #931, #956.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-28 10:01:01 -05:00
jcwalker3andClaude Opus 4.8 c162608175 feat(transport): add transport-neutral MCP bind seam
Closes #931.

The transport was a literal passed once at the bottom of the module,
bind_native_mcp_transport(transport="stdio"), and every guard that asks
"is this a trusted native session" resolves that question through the
value bound there. The string was a constant in the authorization chain
rather than configuration, so a remote transport had no way to be
expressed and no way to be told apart from an untrusted offline import.

mcp_transport_config becomes the single source of truth: it owns the
permitted set, the default, and the resolution of deployment
configuration (GITEA_MCP_TRANSPORT) into a validated identifier. Unset
configuration still yields stdio, so existing deployments are unchanged.
streamable-http is accepted so the bind is pluggable; standing up its
listener remains #938. sse is deliberately not registered.

bind_native_mcp_transport now resolves from configuration when no
transport argument is given, validates against that one permitted set
before writing the runtime record, and pins the result. bound_transport
is the shared accessor every transport-aware guard reads; it reports the
pinned value and never the environment, so a post-bind GITEA_MCP_TRANSPORT
change cannot move what a guard observes -- the rule already applied to
the session-state root under #695 AC2. Rebinding the same identifier is
idempotent; rebinding a different one is refused, so two guards can never
disagree within one process. assert_transport_bound fails closed before
mcp.run, so an absent or invalid bind stops the server instead of serving
tools over a transport no guard can name.

The bound identifier is now recorded in durable provenance:
mutation_provenance_fields gains bound_transport, which reaches the
decision lock through mcp_session_state.save_state and the audit records
that already spread those fields. The pre-existing transport field keeps
its trust-class values, so no durable record changes shape.
assess_transport_for_auth_mint reports the bound transport through the
same accessor; its verdict is unchanged.

#956's anchor fixture and threat model are re-anchored for the lines this
change shifts, and the two boundary claims that #931 makes false -- the
"literal stdio bind" in B1 and the stdio contract at mcp_server.py:4 --
are restated. Without this the #956 guard fails, which is what it is for.

Non-goals, untouched: the remote listener (#938), per-request principal
resolution (#932), client-managed provenance policy (#934), TLS and
remote client authentication, role capability sets, repository-binding
semantics.

Tests: tests/test_issue_931_transport_bind_seam.py, 42 tests, covering
default bind, explicit stdio, the accepted non-stdio identifier, an
unregistered identifier, no bind at all, one-value agreement across
guards, post-bind environment tampering, the durable decision-lock
record, the rebind contract, and the #695 protections that must not
regress.

Full suite from the branches/ worktree: 28 failed, 5843 passed, 6
skipped, 1047 subtests. The pinned base 9b80e75c reports 28 failed, 5801
passed, 6 skipped, 1047 subtests. The failing test IDs are identical sets
-- no failure added, none fixed; the +42 passed are this PR's new module.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-28 09:59:01 -05:00
sysadmin 9b80e75ca3 Merge pull request 'docs(remote-mcp): threat model, trust boundaries, and decomposition ruling (Closes #956)' (#965) from docs/issue-956-remote-mcp-threat-model into master 2026-07-28 08:19:45 -05:00
sysadminandClaude Opus 4.8 b3de9c941c docs(remote-mcp): threat model, trust boundaries, and decomposition ruling (Closes #956)
#930 inventoried stdio coupling; nothing stated what the adversary is, what
each boundary protects, or why one process may hold credentials for several
services. This adds that document as child 2 of epic #929.

Adds docs/remote-mcp/threat-model.md covering 10 assets, the 5 adversaries
#956 names, 9 trust boundaries, the data flows between them, a 14-entry
per-boundary credential inventory, and an explicit decomposition ruling.

Findings established from live native evidence at this commit:

- Four prgs roles (reviewer, merger, reconciler, controller) resolve to one
  Gitea account, so separation of duty between approving and landing is
  enforced only by which process a call reaches. The mdcps tenant has no role
  separation at all: author, reviewer, and merger share one account.
- Any one role process can resolve every other role's credential.
  gitea_list_profiles reports "credentials present" for other roles because it
  calls resolve_token on each one.
- The Gitea server reads Jenkins and GlitchTip secrets out of the keychain to
  produce the "authenticated" word in gitea_audit_config's service summaries.
- Jenkins and GlitchTip were already decomposed into separate MCP servers; the
  credential references were left behind in the Gitea configuration.

Ruling D1 forbids a single integration process from holding credentials for
unrelated services, with one time-boxed dual-run exception for the local fleet
that expires with #939. D2 requires separation of duty to be credential-backed,
D3 scopes credential resolution to the request principal, and D4 gives
coordination state its own authority. Every #929 child from 2 through 10 is
mapped to the boundary it implements.

Anchors are enforced rather than asserted. #930's inventory anchors into
gitea_mcp_server.py had already drifted between 7bf4f125 and aad5c8b4 with
nothing detecting it, so this change ships the guard that was missing:
docs/remote-mcp/threat-model-anchors.json declares all 58 anchors with the
substring each must contain, and tests/test_issue_956_threat_model.py fails if
any anchor does not resolve, if the document cites an anchor the fixture does
not cover, or if the structural obligations regress.

Documentation only. No server behavior changes.

Tests:
- tests/test_issue_956_threat_model.py: 17 passed.
- Four sabotage probes confirm the validator is not passing vacuously
  (shifted anchor, undeclared citation, broken count tally, and a blanked
  boundary owner). The last two probes exposed real weaknesses in the checks
  themselves, which were fixed: the section slice now stops at the next
  heading, and boundary ownership is read only from mapping table rows.
- Docs-sensitive sweep (17 modules referencing docs/): 559 passed.
- Full suite: 30 failed, 5799 passed, 6 skipped. All 30 reproduce on a clean
  base worktree at aad5c8b4; branch failures are a strict subset of base
  failures. No production file is modified by this change.

Closes #956

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-28 04:41:26 -04:00
sysadmin aad5c8b423 Merge pull request 'fix(author): unify the bootstrap and lock_issue issue-lock contract (Closes #953)' (#954) from fix/issue-953-bootstrap-lock-provenance into master
Merges PR #954 at approved head b4c9f55890.

Approval: review 633 APPROVE at b4c9f55890.
Base: master at 82d71b7702.

Closes #953
2026-07-28 02:21:11 -05:00
sysadminandClaude Opus 4.8 b4c9f55890 test(author): execute the native #953 tool functions and compensation paths (#953)
Review 632 F4. Every prior reference to `gitea_recover_incomplete_bootstrap_lock`
and `gitea_inspect_issue_lock_contract` under `tests/` was a string literal — a
`tool=` argument, an assertion on returned prose, or a docs substring check —
and the suite reconstructed the recovery sequence by hand from
`assess_bootstrap_lock_recovery`, `build_canonical_issue_lock`,
`build_recovery_record`, and `bind_session_lock`. A hand-written sequence
validates the decision layer but cannot see a divergence between itself and
the tool body, which is exactly how F1 and F2 — both call-site defects —
survived 61 passing cases.

The new cases drive the registered functions against a real `git init`
repository, a real durable lock file, and config-backed profiles:

* `NamespaceMutationWallOnRecovery` — the gate is reached with this task and
  `author_role_exclusive=True`; its return value aborts the tool rather than
  being computed and discarded; the author namespace succeeds and produces a
  canonical lock; a reviewer namespace is refused with `namespace_block` even
  when the claimant data would otherwise match, and emits the standard BLOCKED
  audit record; merger is refused; a mismatched claimant profile is refused by
  the exact-owner layer with the namespace wall explicitly clear; a mismatched
  head is refused; every refusal leaves the lock bytes, generation, branch,
  worktree, and an unrelated lock untouched. A subtest matrix asserts the
  role-kind wall admits `author` and refuses reviewer, merger, limited, and
  mixed.
* `InspectionToolExecutes` — the registered read-only tool reports the contract
  and the recovery preview while leaving lock bytes, mtime, HEAD, and porcelain
  status unchanged, and reports an absent lock without creating one.
* `Ac7PostCompensationGuidance` — drives the real bootstrap to its AC7 refusal
  with a forced partial lock and asserts the returned action against the state
  the rollback actually left: complete cleanup directs to a bootstrap retry and
  that retry is then executed and succeeds, leaving exactly one canonical lock
  and one branch; partial cleanup with a surviving lock, and with a surviving
  branch and worktree, each get their own executable action; a rollback that
  never completed is distinguished from both; no recommendation names a deleted
  artifact; unrelated locks are byte-identical afterwards. Two cases cover
  `release_session_lock` directly — that the rollback now really removes the
  lock, and that it refuses a lock owned by another session.
* `NativeEndToEndBootstrapToCreatePr` — bootstrap, inspect, heartbeat,
  legitimate divergence (commit and push), pre-mutation ownership re-check, and
  the unchanged #447 create-PR provenance guard, in one sequence against a real
  origin. `gitea_lock_issue` is patched to fail the test if anything reaches for
  it, so the bootstrap lock is proved to carry the whole cycle unrepaired.
* `DeadProvenanceConstantRemoved` — the removed constant stays removed and the
  sanctioned source set stays unwidened.

`_NativeToolBase` clears `role_session_router` route state per test: the sticky
reviewer-stop marker is process-global and, now that this task is registered in
`AUTHOR_TASKS`, an earlier suite leaving it set would make the first gate refuse
before the namespace gate under test is reached.

Also documents both new tools in `docs/mcp-tool-inventory.md`, so this branch
adds no drift to `test_documented_inventory_equals_registered_tools`; the
failure reason there is now identical to the pinned base's.

Suite: 61 -> 92 cases.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-28 02:08:21 -04:00
sysadminandClaude Opus 4.8 1aa351718a fix(author): derive AC7 guidance from the state compensation actually leaves (#953)
Review 632 F2. The bootstrap AC7 read-back refusal calls
`run_compensating_recovery` and *then* reported
`author_lock_contract.recommended_action(contract)` — advice computed from the
malformed lock that provoked the rollback, not from the state the rollback
left. Compensation releases the lock, removes the worktree (always clean
there, no implementation bytes having been written), and deletes the branch,
so an author following that advice got `no_durable_lock` from
`gitea_recover_incomplete_bootstrap_lock` and, had the lock survived,
`worktree_invalid` instead; the `gitea_lock_issue` half of the same sentence
cannot bind a worktree that no longer exists. Two refusals in a row for a
state a plain bootstrap retry fixes — the unexecutable-guidance failure class
this issue exists to remove, reintroduced on the new fail-closed path.

Investigating that path surfaced why the "clean retry" state was in practice
unreachable: `run_compensating_recovery` has called
`issue_lock_store.release_session_lock` since #850, and that function has
never existed. The `AttributeError` landed in a bare `except Exception: pass`,
so every rollback removed the branch and worktree and silently left the lock
behind — precisely the uninspectable, unrecoverable state #953 is about
(`gitea_recover_incomplete_bootstrap_lock` refuses `worktree_invalid`,
`gitea_lock_issue` has no worktree to bind). Confirmed dead at the pinned base
`82d71b77`, not introduced by this branch.

`release_session_lock` is therefore implemented: it removes exactly one
durable lock whose recorded `owner_session` matches the caller's, keyed by
repository when known, refusing on zero or multiple matches so no caller can
delete a lock it does not own and an ambiguous directory is never guessed at.
Bootstrap phase journals and session pointers that share the directory are
excluded by shape. The flock sidecar is deliberately left alone. The caller no
longer swallows a release failure; it records `lock_release_failed:...`.

`assess_post_compensation_state` then classifies from directly observed
durable state — lock file, worktree directory, and branch ref — rather than
from the journal's `rolled_back` list, which records only what compensation
attempted. Three distinct states: `complete` (nothing remains),
`partial` (rollback ran, artifacts survive by design or because a step
errored), `failed` (rollback never completed, so nothing is proven removed).

`post_compensation_action` answers for exactly what survives:

  complete                      -> re-run gitea_bootstrap_author_issue_worktree
  lock + branch + worktree      -> gitea_recover_incomplete_bootstrap_lock
  branch + worktree, no lock    -> gitea_lock_issue (still base-equivalent)
  lock only, worktree gone      -> gitea_inspect_issue_lock_contract
  branch only                   -> gitea_inspect_issue_lock_contract, then retry
  rollback did not complete     -> gitea_inspect_issue_lock_contract

No branch names an artifact the classification says is gone, and a failed
rollback step is stated rather than presented as an intentional outcome. The
refusal payload carries `compensating_recovery` and `post_compensation_state`
alongside the derived `exact_next_action`, still with
`implementation_allowed: false`.

Also removes the dead `SOURCE_BOOTSTRAP_LOCK_RECOVERY` constant (review 632
F3), which had no readers and implied a second lock source; the deliberate
reuse of `SOURCE_LOCK_ISSUE` is now stated as a comment. The #447 guard and
`SANCTIONED_LOCK_SOURCES` remain untouched.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-28 02:08:21 -04:00
sysadminandClaude Opus 4.8 55d66c57e4 fix(author): gate bootstrap-lock recovery on the namespace mutation wall (#953)
Review 632 F1. `gitea_recover_incomplete_bootstrap_lock` writes the same
durable author issue lock as `gitea_recover_dirty_orphaned_issue_worktree`
but gated only on the reviewer-stop router check and the profile permission
block. Every other author state-creating mutation carries a third gate,
`_namespace_mutation_block`, and this tool was the outlier.

The surviving gates did not cover the gap: the task's required permission is
`gitea.issue.comment`, which every configured role holds, and
`_ensure_matching_profile` returns a profile name both call sites discard, so
it refuses nothing. A reviewer-bound session reached the exact-owner claimant
comparison inside `assess_bootstrap_lock_recovery` and was refused there —
one layer too late, with no namespace evaluation and no BLOCKED audit record
of the attempt.

Adding the call alone would have been inert. `check_author_mutation_namespace`
routes through `role_session_router.required_role_for_task`, which reads
`TASK_REQUIRED_ROLE` — not `task_capability_map` — and that table had no entry
for this task, so the check short-circuited to "allowed" for every caller.
The task is therefore registered in `TASK_REQUIRED_ROLE` and `AUTHOR_TASKS`,
matching the role the capability map already records.

The reviewer-namespace check alone is also insufficient for a task gated on
`gitea.issue.comment`, since merger, controller, and reconciler profiles hold
it too. `role_namespace_gate.check_author_role_kind` is added as an opt-in,
additive wall requiring the active profile's derived role kind to be exactly
`author`; `mixed` is refused rather than admitted. `_namespace_mutation_block`
gains a keyword-only `author_role_exclusive` flag, default off, so the six
pre-existing call sites are byte-for-byte unchanged.

Refusals produce the standard structured denial (`namespace_block: true`,
`mcp_namespace`) and the standard BLOCKED audit record, and no lock, lease,
generation, branch, worktree, issue, or PR is touched. Valid `prgs-author`
execution is unaffected. No reviewer permission is broadened and no existing
role wall is weakened.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-28 02:08:21 -04:00
sysadminandClaude Opus 4.8 cdf0daefa9 fix(author): unify the bootstrap and lock_issue issue-lock contract (Closes #953)
gitea_bootstrap_author_issue_worktree wrote a lock no downstream author
operation accepts, then directed the author straight to implementation. Once
the branch carried commits, heartbeat, re-lock, exact-owner renewal, and the
#447 create-PR guard all refused simultaneously and no sanctioned recovery
path remained eligible.

Each of those gates is individually correct. The defect was that two writers
disagreed about what a lock is.

- Add author_lock_contract as the single canonical definition: claimant,
  work_lease, lock_provenance, generation, and an explicit expiration state.
  Both gitea_lock_issue and bootstrap now build through it.
- Promote issue_lock_store.lock_claimant to the one shared claimant reader and
  use it in the ownership check, so a claimant recorded at the lock top level
  is read rather than refused. The values are still compared against
  server-resolved identity and profile, so no legacy placement grants anything
  the canonical placement would not.
- Represent missing expiration explicitly. An absent expires_at previously read
  as "not yet expired", leaving a malformed lock permanently non-expiring and
  permanently ineligible for #760 renewal.
- Bootstrap reads its lock back and verifies it structurally before reporting
  success. A partial lock fails closed while the worktree is still
  base-equivalent, names the missing fields, and never reports
  implementation_allowed. Its exact_next_action now matches the state returned.
- Add gitea_recover_incomplete_bootstrap_lock for locks already written by the
  old bootstrap, including those whose branches carry pushed commits. It never
  moves, resets, or rewinds a branch, never requires base-equivalence, never
  pushes or opens a PR, and touches only the target lock. It proves repository,
  issue, claimant username and profile, branch, worktree, registration, and
  head before writing, refuses healthy foreign-owned locks, and mints
  provenance and authorization server-side.
- Add gitea_inspect_issue_lock_contract, a strictly read-only surface.
- Document the required ordering and the recovery path.

The #447 provenance guard is unchanged and the sanctioned source set was not
widened: bootstrap now satisfies the guard rather than the guard being relaxed
to admit bootstrap.

Tests: 61 new cases covering the canonical schema, immediate heartbeat,
renewal before and after commits, the create_pr guard, executable next actions,
partial and malformed and missing-expiration and expired and same-owner and
foreign-owner and legacy locks, recovery isolation, read-only inspection, and
the full bootstrap-implement-commit-push-create_pr regression against isolated
fixtures.

Full suite at head: 28 failed, 5753 passed, 6 skipped, 1042 subtests.
Full suite at base 82d71b77: 28 failed, 5692 passed, 6 skipped, 1042 subtests.
Failing test-ID sets are identical, so there are zero regressions; the +61
passes are this issue's new suite.

Issue #949 was preserved and not recovered: its branch remains at
92615f474b and its worktree, lock, and PR state
were not touched. The #949-shaped regression uses isolated fixtures only.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-28 00:35:12 -04:00
29 changed files with 6077 additions and 2799 deletions
+134 -24
View File
@@ -23,6 +23,7 @@ import shutil
import subprocess import subprocess
from typing import Any, Mapping from typing import Any, Mapping
import author_lock_contract
import author_mutation_worktree import author_mutation_worktree
import control_plane_db import control_plane_db
import issue_lock_store import issue_lock_store
@@ -278,6 +279,36 @@ def _verify_assignment_and_lease_ids(
return None return None
def _branch_exists(canonical_repo_root: str, branch_name: str) -> bool:
"""Whether *branch_name* still resolves in the canonical checkout (#953 F2).
Used after compensating recovery to observe what survived rather than infer
it from the journal. Fails closed to ``True``: an unobservable branch is
reported as present, so the recommendation stays conservative rather than
telling an author to re-bootstrap over something that may still be there.
"""
if not branch_name:
return False
try:
res = subprocess.run(
[
"git",
"-C",
canonical_repo_root,
"rev-parse",
"--verify",
"--quiet",
f"refs/heads/{branch_name}",
],
capture_output=True,
text=True,
check=False,
)
except Exception:
return True
return res.returncode == 0
def run_compensating_recovery( def run_compensating_recovery(
journal: dict[str, Any], journal: dict[str, Any],
canonical_repo_root: str, canonical_repo_root: str,
@@ -314,10 +345,21 @@ def run_compensating_recovery(
issue_number=issue_num, issue_number=issue_num,
session=session_id, session=session_id,
lock_dir=journal_dir, lock_dir=journal_dir,
remote=journal.get("remote"),
# The same defaults the lock was written under, so the
# rollback targets the exact file bind_session_lock keyed.
org=journal.get("org") or "Scaled-Tech-Consulting",
repo=journal.get("repo") or "Gitea-Tools",
) )
rolled_back.append(f"lock:issue-{issue_num}") rolled_back.append(f"lock:issue-{issue_num}")
except Exception: except Exception as exc:
pass # #953 F2: a swallowed failure here is what made the rollback
# report success while leaving an unrecoverable lock behind.
# Record it so the post-compensation classification can see the
# lock survived and recommend accordingly.
rolled_back.append(
f"lock_release_failed:issue-{issue_num}:{type(exc).__name__}"
)
artifacts["lock_created"] = False artifacts["lock_created"] = False
worktree_created = ( worktree_created = (
@@ -1198,26 +1240,31 @@ def bootstrap_author_issue_worktree(
save_phase_journal(journal, journal_dir=lock_dir) save_phase_journal(journal, journal_dir=lock_dir)
# Phase 6: STATE_ESTABLISHED — Issue Lock Acquisition # Phase 6: STATE_ESTABLISHED — Issue Lock Acquisition
#
# #953: this used to hand-build a thinner record — claimant at the top
# level, no work_lease, no lock_provenance, no expiry — which every
# downstream reader then refused. It now builds through the one shared
# canonical contract, so the lock bootstrap writes is the same lock
# gitea_lock_issue writes.
from datetime import datetime, timezone from datetime import datetime, timezone
try: try:
lock_data = { lock_data = author_lock_contract.build_canonical_issue_lock(
"remote": remote, issue_number=issue_number,
"org": org or "Scaled-Tech-Consulting", branch_name=target_branch,
"repo": repo or "Gitea-Tools", worktree_path=target_worktree,
"issue_number": issue_number, remote=remote,
"branch": target_branch, org=org or "Scaled-Tech-Consulting",
"branch_name": target_branch, repo=repo or "Gitea-Tools",
"worktree_path": target_worktree, identity=identity,
"owner_session": session, profile=profile,
"claimant": { tool="gitea_bootstrap_author_issue_worktree",
"username": identity, source=author_lock_contract.SOURCE_BOOTSTRAP,
"profile": profile, owner_session=session,
}, assignment_id=assignment_id,
"assignment_id": assignment_id, lease_id=lease_id,
"lease_id": lease_id, expected_base_sha=live_master_sha,
"expected_base_sha": live_master_sha, )
"created_at": datetime.now(timezone.utc).isoformat(), lock_data["created_at"] = datetime.now(timezone.utc).isoformat()
}
journal.setdefault("pending_creations", {})["lock"] = True journal.setdefault("pending_creations", {})["lock"] = True
journal["artifacts_created"]["lock_created"] = True journal["artifacts_created"]["lock_created"] = True
save_phase_journal(journal, journal_dir=lock_dir) save_phase_journal(journal, journal_dir=lock_dir)
@@ -1235,9 +1282,65 @@ def bootstrap_author_issue_worktree(
"exact_next_action": "Verify lease/assignment state and retry.", "exact_next_action": "Verify lease/assignment state and retry.",
} }
# ── #953 AC7: verify the lock that was actually written ──
# Reporting "lock_created: true" and then directing the author to
# implement is what produced the unrecoverable state: by the time any
# reader refused the lock, the branch already carried commits and every
# sanctioned recovery path had become ineligible. The lock is therefore
# read back from disk and structurally verified *before* this function
# can report success, and a partial lock fails closed here — while the
# branch is still base-equivalent and recovery is still cheap.
written_lock = issue_lock_store.read_lock_file(lock_res)
contract = author_lock_contract.assess_lock_contract(written_lock)
if not contract["canonical"]:
journal["failure_reason"] = author_lock_contract.format_contract_refusal(
contract
)
compensation = run_compensating_recovery(
journal, root, journal_dir=lock_dir
)
# AC5/AC15: the recommendation must describe the state compensation
# actually left, not the state that provoked it.
# ``run_compensating_recovery`` has by now released the lock, removed
# the worktree, and deleted the branch, so recommending
# incomplete-lock recovery for those exact artifacts would refuse
# twice over. Observe what survived and answer for that.
post_state = author_lock_contract.assess_post_compensation_state(
compensation,
lock_present=bool(lock_res) and os.path.exists(lock_res),
worktree_present=os.path.isdir(target_worktree),
branch_present=_branch_exists(root, target_branch),
)
return {
"success": False,
"reason_code": "incomplete_issue_lock_contract",
"message": author_lock_contract.format_contract_refusal(contract),
"issue_number": issue_number,
"branch_name": target_branch,
"worktree_path": target_worktree,
"lock_state": lock_res,
"lock_contract": contract,
"missing_fields": contract["missing_fields"],
"implementation_allowed": False,
"compensating_recovery": compensation,
"post_compensation_state": post_state,
# AC15: never strand a branch or worktree without a structured
# recovery recommendation — and never name an artifact the
# rollback has already deleted.
"exact_next_action": author_lock_contract.post_compensation_action(
post_state,
issue_number=issue_number,
branch_name=target_branch,
worktree_path=target_worktree,
missing_fields=contract["missing_fields"],
),
"phase_journal": journal,
}
journal["phases"][PHASE_6_STATE_ESTABLISHED] = { journal["phases"][PHASE_6_STATE_ESTABLISHED] = {
"status": "completed", "status": "completed",
"lock": lock_res, "lock": lock_res,
"lock_contract": contract["contract"],
} }
journal["phases"][PHASE_7_TRANSITION_COMPLETED] = { journal["phases"][PHASE_7_TRANSITION_COMPLETED] = {
"status": "completed", "status": "completed",
@@ -1261,9 +1364,16 @@ def bootstrap_author_issue_worktree(
"assignment_id": assignment_id, "assignment_id": assignment_id,
"idempotency_key": key, "idempotency_key": key,
"lock_state": lock_res, "lock_state": lock_res,
"lock_contract": contract,
# #953 AC6: the canonical ownership token for this claim. Never null
# on a successful bootstrap — it is the fencing token every
# subsequent heartbeat and renewal is checked against.
"task_session_id": contract["task_session_id"],
"implementation_allowed": True,
"phase_journal": journal, "phase_journal": journal,
"exact_next_action": ( # #953 AC5: executable under the state actually returned. The lock
"Call gitea_whoami, then gitea_resolve_task_capability(task='work_issue') " # has been read back and verified canonical, so proceeding to
"and proceed with author implementation in the bootstrapped worktree." # implementation is genuinely the correct next step here — which is
), # exactly what the old unconditional wording could not promise.
"exact_next_action": author_lock_contract.recommended_action(contract),
} }
+623
View File
@@ -0,0 +1,623 @@
"""One canonical author issue-lock contract shared by every writer (#953).
Before this module, ``gitea_lock_issue`` and
``gitea_bootstrap_author_issue_worktree`` each wrote their own lock record.
``gitea_lock_issue`` wrote the canonical shape — ``work_lease`` carrying the
claimant plus a sanctioned ``lock_provenance`` — while bootstrap wrote a thinner
record with the claimant at the lock top level, ``lease_id: null``, and no
``work_lease``, ``lock_provenance``, or expiry at all.
Every downstream reader was written against the canonical shape, so a lock that
bootstrap reported as successfully created was simultaneously:
* un-heartbeatable — the ownership check read the claimant only from
``work_lease.claimant``;
* un-renewable — expiry is read only from ``work_lease.expires_at``, so a
missing lease read as "never expires", and #760 exact-owner renewal only ever
assesses an *expired* lease;
* un-re-lockable — the branch had by then advanced past its base;
* and rejected by the #447 create-PR provenance guard.
Each of those gates is individually correct. The defect was that two writers
disagreed about what a lock *is*. This module is the single definition, and both
writers now build through it.
Nothing here weakens a guard. ``build_sanctioned_lock_provenance`` remains the
only provenance source, provenance is never accepted from a caller, and the
#447 guard is untouched — this module simply makes bootstrap satisfy it.
"""
from __future__ import annotations
from datetime import datetime, timedelta, timezone
from typing import Any, Mapping
import issue_lock_provenance
import issue_lock_store
import lease_policy
# Bootstrap writes through the same sanctioned source as gitea_lock_issue: the
# lock it produces *is* a canonical lock, not a second dialect that readers must
# learn. Adding a distinct source would have required widening
# SANCTIONED_LOCK_SOURCES, which is exactly the #447 weakening this issue's
# safety requirements forbid.
SOURCE_BOOTSTRAP = issue_lock_provenance.SOURCE_LOCK_ISSUE
# Recovery of an incomplete bootstrap lock (#953 AC8-AC11) deliberately writes
# through SOURCE_LOCK_ISSUE too, and records its distinctness in
# ``lock_provenance.written_by_tool`` plus the ``bootstrap_lock_recovery``
# transition block instead. There is no distinct recovery *source* constant, for
# the same reason bootstrap has none: minting one would require widening
# SANCTIONED_LOCK_SOURCES, which the #447 safety requirements forbid.
#: Top-level keys every canonical author issue lock must carry.
REQUIRED_LOCK_FIELDS: tuple[str, ...] = (
"remote",
"org",
"repo",
"issue_number",
"branch_name",
"worktree_path",
"work_lease",
"lock_provenance",
)
#: Keys every canonical ``work_lease`` must carry.
REQUIRED_WORK_LEASE_FIELDS: tuple[str, ...] = (
"operation_type",
"issue_number",
"branch",
"worktree_path",
"claimant",
"created_at",
"expires_at",
"last_heartbeat_at",
"task_session_id",
"lifecycle_version",
)
# ── Explicit expiration states (AC12) ──
# The bug this replaces: a lock with no recorded expiry produced
# ``is_lease_expired() -> False``, which reads as "not yet expired" and made the
# lock permanently non-expiring *and* permanently ineligible for the renewal
# path, which only ever assesses an expired lease. "Absent" and "in the future"
# are different facts and are now named differently.
EXPIRATION_RECORDED = "recorded"
EXPIRATION_MISSING = "missing"
EXPIRATION_UNPARSEABLE = "unparseable"
#: Structural verdicts returned by :func:`assess_lock_contract`.
CONTRACT_CANONICAL = "canonical"
CONTRACT_INCOMPLETE = "incomplete"
CONTRACT_LEGACY = "legacy"
CONTRACT_ABSENT = "absent"
def _text(value: Any) -> str:
return str(value or "").strip()
def now_utc() -> datetime:
return datetime.now(timezone.utc)
def format_timestamp(value: datetime) -> str:
"""Serialize in the durable ``...Z`` form already used on disk."""
return (
value.astimezone(timezone.utc)
.replace(microsecond=0)
.isoformat()
.replace("+00:00", "Z")
)
def lock_claimant(lock: Mapping[str, Any] | None) -> dict[str, str]:
"""Read the claimant from either canonical or legacy placement.
``work_lease.claimant`` is canonical and is preferred. A top-level
``claimant`` is the legacy/bootstrap placement and is accepted as a
fallback (AC14) — three separate readers already disagreed about this
(``issue_lock_store``, ``issue_lock_renewal``, ``issue_lock_recovery``),
which is why it now lives in one place.
Reading a legacy placement is *not* a widening: every caller still compares
the values it returns against server-resolved identity and profile. This
only decides where to look, never whether ownership is proven.
Delegates to ``issue_lock_store.lock_claimant`` rather than reimplementing
the rule. A second copy here would be a fourth reader that could drift from
the other three, which is the exact failure #953 exists to end. It lives in
the store because ``author_lock_contract`` imports the store, so defining it
here would make that import circular.
"""
recorded = issue_lock_store.lock_claimant(dict(lock) if isinstance(lock, Mapping) else None)
return {
"username": _text(recorded.get("username")),
"profile": _text(recorded.get("profile")),
}
def claimant_placement(lock: Mapping[str, Any] | None) -> str:
"""Where the claimant was found: ``work_lease``, ``top_level``, or ``absent``."""
if not isinstance(lock, Mapping):
return "absent"
lease = lock.get("work_lease")
if isinstance(lease, Mapping) and isinstance(lease.get("claimant"), Mapping):
return "work_lease"
if isinstance(lock.get("claimant"), Mapping):
return "top_level"
return "absent"
def build_claimant(*, username: str | None, profile: str | None) -> dict[str, str]:
"""Build the canonical claimant pair from server-resolved values."""
return {"username": _text(username), "profile": _text(profile)}
def build_author_issue_work_lease(
*,
issue_number: int,
branch_name: str,
worktree_path: str,
claimant: Mapping[str, Any],
task_session_id: str | None = None,
created: datetime | None = None,
) -> dict[str, Any]:
"""Build the canonical author ``work_lease``.
The single definition behind both writers. The TTL comes from the central
policy rather than a literal, and the window slides from the last valid
heartbeat (#790), so an abandoned task releases its claim within one TTL.
"""
started = created or now_utc()
policy = lease_policy.policy_for(lease_policy.TASK_CLASS_AUTHOR_ISSUE_WORK)
expires = started + timedelta(minutes=policy.initial_ttl_minutes)
session_id = _text(task_session_id) or issue_lock_store.mint_task_session_id(
issue_lock_store.AUTHOR_ISSUE_WORK_LEASE
)
return {
"operation_type": issue_lock_store.AUTHOR_ISSUE_WORK_LEASE,
"issue_number": int(issue_number),
"pr_number": None,
"branch": branch_name,
"worktree_path": worktree_path,
"claimant": dict(claimant),
"created_at": format_timestamp(started),
"expires_at": format_timestamp(expires),
"last_heartbeat_at": format_timestamp(started),
# #790 AC-N1: the ownership key for this task, distinct from the
# recorded PID, which is the shared daemon and identifies no task.
"task_session_id": session_id,
# #790 AC-N8: the explicit lifecycle marker. Its absence — never a
# timestamp comparison — is what makes a lock legacy.
"lifecycle_version": lease_policy.LIFECYCLE_HEARTBEAT_V1,
"heartbeat_count": 1,
}
def build_canonical_issue_lock(
*,
issue_number: int,
branch_name: str,
worktree_path: str,
remote: str,
org: str,
repo: str,
identity: str | None,
profile: str | None,
tool: str,
source: str = issue_lock_provenance.SOURCE_LOCK_ISSUE,
owner_session: str | None = None,
assignment_id: str | None = None,
lease_id: str | None = None,
expected_base_sha: str | None = None,
task_session_id: str | None = None,
created: datetime | None = None,
) -> dict[str, Any]:
"""Build a complete canonical lock record.
``tool`` and ``source`` are server-supplied. There is deliberately no
parameter through which a caller could inject provenance: the #953 safety
requirements forbid caller-manufactured provenance, so provenance is always
minted here from ``build_sanctioned_lock_provenance``.
"""
claimant = build_claimant(username=identity, profile=profile)
work_lease = build_author_issue_work_lease(
issue_number=issue_number,
branch_name=branch_name,
worktree_path=worktree_path,
claimant=claimant,
task_session_id=task_session_id,
created=created,
)
record: dict[str, Any] = {
"remote": remote,
"org": org,
"repo": repo,
"issue_number": int(issue_number),
"branch": branch_name,
"branch_name": branch_name,
"worktree_path": worktree_path,
"work_lease": work_lease,
"lock_provenance": issue_lock_provenance.build_sanctioned_lock_provenance(
tool=tool,
source=source,
claimant=claimant,
),
}
if owner_session is not None:
record["owner_session"] = owner_session
if assignment_id is not None:
record["assignment_id"] = assignment_id
# #953 AC6: a null lease id is recorded only when no workflow lease was
# allocated for this bootstrap. The task-session identifier in the
# work_lease is what downstream ownership checks fence on, and it is never
# null on a canonical lock.
if lease_id is not None:
record["lease_id"] = lease_id
if expected_base_sha is not None:
record["expected_base_sha"] = expected_base_sha
return record
def expiration_state(lock: Mapping[str, Any] | None) -> dict[str, Any]:
"""Classify a lock's recorded expiry explicitly (AC12).
Distinguishes "no expiry was ever recorded" from "an expiry was recorded
and is still in the future". Collapsing those two into a single ``False``
from ``is_lease_expired`` is what let a malformed lock be treated as
permanently live and simultaneously never renewable.
"""
if not isinstance(lock, Mapping):
return {"state": EXPIRATION_MISSING, "expires_at": None, "expired": None}
lease = lock.get("work_lease")
raw = lease.get("expires_at") if isinstance(lease, Mapping) else None
text = _text(raw)
if not text:
return {"state": EXPIRATION_MISSING, "expires_at": None, "expired": None}
try:
parsed = datetime.fromisoformat(text.replace("Z", "+00:00")).astimezone(
timezone.utc
)
except ValueError:
return {"state": EXPIRATION_UNPARSEABLE, "expires_at": text, "expired": None}
return {
"state": EXPIRATION_RECORDED,
"expires_at": text,
"expired": parsed <= now_utc(),
}
def missing_contract_fields(lock: Mapping[str, Any] | None) -> list[str]:
"""Name every canonical field a lock does not carry (AC7)."""
if not isinstance(lock, Mapping):
return ["<no lock record>"]
missing: list[str] = []
for field in REQUIRED_LOCK_FIELDS:
value = lock.get(field)
if value is None or (isinstance(value, str) and not value.strip()):
missing.append(field)
lease = lock.get("work_lease")
if not isinstance(lease, Mapping):
if "work_lease" not in missing:
missing.append("work_lease")
else:
for field in REQUIRED_WORK_LEASE_FIELDS:
value = lease.get(field)
if value is None or (isinstance(value, str) and not value.strip()):
missing.append(f"work_lease.{field}")
provenance = lock.get("lock_provenance")
if isinstance(provenance, Mapping):
if (
_text(provenance.get("source"))
not in issue_lock_provenance.SANCTIONED_LOCK_SOURCES
):
missing.append("lock_provenance.source (not sanctioned)")
if not _text(provenance.get("written_by_tool")):
missing.append("lock_provenance.written_by_tool")
claimant = lock_claimant(lock)
if not claimant["username"]:
missing.append("claimant.username")
if not claimant["profile"]:
missing.append("claimant.profile")
return missing
def assess_lock_contract(lock: Mapping[str, Any] | None) -> dict[str, Any]:
"""Structural, read-only verdict on a durable lock record (AC7, AC16).
Pure inspection: it reads the record it is handed and mutates nothing —
no lock, lease, branch, worktree, issue, or PR. Callers use it both to
verify a lock they just wrote and to report on one they found.
"""
if not isinstance(lock, Mapping) or not lock:
return {
"contract": CONTRACT_ABSENT,
"canonical": False,
"missing_fields": ["<no lock record>"],
"claimant": {"username": "", "profile": ""},
"claimant_placement": "absent",
"expiration": {
"state": EXPIRATION_MISSING,
"expires_at": None,
"expired": None,
},
"heartbeatable": False,
"create_pr_eligible": False,
"lock_generation": None,
"task_session_id": None,
"reasons": ["no durable lock record"],
}
missing = missing_contract_fields(lock)
claimant = lock_claimant(lock)
placement = claimant_placement(lock)
expiration = expiration_state(lock)
provenance_check = issue_lock_provenance.assess_lock_file_for_create_pr(dict(lock))
# Canonical means: every required field present, the claimant in the
# canonical placement, an expiry actually recorded, and the untouched #447
# guard satisfied.
canonical = (
not missing
and placement == "work_lease"
and expiration["state"] == EXPIRATION_RECORDED
and bool(provenance_check.get("proven"))
)
if canonical:
contract = CONTRACT_CANONICAL
elif placement == "top_level" and claimant["username"] and claimant["profile"]:
contract = CONTRACT_LEGACY
else:
contract = CONTRACT_INCOMPLETE
reasons: list[str] = []
if missing:
reasons.append("missing canonical fields: " + ", ".join(missing))
if placement == "top_level":
reasons.append(
"claimant recorded at the lock top level rather than in work_lease "
"(legacy/bootstrap placement)"
)
if expiration["state"] == EXPIRATION_MISSING:
reasons.append(
"no expiration recorded; the lock is neither expirable nor renewable "
"until it is upgraded"
)
elif expiration["state"] == EXPIRATION_UNPARSEABLE:
reasons.append(f"unparseable expires_at '{expiration['expires_at']}'")
if provenance_check.get("block"):
reasons.extend(provenance_check.get("reasons") or [])
# Heartbeat needs the claimant pair (from either placement, post-fix) plus a
# task-session identifier to fence on.
lease = lock.get("work_lease")
task_session_id = (
_text(lease.get("task_session_id")) if isinstance(lease, Mapping) else ""
)
heartbeatable = bool(
claimant["username"] and claimant["profile"] and task_session_id
)
return {
"contract": contract,
"canonical": canonical,
"missing_fields": missing,
"claimant": claimant,
"claimant_placement": placement,
"expiration": expiration,
"heartbeatable": heartbeatable,
"create_pr_eligible": bool(provenance_check.get("proven")),
"lock_generation": lock.get("lock_generation"),
"task_session_id": task_session_id or None,
"reasons": reasons,
}
def format_contract_refusal(assessment: Mapping[str, Any]) -> str:
"""Human-readable refusal naming exactly what the lock is missing."""
missing = ", ".join(assessment.get("missing_fields") or []) or "unknown fields"
return (
"Issue lock contract incomplete (#953): "
f"{missing}. The lock cannot be heartbeated, renewed, or accepted by "
"gitea_create_pr in this state (fail closed)"
)
def recommended_action(assessment: Mapping[str, Any]) -> str:
"""The one executable next step for a lock in this state (AC5, AC15)."""
contract = assessment.get("contract")
if contract == CONTRACT_CANONICAL:
return (
"Lock is canonical. Call gitea_whoami, then "
"gitea_resolve_task_capability(task='work_issue'), then proceed with "
"author implementation in the bootstrapped worktree."
)
if contract == CONTRACT_ABSENT:
return (
"No durable lock exists. Call gitea_lock_issue for this issue and "
"branch before writing any implementation bytes."
)
return (
"Do not begin implementation. Call "
"gitea_recover_incomplete_bootstrap_lock for this exact issue, branch, "
"and worktree to upgrade the lock to the canonical contract, or "
"gitea_lock_issue while the worktree is still base-equivalent."
)
# ── Post-compensation recovery guidance (#953 AC5/AC15, review 632 F2) ──
#
# ``recommended_action`` above answers "what can be done about a lock in this
# shape". That is the wrong question on the bootstrap AC7 refusal path, because
# ``run_compensating_recovery`` has already run by the time the answer is
# reported: it releases the lock, removes the worktree when clean — which it
# always is there, no implementation bytes having been written — and deletes the
# created branch. Recommending incomplete-lock recovery for those artifacts
# hands the author two refusals in a row (``no_durable_lock``, then
# ``worktree_invalid``) for a state that a plain bootstrap retry would fix. The
# advice must describe the state that actually *remains*.
#: Compensation removed every artifact this transition created.
CLEANUP_COMPLETE = "complete"
#: Compensation removed some artifacts; others survive and are still actionable.
CLEANUP_PARTIAL = "partial"
#: Compensation itself failed or could not be observed; nothing is provable.
CLEANUP_FAILED = "failed"
def assess_post_compensation_state(
recovery: Mapping[str, Any] | None,
*,
lock_present: bool,
worktree_present: bool,
branch_present: bool,
) -> dict[str, Any]:
"""Classify what survived compensation, from observed durable state.
Pure. The caller observes the filesystem and git; this decides. Observation
is authoritative over the journal's ``rolled_back`` list, which records what
compensation *attempted*: ``run_compensating_recovery`` swallows a failed
lock release and appends nothing, so an absent marker proves nothing either
way. The list is still carried through as corroborating evidence.
The three states are distinct facts, not degrees of the same one:
* ``CLEANUP_COMPLETE`` — compensation ran and nothing it created remains.
* ``CLEANUP_PARTIAL`` — compensation ran and artifacts survive, whether by
design (a worktree dirty at rollback time, a branch carrying commits) or
because a rollback step errored. Either way the surviving set was observed
directly, so it is known and actionable; ``failed_rollback_steps`` records
which cause applies.
* ``CLEANUP_FAILED`` — compensation never ran to completion, so nothing it
would have removed can be assumed removed.
"""
rolled_back = list((recovery or {}).get("rolled_back") or [])
executed = bool((recovery or {}).get("executed"))
failed_steps = [entry for entry in rolled_back if "_failed" in entry]
surviving: list[str] = []
if lock_present:
surviving.append("lock")
if worktree_present:
surviving.append("worktree")
if branch_present:
surviving.append("branch")
if not executed:
state = CLEANUP_FAILED
elif surviving:
state = CLEANUP_PARTIAL
else:
state = CLEANUP_COMPLETE
return {
"cleanup_state": state,
"compensation_executed": executed,
"lock_present": bool(lock_present),
"worktree_present": bool(worktree_present),
"branch_present": bool(branch_present),
"surviving_artifacts": surviving,
"removed_artifacts": [
name
for name, present in (
("lock", lock_present),
("worktree", worktree_present),
("branch", branch_present),
)
if not present
],
"failed_rollback_steps": failed_steps,
"rolled_back": rolled_back,
}
def post_compensation_action(
state: Mapping[str, Any],
*,
issue_number: int,
branch_name: str,
worktree_path: str,
missing_fields: list[str] | None = None,
) -> str:
"""The one executable next step for the state compensation actually left.
Every branch names only artifacts the classification says still exist, so no
recommendation can point at something the rollback deleted.
"""
missing = ", ".join(missing_fields or []) or "the reported missing fields"
cleanup_state = state.get("cleanup_state")
lock_present = bool(state.get("lock_present"))
worktree_present = bool(state.get("worktree_present"))
branch_present = bool(state.get("branch_present"))
if not state.get("compensation_executed"):
# Compensation never ran, so nothing was rolled back and nothing about
# the remaining state was decided. The read-only surface is the only
# action executable under any state.
return (
"Compensating rollback did not complete, so the remaining state is "
f"not proven. Call gitea_inspect_issue_lock_contract for issue "
f"#{issue_number} (read-only) to establish what survives before any "
"further action. Do not retry bootstrap until it is known."
)
prefix = ""
failed_steps = state.get("failed_rollback_steps") or []
if failed_steps:
prefix = (
"Compensating rollback reported a failed step "
f"({', '.join(failed_steps)}); what survives was observed directly "
"and the action below is scoped to exactly that. "
)
if cleanup_state == CLEANUP_COMPLETE:
return (
"Compensating rollback removed the malformed lock, the branch, and "
f"the worktree, so nothing from this attempt remains. Resolve "
f"{missing} and re-run gitea_bootstrap_author_issue_worktree for "
f"issue #{issue_number} from the clean pre-bootstrap state. Do not "
"call gitea_recover_incomplete_bootstrap_lock: there is no lock, "
"branch, or worktree left for it to act on."
)
if lock_present and worktree_present and branch_present:
return prefix + (
"The lock, branch, and worktree all survive. Call "
"gitea_recover_incomplete_bootstrap_lock for issue "
f"#{issue_number}, branch '{branch_name}', and worktree "
f"'{worktree_path}', passing the worktree's current head as "
"expected_head, to upgrade the lock to the canonical contract."
)
if not lock_present and worktree_present and branch_present:
return prefix + (
"The malformed lock was released but the branch and worktree "
"survive. No implementation bytes were written, so the worktree is "
f"still base-equivalent: call gitea_lock_issue for issue "
f"#{issue_number} on branch '{branch_name}' from worktree "
f"'{worktree_path}' to acquire a canonical lock."
)
if lock_present and not worktree_present:
return prefix + (
f"The worktree for issue #{issue_number} is gone but the durable "
"lock survived, so neither gitea_recover_incomplete_bootstrap_lock "
"(it would refuse worktree_invalid) nor gitea_lock_issue (it has no "
"worktree to bind) is executable. Call "
"gitea_inspect_issue_lock_contract for issue "
f"#{issue_number} (read-only) to confirm the surviving lock; it "
"must be released by its recorded owner before bootstrap is "
"retried."
)
# Lock gone, worktree gone, some git artifact left (a branch with commits,
# or a branch this transition did not create).
return prefix + (
"Compensating rollback removed the lock and worktree; branch "
f"'{branch_name}' survives and was not deleted. Call "
f"gitea_inspect_issue_lock_contract for issue #{issue_number} "
"(read-only) to confirm no durable lock remains, then re-run "
"gitea_bootstrap_author_issue_worktree, which will adopt the existing "
"branch rather than recreating it."
)
+304
View File
@@ -0,0 +1,304 @@
"""Target-specific recovery for incomplete bootstrap issue locks (#953).
The situation this exists for: ``gitea_bootstrap_author_issue_worktree``
reported success, wrote an incomplete lock, and told the author to implement.
The author did — legitimately, following the tool's own reported next action —
and the branch now carries real committed and pushed work. At that point every
pre-existing recovery path is simultaneously ineligible:
* heartbeat refuses, because the claimant is not where it looks;
* ``gitea_lock_issue`` refuses, because the branch is no longer base-equivalent;
* #760 exact-owner renewal never engages, because a lock with no recorded
expiry is never *expired*;
* the #447 create-PR guard refuses, because there is no provenance.
Distinct from every neighbouring path: #753 ``issue_lock_recovery`` requires a
dead owner PID, #760 ``issue_lock_renewal`` requires an *expired* lease, and
#442 ``issue_lock_adoption`` decides branch adoption. None of them addresses a
lock that is structurally incomplete and therefore never expires at all.
**What this will not do.** It never moves, resets, or rewinds a branch, and
never requires base-equivalence — the committed work is the thing being
preserved. It never pushes and never opens a pull request. It touches only the
one lock file named by (remote, org, repo, issue). It accepts no caller-supplied
provenance and no caller-supplied authorization flag; both are minted
server-side. It refuses a healthy foreign-owned lock outright, and a matching
username alone is never accepted as proof of ownership — the profile must match
too, and the lock's recorded binding must agree with the observed branch,
worktree, and head.
"""
from __future__ import annotations
import os
from typing import Any, Mapping
import author_lock_contract
import issue_lock_store
#: Refusal codes, so callers can branch on cause rather than parse prose.
REFUSAL_NO_LOCK = "no_durable_lock"
REFUSAL_ALREADY_CANONICAL = "already_canonical"
REFUSAL_FOREIGN_CLAIMANT = "foreign_claimant"
REFUSAL_HEALTHY_FOREIGN = "healthy_foreign_lock"
REFUSAL_IDENTITY_UNRESOLVED = "identity_unresolved"
REFUSAL_BINDING_MISMATCH = "binding_mismatch"
REFUSAL_WORKTREE_INVALID = "worktree_invalid"
REFUSAL_HEAD_MISMATCH = "head_mismatch"
def _text(value: Any) -> str:
return str(value or "").strip()
def _same_realpath(left: str | None, right: str | None) -> bool:
lhs, rhs = _text(left), _text(right)
if not lhs or not rhs:
return False
try:
return os.path.realpath(lhs) == os.path.realpath(rhs)
except OSError:
return lhs == rhs
def assess_bootstrap_lock_recovery(
existing_lock: Mapping[str, Any] | None,
*,
issue_number: int,
branch_name: str,
worktree_path: str,
remote: str,
org: str,
repo: str,
identity: str | None,
profile: str | None,
observed_head: str | None,
declared_head: str | None,
worktree_exists: bool,
worktree_registered: bool,
current_branch: str | None,
now: Any = None,
) -> dict[str, Any]:
"""Decide whether this exact lock may be upgraded by this exact caller.
Pure: every input is an observation the caller already made, and nothing
here reads or writes the filesystem, git, or Gitea. That is what makes the
same decision testable in isolation and reusable by the read-only
inspection surface, which must not mutate anything (AC16).
Returns a dict with ``recovery_sanctioned`` plus the full evidence set. A
refusal never raises — it reports, so the caller can surface exactly which
piece of evidence was missing.
"""
reasons: list[str] = []
refusal_code: str | None = None
contract = author_lock_contract.assess_lock_contract(existing_lock)
if not existing_lock:
return {
"recovery_sanctioned": False,
"refusal_code": REFUSAL_NO_LOCK,
"reasons": [
f"no durable issue lock exists for issue #{issue_number}; there is "
"nothing to recover (fail closed)"
],
"contract": contract,
"evidence": {},
"expected_generation": None,
}
active_identity = _text(identity)
active_profile = _text(profile)
recorded = author_lock_contract.lock_claimant(existing_lock)
freshness = issue_lock_store.assess_lock_freshness(dict(existing_lock), now=now)
generation = issue_lock_store.lock_generation(existing_lock)
evidence: dict[str, Any] = {
"recorded_claimant": recorded,
"active_identity": active_identity,
"active_profile": active_profile,
"recorded_branch": existing_lock.get("branch_name"),
"recorded_worktree": existing_lock.get("worktree_path"),
"recorded_owner_session": existing_lock.get("owner_session"),
"recorded_generation": generation,
"recorded_remote": existing_lock.get("remote"),
"recorded_org": existing_lock.get("org"),
"recorded_repo": existing_lock.get("repo"),
"observed_head": _text(observed_head),
"declared_head": _text(declared_head),
"current_branch": _text(current_branch),
"worktree_exists": bool(worktree_exists),
"worktree_registered": bool(worktree_registered),
"freshness": freshness,
"claimant_placement": contract.get("claimant_placement"),
"expiration_state": contract.get("expiration", {}).get("state"),
}
# ── Repository and issue identity (AC10) ──
if _text(existing_lock.get("remote")) != _text(remote):
reasons.append(
f"recorded remote '{existing_lock.get('remote')}' does not match '{remote}'"
)
refusal_code = refusal_code or REFUSAL_BINDING_MISMATCH
if _text(existing_lock.get("org")) != _text(org):
reasons.append(
f"recorded org '{existing_lock.get('org')}' does not match '{org}'"
)
refusal_code = refusal_code or REFUSAL_BINDING_MISMATCH
if _text(existing_lock.get("repo")) != _text(repo):
reasons.append(
f"recorded repo '{existing_lock.get('repo')}' does not match '{repo}'"
)
refusal_code = refusal_code or REFUSAL_BINDING_MISMATCH
if existing_lock.get("issue_number") != issue_number:
reasons.append(
f"lock targets issue #{existing_lock.get('issue_number')}, not "
f"#{issue_number}"
)
refusal_code = refusal_code or REFUSAL_BINDING_MISMATCH
# ── Branch and worktree binding (AC10) ──
if _text(existing_lock.get("branch_name")) != _text(branch_name):
reasons.append(
f"recorded branch '{existing_lock.get('branch_name')}' does not match "
f"'{branch_name}'"
)
refusal_code = refusal_code or REFUSAL_BINDING_MISMATCH
if not _same_realpath(existing_lock.get("worktree_path"), worktree_path):
reasons.append(
f"recorded worktree '{existing_lock.get('worktree_path')}' does not "
f"match '{worktree_path}'"
)
refusal_code = refusal_code or REFUSAL_BINDING_MISMATCH
# ── The worktree is real, registered, and on the branch (AC10) ──
# Deliberately no base-equivalence requirement and no constraint on how far
# the branch has advanced: the whole point is that it already carries the
# author's legitimate commits (AC9).
if not worktree_exists:
reasons.append(f"declared worktree '{worktree_path}' does not exist")
refusal_code = refusal_code or REFUSAL_WORKTREE_INVALID
if not worktree_registered:
reasons.append(f"worktree '{worktree_path}' is not a registered git worktree")
refusal_code = refusal_code or REFUSAL_WORKTREE_INVALID
if _text(current_branch) != _text(branch_name):
reasons.append(
f"worktree is on branch '{_text(current_branch) or 'unknown'}', not "
f"'{branch_name}'"
)
refusal_code = refusal_code or REFUSAL_WORKTREE_INVALID
# ── Current head fencing (AC10) ──
# The caller names the commit it believes it is recovering. A mismatch means
# the worktree moved under the caller, so the decision is stale.
if not _text(observed_head):
reasons.append("could not observe the worktree head")
refusal_code = refusal_code or REFUSAL_HEAD_MISMATCH
elif _text(declared_head) and _text(declared_head) != _text(observed_head):
reasons.append(
f"declared head '{_text(declared_head)}' does not match observed head "
f"'{_text(observed_head)}'"
)
refusal_code = refusal_code or REFUSAL_HEAD_MISMATCH
# ── Ownership (AC10, AC11) ──
# A matching username alone is never sufficient: the profile must match too,
# and both are compared against server-resolved values the caller cannot set.
if not active_identity or not active_profile:
reasons.append(
"active identity and profile could not both be resolved; ownership "
"cannot be proven"
)
refusal_code = refusal_code or REFUSAL_IDENTITY_UNRESOLVED
if not recorded["username"] or not recorded["profile"]:
reasons.append(
"durable lock does not record both a claimant username and profile"
)
refusal_code = refusal_code or REFUSAL_FOREIGN_CLAIMANT
elif (
recorded["username"] != active_identity
or recorded["profile"] != active_profile
):
# AC11: a foreign-owned lock is never recoverable through this path,
# healthy or not. The healthy case is reported distinctly so the refusal
# is legible, but both refuse.
if freshness.get("live"):
reasons.append(
f"lock is owned by a healthy foreign claimant "
f"'{recorded['username']}/{recorded['profile']}'; takeover is not "
"a recovery path"
)
refusal_code = REFUSAL_HEALTHY_FOREIGN
else:
reasons.append(
f"lock claimant '{recorded['username']}/{recorded['profile']}' "
f"does not match active '{active_identity}/{active_profile}'"
)
refusal_code = refusal_code or REFUSAL_FOREIGN_CLAIMANT
# ── Nothing to recover ──
# A lock that is already canonical is left strictly alone. Rewriting it would
# mint a new task-session identifier and invalidate the heartbeat token the
# legitimate owner is already using.
if contract.get("canonical") and not reasons:
return {
"recovery_sanctioned": False,
"refusal_code": REFUSAL_ALREADY_CANONICAL,
"reasons": [
"lock already satisfies the canonical contract; no recovery is "
"required"
],
"contract": contract,
"evidence": evidence,
"expected_generation": generation,
}
sanctioned = not reasons
return {
"recovery_sanctioned": sanctioned,
"refusal_code": None if sanctioned else refusal_code,
"reasons": reasons,
"contract": contract,
"evidence": evidence,
"expected_generation": generation,
}
def build_recovery_record(
assessment: Mapping[str, Any],
*,
recovered_at: str,
new_task_session_id: str,
) -> dict[str, Any]:
"""Auditable record of the ownership and generation transition (AC10).
A recovered lock must never read as an original claim, so both sides of the
transition are preserved: what the incomplete lock recorded, and what
replaced it.
"""
evidence = dict(assessment.get("evidence") or {})
contract = dict(assessment.get("contract") or {})
return {
"recovery_kind": "incomplete_bootstrap_lock",
"recovered_at": recovered_at,
"prior_contract": contract.get("contract"),
"prior_missing_fields": list(contract.get("missing_fields") or []),
"prior_claimant_placement": evidence.get("claimant_placement"),
"prior_expiration_state": evidence.get("expiration_state"),
"prior_generation": evidence.get("recorded_generation"),
"prior_owner_session": evidence.get("recorded_owner_session"),
"prior_freshness": (evidence.get("freshness") or {}).get("status"),
"replacement_task_session_id": new_task_session_id,
"preserved_head": evidence.get("observed_head"),
"branch_reset": False,
"base_equivalence_required": False,
}
def format_recovery_refusal(assessment: Mapping[str, Any]) -> str:
reasons = "; ".join(
assessment.get("reasons") or ["unknown bootstrap lock recovery refusal"]
)
code = assessment.get("refusal_code") or "refused"
return f"Bootstrap lock recovery refused ({code}): {reasons} (fail closed)"
+1 -164
View File
@@ -32,7 +32,7 @@ from typing import Any, Iterator, Sequence
import dependency_graph import dependency_graph
import gitea_audit import gitea_audit
SCHEMA_VERSION = 6 SCHEMA_VERSION = 5
# Assignable work kinds only — raw monitoring incidents are never work items. # Assignable work kinds only — raw monitoring incidents are never work items.
WORK_KINDS = frozenset({"issue", "pr"}) WORK_KINDS = frozenset({"issue", "pr"})
@@ -65,33 +65,6 @@ CREATE TABLE IF NOT EXISTS sessions (
status TEXT NOT NULL DEFAULT 'active' status TEXT NOT NULL DEFAULT 'active'
); );
-- Fleet-level MCP server runtime registry (#949). Distinct from ``sessions``:
-- a session is one allocator *task*, while a row here is one *server process*
-- that bound native MCP transport. Each server writes its own row and no other,
-- so a caller can never supply this evidence. Creating the table is itself the
-- v5→v6 migration: additive, idempotent, and it touches no existing table.
CREATE TABLE IF NOT EXISTS mcp_server_runtimes (
runtime_id TEXT PRIMARY KEY,
namespace TEXT NOT NULL,
profile TEXT,
role TEXT,
remote TEXT,
org TEXT,
repo TEXT,
repository_root TEXT,
pid INTEGER NOT NULL,
cohort_id TEXT,
cohort_source TEXT,
client_provenance TEXT,
boot_id TEXT,
startup_head TEXT,
daemon_start_head TEXT,
transport TEXT,
registered_at TEXT NOT NULL,
last_heartbeat_at TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'running'
);
CREATE TABLE IF NOT EXISTS work_items ( CREATE TABLE IF NOT EXISTS work_items (
work_item_id INTEGER PRIMARY KEY AUTOINCREMENT, work_item_id INTEGER PRIMARY KEY AUTOINCREMENT,
remote TEXT NOT NULL, remote TEXT NOT NULL,
@@ -937,142 +910,6 @@ class ControlPlaneDB:
rows = conn.execute(sql, params).fetchall() rows = conn.execute(sql, params).fetchall()
return [dict(r) for r in rows] return [dict(r) for r in rows]
# ── MCP server runtime registry (#949) ────────────────────────────────
_RUNTIME_COLUMNS: tuple[str, ...] = (
"runtime_id",
"namespace",
"profile",
"role",
"remote",
"org",
"repo",
"repository_root",
"pid",
"cohort_id",
"cohort_source",
"client_provenance",
"boot_id",
"startup_head",
"daemon_start_head",
"transport",
"registered_at",
"last_heartbeat_at",
"status",
)
def register_mcp_server_runtime(
self,
record: dict[str, Any],
*,
retention_seconds: int | None = None,
now: str | None = None,
) -> dict[str, Any]:
"""Record that *this* server process is running (#949).
Only the described process ever calls this, at native transport bind.
Two housekeeping deletions happen here — on the startup write path, never
on the read path — so the registry cannot grow without bound:
* rows that recorded the **same PID**, which the calling process now
owns and which therefore cannot still be a live server;
* rows older than *retention_seconds*.
Neither can hide a live duplicate: a concurrently running server holds a
different PID and writes a fresher row.
"""
stamp = now or _ts()
row = {key: record.get(key) for key in self._RUNTIME_COLUMNS}
row["runtime_id"] = str(record.get("runtime_id") or "").strip()
if not row["runtime_id"]:
raise ValueError("runtime_id is required (fail closed)")
row["namespace"] = str(record.get("namespace") or "").strip()
if not row["namespace"]:
raise ValueError("namespace is required (fail closed)")
if record.get("pid") is None:
raise ValueError("pid is required (fail closed)")
row["pid"] = int(record["pid"])
row["registered_at"] = record.get("registered_at") or stamp
row["last_heartbeat_at"] = record.get("last_heartbeat_at") or stamp
row["status"] = record.get("status") or "running"
columns = ", ".join(self._RUNTIME_COLUMNS)
placeholders = ", ".join("?" for _ in self._RUNTIME_COLUMNS)
values = [row[key] for key in self._RUNTIME_COLUMNS]
with self._tx() as conn:
conn.execute(
"DELETE FROM mcp_server_runtimes WHERE pid = ? AND runtime_id != ?",
(row["pid"], row["runtime_id"]),
)
if retention_seconds and retention_seconds > 0:
cutoff = _ts(
_utc_now() - timedelta(seconds=int(retention_seconds))
)
conn.execute(
"DELETE FROM mcp_server_runtimes WHERE registered_at < ?",
(cutoff,),
)
conn.execute(
f"INSERT OR REPLACE INTO mcp_server_runtimes({columns}) "
f"VALUES ({placeholders})",
values,
)
stored = conn.execute(
"SELECT * FROM mcp_server_runtimes WHERE runtime_id = ?",
(row["runtime_id"],),
).fetchone()
return dict(stored)
def heartbeat_mcp_server_runtime(self, runtime_id: str) -> None:
"""Refresh a runtime row's liveness timestamp.
Never called by the inventory read path: liveness is established from
process evidence so that reading the fleet mutates nothing (#949 AC11).
"""
with self._tx() as conn:
conn.execute(
"UPDATE mcp_server_runtimes SET last_heartbeat_at = ? "
"WHERE runtime_id = ?",
(_ts(), runtime_id),
)
def mark_mcp_server_runtime_stopped(self, runtime_id: str) -> None:
"""Mark a runtime row stopped (graceful shutdown bookkeeping)."""
with self._tx() as conn:
conn.execute(
"UPDATE mcp_server_runtimes SET status = 'stopped', "
"last_heartbeat_at = ? WHERE runtime_id = ?",
(_ts(), runtime_id),
)
def list_mcp_server_runtimes(
self,
*,
statuses: Sequence[str] | None = None,
limit: int = 500,
) -> list[dict[str, Any]]:
"""List runtime registry rows, newest registration first (read-only).
Ordering is fully deterministic: ties on ``registered_at`` break on
``runtime_id`` so two callers reading one snapshot see one order.
"""
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 ""
params.append(max(1, int(limit)))
with self._tx(immediate=False) as conn:
rows = conn.execute(
f"SELECT * FROM mcp_server_runtimes {where} "
"ORDER BY registered_at DESC, runtime_id ASC LIMIT ?",
params,
).fetchall()
return [dict(r) for r in rows]
# ── work items ──────────────────────────────────────────────────────── # ── work items ────────────────────────────────────────────────────────
def upsert_work_item( def upsert_work_item(
+209
View File
@@ -0,0 +1,209 @@
# The canonical author issue-lock contract (#953)
Every author issue lock has exactly one shape. Both writers —
`gitea_bootstrap_author_issue_worktree` and `gitea_lock_issue` — build it
through `author_lock_contract.build_canonical_issue_lock`, and every reader
consumes that same shape.
Before #953 the two writers disagreed. `gitea_lock_issue` wrote the canonical
record; bootstrap wrote a thinner one with the claimant at the lock top level,
`lease_id: null`, and no `work_lease`, `lock_provenance`, or expiry. Because
every reader was written against the canonical shape, a lock that bootstrap
reported as successfully created could not be heartbeated, renewed, re-locked,
or accepted by `gitea_create_pr`. Each of those gates was individually correct;
the defect was that two writers disagreed about what a lock *is*.
## Required ordering
**Finalize the lock before writing any implementation bytes.** This ordering is
what keeps recovery cheap: while the worktree is still base-equivalent, a lock
problem can be fixed by simply calling `gitea_lock_issue` again. Once the branch
carries commits, base-equivalence is gone and the ordinary re-lock path is no
longer available.
1. `gitea_whoami` — resolve identity and profile.
2. `gitea_resolve_task_capability(task='work_issue')`.
3. `gitea_bootstrap_author_issue_worktree` — creates the branch, the registered
worktree under `branches/`, and a **canonical** lock. It reads the lock back
and verifies it structurally before reporting success; a partial lock fails
closed here, with the missing fields named, and never reports
`implementation_allowed: true`.
4. `gitea_heartbeat_issue_lock` — prove the lock is usable, using the
`task_session_id` bootstrap returned.
5. Implement, commit, push.
6. `gitea_create_pr`.
If bootstrap returns `success: false` with
`reason_code: incomplete_issue_lock_contract`, **do not implement**. Its
`exact_next_action` names the executable recovery step. Bootstrap's reported
next action always matches the state it actually returned.
### What that refusal leaves behind
The AC7 refusal runs `run_compensating_recovery` *before* it reports, so the
advice has to describe the post-rollback state rather than the shape of the lock
that provoked it. Recommending incomplete-lock recovery for artifacts the
rollback already deleted would produce `no_durable_lock` and then
`worktree_invalid` — two refusals for a state a plain retry fixes.
The refusal therefore carries `compensating_recovery` and
`post_compensation_state`, and derives `exact_next_action` from what was
observed on disk. `cleanup_state` is one of:
| `cleanup_state` | Meaning | Next action |
| --- | --- | --- |
| `complete` | lock, branch, and worktree all removed | resolve `missing_fields` and re-run `gitea_bootstrap_author_issue_worktree` |
| `partial` | rollback ran; some artifacts survive, by design or because a step errored | scoped to exactly what survives — see below |
| `failed` | rollback never completed, so nothing is proven removed | `gitea_inspect_issue_lock_contract` (read-only) before anything else |
Within `partial`, the surviving set decides the action:
| Survives | Next action |
| --- | --- |
| lock + branch + worktree | `gitea_recover_incomplete_bootstrap_lock` for that exact issue, branch, and worktree |
| branch + worktree (lock released) | `gitea_lock_issue` — no implementation bytes were written, so the worktree is still base-equivalent |
| lock only (worktree removed) | `gitea_inspect_issue_lock_contract`; the surviving lock must be released by its recorded owner before bootstrap is retried |
| branch only | `gitea_inspect_issue_lock_contract`, then re-run bootstrap, which adopts the existing branch |
`failed_rollback_steps` names any rollback step that errored, and the returned
action says so rather than presenting the surviving state as intentional.
> The lock half of that rollback was dead code until #953 review 632 F2:
> `run_compensating_recovery` called `issue_lock_store.release_session_lock`,
> which did not exist, inside a bare `except Exception: pass`. Every rollback
> removed the branch and worktree and silently left the lock — the exact
> uninspectable, unrecoverable state this issue exists to eliminate. The
> function now exists, releases only a lock whose recorded `owner_session`
> matches, and its failures are recorded rather than swallowed.
## The contract
A canonical lock carries every field in
`author_lock_contract.REQUIRED_LOCK_FIELDS`:
| Field | Meaning |
| --- | --- |
| `remote`, `org`, `repo`, `issue_number` | repository and issue identity |
| `branch_name`, `worktree_path` | the binding this claim owns |
| `work_lease` | the canonical lease block, below |
| `lock_provenance` | sanctioned source, minted server-side |
| `lock_generation` | monotonic; every write advances it |
`work_lease` carries every field in
`author_lock_contract.REQUIRED_WORK_LEASE_FIELDS`, notably:
| Field | Meaning |
| --- | --- |
| `claimant.{username,profile}` | **canonical** claimant placement |
| `expires_at` | sliding TTL from `lease_policy` |
| `last_heartbeat_at`, `heartbeat_count` | liveness evidence |
| `task_session_id` | the ownership fencing token — never null |
| `lifecycle_version` | `heartbeat-v1`; its absence is what makes a lock legacy |
### Claimant placement and legacy compatibility
`work_lease.claimant` is canonical. A top-level `claimant` is the legacy
placement written by pre-#953 bootstrap and is still **read** — through the one
shared reader, `issue_lock_store.lock_claimant` — so an existing lock is not
refused for "not recording a claimant" when it plainly records one.
Tolerating the placement is not a widening. Every caller still compares the
values against server-resolved identity and profile, so a legacy placement
grants nothing the canonical placement would not. When both are present, the
`work_lease` copy wins: after an upgrade, a stale top-level copy must never
decide ownership.
### Expiration is explicit
A lock with no recorded expiry is **not** "not yet expired". `is_lease_expired`
returns `False` for it, which used to make such a lock permanently non-expiring
*and* permanently ineligible for #760 exact-owner renewal, which only ever
assesses an expired lease. `author_lock_contract.expiration_state` names the
real fact: `recorded`, `missing`, or `unparseable`. A `missing` expiry makes the
lock eligible for the recovery path below rather than stranding it.
## Recovering an existing incomplete bootstrap lock
For locks already written by the old bootstrap — including those whose branches
already carry legitimate committed and pushed work — use:
```text
gitea_inspect_issue_lock_contract(issue_number, branch_name, worktree_path, remote=...)
gitea_recover_incomplete_bootstrap_lock(issue_number, branch_name, worktree_path, expected_head, remote=...)
```
`gitea_inspect_issue_lock_contract` is strictly read-only: it performs no lock,
lease, branch, worktree, issue, or pull-request mutation. Use it first to see
which fields are missing and what the recommended action is; pass `dry_run=True`
to the recovery tool to preview the decision without writing.
`gitea_recover_incomplete_bootstrap_lock` upgrades that one lock to the
canonical contract. Before writing anything it verifies:
* repository (`remote`, `org`, `repo`) and issue number
* claimant username **and** profile against the server-resolved values — a
matching username alone is never accepted
* branch, worktree path, worktree existence, and worktree registration
* the worktree is on the recorded branch
* the observed head equals the caller's `expected_head`
* the existing lock's generation and provenance state
* the absence of healthy foreign ownership
What it deliberately does **not** do:
* it never moves, resets, or rewinds the branch, and never requires
base-equivalence — preserving the committed work is the entire point;
* it never pushes and never creates a pull request;
* it touches only the single lock file for that exact remote/org/repo/issue;
* it accepts no caller-supplied provenance and no caller-supplied authorization
flag — both are minted server-side.
A recovered lock records a `bootstrap_lock_recovery` block holding both sides of
the transition — prior contract, prior missing fields, prior generation, prior
owning session, the replacement `task_session_id`, and the preserved head — so a
recovered claim never reads as an original one.
### Gates, in order
`gitea_recover_incomplete_bootstrap_lock` is an author-only durable-lock
mutation and carries the same three gates as every comparable author operation,
in this order:
1. `role_session_router.check_author_mutation_after_reviewer_stop` — no author
fallback after a reviewer `wrong_role_stop`.
2. `_namespace_mutation_block(task, remote=remote, author_role_exclusive=True)`
the namespace wall. It refuses a reviewer-bound session and, because this
task's required permission is `gitea.issue.comment` (which merger,
controller, and reconciler profiles also hold), additionally requires the
active profile's derived role kind to be exactly `author`. A refusal carries
`namespace_block: true` and emits a `BLOCKED` audit record naming the
namespace and profile.
3. `_profile_permission_block` — operation, provenance, and session-context
gates.
Exact-owner claimant matching inside `assess_bootstrap_lock_recovery` runs
*after* all three. It is a further layer, never a substitute for them: on its
own it refuses one step too late and leaves the audit trail silent about the
attempt.
### Refusals
| `refusal_code` | Meaning |
| --- | --- |
| `no_durable_lock` | nothing to recover |
| `already_canonical` | lock is fine; rewriting would invalidate a live heartbeat token |
| `foreign_claimant` | recorded claimant is not the active identity/profile pair |
| `healthy_foreign_lock` | a live foreign-owned lock; takeover is not a recovery path |
| `identity_unresolved` | identity or profile could not be resolved |
| `binding_mismatch` | repository, issue, branch, or worktree does not match |
| `worktree_invalid` | worktree missing, unregistered, or on another branch |
| `head_mismatch` | the worktree moved under the caller |
## The #447 create-PR provenance guard is unchanged
`issue_lock_provenance.assess_lock_file_for_create_pr` still requires both a
sanctioned `lock_provenance` and a `work_lease`, and the sanctioned source set
was **not** widened. Bootstrap writes through
`issue_lock_provenance.SOURCE_LOCK_ISSUE` — the lock it produces *is* a
canonical lock, not a second dialect with its own exemption. Bootstrap now
satisfies the guard rather than the guard being relaxed to admit bootstrap.
-134
View File
@@ -1,134 +0,0 @@
# Authoritative MCP fleet inventory (#949)
`gitea_assess_fleet_inventory` is the read-only native capability that answers a
question no other surface could: **is exactly one server running for each
configured PRGS profile, and do they all belong to one client cohort?**
## Why the existing surfaces were not enough
| Surface | What it proves | Why it cannot prove the fleet |
| --- | --- | --- |
| `gitea_get_runtime_context` | Profile, identity, workspace binding of **the process answering the call** | Says nothing about the other four namespaces |
| `gitea_assess_master_parity` | Startup vs current vs live revision of **that same process** | Five self-reports of one revision do not establish five processes, nor the absence of a sixth |
| `gitea_assess_mcp_namespace_health` | Whether one named namespace can invoke one tool | Accepts `process`, `probe_result` and `registered_tools` **from the caller** — a capability whose inputs come from the party it constrains is not evidence |
| control-plane `sessions` table | Allocator **task** sessions | A session is a unit of work, not a server process; nothing recorded that a server exists |
Before this capability, satisfying a strict five-process/single-cohort gate
required shell process inspection, cached JSON, or source reading — none of
which are sanctioned workflow evidence.
## Evidence model
Two independent sources must agree before a member counts as running.
**1. Control-plane runtime registry** — table `mcp_server_runtimes`.
Each server writes exactly one row *about itself*, from the official entrypoint,
immediately after the native MCP transport bind. Only a transport-bound process
reaches that line, and no MCP caller can reach it at all. The row is
authoritative for identity: namespace, profile, role, repository binding,
cohort, startup revision, transport, and PID.
**2. Server-side process observation** — a process listing performed by the
server answering the inventory call, never by the caller. It is authoritative
for existence and liveness, and it is the only source that can reveal a running
server the registry does not know about.
A member is `live` only when a registry row has a matching running process whose
start time precedes the registration — so a recycled PID cannot impersonate a
server that has since exited.
### Deliberate non-inferences
* **Configuration is not existence.** A configured profile with no live
corroborated row is `missing`, never `running` (AC8).
* **Matching revisions are not a cohort.** `single_cohort` is derived only from
recorded cohort identity. Five members at one revision with an unknown cohort
yield `single_cohort: null` and a closed gate, never `true` (AC7).
* **Parent-client status is not member health.** Each member is classified from
its own evidence.
Cohort identity comes from `GITEA_MCP_CLIENT_COHORT_ID` when the client sets it,
otherwise from the parent process that launched the server. A server whose
parent has gone away (reparented to init) reports an unknown cohort rather than
guessing.
## Result shape
Top-level verdict fields:
| Field | Meaning |
| --- | --- |
| `inventory_complete` | All evidence was obtainable. False whenever anything below is unknown. |
| `incomplete_reasons` | Every distinct reason completeness failed. |
| `configured_members` | The five expected members with their instance counts and health. |
| `running_members` | Live, corroborated instances of expected profiles. |
| `missing_members` | Expected profiles with no live instance. |
| `duplicate_members` | Expected profiles with more than one live instance, with every PID. |
| `unexpected_members` | Live members outside the expected roster. Never folded into duplicates. |
| `stale_members` | Registry rows whose process is dead, unobserved, PID-recycled, or unknown. |
| `unregistered_processes` | Running MCP server processes with no registry row. |
| `repository_binding_mismatches` | Live members bound to another repository, or with an incomplete binding. |
| `role_mismatches` | Live members whose declared role or namespace contradicts the configured profile. |
| `single_cohort` / `mixed_cohort` | `true`/`false`, or `null` when cohort evidence is unknown. |
| `mixed_revision` / `startup_revisions` | Revision spread across live members. |
| `exactly_one_per_profile` | No missing and no duplicate expected members. |
| `mutation_gate_satisfied` | The full invariant held. |
| `blocked_reason` / `blocked_reasons` | Why the gate is closed; the first is the headline. |
| `mutations_performed` | Always `[]`. |
Ordering is deterministic — every list sorts by namespace, profile, PID, then
runtime id — so two callers reading one snapshot see identical structure.
## How controller and reconciler gates consume it
Classification is a **pure function of the snapshot**. The answering namespace is
reported as `answering_namespace` metadata and never affects the verdict, so
`gitea-controller` and `gitea-reconciler` return the same result for the same
fleet. Neither is privileged over the other.
Consume it like this:
1. Call `gitea_assess_fleet_inventory` from `gitea-controller` **or**
`gitea-reconciler` (any namespace holding `gitea.read` may call it).
2. If `mutation_gate_satisfied` is `true`, the exact-one-instance-per-profile,
single-cohort, single-revision invariant is proven; proceed.
3. Otherwise **stop and report `blocked_reason` verbatim**. Distinguish the
cases — they need different operator actions:
* `missing_members` — the client did not launch that namespace; reconnect it.
* `duplicate_members` — a second client or a manual launch is running that
profile; the **operator** quits the extra client. This capability never
terminates a process.
* `unexpected_members` — an unconfigured PRGS server is live; investigate
before trusting any gate.
* `unregistered_processes` — a server is running code that predates this
capability, or failed to register; the inventory is incomplete by
construction and must not be reported as healthy.
* `mixed_cohort` / `mixed_revision` — the fleet is not one coherent unit.
* `single_cohort: null` — cohort evidence is missing; this is *unknown*, not
*healthy*.
Never treat `inventory_complete: false` as a soft warning. It means the
inventory cannot describe the whole fleet, which is exactly the state the
mutation gate exists to refuse.
## Read-only guarantees
The capability performs no restart, reconnect, drain, lease mutation, issue
mutation, repository write, or process termination, and it never manufactures
historical restart evidence. The only signal it sends is `signal 0`, a
liveness/permission check that delivers nothing to the target process. Registry
writes happen exclusively on the startup path of the process being described,
never on this read path.
## Scope boundaries
This capability is observation only. Related concerns live elsewhere and are
deliberately not absorbed here:
* **#950** — controller role metadata and capability-routing consistency.
* **#951** — durable synchronization and restart receipts.
* **#952** — stale-lease inspection, dashboard, and executor consistency.
* **#900** — cohort lifecycle supervision (draining and reaping superseded
cohorts) is a *mutating* transition that would consume this evidence.
* **#948** — the client/session/generation ownership model that this evidence
feeds.
+2 -1
View File
@@ -54,7 +54,6 @@ that gates each call, not which tools exist.
- `gitea_assess_already_landed_reconciliation` - `gitea_assess_already_landed_reconciliation`
- `gitea_assess_conflict_fix_classification` - `gitea_assess_conflict_fix_classification`
- `gitea_assess_conflict_fix_push` - `gitea_assess_conflict_fix_push`
- `gitea_assess_fleet_inventory`
- `gitea_assess_gitea_operation_path` - `gitea_assess_gitea_operation_path`
- `gitea_assess_master_parity` - `gitea_assess_master_parity`
- `gitea_assess_mcp_namespace_health` - `gitea_assess_mcp_namespace_health`
@@ -104,6 +103,7 @@ that gates each call, not which tools exist.
- `gitea_get_shell_health` - `gitea_get_shell_health`
- `gitea_heartbeat_issue_lock` - `gitea_heartbeat_issue_lock`
- `gitea_heartbeat_reviewer_pr_lease` - `gitea_heartbeat_reviewer_pr_lease`
- `gitea_inspect_issue_lock_contract`
- `gitea_inspect_workflow_lease` - `gitea_inspect_workflow_lease`
- `gitea_issue_irrecoverable_provenance_authorization` - `gitea_issue_irrecoverable_provenance_authorization`
- `gitea_list_dependency_edges` - `gitea_list_dependency_edges`
@@ -135,6 +135,7 @@ that gates each call, not which tools exist.
- `gitea_record_pre_review_command` - `gitea_record_pre_review_command`
- `gitea_record_shell_spawn_outcome` - `gitea_record_shell_spawn_outcome`
- `gitea_record_stable_branch_push_attempt` - `gitea_record_stable_branch_push_attempt`
- `gitea_recover_incomplete_bootstrap_lock`
- `gitea_release_merger_pr_lease` - `gitea_release_merger_pr_lease`
- `gitea_release_reviewer_pr_lease` - `gitea_release_reviewer_pr_lease`
- `gitea_release_workflow_lease` - `gitea_release_workflow_lease`
+80
View File
@@ -0,0 +1,80 @@
{
"_comment": [
"Machine-checkable anchor table for docs/remote-mcp/threat-model.md (#956).",
"Every file:line anchor cited in the threat model must appear here, and the",
"source line at that anchor must contain the 'expect' substring.",
"tests/test_issue_956_threat_model.py enforces both directions, so a refactor",
"that shifts a line number fails the suite instead of silently rotting the",
"document. #930's inventory had no such guard and its gitea_mcp_server.py",
"anchors drifted between 7bf4f125 and aad5c8b4."
],
"generated_against_commit": "a143cd065ba06e1a2bdc5143a19ec156e53650ef",
"anchors": [
{"anchor": "gitea_mcp_server.py:24728", "expect": "mcp_daemon_guard.bind_native_mcp_transport()"},
{"anchor": "mcp_daemon_guard.py:49", "expect": "_PRODUCTION_TRANSPORTS = mcp_transport_config.SUPPORTED_TRANSPORTS"},
{"anchor": "mcp_daemon_guard.py:195", "expect": "def bind_native_mcp_transport"},
{"anchor": "irrecoverable_provenance.py:497", "expect": "def assess_transport_for_auth_mint"},
{"anchor": "gitea_mcp_server.py:9132", "expect": "assess_transport_for_auth_mint()"},
{"anchor": "gitea_mcp_server.py:9381", "expect": "assess_transport_for_auth_mint()"},
{"anchor": "mcp_server.py:4", "expect": "The transport is selected by deployment configuration"},
{"anchor": "gitea_mcp_server.py:15415", "expect": "def _is_client_managed_process"},
{"anchor": "gitea_mcp_server.py:15445", "expect": "def _provenance_mutation_block"},
{"anchor": "gitea_mcp_server.py:15453", "expect": "unsupported_manual_launch"},
{"anchor": "gitea_mcp_server.py:19004", "expect": "server_provenance"},
{"anchor": "gitea_mcp_server.py:21445", "expect": "def _check_mcp_runtimes_diagnostics"},
{"anchor": "gitea_mcp_server.py:21465", "expect": "\"ps\", \"-o\", \"pid,lstart,command\""},
{"anchor": "gitea_mcp_server.py:21509", "expect": "\"ps\", \"eww\""},
{"anchor": "gitea_config.py:1172", "expect": "RECOGNIZED_GITEA_ENV_KEYS"},
{"anchor": "gitea_config.py:1233", "expect": "GITEA_CLIENT_MANAGED"},
{"anchor": "gitea_config.py:54", "expect": "ENV_PROFILE = \"GITEA_MCP_PROFILE\""},
{"anchor": "gitea_config.py:97", "expect": "_REVIEW_MERGE_OPS"},
{"anchor": "gitea_config.py:499", "expect": "repository authorization scope"},
{"anchor": "gitea_config.py:956", "expect": "def _keychain_token"},
{"anchor": "gitea_config.py:974", "expect": "def resolve_token"},
{"anchor": "gitea_config.py:1015", "expect": "def keychain_auth"},
{"anchor": "gitea_config.py:294", "expect": "def _validate_identity_auth"},
{"anchor": "mcp_daemon_guard.py:583", "expect": "def assert_keychain_access_allowed"},
{"anchor": "gitea_mcp_server.py:19261", "expect": "def gitea_list_profiles"},
{"anchor": "gitea_mcp_server.py:19312", "expect": "gitea_config.resolve_token(p)"},
{"anchor": "gitea_mcp_server.py:19555", "expect": "def gitea_audit_config"},
{"anchor": "gitea_mcp_server.py:19577", "expect": "service_summaries(config)"},
{"anchor": "gitea_config.py:704", "expect": "def resolve_service"},
{"anchor": "gitea_config.py:837", "expect": "def service_summaries"},
{"anchor": "gitea_config.py:851", "expect": "_keychain_token(auth.get(\"id\"))"},
{"anchor": "gitea_mcp_server.py:17710", "expect": "\"jenkins-mcp\""},
{"anchor": "gitea_mcp_server.py:17716", "expect": "external-mcp"},
{"anchor": "gitea_mcp_server.py:17737", "expect": "\"glitchtip-mcp\""},
{"anchor": "gitea_mcp_server.py:17742", "expect": "external-mcp"},
{"anchor": "mcp_discoverability.py:9", "expect": "EXPECTED_JENKINS_TOOLS"},
{"anchor": "mcp_discoverability.py:17", "expect": "EXPECTED_GLITCHTIP_TOOLS"},
{"anchor": "sentry_incident_bridge.py:36", "expect": "SENTRY_AUTH_TOKEN"},
{"anchor": "sentry_incident_bridge.py:190", "expect": "def resolve_token"},
{"anchor": "sentry_incident_bridge.py:289", "expect": "Authorization"},
{"anchor": "sentry_observability.py:55", "expect": "SENTRY_DSN"},
{"anchor": "master_parity_gate.py:168", "expect": "def capture_startup_parity"},
{"anchor": "master_parity_gate.py:255", "expect": "mutation_safe"},
{"anchor": "gitea_mcp_server.py:19105", "expect": "def gitea_assess_master_parity"},
{"anchor": "gitea_mcp_server.py:193", "expect": "ACTIVE_WORKTREE_ENV"},
{"anchor": "gitea_mcp_server.py:194", "expect": "AUTHOR_WORKTREE_ENV"},
{"anchor": "gitea_mcp_server.py:2351", "expect": "/tmp/gitea_issue_lock.json"},
{"anchor": "gitea_mcp_server.py:10897", "expect": "def gitea_bootstrap_author_issue_worktree"},
{"anchor": "mcp_server.py:13", "expect": "/tmp/mcp_server_stderr.log"},
{"anchor": "issue_lock_store.py:26", "expect": "DEFAULT_LOCK_DIR"},
{"anchor": "issue_lock_store.py:83", "expect": "def session_pointer_path"},
{"anchor": "issue_lock_store.py:98", "expect": "def is_process_alive"},
{"anchor": "mcp_session_state.py:27", "expect": "DEFAULT_STATE_DIR"},
{"anchor": "control_plane_db.py:47", "expect": "DEFAULT_DB_PATH"},
{"anchor": "control_plane_db.py:380", "expect": "mode=0o700"},
{"anchor": "control_plane_db.py:386", "expect": "sqlite3.connect"},
{"anchor": "control_plane_db.py:1145", "expect": "os.getpid()"},
{"anchor": "gitea_mcp_server.py:12804", "expect": "owner_pid_alive"}
]
}
+407
View File
@@ -0,0 +1,407 @@
# Remote-MCP threat model, trust boundaries, and service decomposition
What the adversary is, what each boundary protects, and which services may share a process.
- **Issue:** #956 (Remote-MCP threat model), child of epic #929, cross-linked to #955.
- **Depends on:** #930 (closed) — `docs/remote-mcp/coupling-inventory.md`.
- **Blocks:** #932, #933, #934, #938.
- **Generated against commit:** `a143cd065ba06e1a2bdc5143a19ec156e53650ef` (#931's transport
bind seam). Originally generated against `aad5c8b42361d380a8eeb07b94b90815e594c2c5`
(`master`) and re-anchored when #931 shifted the cited lines.
- **Scope:** documentation only. This child changes no server behavior. It adds one
document, one anchor fixture, and the test that enforces them.
## Relationship to #930
#930 asked *what breaks when the process stops being local*. This document asks *what an
attacker gets, and where we stop them*. The two are deliberately different axes: #930
classifies each coupling as portable, seam, replacement, or cannot-be-remote; this document
classifies each **credential** by blast radius and each **boundary** by what crossing it
requires. An entry can be perfectly portable and still be a trust disaster —
`gitea_config.py:851` is portable Python that reads a CI secret from inside the Gitea server.
### Anchors are enforced, not asserted
Every `file:line` in this document is declared in `docs/remote-mcp/threat-model-anchors.json`
with the substring that must appear at that line, and
`tests/test_issue_956_threat_model.py` fails if any anchor does not resolve or if the
document cites an anchor the fixture does not cover.
This guard exists because #930 did not have one. Its inventory was generated at
`7bf4f125`; by `aad5c8b4` its `gitea_mcp_server.py` anchors had drifted — the transport
bind it cited at line 23750 now lives at `gitea_mcp_server.py:24728`, and its
client-managed provenance anchor at 14588 now lands in an unrelated function. Nothing
failed, because nothing checked. Anchors into a ~24,700-line module rot silently, and a
security document that cannot prove its own citations is worse than none, because it is
trusted.
---
## 1. Assets
What an adversary wants. Ordered by consequence, not by likelihood.
| ID | Asset | Why it matters |
| -- | ----- | -------------- |
| A1 | Merge authority on `Scaled-Tech-Consulting/Gitea-Tools` | This repository *is* the control plane. Code merged here becomes the gate that authorizes every future mutation, so merge authority is self-amplifying: one merge can disable every other control in this document. |
| A2 | Write authority on the `mdcps` tenant | A second, unrelated organization reachable from the same configuration. Compromise here is a cross-organization incident, not an internal one. |
| A3 | The eight Gitea role credentials | Long-lived bearer tokens. Possession is authority; there is no second factor at the API. |
| A4 | Jenkins read access (`mdcps`, enabled) | Build logs routinely carry deployment topology, internal hostnames, and accidentally-echoed secrets. |
| A5 | Error-tracking read access (GlitchTip / Sentry) | Event payloads carry stack frames, request context, and production user data. |
| A6 | Coordination-state integrity | The locks, leases, and review-decision records that make "exactly one owner" true. Corrupting them needs no Gitea credential and produces duplicate or lost work. |
| A7 | The operator's checkout and worktrees | Unmerged code, branch state, and the filesystem the author tools write to. |
| A8 | The macOS login keychain | The meta-credential. Everything in A3, A4, and A5 resolves from it. |
| A9 | Separation of duty between review and merge | The property that no single actor both approves and lands a change. An *asset*, not a control, because it is what the controls exist to produce. |
| A10 | Audit and provenance records | Determine whether an incident is reconstructable. An attacker who can forge provenance makes an intrusion indistinguishable from normal work. |
## 2. Adversaries
| ID | Adversary | Capability assumed | Not assumed |
| -- | --------- | ------------------ | ----------- |
| ADV1 | **Compromised LLM client** | Full control of one MCP client. Issues arbitrary tool calls, in any order, with any arguments, at machine speed. Sees every tool result. | Cannot read the operator's disk except through tools; cannot execute arbitrary local code outside the tool surface. |
| ADV2 | **Prompt injection** via repository content | Controls text the model reads and treats as instruction — issue bodies, PR descriptions, review comments, commit messages, file contents. Reaches the model on any read of untrusted content. | Holds no credential and issues no call directly. Its entire power is causing an *authorized* client to act. |
| ADV3 | **Malicious tool arguments** | Supplies hostile values to any parameter — paths, branch names, session identifiers, worktree paths, issue numbers — including traversal, injection, and confusion between look-alike identifiers. | Cannot bypass a gate that actually validates its input. |
| ADV4 | **Network attacker** | Observes and modifies traffic between client, server, and Gitea. Attempts downgrade, replay, and endpoint impersonation. | Does not hold a valid credential at the start. |
| ADV5 | **Curious operator** | Legitimate local access to the workstation: process table, `/tmp`, home directory, keychain prompts. Not malicious, but not authorized for every role either. | Does not defeat the OS keychain's own access control without a prompt. |
ADV2 is the adversary this architecture most under-models. Every other adversary must first
obtain something. Prompt injection obtains nothing: it borrows authority the client already
holds and is indistinguishable at the tool boundary from legitimate work. Each boundary
below therefore states whether it constrains ADV2 at all — and most do not, because they
authenticate the *caller*, not the *intent*.
## 3. Trust boundaries
"Crossing requires today" is what the code actually enforces at
`a143cd065ba06e1a2bdc5143a19ec156e53650ef`, not what the design intends.
| ID | Boundary | Protects | Crossing requires today | Crossing must require remotely |
| -- | -------- | -------- | ----------------------- | ------------------------------ |
| B1 | LLM client ↔ MCP server session | A1, A3, A10 — that a mutating session was established through the sanctioned client path | A single configured bind (`gitea_mcp_server.py:24728`) validated against one closed allowlist (`mcp_daemon_guard.py:49`, `mcp_daemon_guard.py:195`) — since #931 the identifier comes from deployment configuration and defaults to the local transport, so the boundary no longer rests on a literal, but it still rests on the *bind* rather than on an authenticated caller; client-managed provenance (`gitea_mcp_server.py:15415`) or a refusal (`gitea_mcp_server.py:15453`); production transport before recovery-authorization mint (`irrecoverable_provenance.py:497`, consumed at `gitea_mcp_server.py:9132` and `gitea_mcp_server.py:9381`) | An authenticated handshake issuing a server-side session identity bound to a principal, with the transport recorded in provenance. The physical proof (a pipe) must become a cryptographic one. |
| B2 | Role ↔ role | A9 — that author, reviewer, merger, and reconciler are distinct authorities | **The process boundary only.** The role is a property of the process, read once from `GITEA_MCP_PROFILE` (`gitea_config.py:54`). A caller gets author permissions by connecting to the author process. Review and merge are the operations singled out for extra care (`gitea_config.py:97`) | A per-request principal, so the role follows from the credential presented and cannot be selected by reaching a different endpoint. |
| B3 | MCP server ↔ credential store | A3, A8 — that only sanctioned code turns a profile into a token | `_keychain_token` shelling out to the login keychain (`gitea_config.py:956`), dispatched by `resolve_token` (`gitea_config.py:974`) with the reference type built at `gitea_config.py:1015`, gated by `assert_keychain_access_allowed` (`mcp_daemon_guard.py:583`). Inline secrets are rejected at config load (`gitea_config.py:294`) | A credential provider keyed by the *request* principal, returning only that principal's credential, with the source recorded and the value never returned. |
| B4 | MCP server ↔ Gitea | A1, A2 — that only authorized calls reach the forge | A bearer token over TLS. Server-side, nothing distinguishes one role's token from another beyond the account it belongs to | Unchanged at the forge; the endpoint in front of it must refuse unauthenticated and plaintext connections before tool dispatch. |
| B5 | MCP server ↔ caller's filesystem | A7 — that a tool acts on the *caller's* disk or refuses | Nothing. The server's disk *is* the caller's disk. Worktree bootstrap writes directly (`gitea_mcp_server.py:10897`); the active workspace is process-global (`gitea_mcp_server.py:193`, `gitea_mcp_server.py:194`) | An explicit per-tool classification, enforced at dispatch, refusing filesystem tools over a transport that cannot reach the caller's disk. A green verdict about the wrong disk is the failure to prevent. |
| B6 | MCP server ↔ coordination state | A6, A9 — mutual exclusion | Local files and a local SQLite database, with liveness judged from the local process table (`issue_lock_store.py:98`), keyed on paths under one user's home (`issue_lock_store.py:26`, `mcp_session_state.py:27`, `control_plane_db.py:47`) and on `os.getpid()` (`control_plane_db.py:1145`, `gitea_mcp_server.py:12804`). A legacy global slot still exists at `gitea_mcp_server.py:2351`, and the session-pointer file is named per PID (`issue_lock_store.py:83`) | One authority per ownership question, with liveness from session identity and expiry, and atomic acquire, renew, and release across hosts. |
| B7 | Gitea integration ↔ unrelated integrations | A4, A5 — that a Gitea compromise is not a CI and observability compromise | **Nothing.** See §5. The Gitea server reads Jenkins and GlitchTip secrets (`gitea_config.py:851`, reached from `gitea_config.py:837`) and holds the Sentry token (`sentry_incident_bridge.py:190`) | A hard process boundary. This is the boundary #956 exists to create. |
| B8 | Tenant ↔ tenant (`prgs` / `mdcps` / `local-lab`) | A2 — that one organization's compromise is not another's | Convention. One configuration declares all three contexts; `resolve_service` fails closed on a *disabled* context (`gitea_config.py:704`) but the credentials of enabled ones remain reachable in-process. A per-profile repository scope exists (`gitea_config.py:499`) | Separate deployments, or at minimum per-tenant credential scopes with no process able to resolve both. |
| B9 | Deployed code ↔ merged policy | A1, A10 — that the running server enforces the rules that were actually merged | Comparing this process's startup commit against this disk (`master_parity_gate.py:168`), conjoined into a single verdict (`master_parity_gate.py:255`) published by `gitea_mcp_server.py:19105` | Freshness defined against the deployed build identity, with an explicit fail-closed verdict when undeterminable. |
### What no boundary constrains
None of B1B9 constrains **ADV2**. Every one authenticates a caller or a process; prompt
injection supplies neither. An injected instruction that reaches an authorized author
session crosses B1, B2, B3, and B5 legitimately, because at each of those boundaries it *is*
the author. The only controls that bite ADV2 are those constraining what an authenticated
principal may do regardless of what it asks for — the per-role permission split (B2), the
repository scope at `gitea_config.py:499`, and separation of duty (A9). Sizing those
controls correctly matters more after the migration, not less, because a remote endpoint
raises the number of clients that can be injected into.
## 4. Data flows
Flows that cross a boundary. `==>` carries a credential; `-->` does not.
```
B1 B4
[LLM client] ====================> [MCP server] ========> [Gitea]
^ stdio pipe today | ^ (A1,A2)
| session identity | |
| after migration | |
| | | B3
untrusted repository content | +======> [macOS login keychain] (A8)
read back into the model (ADV2) | resolves A3, A4, A5
^ |
+----------------------------------+
|
B5 | B6
[operator checkout / worktrees] <--------+-------> [locks · leases · sqlite]
(A7) | (A6)
|
B7 <-- boundary does not exist today
|
+========================+========================+
| | |
[Jenkins] (A4) [GlitchTip] (A5) [Sentry] (A5)
external MCP server external MCP server in-process bridge
```
Two flows deserve attention because neither is obvious from the code:
1. **The keychain flow fans out.** B3 is drawn once but resolves credentials for *every*
configured profile and service, not only the active one. `gitea_list_profiles`
(`gitea_mcp_server.py:19261`) reports each profile's credential status by calling
`resolve_token` on it (`gitea_mcp_server.py:19312`), and `gitea_audit_config`
(`gitea_mcp_server.py:19555`) reports service credential status through
`service_summaries` (`gitea_mcp_server.py:19577`).
2. **The return path is a flow too.** Content read from Gitea travels back into the model
and is treated as instruction. This is the ADV2 edge, and it is the only edge in the
diagram with no authentication on it, because it is not a request.
## 5. Per-boundary credential inventory
**14 credentials in total.** Blast radius is stated as what the credential yields *on its
own*, assuming every gate not backed by the credential itself has been bypassed — because
an attacker holding a token calls the API, not our tools.
| ID | Credential | Holder | Boundary | Blast radius |
| -- | ---------- | ------ | -------- | ------------ |
| CR1 | `prgs-author` Gitea token — account `jcwalker3` | macOS keychain; resolved in-process (`gitea_config.py:974`) | B3 → B4 | Create branches, push, commit, open PRs, create/close/comment issues on the control-plane repo. Cannot approve or merge. The one credential whose identity is genuinely distinct. |
| CR2 | `prgs-reviewer` Gitea token — account `sysadmin` | macOS keychain | B3 → B4 | Approve and request changes. **Shares one Gitea account with CR3, CR4, CR5.** |
| CR3 | `prgs-merger` Gitea token — account `sysadmin` | macOS keychain | B3 → B4 | Merge to `master` — A1 in full. Same account as CR2. |
| CR4 | `prgs-reconciler` Gitea token — account `sysadmin` | macOS keychain | B3 → B4 | Close PRs, delete branches, irrecoverable decision-lock recovery. Same account as CR2. |
| CR5 | `prgs-controller` Gitea token — account `sysadmin` | macOS keychain | B3 → B4 | Same operation set as CR4. Same account as CR2. |
| CR6 | `mdcps-author` Gitea token — account `913443` | macOS keychain | B3 → B4, B8 | Author operations on a second organization. **Shares one account with CR7 and CR8.** |
| CR7 | `mdcps-reviewer` Gitea token — account `913443` | macOS keychain | B3 → B4, B8 | Approve and request changes on `mdcps`. Same account as CR6. |
| CR8 | `mdcps-merger` Gitea token — account `913443` | macOS keychain | B3 → B4, B8 | Merge on `mdcps` — A2 in full. Same account as CR6. |
| CR9 | MDCPS Jenkins read credential | macOS keychain, read from the Gitea server process (`gitea_config.py:851`) | B7 | Read CI jobs, builds, and logs (A4). Enabled today. |
| CR10 | MDCPS GlitchTip read credential | macOS keychain, read from the Gitea server process (`gitea_config.py:851`) | B7 | Read error events and their payloads (A5). Enabled today. |
| CR11 | `SENTRY_AUTH_TOKEN` | Process environment, read in-process (`sentry_incident_bridge.py:36`, `sentry_incident_bridge.py:190`), sent as a bearer header (`sentry_incident_bridge.py:289`) | B7 | Read and reconcile Sentry issues (A5). Not a keychain credential — an env var, so it is inherited by anything the process spawns. |
| CR12 | `SENTRY_DSN` | Process environment (`sentry_observability.py:55`) | B7 | Write events into the observability project. Low read value, real forgery value: an attacker can inject fabricated events into the record (A10). |
| CR13 | macOS login keychain access | The operator's login session; gated by `assert_keychain_access_allowed` (`mcp_daemon_guard.py:583`) | B3, ADV5 | **Every other credential in this table except CR11 and CR12.** This is the aggregation point. |
| CR14 | Coordination-store access (no secret) | Filesystem permissions — `control_plane_db.py:47`, created `0o700` (`control_plane_db.py:380`), opened with a local file lock (`control_plane_db.py:386`) | B6, ADV5 | Full read/write of locks, leases, and decision records (A6). **There is no credential here at all** — anything running as the operator can rewrite ownership. |
### Findings
**Finding 1 — Role separation is not credential separation.** Four `prgs` roles resolve to
one Gitea account (`sysadmin`): reviewer, merger, reconciler, and controller. A stolen
reviewer credential *is* a merger credential. A9 — separation of duty between approving and
landing — is therefore enforced entirely by which local process a call reaches (B2), and not
at all by the forge. It survives exactly as long as B2 does, and B2 is the boundary the
migration dissolves.
**Finding 2 — The `mdcps` tenant has no role separation at all.** Author, reviewer, and
merger all resolve to account `913443`. One credential can open a PR, approve it, and merge
it. The in-process self-review check compares the authenticated username against the PR
author and would refuse — but that check runs on our side of B4. It is not a property of
the credential, and an attacker holding the token does not call our tools.
**Finding 3 — Any one role process can resolve every other role's credential.** This is not
inferred; it is demonstrated by tool output. `gitea_list_profiles`
(`gitea_mcp_server.py:19261`) called from the **author** session reports
`identity_status: "credentials present"` for `prgs-merger`, `prgs-reviewer`,
`prgs-reconciler`, and every `mdcps` profile, because it calls `resolve_token` on each one
(`gitea_mcp_server.py:19312`). The author process does not merely *have access to* the
merger's credential — it reads it to answer a status query. B2 is not a credential boundary
in either direction.
**Finding 4 — The Gitea server reads CI and observability secrets.** `gitea_audit_config`
(`gitea_mcp_server.py:19555`) reports `MDCPS Jenkins: enabled, read-only, authenticated`.
That word `authenticated` is produced by `service_summaries` (`gitea_mcp_server.py:19577`,
defined at `gitea_config.py:837`), whose default check calls `_keychain_token` on the
service's own keychain reference (`gitea_config.py:851`). Producing that one line requires
the Gitea MCP server to read the Jenkins secret and the GlitchTip secret out of the
keychain. B7 does not exist.
**Finding 5 — Jenkins and GlitchTip are already decomposed; the reach is residual.** Their
tools live in separately registered servers, marked `external-mcp`
(`gitea_mcp_server.py:17710`, `gitea_mcp_server.py:17716`, `gitea_mcp_server.py:17737`,
`gitea_mcp_server.py:17742`) with their own expected tool sets (`mcp_discoverability.py:9`,
`mcp_discoverability.py:17`). The correct decomposition was already chosen. What remains is
a leak across it: the credential *references* still live in the Gitea configuration and are
still resolved by the Gitea process. #75 bundled these services into one control-plane
umbrella; the tools were separated afterwards, the credentials were not.
**Finding 6 — Sentry is the exception that is not decomposed.** Unlike Jenkins and
GlitchTip, the Sentry bridge runs *inside* the Gitea server, resolving its token from the
process environment (`sentry_incident_bridge.py:190`) and sending it as a bearer header
(`sentry_incident_bridge.py:289`). Being an environment variable rather than a keychain item
makes it strictly worse: it needs no keychain prompt and is inherited by every subprocess the
server spawns — including the `ps` invocations at `gitea_mcp_server.py:21465` and
`gitea_mcp_server.py:21509`, reached from `gitea_mcp_server.py:21445`.
**Finding 7 — The highest-value coordination asset has the weakest gate.** A6 is protected
by filesystem permissions alone (CR14). Corrupting a lease requires no Gitea credential,
produces no forge-side audit record, and breaks the mutual exclusion the entire workflow
assumes. Every other asset costs an attacker a credential; this one costs nothing beyond
local access, which is exactly ADV5's position.
**Finding 8 — Provenance authenticates the launch, not the caller.** `server_provenance` is
reported as exactly `client_managed` or `manual_launch` (`gitea_mcp_server.py:19004`),
derived from environment inspection (`gitea_mcp_server.py:15415`) with the recognized-key
allowlist at `gitea_config.py:1172` and the generator that emits the marker at
`gitea_config.py:1233`. Every one of those facts is fixed at process start. A client that is
trustworthy at launch and compromised a minute later remains `client_managed` for the life
of the process, and the transport contract that underwrites it is stated as a property of
the server itself (`mcp_server.py:4`). Since #931 that contract names the configured
transport rather than asserting stdio, but it is still fixed once, at bind, for the life of
the process.
## 6. Decomposition ruling
This section is the ruling #956 requires. It is a decision, not a recommendation.
**D1 — No unrelated co-residency.** A single integration process **must not** hold, resolve,
or be able to resolve credentials for services it does not itself integrate with.
Concretely: the Gitea MCP service may hold Gitea credentials and nothing else. Jenkins,
GlitchTip, Sentry, and any database credential are **not permitted** to co-reside with Gitea
credentials in one process.
*Rationale.* A process is the smallest unit an attacker takes whole. Once ADV1 or ADV2
controls execution in a process, every credential that process can resolve is theirs, and no
in-process check helps, because the checks are in the process too. Blast radius is therefore
a property of the process boundary and nothing finer. Findings 4 and 6 show that today one
compromise of the Gitea server yields CI read access, error-tracking read access, and — via
CR13 — every role credential on both tenants. That is the single largest reduction in blast
radius available anywhere in epic #929, and it costs no new mechanism: the decomposition
already exists (Finding 5) and is merely leaked across.
**D2 — Separation of duty must be backed by credentials.** Two roles whose separation is a
security property must not resolve to the same forge account. Specifically, reviewer and
merger must be distinct accounts. Today they are not, on either tenant (Findings 1 and 2).
*Rationale.* B2 is a process boundary, and the migration's entire purpose is to replace
process boundaries with request-level ones. A separation enforced only by which process a
call reaches does not survive that replacement — and it is already bypassable by anyone who
holds the token and calls the API instead of the tool.
**D3 — Credential resolution is scoped to the request principal.** A session must resolve its
own credential and must have no path to any other principal's. The resolve-every-profile
behavior behind `gitea_mcp_server.py:19312` and `gitea_mcp_server.py:19577` must report
configured-or-not from configuration alone, without resolving the secret.
*Rationale.* Finding 3. An audit surface that proves a credential exists by fetching it is a
credential-aggregation primitive wearing a diagnostic's clothes.
**D4 — Coordination state is a protected asset with its own authority.** Access to locks,
leases, and decision records must require an authenticated session, not merely local
filesystem access.
*Rationale.* Finding 7. #937 already moves this store for concurrency reasons; the
authorization requirement must land with it, or the store becomes remotely reachable while
still being authorized by nothing.
### Exceptions
**One, time-boxed.** During the dual-run window defined by #939, the **local** stdio fleet
may continue to resolve Jenkins and GlitchTip credential *references* from the shared
configuration, because removing them from the local configuration is not a prerequisite for
standing up the remote endpoint and would strand the operator's existing local workflow.
This exception is bounded by all of:
- It applies to the local stdio deployment only. The remote endpoint (#938) must be
configured with Gitea credentials and no others from its first day.
- It expires when #939 completes. It does not survive cutover.
- It does not extend to Sentry: CR11 and CR12 are process-environment credentials in the
Gitea server (Finding 6) and must be absent from the remote deployment's environment
regardless of dual-run state.
No exception is granted to D2, D3, or D4.
### Consequences for the target architecture
- The remote endpoint serves **Gitea only**. It is not a general control-plane endpoint.
- Jenkins and GlitchTip keep their existing separate servers, and their credential
references move out of the Gitea configuration.
- The Sentry bridge either moves behind its own service boundary or is absent from the
remote deployment. It does not travel with the Gitea server.
- Reviewer and merger accounts diverge before the endpoint is trusted for merges, or A9 is
recorded as unenforced.
## 7. Child-to-boundary mapping
Every #929 child from 2 through 10, mapped to the boundary it implements. A child
implementing more than one boundary names its primary first.
| Child | Issue | Boundaries | What it must establish | Rulings it must honor |
| ----: | ----- | ---------- | ---------------------- | --------------------- |
| 2 | #931 | B1, B9 | The bound transport becomes a validated value that provenance and freshness can both key on. Without it neither B1 nor B9 has an input. | — |
| 3 | #932 | B2 | The role becomes a property of the request, not the process — the boundary the migration otherwise deletes. | D2, D3 |
| 4 | #933 | B3, B7 | Credentials come from a provider keyed by principal. This is where D1 and D3 are either enforced or permanently lost. | D1, D3 |
| 5 | #934 | B1 | Session provenance replaces pipe-and-process-table proof with an authenticated session identity. | — |
| 6 | #935 | B9 | Freshness redefined against deployed build identity, with an explicit undeterminable verdict. | — |
| 7 | #936 | B5 | Every tool classified and the filesystem boundary enforced at dispatch, so a tool cannot return green about the wrong disk. | — |
| 8 | #937 | B6 | One authority per ownership question, with session-identity liveness and atomic transitions. | D4 |
| 9 | #938 | B4, B1, B8 | The endpoint: authentication, principal binding, transport security, and — critically — the deployed credential set. | D1, D2, D3 |
| 10 | #939 | B6 | Dual-run with exactly one coordination authority at every instant, and the rollback that proves the way back. | D1 exception expiry |
Boundary coverage: B1 (#931, #934, #938), B2 (#932), B3 (#933), B4 (#938), B5 (#936),
B6 (#937, #939), B7 (#933), B8 (#938), B9 (#931, #935).
B7 has exactly one owner, #933, and that is deliberate. B7 is not created by standing up an
endpoint; it is created by deciding which credentials a process may resolve, which is
precisely what the credential-provider child does.
## 8. Adversarial walkthrough: a fully compromised LLM client
The required negative case. ADV1 holds one MCP client and issues arbitrary tool calls. What
does it reach at each boundary?
### Before the migration
The client is attached to the local fleet over stdio.
| Boundary | What ADV1 reaches | Stopped by |
| -------- | ----------------- | ---------- |
| B1 | Everything the fleet serves. The client *is* the sanctioned launcher: it satisfies the client-managed check (`gitea_mcp_server.py:15415`) by construction, and provenance is never re-verified after launch (Finding 8). | Nothing. The guard authenticates the launch, not the caller. |
| B2 | All five roles — it is attached to all five namespaces. It can author a PR, approve it from the reviewer namespace, and merge it from the merger namespace. | Only the in-process self-review check, which compares `jcwalker3` (author) against `sysadmin` (reviewer) and **passes**, because Finding 1 made them different accounts while leaving reviewer and merger identical. A9 falls in one sequence of legitimate calls. |
| B3 | Every credential in CR1CR10 via CR13, with no additional prompt — the daemon is already sanctioned, so `assert_keychain_access_allowed` (`mcp_daemon_guard.py:583`) returns immediately. | Nothing. |
| B4 | A1 and A2 in full. | Branch protection at the forge, to the extent configured. |
| B5 | The operator's checkout and every worktree, through the author tools (`gitea_mcp_server.py:10897`), plus the shared stderr path at `mcp_server.py:13`. | Nothing; the server's disk is the target disk. |
| B6 | All coordination state — no credential required (CR14). It can forge lease ownership and clear decision locks. | Filesystem permissions, which it already satisfies. |
| B7 | Jenkins (A4) and GlitchTip (A5) secrets via Finding 4, and CR11/CR12 from its own environment. | Nothing. |
| B8 | Both tenants. | Nothing in-process; only the disabled-context check (`gitea_config.py:704`), which does not apply to enabled contexts. |
| B9 | Can it run stale code? Only by advancing the checkout, which parity detects (`master_parity_gate.py:255`). | This gate holds — the one boundary that behaves as designed against ADV1. |
**Result: total.** One compromised client yields A1 through A8 and A10. The only asset with
real resistance is A1 via branch protection, and the client holds the merger credential
anyway. Nine boundaries, one meaningful stop.
### After the migration
The same client authenticates to the remote endpoint with one role's credential, assuming
#931#939 land **and honor D1D4**.
| Boundary | What ADV1 reaches | Stopped by |
| -------- | ----------------- | ---------- |
| B1 | One authenticated session, bound to one principal. | #934: a forged or expired session identity is refused; the client cannot mint one. |
| B2 | **One role.** Presenting the author credential yields author permissions only. | #932: the principal comes from the credential, not from which endpoint was reached. |
| B3 | **One credential — its own.** | #933 with D3: the provider resolves by principal, and no diagnostic resolves the others. |
| B4 | That role's authority on the forge. | Endpoint authentication (#938); plaintext and unauthenticated attempts refused before dispatch. |
| B5 | **Nothing.** Filesystem tools are refused over the remote transport with a named blocker. | #936. |
| B6 | Its own leases; contention resolves to exactly one winner. | #937 with D4: authenticated session required, not filesystem access. |
| B7 | **Nothing.** No CI or observability credential exists in the process. | D1 — the single largest reduction on this table. |
| B8 | One tenant. | D1 and #938: the deployment carries one tenant's credentials. |
| B9 | Cannot induce stale enforcement. | #935: explicit fail-closed verdict, including undeterminable. |
**Result: bounded.** The compromise is contained to one role on one tenant, with no
filesystem reach and no lateral credential access. A9 survives *only if D2 lands* — if
reviewer and merger still share `sysadmin`, a compromised reviewer session still merges, and
this row reads the same after the migration as before it.
### What the migration does not fix
Against **ADV2**, both tables are identical. Prompt injection does not need to cross a
boundary: it arrives inside an authorized session and asks that session to do what it is
already permitted to do. Every "stopped by" above authenticates a principal, and the
injected instruction has the correct principal. The migration reduces ADV1's blast radius by
roughly an order of magnitude and reduces ADV2's by nothing.
The controls that do constrain ADV2 are per-principal permission scope (#932), repository
scope (`gitea_config.py:499`), and credential-backed separation of duty (D2) — each limiting
what an authenticated session may do *regardless of what it is asked for*. #955's
secure-isolation end state should be read with that distinction in mind: removing credentials
from clients defeats ADV1 and ADV5, and does not by itself defeat ADV2.
Two further items are explicitly out of scope here and unowned by #929:
- **Session-credential rotation and revocation.** #938 names rotation as documentation, but
no child owns proving that a revoked credential stops an in-flight session.
- **ADV3** (malicious tool arguments) is diffused across every child rather than owned. The
per-request principal work in #932 is the natural place to assert that identifiers taken
from the request never authorize anything on their own.
## 9. How to verify this document
1. `PYTHONPATH=. pytest tests/test_issue_956_threat_model.py` — resolves every anchor
against the working tree and checks the document's structural obligations.
2. Pick any five anchors at random and read them; the fixture states what each line must
contain.
3. Reproduce Findings 3 and 4 live: call `gitea_list_profiles` and `gitea_audit_config`
from the **author** namespace. Credential presence reported for roles other than the
active one is Finding 3; `MDCPS Jenkins: enabled, read-only, authenticated` is Finding 4.
If the anchor test fails after an unrelated refactor, the anchors moved and the fixture
needs regenerating — the claims are still true, but they are no longer traceable, which
#956 treats as the same defect.
+412 -214
View File
@@ -1,7 +1,10 @@
#!/usr/bin/env python3 #!/usr/bin/env python3
"""Gitea MCP Server — exposes Gitea operations as MCP tools. """Gitea MCP Server — exposes Gitea operations as MCP tools.
Runs over stdio. All tools authenticate via macOS keychain (git credential fill). The transport is selected by deployment configuration (GITEA_MCP_TRANSPORT) and
defaults to the local client-spawned transport when unset (#931); the permitted
set lives in mcp_transport_config. All tools authenticate via macOS keychain
(git credential fill).
Usage (standalone test): Usage (standalone test):
python3 mcp_server.py python3 mcp_server.py
@@ -196,7 +199,6 @@ RECONCILER_WORKTREE_ENV = "GITEA_RECONCILER_WORKTREE"
import namespace_workspace_binding as nwb # noqa: E402 import namespace_workspace_binding as nwb # noqa: E402
import canonical_repository_root as crr # noqa: E402 # #706 cross-repo canonical root import canonical_repository_root as crr # noqa: E402 # #706 cross-repo canonical root
import mcp_namespace_health # noqa: E402 import mcp_namespace_health # noqa: E402
import mcp_fleet_inventory # noqa: E402 # #949 authoritative fleet inventory
import stale_binding_recovery # noqa: E402 import stale_binding_recovery # noqa: E402
# Worktree env bindings inherited from the parent environment at daemon boot # Worktree env bindings inherited from the parent environment at daemon boot
@@ -2111,6 +2113,8 @@ import issue_lock_store # noqa: E402
import issue_lock_adoption # noqa: E402 import issue_lock_adoption # noqa: E402
import issue_lock_recovery # noqa: E402 import issue_lock_recovery # noqa: E402
import issue_lock_renewal # noqa: E402 import issue_lock_renewal # noqa: E402
import author_lock_contract # noqa: E402
import bootstrap_lock_recovery # noqa: E402
import dirty_orphan_worktree_recovery # noqa: E402 # #860 dirty orphan recovery import dirty_orphan_worktree_recovery # noqa: E402 # #860 dirty orphan recovery
import dirty_same_claimant_session_rebind # noqa: E402 # #864 import dirty_same_claimant_session_rebind # noqa: E402 # #864
import stacked_pr_support # noqa: E402 import stacked_pr_support # noqa: E402
@@ -2659,33 +2663,17 @@ def _build_author_issue_work_lease(
worktree_path: str, worktree_path: str,
host: str | None, host: str | None,
) -> dict: ) -> dict:
created = _work_lease_now() # #953: the lease shape now lives in author_lock_contract so that bootstrap
# #790 Slice A: the window comes from the central policy, not a literal here. # and gitea_lock_issue cannot drift apart again. The policy-derived sliding
# It is also now a *sliding* window — the lease lives ``initial_ttl_minutes`` # TTL (#790 Slice A) and the task-session ownership key (#790 AC-N1) are
# past its last valid heartbeat rather than a fixed four hours past its # unchanged — they simply have one definition instead of two.
# creation, so an abandoned task stops holding the claim within one TTL. return author_lock_contract.build_author_issue_work_lease(
policy = lease_policy.policy_for(lease_policy.TASK_CLASS_AUTHOR_ISSUE_WORK) issue_number=issue_number,
expires = created + timedelta(minutes=policy.initial_ttl_minutes) branch_name=branch_name,
return { worktree_path=worktree_path,
"operation_type": AUTHOR_ISSUE_WORK_LEASE, claimant=_work_lease_claimant(host),
"issue_number": issue_number, created=_work_lease_now(),
"pr_number": None, )
"branch": branch_name,
"worktree_path": worktree_path,
"claimant": _work_lease_claimant(host),
"created_at": _work_lease_timestamp(created),
"expires_at": _work_lease_timestamp(expires),
"last_heartbeat_at": _work_lease_timestamp(created),
# #790 AC-N1: the ownership key for this task. Distinct from the recorded
# PID, which is the shared daemon and identifies no individual task.
"task_session_id": issue_lock_store.mint_task_session_id(
AUTHOR_ISSUE_WORK_LEASE
),
# #790 AC-N8: the explicit lifecycle marker. Its absence — never a
# timestamp comparison — is what makes a lock legacy.
"lifecycle_version": lease_policy.LIFECYCLE_HEARTBEAT_V1,
"heartbeat_count": 1,
}
def _active_work_lease_block( def _active_work_lease_block(
@@ -4955,6 +4943,370 @@ def gitea_heartbeat_issue_lock(
return outcome return outcome
def _observe_recovery_worktree(worktree_path: str) -> dict:
"""Observe head, branch, existence, and registration for lock recovery.
Read-only: it runs ``git`` queries and touches nothing. Kept separate from
the decision so the decision stays a pure function of observations (#953).
"""
observation = {
"worktree_exists": os.path.isdir(worktree_path),
"worktree_registered": False,
"current_branch": "",
"observed_head": "",
}
if not observation["worktree_exists"]:
return observation
try:
observation["current_branch"] = subprocess.run(
["git", "-C", worktree_path, "rev-parse", "--abbrev-ref", "HEAD"],
capture_output=True,
text=True,
check=False,
).stdout.strip()
observation["observed_head"] = subprocess.run(
["git", "-C", worktree_path, "rev-parse", "HEAD"],
capture_output=True,
text=True,
check=False,
).stdout.strip()
listing = subprocess.run(
["git", "-C", worktree_path, "worktree", "list", "--porcelain"],
capture_output=True,
text=True,
check=False,
).stdout
real = os.path.realpath(worktree_path)
observation["worktree_registered"] = any(
os.path.realpath(line.split(" ", 1)[1].strip()) == real
for line in listing.splitlines()
if line.startswith("worktree ")
)
except Exception: # fail closed: unobservable is not provable
return observation
return observation
@mcp.tool()
def gitea_inspect_issue_lock_contract(
issue_number: int,
branch_name: str | None = None,
worktree_path: str | None = None,
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
) -> dict:
"""Inspect a durable author issue lock against the canonical contract (#953 AC8/AC16).
Strictly read-only. It performs no lock, lease, branch, worktree, issue, or
pull-request mutation of any kind it reads the durable lock record and
reports. Use it to find out *why* a lock is being refused before choosing a
recovery path, and to confirm afterwards that recovery produced a canonical
lock.
Reports which canonical fields are missing, where the claimant is recorded
(``work_lease`` is canonical, top level is the legacy/bootstrap placement),
whether an expiration is actually recorded as opposed to absent, which
used to masquerade as "not yet expired" whether the lock can be
heartbeated, and whether it satisfies the untouched #447 create-PR
provenance guard.
Args:
issue_number: The issue whose lock to inspect.
branch_name: Optional; when given, the recovery eligibility preview is
evaluated against this branch.
worktree_path: Optional; when given, the recovery eligibility preview is
evaluated against this worktree.
remote: Known instance 'dadeschools' or 'prgs'.
host: Override the Gitea host.
org: Override the owner/organization.
repo: Override the repository name.
Returns:
dict with 'success', 'lock_present', 'lock_contract' (the structural
verdict), 'recommended_action', and when branch_name and
worktree_path are supplied a non-mutating 'recovery_preview'.
"""
blocked = _profile_permission_block(
"gitea.read",
issue_number=issue_number,
remote=remote,
host=host,
org=org,
repo=repo,
org_explicit=org is not None,
repo_explicit=repo is not None,
)
if blocked:
return blocked
h, o, r = _resolve(remote, host, org, repo)
existing = _load_existing_issue_lock(
remote=remote, org=o, repo=r, issue_number=issue_number
)
contract = author_lock_contract.assess_lock_contract(existing)
result = {
"success": True,
"performed": False,
"mutation_performed": False,
"read_only": True,
"issue_number": issue_number,
"lock_present": bool(existing),
"lock_contract": contract,
"lock_freshness": (
issue_lock_store.assess_lock_freshness(dict(existing))
if existing
else {"status": issue_lock_store.STATUS_ABSENT, "live": False}
),
"recommended_action": author_lock_contract.recommended_action(contract),
}
if branch_name and worktree_path:
resolved = issue_lock_worktree.resolve_author_worktree_path(
worktree_path, _canonical_local_git_root()
)
observation = _observe_recovery_worktree(resolved)
claimant = _work_lease_claimant(h)
result["recovery_preview"] = bootstrap_lock_recovery.assess_bootstrap_lock_recovery(
existing,
issue_number=issue_number,
branch_name=branch_name,
worktree_path=resolved,
remote=remote,
org=o,
repo=r,
identity=claimant.get("username"),
profile=claimant.get("profile"),
observed_head=observation["observed_head"],
declared_head=None,
worktree_exists=observation["worktree_exists"],
worktree_registered=observation["worktree_registered"],
current_branch=observation["current_branch"],
)
return result
@mcp.tool()
def gitea_recover_incomplete_bootstrap_lock(
issue_number: int,
branch_name: str,
worktree_path: str,
expected_head: str,
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
dry_run: bool = False,
) -> dict:
"""Upgrade an incomplete bootstrap issue lock to the canonical contract (#953 AC8-AC11).
Explicit, target-specific recovery. It does **not** widen
``gitea_lock_issue``, and it is not a takeover path.
The state it repairs: ``gitea_bootstrap_author_issue_worktree`` reported
success but wrote a lock with the claimant at the top level, no
``work_lease``, no ``lock_provenance``, and no expiry. The author then
implemented, committed, and pushed following bootstrap's own reported next
action after which heartbeat, re-lock, exact-owner renewal, and the #447
create-PR guard all refuse simultaneously.
Deliberate non-behaviours: the branch is never moved, reset, or rewound, and
base-equivalence is never required preserving the already-committed and
pushed work is the entire point. Nothing is pushed and no pull request is
created. Only the single lock file for this exact (remote, org, repo, issue)
is written.
Ownership is proven, never asserted. The claimant recorded on the durable
lock must match **both** the server-resolved identity and the active
profile; a matching username alone is refused. Repository, issue, branch,
worktree, registration, current branch, and head are all verified before any
write, and the declared ``expected_head`` must equal the observed head. A
healthy foreign-owned lock is refused outright. Provenance and authorization
are minted server-side there is no parameter through which a caller can
supply either.
Args:
issue_number: The issue whose incomplete lock is being recovered.
branch_name: The branch recorded on the lock; must match.
worktree_path: The registered worktree recorded on the lock; must match.
expected_head: Full SHA the caller believes the worktree is at. A
mismatch fails closed, so a worktree that moved underneath the
caller cannot be recovered against stale evidence.
remote: Known instance 'dadeschools' or 'prgs'.
host: Override the Gitea host.
org: Override the owner/organization.
repo: Override the repository name.
dry_run: Report the decision and evidence, mutate nothing.
Returns:
dict with 'success', 'performed', the resulting canonical
'lock_contract' and 'work_lease', the auditable
'bootstrap_lock_recovery' transition record, and 'exact_next_action'; on
refusal 'success'/'performed' False with 'refusal_code' and 'reasons'
naming exactly which evidence was missing, and no write performed.
"""
task = "recover_incomplete_bootstrap_lock"
ok, block_reasons = role_session_router.check_author_mutation_after_reviewer_stop(
task
)
if not ok:
return _author_mutation_block(block_reasons)
# #953 F1: the namespace/session wall every author state-creating mutation
# carries, and which this tool — the structural neighbour of
# gitea_recover_dirty_orphaned_issue_worktree, writing the same durable
# lock — was the only one to omit. Exact-owner claimant matching inside
# assess_bootstrap_lock_recovery is a later layer, not a substitute: it
# refuses one commit too late and leaves no BLOCKED audit record of the
# attempt. author_role_exclusive is required here because this task is gated
# on gitea.issue.comment, which merger, controller, and reconciler profiles
# also hold.
blocked = _namespace_mutation_block(
task, remote=remote, author_role_exclusive=True
)
if blocked:
return blocked
blocked = _profile_permission_block(
task_capability_map.required_permission(task),
issue_number=issue_number,
remote=remote,
host=host,
org=org,
repo=repo,
org_explicit=org is not None,
repo_explicit=repo is not None,
)
if blocked:
return blocked
h, o, r = _resolve(remote, host, org, repo)
resolved_worktree = issue_lock_worktree.resolve_author_worktree_path(
worktree_path, _canonical_local_git_root()
)
existing = _load_existing_issue_lock(
remote=remote, org=o, repo=r, issue_number=issue_number
)
observation = _observe_recovery_worktree(resolved_worktree)
# The claimant pair is resolved server-side from the live session; the
# caller cannot influence which identity or profile recovery compares
# against.
claimant = _work_lease_claimant(h)
assessment = bootstrap_lock_recovery.assess_bootstrap_lock_recovery(
existing,
issue_number=issue_number,
branch_name=branch_name,
worktree_path=resolved_worktree,
remote=remote,
org=o,
repo=r,
identity=claimant.get("username"),
profile=claimant.get("profile"),
observed_head=observation["observed_head"],
declared_head=expected_head,
worktree_exists=observation["worktree_exists"],
worktree_registered=observation["worktree_registered"],
current_branch=observation["current_branch"],
)
if not assessment["recovery_sanctioned"]:
return {
"success": False,
"performed": False,
"mutation_performed": False,
"issue_number": issue_number,
"refusal_code": assessment["refusal_code"],
"reasons": assessment["reasons"],
"message": bootstrap_lock_recovery.format_recovery_refusal(assessment),
"lock_contract": assessment["contract"],
"evidence": assessment["evidence"],
}
if dry_run:
return {
"success": True,
"performed": False,
"mutation_performed": False,
"dry_run": True,
"issue_number": issue_number,
"would_recover": True,
"lock_contract": assessment["contract"],
"evidence": assessment["evidence"],
"exact_next_action": (
"Re-run without dry_run=True to upgrade this lock to the "
"canonical contract."
),
}
recovered = author_lock_contract.build_canonical_issue_lock(
issue_number=issue_number,
branch_name=branch_name,
worktree_path=resolved_worktree,
remote=remote,
org=o,
repo=r,
identity=claimant.get("username"),
profile=claimant.get("profile"),
tool="gitea_recover_incomplete_bootstrap_lock",
source=issue_lock_provenance.SOURCE_LOCK_ISSUE,
owner_session=(existing or {}).get("owner_session"),
expected_base_sha=(existing or {}).get("expected_base_sha"),
)
# AC10: preserve both sides of the transition so a recovered lock never
# reads as an original claim.
recovered["bootstrap_lock_recovery"] = bootstrap_lock_recovery.build_recovery_record(
assessment,
recovered_at=_work_lease_timestamp(_work_lease_now()),
new_task_session_id=recovered["work_lease"]["task_session_id"],
)
try:
lock_path = issue_lock_store.bind_session_lock(
recovered,
expected_generation=assessment["expected_generation"],
recovery_sanctioned=True,
)
except Exception as exc:
return {
"success": False,
"performed": False,
"mutation_performed": False,
"issue_number": issue_number,
"refusal_code": "lock_write_failed",
"reasons": [str(exc)],
"message": f"Recovered lock could not be persisted: {exc} (fail closed)",
}
written = issue_lock_store.read_lock_file(lock_path)
contract = author_lock_contract.assess_lock_contract(written)
return {
"success": True,
"performed": True,
"mutation_performed": True,
"issue_number": issue_number,
"branch_name": branch_name,
"worktree_path": resolved_worktree,
"lock_file_path": lock_path,
"lock_contract": contract,
"work_lease": (written or {}).get("work_lease"),
"task_session_id": contract["task_session_id"],
"lock_generation": (written or {}).get("lock_generation"),
"prior_generation": assessment["expected_generation"],
"bootstrap_lock_recovery": (written or {}).get("bootstrap_lock_recovery"),
"preserved_head": observation["observed_head"],
"branch_reset": False,
"pushed": False,
"pr_created": False,
"exact_next_action": (
"Lock is canonical. Heartbeat it with the returned task_session_id, "
"then continue the author workflow; publish and create the pull "
"request through the normal sanctioned calls."
),
}
@mcp.tool() @mcp.tool()
def gitea_recover_dirty_orphaned_issue_worktree( def gitea_recover_dirty_orphaned_issue_worktree(
issue_number: int, issue_number: int,
@@ -15156,8 +15508,21 @@ def _profile_permission_block(required_operation: str, **extra_fields) -> dict |
) )
def _namespace_mutation_block(mutation_task: str, **extra_fields) -> dict | None: def _namespace_mutation_block(
"""Reviewer/author namespace alignment gate (#209).""" mutation_task: str,
*,
author_role_exclusive: bool = False,
**extra_fields,
) -> dict | None:
"""Reviewer/author namespace alignment gate (#209).
``author_role_exclusive`` additionally requires the active profile's derived
role kind to be exactly ``author`` (#953 F1). Off by default, so the six
pre-existing call sites are unchanged. Tools whose required permission is
``gitea.issue.comment`` which every configured role holds opt in, since
the reviewer-namespace check alone would let a merger, controller, or
reconciler session through to a durable author lock write.
"""
required_permission = task_capability_map.required_permission(mutation_task) required_permission = task_capability_map.required_permission(mutation_task)
required_role = task_capability_map.required_role(mutation_task) required_role = task_capability_map.required_role(mutation_task)
# #714: evaluate active profile only — never auto-switch. # #714: evaluate active profile only — never auto-switch.
@@ -15179,6 +15544,9 @@ def _namespace_mutation_block(mutation_task: str, **extra_fields) -> dict | None
} }
ok, reasons = role_namespace_gate.check_author_mutation_namespace( ok, reasons = role_namespace_gate.check_author_mutation_namespace(
mutation_task, profile) mutation_task, profile)
if ok and author_role_exclusive:
ok, reasons = role_namespace_gate.check_author_role_kind(
mutation_task, profile)
if ok: if ok:
return None return None
blocked = { blocked = {
@@ -18812,116 +19180,6 @@ def gitea_assess_master_parity(
return out return out
def _fleet_namespace_for_profile(profile_name: str | None) -> str | None:
"""Fleet namespace for *profile_name* (#949).
``role_namespace_gate.infer_mcp_namespace`` recognises only author and
reviewer and returns the profile name for everything else, so it cannot
name the controller, merger, or reconciler namespaces. The fleet roster
carries that mapping; this defers to it and keeps the existing inference as
the fallback for profiles outside the roster. Widening the shared helper is
controller role-metadata work owned by #950.
"""
return mcp_fleet_inventory.namespace_for_profile(
profile_name,
default=role_namespace_gate.infer_mcp_namespace(profile_name),
)
@mcp.tool()
def gitea_assess_fleet_inventory(
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
) -> dict:
"""Read-only: authoritative inventory of the running PRGS MCP fleet (#949).
Every other runtime surface is per-process. ``gitea_get_runtime_context``
and ``gitea_assess_master_parity`` describe only the server answering the
call; ``gitea_assess_mcp_namespace_health`` takes ``process`` and
``probe_result`` from the caller and so cannot constrain it. This tool takes
**no evidence parameters at all**: it combines the control-plane runtime
registry, which each server writes about itself at native transport bind,
with a process observation performed by this server. A member counts as
running only when both agree, so configuration alone never counts as a
running server and a caller cannot supply the answer.
Classification is a pure function of the snapshot, so ``gitea-controller``
and ``gitea-reconciler`` return the same verdict for the same fleet;
``answering_namespace`` is metadata and never changes it.
Fails closed. Any unreadable registry, unavailable process listing,
unregistered server process, unknown cohort, or unknown startup revision
sets ``inventory_complete`` false and ``mutation_gate_satisfied`` false with
a specific ``blocked_reason``. Matching Git revisions never establish a
single cohort ``single_cohort`` is derived only from recorded cohort
identity and stays ``null`` when that is unknown.
Strictly read-only: no restart, reconnect, lease mutation, issue mutation,
repository write, or process termination, and duplicates are reported and
never terminated. The only signal sent is ``signal 0`` liveness probing,
which delivers nothing to the target process. ``mutations_performed`` is
always an empty list.
Args:
remote: Known instance 'dadeschools' or 'prgs'. Declares which
repository binding fleet members are expected to carry.
host: Override the Gitea host.
org: Override the expected owner/organization binding.
repo: Override the expected repository binding.
Returns:
dict with 'inventory_complete', 'configured_members', 'running_members',
'missing_members', 'duplicate_members', 'unexpected_members',
'stale_members', 'unregistered_processes', 'single_cohort',
'mixed_cohort', 'mixed_revision', 'exactly_one_per_profile',
'mutation_gate_satisfied', 'blocked_reason', and per-member evidence.
"""
read_block = _profile_operation_gate("gitea.read")
if read_block:
return {
"success": False,
"read_only": True,
"inventory_complete": False,
"mutation_gate_satisfied": False,
"blocked_reason": "the active profile may not read control-plane state",
"reasons": read_block,
"permission_report": _permission_block_report("gitea.read"),
"mutations_performed": [],
}
ctx = session_ctx.get_session_context() or {}
expected_binding = {
"remote": remote or ctx.get("remote"),
"org": org or ctx.get("org"),
"repo": repo or ctx.get("repository"),
}
db, db_errors = _control_plane_db_or_error()
runtime_rows: list[dict] = []
registry_error: str | None = None
if db is None:
registry_error = "; ".join(db_errors) or "control-plane DB unavailable"
else:
try:
runtime_rows = db.list_mcp_server_runtimes(statuses=("running",))
except Exception as exc: # noqa: BLE001 - unreadable registry is unknown evidence
registry_error = f"runtime registry could not be read: {_redact(str(exc))}"
result = mcp_fleet_inventory.classify_fleet_inventory(
runtime_rows=runtime_rows,
process_scan=mcp_fleet_inventory.scan_mcp_server_processes(),
expected_binding=expected_binding,
registry_available=registry_error is None,
registry_error=registry_error,
answering_namespace=_fleet_namespace_for_profile(_active_profile_name(host)),
)
result["expected_repository_binding"] = expected_binding
result["summary"] = mcp_fleet_inventory.summarize(result)
return result
# #781: documented in the canonical review workflow as a tool reviewers call # #781: documented in the canonical review workflow as a tool reviewers call
# before workflow load, but the registration decorator had been lost, so the # before workflow load, but the registration decorator had been lost, so the
# documented inventory named something no namespace could reach. # documented inventory named something no namespace could reach.
@@ -24455,72 +24713,6 @@ def gitea_quarantine_contaminated_review(
} }
def _register_fleet_runtime(transport: str = "stdio") -> dict | None:
"""Record this server process in the control-plane runtime registry (#949).
Called once from the official entrypoint immediately after the native
transport bind, so the row can only ever describe the process writing it.
This is the evidence ``gitea_assess_fleet_inventory`` reads; without it the
fleet is unprovable, but a failure here must never prevent the server from
serving, so every error is swallowed after being logged to stderr.
"""
try:
profile_name = gitea_config.selected_profile_name()
try:
profile = get_profile()
except Exception: # noqa: BLE001 - identity resolution must not block boot
profile = {}
allowed = (profile or {}).get("allowed_operations") or []
forbidden = (profile or {}).get("forbidden_operations") or []
role = _role_kind(allowed, forbidden) if allowed else None
ctx = session_ctx.get_session_context() or {}
env = os.environ
provenance = (
"client_managed"
if (
(env.get("GITEA_CLIENT_MANAGED") or "").strip().lower()
in {"1", "true", "yes", "client_managed"}
or (env.get("GITEA_MCP_CLIENT_MANAGED") or "").strip().lower()
in {"1", "true", "yes", "client_managed"}
or (env.get("GITEA_SERVER_PROVENANCE") or "").strip()
== "client_managed"
)
else "manual_launch"
)
record = mcp_fleet_inventory.build_process_runtime_record(
namespace=_fleet_namespace_for_profile(profile_name),
profile=profile_name,
role=role,
remote=ctx.get("remote"),
org=ctx.get("org"),
repo=ctx.get("repository"),
repository_root=PROJECT_ROOT,
startup_head=_process_boot_head_sha,
daemon_start_head=_process_boot_head_sha,
transport=transport,
client_provenance=provenance,
env=env,
)
db, errors = _control_plane_db_or_error()
if db is None:
sys.stderr.write(
f"--- fleet runtime registration skipped: {'; '.join(errors)} ---\n"
)
return None
db.register_mcp_server_runtime(
record,
retention_seconds=mcp_fleet_inventory.RUNTIME_RETENTION_SECONDS,
)
sys.stderr.write(
f"--- fleet runtime registered: {record['runtime_id']} "
f"(cohort {record['cohort_id']}) ---\n"
)
return record
except Exception as exc: # noqa: BLE001 - never block startup
sys.stderr.write(f"--- fleet runtime registration failed: {exc} ---\n")
return None
# ── Entry point ─────────────────────────────────────────────────────────────── # ── Entry point ───────────────────────────────────────────────────────────────
if __name__ == "__main__": if __name__ == "__main__":
@@ -24529,12 +24721,11 @@ if __name__ == "__main__":
# basename-only stack frames, and import-only launch cannot reconstruct # basename-only stack frames, and import-only launch cannot reconstruct
# native transport; offline imports / standalone scripts fail closed. # native transport; offline imports / standalone scripts fail closed.
mcp_daemon_guard.mark_sanctioned_daemon() mcp_daemon_guard.mark_sanctioned_daemon()
mcp_daemon_guard.bind_native_mcp_transport(transport="stdio") # #931: the transport is no longer a literal here. It comes from deployment
# #949: record this process in the control-plane runtime registry. Only a # configuration (GITEA_MCP_TRANSPORT), defaults to stdio when unset, and is
# transport-bound server reaches this point, so the row is authoritative # validated against the single permitted set in mcp_transport_config. An
# evidence that this namespace is actually running — evidence no caller of # unregistered identifier raises here, before any tool can dispatch.
# gitea_assess_fleet_inventory can supply. mcp_daemon_guard.bind_native_mcp_transport()
_register_fleet_runtime(transport="stdio")
# Lock this session's launch profile into the environment so child CLI # Lock this session's launch profile into the environment so child CLI
# processes (e.g. review_pr.py) can detect and refuse profile # processes (e.g. review_pr.py) can detect and refuse profile
# side-channel overrides (#199). # side-channel overrides (#199).
@@ -24552,4 +24743,11 @@ if __name__ == "__main__":
sys.stderr.write( sys.stderr.write(
f"--- Sentry observability: {_sentry_status.get('reason')} ---\n" f"--- Sentry observability: {_sentry_status.get('reason')} ---\n"
) )
mcp.run(transport="stdio") # #931: serve over exactly the transport that was bound and pinned, and only
# if this entrypoint is commissioned to execute it. Recognition is not
# execution authorization (review 635): a registered remote identifier is
# bound, pinned and recorded, but serving it would start a listener with no
# authentication or per-request principal, which #938 owns. An absent bind
# keeps its pre-existing #695 failure; an uncommissioned transport fails
# closed here, before any listener exists and before any tool can dispatch.
mcp.run(transport=mcp_daemon_guard.authorize_transport_execution("tool service"))
+6
View File
@@ -500,6 +500,11 @@ def assess_transport_for_auth_mint() -> dict[str, Any]:
native = mcp_daemon_guard.is_native_mcp_transport() native = mcp_daemon_guard.is_native_mcp_transport()
pytest = mcp_daemon_guard.is_pytest_runtime() pytest = mcp_daemon_guard.is_pytest_runtime()
production = mcp_daemon_guard.is_production_native_mcp_transport() production = mcp_daemon_guard.is_production_native_mcp_transport()
# #931: report which transport underwrites the verdict, read through the
# one shared accessor rather than assumed to be stdio. The gate's decision
# is unchanged here; naming the transport is what lets #932 re-derive the
# guarantee from an authenticated session instead of from the bind.
bound = mcp_daemon_guard.bound_transport()
if not native and not pytest: if not native and not pytest:
reasons.append( reasons.append(
"irrecoverable provenance authorization requires production native " "irrecoverable provenance authorization requires production native "
@@ -510,6 +515,7 @@ def assess_transport_for_auth_mint() -> dict[str, Any]:
"allowed": not reasons, "allowed": not reasons,
"native_mcp_transport": native, "native_mcp_transport": native,
"production_native_mcp_transport": production, "production_native_mcp_transport": production,
"bound_transport": bound,
"pytest": pytest, "pytest": pytest,
"reasons": reasons, "reasons": reasons,
} }
+128 -9
View File
@@ -299,9 +299,13 @@ def _ownership_refusals(
f"lock worktree '{lock.get('worktree_path')}' does not match " f"lock worktree '{lock.get('worktree_path')}' does not match "
f"'{worktree_path}'" f"'{worktree_path}'"
) )
lease = lock.get("work_lease") if isinstance(lock, dict) else None # #953 AC2/AC13/AC14: read through the shared claimant reader so a lock
claimant = lease.get("claimant") if isinstance(lease, dict) else None # written by bootstrap — which records the claimant at the top level — is
claimant = claimant if isinstance(claimant, dict) else {} # not refused for "not recording a claimant" when it plainly records one.
# This is not a widening: the values are still compared against the
# server-resolved identity and profile immediately below, so a legacy
# placement grants nothing that the canonical placement would not.
claimant = lock_claimant(lock) if isinstance(lock, dict) else {}
recorded_identity = str(claimant.get("username") or "").strip() recorded_identity = str(claimant.get("username") or "").strip()
recorded_profile = str(claimant.get("profile") or "").strip() recorded_profile = str(claimant.get("profile") or "").strip()
if not recorded_identity or not recorded_profile: if not recorded_identity or not recorded_profile:
@@ -636,6 +640,99 @@ def iter_lock_files(lock_dir: str | None = None) -> list[str]:
return sorted(paths) return sorted(paths)
def release_session_lock(
*,
issue_number: int,
session: str,
lock_dir: str | None = None,
remote: str | None = None,
org: str | None = None,
repo: str | None = None,
) -> str:
"""Remove exactly the durable lock *session* created for *issue_number*.
``author_issue_bootstrap.run_compensating_recovery`` has called this name
since #850, but it was never defined: the call raised ``AttributeError``
into a bare ``except Exception: pass``, so the lock half of every
compensating rollback silently did nothing. The branch and worktree were
removed and the lock was left behind — a state no sanctioned tool can act
on, since recovery refuses ``worktree_invalid`` and ``gitea_lock_issue`` has
no worktree to bind (#953 review 632 F2).
Ownership is proven, not asserted. A record is removed only when its
recorded ``owner_session`` equals *session* and its issue number matches;
``remote``/``org``/``repo`` narrow it further when supplied. Zero matches or
more than one both raise, so a caller can never delete a lock it does not
own and an ambiguous directory is never guessed at. The ``.json.lock`` flock
sidecar is deliberately left in place — it is a zero-byte mutex another
process may hold, and removing it under contention would be a race.
Returns the removed lock file path.
"""
target_issue = int(issue_number)
owner = str(session or "").strip()
if not owner:
raise ValueError(
"release_session_lock requires the owning session id (fail closed)"
)
def _is_owned_durable_lock(record: dict[str, Any] | None) -> bool:
# A durable lock, not a bootstrap phase journal or a session pointer,
# both of which can share a directory and carry the same issue number
# and owner_session.
if not record or "lock_generation" not in record:
return False
if not str(record.get("branch_name") or "").strip():
return False
if not str(record.get("worktree_path") or "").strip():
return False
try:
if int(record.get("issue_number") or 0) != target_issue:
return False
except (TypeError, ValueError):
return False
return str(record.get("owner_session") or "").strip() == owner
# Prefer the exact keyed path when the caller knows the repository; scanning
# is the fallback for callers that only carry the issue number.
if remote and org and repo:
exact = lock_file_path(
remote=remote,
org=org,
repo=repo,
issue_number=target_issue,
lock_dir=lock_dir,
)
if not _is_owned_durable_lock(read_lock_file(exact)):
raise FileNotFoundError(
f"durable issue lock '{exact}' is absent or is not owned by "
f"session '{owner}' (fail closed; nothing released)"
)
os.remove(exact)
return exact
matches: list[str] = []
for path in iter_lock_files(lock_dir):
if _is_owned_durable_lock(read_lock_file(path)):
matches.append(path)
if not matches:
raise FileNotFoundError(
f"no durable issue lock for issue #{target_issue} is owned by "
f"session '{owner}' (fail closed; nothing released)"
)
if len(matches) > 1:
raise RuntimeError(
f"{len(matches)} durable locks for issue #{target_issue} claim "
f"session '{owner}'; refusing to guess which to release "
"(fail closed)"
)
path = matches[0]
os.remove(path)
return path
def find_lock_for_branch( def find_lock_for_branch(
*, *,
remote: str, remote: str,
@@ -1112,21 +1209,43 @@ def assess_same_issue_lease_conflict(
) )
def _lock_claimant(lock: dict[str, Any] | None) -> dict[str, str]: def lock_claimant(lock: dict[str, Any] | None) -> dict[str, str]:
"""Read the claimant from either canonical or legacy placement (#953 AC13/AC14).
``work_lease.claimant`` is the canonical placement and is preferred; a
top-level ``claimant`` is the legacy/bootstrap placement and is accepted as
a fallback. This is the single definition. Before #953 the readers
disagreed: this module, ``issue_lock_renewal``, and ``issue_lock_recovery``
tolerated both placements, while ``_ownership_refusals`` looked only in
``work_lease`` — which is what made a bootstrap-written lock
un-heartbeatable.
Preferring ``work_lease`` over the top level is deliberate: once a legacy
lock is upgraded, the canonical placement is authoritative and a stale
top-level copy must never win.
This decides *where to look*, never whether ownership is proven — every
caller still compares these values against server-resolved identity and
profile.
"""
if not isinstance(lock, dict): if not isinstance(lock, dict):
return {} return {}
claimant = lock.get("claimant") lease = lock.get("work_lease")
claimant = lease.get("claimant") if isinstance(lease, dict) else None
if not isinstance(claimant, dict): if not isinstance(claimant, dict):
lease = lock.get("work_lease") claimant = lock.get("claimant")
claimant = lease.get("claimant") if isinstance(lease, dict) else None
if not isinstance(claimant, dict): if not isinstance(claimant, dict):
return {} return {}
return { return {
"username": str(claimant.get("username") or ""), "username": str(claimant.get("username") or "").strip(),
"profile": str(claimant.get("profile") or ""), "profile": str(claimant.get("profile") or "").strip(),
} }
#: Back-compatible alias for the pre-#953 private name.
_lock_claimant = lock_claimant
def assess_foreign_lock_overwrite( def assess_foreign_lock_overwrite(
existing_lock: dict[str, Any] | None, existing_lock: dict[str, Any] | None,
incoming_lock: dict[str, Any], incoming_lock: dict[str, Any],
+176 -13
View File
@@ -32,6 +32,8 @@ import time
from pathlib import Path from pathlib import Path
from typing import Any from typing import Any
import mcp_transport_config
SANCTIONED_DAEMON_ENV = "GITEA_MCP_SANCTIONED_DAEMON" SANCTIONED_DAEMON_ENV = "GITEA_MCP_SANCTIONED_DAEMON"
ALLOW_DIRECT_IMPORT_ENV = "GITEA_ALLOW_DIRECT_MCP_IMPORT" ALLOW_DIRECT_IMPORT_ENV = "GITEA_ALLOW_DIRECT_MCP_IMPORT"
ALLOW_KEYCHAIN_CLI_ENV = "GITEA_ALLOW_KEYCHAIN_CLI" ALLOW_KEYCHAIN_CLI_ENV = "GITEA_ALLOW_KEYCHAIN_CLI"
@@ -42,7 +44,9 @@ FORCE_PROVENANCE_FAIL_ENV = "GITEA_TEST_FORCE_UNSANCTIONED"
_NATIVE_RUNTIME: dict[str, Any] | None = None _NATIVE_RUNTIME: dict[str, Any] | None = None
# Production transport identifiers accepted by bind_native_mcp_transport. # Production transport identifiers accepted by bind_native_mcp_transport.
_PRODUCTION_TRANSPORTS = frozenset({"stdio"}) # #931: the permitted set is defined once, in mcp_transport_config. This name
# is kept as an alias so the guard never restates a transport identifier.
_PRODUCTION_TRANSPORTS = mcp_transport_config.SUPPORTED_TRANSPORTS
_RUNTIME_MODE_PRODUCTION = "production" _RUNTIME_MODE_PRODUCTION = "production"
_RUNTIME_MODE_TEST = "test" _RUNTIME_MODE_TEST = "test"
_PHASE_ENTRYPOINT_CLAIMED = "entrypoint_claimed" _PHASE_ENTRYPOINT_CLAIMED = "entrypoint_claimed"
@@ -58,6 +62,23 @@ class UnsanctionedRuntimeError(RuntimeError):
"""Raised when mutation/credential code runs outside a native MCP daemon.""" """Raised when mutation/credential code runs outside a native MCP daemon."""
class TransportExecutionError(UnsanctionedRuntimeError):
"""Raised when a bound transport may not be served by this entrypoint (#931).
Subclasses :class:`UnsanctionedRuntimeError` so every existing fail-closed
handler still catches it, while letting a caller that cares distinguish
"nothing is bound" from "something valid is bound but its listener has not
been commissioned". Carries the structured verdict on ``.assessment``.
"""
def __init__(self, message: str, assessment: dict[str, Any] | None = None):
super().__init__(message)
self.assessment = assessment or {}
self.blocker_kind = self.assessment.get("blocker_kind")
self.owner_issue = self.assessment.get("owner_issue")
self.transport = self.assessment.get("transport")
def is_pytest_runtime() -> bool: def is_pytest_runtime() -> bool:
if (os.environ.get(FORCE_PROVENANCE_FAIL_ENV) or "").strip() in { if (os.environ.get(FORCE_PROVENANCE_FAIL_ENV) or "").strip() in {
"1", "1",
@@ -171,23 +192,46 @@ def mark_sanctioned_daemon() -> dict[str, Any]:
return native_runtime_status() return native_runtime_status()
def bind_native_mcp_transport(*, transport: str) -> dict[str, Any]: def bind_native_mcp_transport(*, transport: str | None = None) -> dict[str, Any]:
"""Bind the live native MCP transport lifecycle (#695). """Bind the live native MCP transport lifecycle (#695 / #931).
Must be called from the resolved canonical entrypoint immediately before Must be called from the resolved canonical entrypoint immediately before
the real MCP server transport loop (e.g. ``mcp.run(transport=\"stdio\")``). the real MCP server transport loop (``mcp.run``). Requires a prior
Requires a prior successful :func:`mark_sanctioned_daemon` claim in this successful :func:`mark_sanctioned_daemon` claim in this process.
process. Import-only or offline launch without this bind leaves Import-only or offline launch without this bind leaves
:func:`is_native_mcp_transport` false. :func:`is_native_mcp_transport` false.
#931: ``transport`` is now optional. Omitting it — which is what the
production entrypoint does resolves the identifier from deployment
configuration via :func:`mcp_transport_config.resolve_configured_transport`,
yielding :data:`mcp_transport_config.DEFAULT_TRANSPORT` when nothing is
configured. An explicit argument remains supported for tests and for a
launcher that has already resolved the value. Either way the identifier is
validated against the single permitted set before the runtime record is
written, so no tool can dispatch over an unregistered transport.
The resolved value is pinned into the process-local record and is read back
only through :func:`bound_transport`. Rebinding to a different transport is
refused, so two guards can never observe different values in one process.
""" """
global _NATIVE_RUNTIME global _NATIVE_RUNTIME
transport_name = (transport or "").strip().lower() if transport is None:
if transport_name not in _PRODUCTION_TRANSPORTS: resolution = mcp_transport_config.resolve_configured_transport()
raise UnsanctionedRuntimeError( transport_name = str(resolution["transport"])
f"bind_native_mcp_transport rejected: transport {transport!r} is " if not resolution["supported"]:
f"not a production MCP transport (#695). Allowed: " raise UnsanctionedRuntimeError(
f"{sorted(_PRODUCTION_TRANSPORTS)}." "bind_native_mcp_transport rejected: "
) + "; ".join(resolution["reasons"])
+ " No tool is served over an unregistered transport."
)
else:
transport_name = mcp_transport_config.normalize_transport(transport)
if transport_name not in _PRODUCTION_TRANSPORTS:
raise UnsanctionedRuntimeError(
f"bind_native_mcp_transport rejected: transport {transport!r} is "
f"not a production MCP transport (#695). Allowed: "
f"{sorted(_PRODUCTION_TRANSPORTS)}."
)
entrypoint_path = _caller_official_entrypoint_path() entrypoint_path = _caller_official_entrypoint_path()
if entrypoint_path is None: if entrypoint_path is None:
@@ -216,6 +260,23 @@ def bind_native_mcp_transport(*, transport: str) -> dict[str, Any]:
"between mark and bind (#695)." "between mark and bind (#695)."
) )
# #931: one process binds one transport. Re-binding the same identifier is
# idempotent (a retried launch step must not fail); re-binding a different
# one is refused, because a guard that already read the first value would
# otherwise disagree with a guard that reads the second.
already_bound = (_NATIVE_RUNTIME.get("transport") or "").strip()
if (
already_bound
and _NATIVE_RUNTIME.get("phase") == _PHASE_TRANSPORT_BOUND
and already_bound != transport_name
):
raise UnsanctionedRuntimeError(
"bind_native_mcp_transport rejected: transport is already bound to "
f"{already_bound!r} in this process; rebinding to "
f"{transport_name!r} is forbidden (#931). Restart the server to "
"change the deployment transport."
)
# Pin session-state root for this server lifetime (#695 AC2 / PR #701). # Pin session-state root for this server lifetime (#695 AC2 / PR #701).
# Changing GITEA_MCP_SESSION_STATE_DIR after bind must not manufacture a # Changing GITEA_MCP_SESSION_STATE_DIR after bind must not manufacture a
# second authority domain for decision locks / workflow proofs. # second authority domain for decision locks / workflow proofs.
@@ -353,6 +414,88 @@ def is_production_native_mcp_transport() -> bool:
return (_NATIVE_RUNTIME or {}).get("mode") == _RUNTIME_MODE_PRODUCTION return (_NATIVE_RUNTIME or {}).get("mode") == _RUNTIME_MODE_PRODUCTION
def bound_transport() -> str | None:
"""The one authoritative bound transport identifier, or ``None`` (#931).
This is the shared accessor every transport-aware guard reads. It reports
the value pinned at bind time, never the environment, so changing
``GITEA_MCP_TRANSPORT`` after the bind cannot move what a guard observes
the same rule :func:`pinned_session_state_dir` applies to session state.
``None`` means unbound: an offline import or a launch that never reached
the bind. Callers must treat that as fail-closed, exactly as they already
treat :func:`is_native_mcp_transport` returning false.
"""
if not is_native_mcp_transport():
return None
return (_NATIVE_RUNTIME or {}).get("transport") or None
def assert_transport_bound(context: str = "tool service") -> str:
"""Return the bound transport, or fail closed before *context* (#931).
Called immediately before the server enters its transport loop so an
invalid or absent bind stops the process rather than serving tools over a
transport no guard can name.
"""
transport = bound_transport()
if transport:
return transport
raise UnsanctionedRuntimeError(
f"No MCP transport is bound; refusing {context} (#931). "
"bind_native_mcp_transport must succeed from the canonical entrypoint "
"before any tool is served. Offline import and standalone launch "
"cannot reconstruct a bind."
)
def assess_serve_authorization() -> dict[str, Any]:
"""Structured verdict on whether this process may serve tools (#931).
This is the decision that consumes :func:`bound_transport`. It is what stops
the bound identifier from being reporting-only metadata: the serve path
cannot proceed unless the value pinned at bind is one this entrypoint is
commissioned to execute.
Never raises; returns the verdict so callers and diagnostics can inspect it.
"""
return mcp_transport_config.assess_transport_execution(bound_transport())
def authorize_transport_execution(context: str = "tool service") -> str:
"""Return the transport this process may serve, or fail closed (#931).
Two distinct boundaries, in order:
1. **Bind presence** :func:`assert_transport_bound` enforces the
pre-existing #695 contract, so an unbound runtime keeps its established
failure and reason code.
2. **Execution authorization** the bound identifier must be one this
entrypoint is commissioned to serve. A registered transport whose
listener has not been commissioned is refused here, before any listener
is created and before any tool can dispatch.
That ordering matters: recognition, validation and durable recording all
still happen for a remote identifier, so #931's seam is intact; only the act
of *serving* it is withheld until its owning issue commissions it.
"""
# Boundary 1: unbound stays exactly as fail-closed as it was under #695.
assert_transport_bound(context)
# Boundary 2: bound, but is this entrypoint allowed to serve it?
assessment = assess_serve_authorization()
if assessment.get("allowed"):
return str(assessment["transport"])
reasons = "; ".join(assessment.get("reasons") or []) or "not authorized"
next_action = assessment.get("exact_next_action") or ""
raise TransportExecutionError(
f"Refusing {context} (#931) [{assessment.get('blocker_kind')}]: "
f"{reasons} {next_action}".strip(),
assessment,
)
def is_sanctioned_mcp_daemon() -> bool: def is_sanctioned_mcp_daemon() -> bool:
"""Backward-compatible name; #695 requires native transport, not env alone.""" """Backward-compatible name; #695 requires native transport, not env alone."""
if is_production_native_mcp_transport(): if is_production_native_mcp_transport():
@@ -476,6 +619,21 @@ def native_runtime_status() -> dict[str, Any]:
"entrypoint_path": rt.get("entrypoint_path"), "entrypoint_path": rt.get("entrypoint_path"),
"phase": rt.get("phase"), "phase": rt.get("phase"),
"transport": rt.get("transport"), "transport": rt.get("transport"),
# #931: the authoritative bound identifier, plus the seam that defines
# what may be bound. ``bound_transport`` is None until a bind succeeds,
# so an offline import is distinguishable from a stdio session.
"bound_transport": bound_transport(),
"transport_bound": bound_transport() is not None,
"default_transport": mcp_transport_config.DEFAULT_TRANSPORT,
"supported_transports": list(mcp_transport_config.supported_transports()),
"transport_env": mcp_transport_config.TRANSPORT_ENV,
# #931 review 635: recognition and execution authorization are distinct.
# ``supported`` is what may be bound; ``executable`` is what this
# entrypoint may actually serve. A recognized-but-uncommissioned
# transport reports serve_authorized False with a named blocker.
"executable_transports": list(mcp_transport_config.executable_transports()),
"serve_authorized": bool(assess_serve_authorization().get("allowed")),
"serve_authorization": assess_serve_authorization(),
"mode": rt.get("mode"), "mode": rt.get("mode"),
"session_state_dir": pinned_session_state_dir() or rt.get("session_state_dir"), "session_state_dir": pinned_session_state_dir() or rt.get("session_state_dir"),
"session_state_dir_pinned": pinned_session_state_dir() is not None, "session_state_dir_pinned": pinned_session_state_dir() is not None,
@@ -502,7 +660,12 @@ def mutation_provenance_fields() -> dict[str, Any]:
if st.get("mode") == _RUNTIME_MODE_TEST and st["native_mcp_transport"]: if st.get("mode") == _RUNTIME_MODE_TEST and st["native_mcp_transport"]:
transport = "test_native_mcp" transport = "test_native_mcp"
return { return {
# ``transport`` stays the trust *class* it has always been, so existing
# durable records keep their shape. ``bound_transport`` (#931) adds the
# bound identifier itself, which is what lets an operator tell from a
# durable record which transport performed a mutation.
"transport": transport, "transport": transport,
"bound_transport": st.get("bound_transport"),
"native_mcp_transport": bool(st["native_mcp_transport"]), "native_mcp_transport": bool(st["native_mcp_transport"]),
"production_native_mcp_transport": bool( "production_native_mcp_transport": bool(
st.get("production_native_mcp_transport") st.get("production_native_mcp_transport")
-785
View File
@@ -1,785 +0,0 @@
"""Authoritative, read-only PRGS MCP fleet inventory (#949).
Every pre-existing runtime surface is *per process*. ``gitea_get_runtime_context``
and ``gitea_assess_master_parity`` describe only the server answering the call.
``gitea_assess_mcp_namespace_health`` accepts ``process``, ``probe_result`` and
``registered_tools`` **from the caller**, so it cannot constrain the caller. The
control-plane ``sessions`` table records allocator *task* sessions, not server
processes. Five independent self-reports of the same revision therefore never
proved that exactly five processes exist, that no sixth exists, or that all five
belong to one client cohort.
Evidence model
--------------
Two independent sources must agree before a fleet member counts as running:
``control-plane runtime registry``
A row each server writes **about itself** at native transport bind
(:func:`build_process_runtime_record`). No caller can supply it. It is
authoritative for identity: namespace, profile, role, repository binding,
cohort, and the revision the process started at.
``server-side process observation``
A process listing performed by the server answering the inventory call
(:func:`scan_mcp_server_processes`), never by the caller. It is
authoritative for existence and liveness, and it is the only source that
can show a process the registry does not know about.
A member is ``live`` only when a registry row has a matching, still-running
process whose start time precedes the registration (so a recycled PID cannot
impersonate a dead server). Anything the two sources cannot jointly establish
is reported as unknown and fails the mutation gate closed configuration alone
never counts as a running member, and matching Git revisions never establish a
single cohort.
This module performs no restart, reconnect, lease mutation, issue mutation, or
process termination. The only signal it ever sends is ``signal 0`` liveness
probing, which delivers nothing to the target process.
"""
from __future__ import annotations
import os
import secrets
import subprocess
from datetime import datetime, timezone
from typing import Any, Iterable, Mapping, Sequence
# ── expected fleet ────────────────────────────────────────────────────────────
# The configured PRGS fleet. Each entry is one expected member; the roster is
# the definition of "expected" for missing/unexpected classification.
EXPECTED_PRGS_FLEET: tuple[dict[str, str], ...] = (
{"namespace": "gitea-author", "profile": "prgs-author", "role": "author"},
{
"namespace": "gitea-controller",
"profile": "prgs-controller",
"role": "controller",
},
{"namespace": "gitea-reviewer", "profile": "prgs-reviewer", "role": "reviewer"},
{"namespace": "gitea-merger", "profile": "prgs-merger", "role": "merger"},
{
"namespace": "gitea-reconciler",
"profile": "prgs-reconciler",
"role": "reconciler",
},
)
# Cohort identity supplied by a client that manages the whole fleet.
COHORT_ID_ENV = "GITEA_MCP_CLIENT_COHORT_ID"
# Registry rows older than this are pruned at *registration* time (a startup
# write), never on the read path. Dead rows inside the window are still reported
# as stale evidence rather than silently dropped.
RUNTIME_RETENTION_SECONDS = 7 * 24 * 3600
# Liveness classifications for a registry row.
LIVENESS_LIVE = "live"
LIVENESS_DEAD = "dead"
LIVENESS_PID_RECYCLED = "pid_recycled"
LIVENESS_UNOBSERVED = "unobserved"
LIVENESS_UNKNOWN = "unknown"
# Per-member health classifications.
HEALTH_RUNNING = "running"
HEALTH_MISSING = "missing"
HEALTH_DUPLICATE = "duplicate"
HEALTH_UNEXPECTED = "unexpected"
HEALTH_STALE = "stale"
HEALTH_UNKNOWN = "unknown"
EVIDENCE_AUTHORITY = "control_plane_runtime_registry+server_process_observation"
_MCP_PROCESS_MARKER = "mcp_server.py"
_LSTART_FORMAT = "%a %b %d %H:%M:%S %Y"
_ISO_FORMAT = "%Y-%m-%dT%H:%M:%SZ"
# ── time helpers ──────────────────────────────────────────────────────────────
def _utcnow() -> datetime:
return datetime.now(timezone.utc)
def iso_now() -> str:
return _utcnow().strftime(_ISO_FORMAT)
def _parse_iso(value: Any) -> datetime | None:
text = (str(value) if value is not None else "").strip()
if not text:
return None
if text.endswith("Z"):
text = text[:-1] + "+00:00"
try:
parsed = datetime.fromisoformat(text)
except ValueError:
return None
if parsed.tzinfo is None:
parsed = parsed.replace(tzinfo=timezone.utc)
return parsed.astimezone(timezone.utc)
def _int_or_none(value: Any) -> int | None:
try:
return int(value)
except (TypeError, ValueError):
return None
def _clean(value: Any) -> str | None:
text = (str(value) if value is not None else "").strip()
return text or None
# ── process-level evidence (server side only) ─────────────────────────────────
def probe_pid_alive(pid: int | None) -> bool | None:
"""Return whether *pid* exists. ``None`` when it cannot be determined.
Uses ``signal 0``, which performs a permission/existence check and delivers
nothing to the target. This module never sends a terminating signal.
"""
resolved = _int_or_none(pid)
if resolved is None or resolved <= 0:
return None
try:
os.kill(resolved, 0)
except ProcessLookupError:
return False
except PermissionError:
# The process exists but belongs to another user.
return True
except OSError:
return None
return True
def scan_mcp_server_processes(*, runner=subprocess.run) -> dict[str, Any]:
"""Observe running Gitea MCP server processes from this server process.
This is deliberately performed by the answering server, never by the caller:
a caller-supplied process list is exactly the input a fleet gate must not
trust. When the listing cannot be obtained the result reports
``available=False`` so the inventory fails closed instead of assuming that
no unregistered process exists.
"""
try:
proc = runner(
["ps", "-o", "pid,lstart,command", "-ax"],
capture_output=True,
text=True,
check=True,
)
except Exception as exc: # noqa: BLE001 - any failure means unknown evidence
return {
"available": False,
"processes": [],
"reason": f"process listing unavailable: {exc}",
}
processes: list[dict[str, Any]] = []
for raw_line in (proc.stdout or "").splitlines()[1:]:
line = raw_line.strip()
if not line or _MCP_PROCESS_MARKER not in line:
continue
parts = line.split(None, 6)
if len(parts) < 7:
continue
pid = _int_or_none(parts[0])
if pid is None:
continue
try:
naive = datetime.strptime(" ".join(parts[1:6]), _LSTART_FORMAT)
started_at = naive.astimezone(timezone.utc)
except (ValueError, OSError):
started_at = None
processes.append(
{
"pid": pid,
"started_at": started_at.strftime(_ISO_FORMAT) if started_at else None,
"command": parts[6],
}
)
processes.sort(key=lambda item: item["pid"])
return {"available": True, "processes": processes, "reason": None}
# ── cohort / registration record ──────────────────────────────────────────────
def namespace_for_profile(profile: str | None, *, default: str | None = None) -> str | None:
"""Map a configured profile to its fleet namespace.
Resolved from the expected roster rather than from name-shape heuristics.
``role_namespace_gate.infer_mcp_namespace`` only recognises author and
reviewer, so it returns the *profile* name for controller, merger, and
reconciler which would label three of the five members with a namespace
that does not exist. Widening that helper is controller role-metadata work
and belongs to #950; the roster already carries the mapping this inventory
needs, so it is read from there.
A profile outside the roster falls back to *default* (typically the
caller's existing inference), so an unexpected member is still described
rather than dropped.
"""
cleaned = _clean(profile)
for entry in EXPECTED_PRGS_FLEET:
if entry["profile"] == cleaned:
return entry["namespace"]
return default if default is not None else cleaned
def derive_cohort_identity(env: Mapping[str, str] | None = None) -> dict[str, Any]:
"""Derive the client/cohort identity of this server process.
A cohort is the set of servers a single client launched together. The
parent process is the durable expression of that: an IDE/CLI client spawns
every namespace as its own child. A process whose parent has gone away
(reparented to init) cannot prove which cohort it belongs to, and says so
rather than guessing.
Revisions are deliberately not consulted here. Two servers built from the
same commit are not thereby one cohort (#949 AC7).
"""
source_env = os.environ if env is None else env
explicit = _clean(source_env.get(COHORT_ID_ENV))
if explicit:
return {"cohort_id": explicit, "cohort_source": "explicit_env"}
try:
ppid = os.getppid()
except OSError:
ppid = 0
if ppid and ppid > 1:
return {"cohort_id": f"ppid:{ppid}", "cohort_source": "parent_process"}
return {
"cohort_id": None,
"cohort_source": "unknown",
"cohort_reason": (
"parent process is unavailable or reparented to init; this server "
"cannot prove which client cohort launched it"
),
}
def build_process_runtime_record(
*,
namespace: str,
profile: str | None,
role: str | None,
remote: str | None = None,
org: str | None = None,
repo: str | None = None,
repository_root: str | None = None,
pid: int | None = None,
startup_head: str | None = None,
daemon_start_head: str | None = None,
transport: str | None = None,
client_provenance: str | None = None,
env: Mapping[str, str] | None = None,
boot_id: str | None = None,
registered_at: str | None = None,
) -> dict[str, Any]:
"""Build the row a server writes about itself at native transport bind.
Every field describes the *calling* process. Nothing here is caller-supplied
in the MCP sense: the only code that reaches this function is the official
entrypoint of the process being described.
"""
cohort = derive_cohort_identity(env)
resolved_pid = _int_or_none(pid)
if resolved_pid is None:
resolved_pid = os.getpid()
token = boot_id or secrets.token_hex(8)
return {
"runtime_id": f"{namespace}:{resolved_pid}:{token}",
"namespace": namespace,
"profile": _clean(profile),
"role": _clean(role),
"remote": _clean(remote),
"org": _clean(org),
"repo": _clean(repo),
"repository_root": _clean(repository_root),
"pid": resolved_pid,
"cohort_id": cohort["cohort_id"],
"cohort_source": cohort["cohort_source"],
"client_provenance": _clean(client_provenance) or "unknown",
"boot_id": token,
"startup_head": _clean(startup_head),
"daemon_start_head": _clean(daemon_start_head),
"transport": _clean(transport),
"registered_at": registered_at or iso_now(),
"status": "running",
}
# ── classification ────────────────────────────────────────────────────────────
def _expected_index(
expected_fleet: Sequence[Mapping[str, str]],
) -> dict[str, dict[str, str]]:
index: dict[str, dict[str, str]] = {}
for entry in expected_fleet:
profile = _clean(entry.get("profile"))
if profile:
index[profile] = dict(entry)
return index
def _sort_key(member: Mapping[str, Any]) -> tuple:
return (
str(member.get("namespace") or ""),
str(member.get("profile") or ""),
_int_or_none(member.get("pid")) or 0,
str(member.get("runtime_id") or ""),
)
def _normalize_row(
row: Mapping[str, Any],
*,
observed_by_pid: Mapping[int, Mapping[str, Any]],
process_scan_available: bool,
now: datetime,
) -> dict[str, Any]:
pid = _int_or_none(row.get("pid"))
member: dict[str, Any] = {
"runtime_id": _clean(row.get("runtime_id")),
"namespace": _clean(row.get("namespace")),
"profile": _clean(row.get("profile")),
"role": _clean(row.get("role")),
"remote": _clean(row.get("remote")),
"org": _clean(row.get("org")),
"repo": _clean(row.get("repo")),
"repository_root": _clean(row.get("repository_root")),
"pid": pid,
"cohort_id": _clean(row.get("cohort_id")),
"cohort_source": _clean(row.get("cohort_source")) or "unknown",
"client_provenance": _clean(row.get("client_provenance")) or "unknown",
"boot_id": _clean(row.get("boot_id")),
"startup_head": _clean(row.get("startup_head")),
"daemon_start_head": _clean(row.get("daemon_start_head")),
"transport": _clean(row.get("transport")),
"registered_at": _clean(row.get("registered_at")),
"last_heartbeat_at": _clean(row.get("last_heartbeat_at")),
"recorded_status": _clean(row.get("status")) or "unknown",
}
pid_alive = probe_pid_alive(pid)
member["pid_alive"] = pid_alive
observed = observed_by_pid.get(pid) if pid is not None else None
member["process_observed"] = bool(observed) if process_scan_available else None
if not process_scan_available:
# Existence cannot be corroborated; never upgrade to live on the
# registry's word alone.
member["liveness"] = LIVENESS_UNKNOWN
member["liveness_reason"] = (
"process observation unavailable; registry rows cannot be corroborated"
)
elif pid_alive is False:
member["liveness"] = LIVENESS_DEAD
member["liveness_reason"] = "recorded PID is not running"
elif pid_alive is None:
member["liveness"] = LIVENESS_UNKNOWN
member["liveness_reason"] = "PID liveness could not be determined"
elif observed is None:
member["liveness"] = LIVENESS_UNOBSERVED
member["liveness_reason"] = (
"recorded PID is not a running Gitea MCP server process"
)
else:
started_at = _parse_iso(observed.get("started_at"))
registered_at = _parse_iso(member["registered_at"])
if started_at and registered_at and started_at > registered_at:
member["liveness"] = LIVENESS_PID_RECYCLED
member["liveness_reason"] = (
"the process now holding this PID started after the registry row "
"was written; the registered server is gone"
)
else:
member["liveness"] = LIVENESS_LIVE
member["liveness_reason"] = None
heartbeat = _parse_iso(member["last_heartbeat_at"])
member["heartbeat_age_seconds"] = (
int((now - heartbeat).total_seconds()) if heartbeat else None
)
return member
def _binding_matches(
member: Mapping[str, Any], expected_binding: Mapping[str, Any] | None
) -> bool | None:
if not expected_binding:
return None
for field in ("remote", "org", "repo"):
expected = _clean(expected_binding.get(field))
if expected is None:
continue
actual = _clean(member.get(field))
if actual is None:
return None
if actual != expected:
return False
return True
def classify_fleet_inventory(
*,
runtime_rows: Iterable[Mapping[str, Any]],
process_scan: Mapping[str, Any] | None = None,
expected_fleet: Sequence[Mapping[str, str]] = EXPECTED_PRGS_FLEET,
expected_binding: Mapping[str, Any] | None = None,
registry_available: bool = True,
registry_error: str | None = None,
now: datetime | None = None,
answering_namespace: str | None = None,
) -> dict[str, Any]:
"""Classify a fleet snapshot. Pure: identical input yields identical output.
The verdict never depends on which namespace asked, so controller and
reconciler agree by construction; ``answering_namespace`` is reported as
metadata only.
"""
moment = now or _utcnow()
scan = dict(process_scan or {"available": False, "processes": [], "reason": None})
scan_available = bool(scan.get("available"))
observed_processes = list(scan.get("processes") or [])
observed_by_pid: dict[int, Mapping[str, Any]] = {}
for proc in observed_processes:
observed_pid = _int_or_none(proc.get("pid"))
if observed_pid is not None:
observed_by_pid[observed_pid] = proc
expected_index = _expected_index(expected_fleet)
members = [
_normalize_row(
row,
observed_by_pid=observed_by_pid,
process_scan_available=scan_available,
now=moment,
)
for row in (runtime_rows or [])
]
live = [m for m in members if m["liveness"] == LIVENESS_LIVE]
not_live = [m for m in members if m["liveness"] != LIVENESS_LIVE]
live_by_profile: dict[str, list[dict[str, Any]]] = {}
for member in live:
live_by_profile.setdefault(member["profile"] or "", []).append(member)
running_members: list[dict[str, Any]] = []
missing_members: list[dict[str, Any]] = []
duplicate_members: list[dict[str, Any]] = []
unexpected_members: list[dict[str, Any]] = []
binding_mismatches: list[dict[str, Any]] = []
role_mismatches: list[dict[str, Any]] = []
configured_members: list[dict[str, Any]] = []
for entry in expected_fleet:
profile = _clean(entry.get("profile")) or ""
instances = sorted(live_by_profile.get(profile, []), key=_sort_key)
configured_members.append(
{
"namespace": _clean(entry.get("namespace")),
"profile": profile,
"role": _clean(entry.get("role")),
"instance_count": len(instances),
"health": (
HEALTH_MISSING
if not instances
else HEALTH_RUNNING
if len(instances) == 1
else HEALTH_DUPLICATE
),
}
)
if not instances:
missing_members.append(
{
"namespace": _clean(entry.get("namespace")),
"profile": profile,
"role": _clean(entry.get("role")),
"health": HEALTH_MISSING,
"reason": (
"no live registry row corroborated by a running server "
"process"
),
}
)
continue
for instance in instances:
instance["health"] = (
HEALTH_RUNNING if len(instances) == 1 else HEALTH_DUPLICATE
)
instance["expected"] = True
running_members.append(instance)
if len(instances) > 1:
duplicate_members.append(
{
"namespace": _clean(entry.get("namespace")),
"profile": profile,
"role": _clean(entry.get("role")),
"health": HEALTH_DUPLICATE,
"instance_count": len(instances),
"pids": sorted(
i["pid"] for i in instances if i["pid"] is not None
),
"instances": instances,
"reason": (
"more than one live server is registered for this profile"
),
}
)
for member in sorted(live, key=_sort_key):
profile = member["profile"] or ""
expected_entry = expected_index.get(profile)
if expected_entry is not None:
expected_role = _clean(expected_entry.get("role"))
actual_role = member["role"]
if expected_role and actual_role and actual_role != expected_role:
role_mismatches.append(
{
"namespace": member["namespace"],
"profile": profile,
"pid": member["pid"],
"expected_role": expected_role,
"declared_role": actual_role,
"reason": "declared role does not match the configured profile",
}
)
expected_namespace = _clean(expected_entry.get("namespace"))
if (
expected_namespace
and member["namespace"]
and member["namespace"] != expected_namespace
):
role_mismatches.append(
{
"namespace": member["namespace"],
"profile": profile,
"pid": member["pid"],
"expected_namespace": expected_namespace,
"declared_role": member["role"],
"reason": (
"profile is served from a namespace it is not "
"configured for"
),
}
)
else:
member["health"] = HEALTH_UNEXPECTED
member["expected"] = False
unexpected_members.append(member)
match = _binding_matches(member, expected_binding)
member["repository_binding_matches"] = match
if match is False:
binding_mismatches.append(
{
"namespace": member["namespace"],
"profile": profile,
"pid": member["pid"],
"remote": member["remote"],
"org": member["org"],
"repo": member["repo"],
"expected": dict(expected_binding or {}),
"reason": "member is bound to a different repository",
}
)
elif match is None and expected_binding:
binding_mismatches.append(
{
"namespace": member["namespace"],
"profile": profile,
"pid": member["pid"],
"remote": member["remote"],
"org": member["org"],
"repo": member["repo"],
"expected": dict(expected_binding or {}),
"reason": "member did not record a complete repository binding",
}
)
stale_members: list[dict[str, Any]] = []
for member in sorted(not_live, key=_sort_key):
member["health"] = (
HEALTH_UNKNOWN if member["liveness"] == LIVENESS_UNKNOWN else HEALTH_STALE
)
stale_members.append(member)
registered_pids = {m["pid"] for m in live if m["pid"] is not None}
unregistered_processes: list[dict[str, Any]] = []
if scan_available:
for proc in observed_processes:
observed_pid = _int_or_none(proc.get("pid"))
if observed_pid is None or observed_pid in registered_pids:
continue
unregistered_processes.append(
{"pid": observed_pid, "started_at": proc.get("started_at")}
)
unregistered_processes.sort(key=lambda item: item["pid"])
# Cohort. Derived only from recorded cohort identity — never from revisions.
cohort_ids = {m["cohort_id"] for m in live}
cohort_unknown = any(cohort_id is None for cohort_id in cohort_ids)
known_cohorts = sorted(c for c in cohort_ids if c is not None)
if not live or cohort_unknown:
single_cohort: bool | None = None
mixed_cohort: bool | None = None
else:
single_cohort = len(known_cohorts) == 1
mixed_cohort = len(known_cohorts) > 1
# Revision spread. Reported independently of cohort; never used to infer it.
heads = {m["startup_head"] for m in live}
head_unknown = any(head is None for head in heads)
known_heads = sorted(h for h in heads if h is not None)
mixed_revision = (len(known_heads) > 1) if known_heads else None
exactly_one_per_profile = not missing_members and not duplicate_members
no_unexpected_members = not unexpected_members
incomplete_reasons: list[str] = []
if not registry_available:
incomplete_reasons.append(
registry_error or "the control-plane runtime registry could not be read"
)
if not scan_available:
incomplete_reasons.append(
str(scan.get("reason") or "server-side process observation unavailable")
)
if unregistered_processes:
pids = ", ".join(str(p["pid"]) for p in unregistered_processes)
incomplete_reasons.append(
f"running Gitea MCP server process(es) with no runtime registry row "
f"(PIDs: {pids}); the fleet contains members this inventory cannot "
f"describe"
)
if live and cohort_unknown:
incomplete_reasons.append(
"one or more live members did not record a client cohort identity; "
"matching revisions do not establish a single cohort"
)
if live and head_unknown:
incomplete_reasons.append(
"one or more live members did not record a startup revision"
)
if any(m["liveness"] == LIVENESS_UNKNOWN for m in members):
incomplete_reasons.append(
"liveness of one or more registry rows could not be determined"
)
inventory_complete = not incomplete_reasons
blocked_reasons: list[str] = list(incomplete_reasons)
if missing_members:
names = ", ".join(sorted(m["profile"] for m in missing_members))
blocked_reasons.append(f"expected fleet member(s) not running: {names}")
if duplicate_members:
names = ", ".join(sorted(d["profile"] for d in duplicate_members))
blocked_reasons.append(
f"duplicate server(s) registered for profile(s): {names}"
)
if unexpected_members:
names = ", ".join(
sorted(
str(m["profile"] or m["namespace"] or "?") for m in unexpected_members
)
)
blocked_reasons.append(f"unexpected fleet member(s) running: {names}")
if mixed_cohort:
blocked_reasons.append(
"live members span more than one client cohort: "
+ ", ".join(known_cohorts)
)
if mixed_revision:
blocked_reasons.append(
"live members started at different revisions: " + ", ".join(known_heads)
)
if binding_mismatches:
blocked_reasons.append(
"one or more live members are not bound to the expected repository"
)
if role_mismatches:
blocked_reasons.append(
"one or more live members declare a role or namespace that does not "
"match the configured profile"
)
mutation_gate_satisfied = bool(
inventory_complete
and exactly_one_per_profile
and no_unexpected_members
and single_cohort is True
and mixed_revision is False
and not binding_mismatches
and not role_mismatches
)
if not mutation_gate_satisfied and not blocked_reasons:
blocked_reasons.append(
"the fleet snapshot did not establish the exact-one-instance-per-"
"profile, single-cohort invariant"
)
return {
"success": True,
"read_only": True,
"evidence_authority": EVIDENCE_AUTHORITY,
"answering_namespace": _clean(answering_namespace),
"generated_at": moment.strftime(_ISO_FORMAT),
"inventory_complete": inventory_complete,
"incomplete_reasons": incomplete_reasons,
"configured_members": configured_members,
"running_members": sorted(running_members, key=_sort_key),
"missing_members": sorted(missing_members, key=_sort_key),
"duplicate_members": sorted(duplicate_members, key=_sort_key),
"unexpected_members": sorted(unexpected_members, key=_sort_key),
"stale_members": stale_members,
"unregistered_processes": unregistered_processes,
"repository_binding_mismatches": sorted(binding_mismatches, key=_sort_key),
"role_mismatches": sorted(role_mismatches, key=_sort_key),
"cohort_ids": known_cohorts,
"single_cohort": single_cohort,
"mixed_cohort": mixed_cohort,
"startup_revisions": known_heads,
"mixed_revision": mixed_revision,
"exactly_one_per_profile": exactly_one_per_profile,
"no_unexpected_members": no_unexpected_members,
"expected_member_count": len(expected_fleet),
"running_member_count": len(running_members),
"mutation_gate_satisfied": mutation_gate_satisfied,
"blocked_reason": blocked_reasons[0] if blocked_reasons else None,
"blocked_reasons": blocked_reasons,
"process_observation": {
"available": scan_available,
"observed_process_count": len(observed_processes),
"reason": scan.get("reason"),
},
"registry": {
"available": registry_available,
"row_count": len(members),
"error": registry_error,
},
"mutations_performed": [],
}
def summarize(result: Mapping[str, Any]) -> str:
"""One-line human summary of a classification result."""
if result.get("mutation_gate_satisfied"):
return (
f"fleet healthy: {result.get('running_member_count')} of "
f"{result.get('expected_member_count')} members running in a single "
f"cohort at one revision"
)
return f"fleet not provable: {result.get('blocked_reason')}"
+7 -3
View File
@@ -1,7 +1,10 @@
#!/usr/bin/env python3 #!/usr/bin/env python3
"""Gitea MCP Server — exposes Gitea operations as MCP tools. """Gitea MCP Server — exposes Gitea operations as MCP tools.
Runs over stdio. All tools authenticate via macOS keychain (git credential fill). The transport is selected by deployment configuration (GITEA_MCP_TRANSPORT) and
defaults to the local client-spawned transport when unset (#931); the permitted
set lives in mcp_transport_config. All tools authenticate via macOS keychain
(git credential fill).
""" """
import os import os
import sys import sys
@@ -43,8 +46,9 @@ check_conflict_markers()
# #558 / #695: claim the official entrypoint before loading mutation modules. # #558 / #695: claim the official entrypoint before loading mutation modules.
# This alone does NOT authorize mutations — gitea_mcp_server binds the live # This alone does NOT authorize mutations — gitea_mcp_server binds the live
# native MCP transport (stdio) immediately before mcp.run. Import-only or # native MCP transport immediately before mcp.run, over the configured
# offline launch without that bind fails closed on mutations. # transport (#931). Import-only or offline launch without that bind fails
# closed on mutations.
try: try:
import mcp_daemon_guard import mcp_daemon_guard
+4
View File
@@ -588,9 +588,13 @@ def save_state(
bool(prov.get("production_native_mcp_transport")), bool(prov.get("production_native_mcp_transport")),
) )
body.setdefault("transport", prov.get("transport")) body.setdefault("transport", prov.get("transport"))
# #931: record the bound transport identifier itself, so a durable
# decision lock names which transport performed the mutation.
body.setdefault("bound_transport", prov.get("bound_transport"))
except Exception: except Exception:
body.setdefault("native_mcp_transport", False) body.setdefault("native_mcp_transport", False)
body.setdefault("transport", "untrusted") body.setdefault("transport", "untrusted")
body.setdefault("bound_transport", None)
envelope = { envelope = {
"kind": kind, "kind": kind,
-1
View File
@@ -40,7 +40,6 @@ NON_TOOL_IDENTIFIERS: frozenset[str] = frozenset(
"gitea_auth", "gitea_auth",
"gitea_config", "gitea_config",
"gitea_mcp_server", "gitea_mcp_server",
"mcp_fleet_inventory",
"mcp_server", "mcp_server",
"offline_mcp_helper", "offline_mcp_helper",
"offline_mcp_runner", "offline_mcp_runner",
+292
View File
@@ -0,0 +1,292 @@
"""Single authoritative source for the bound MCP transport identifier (#931).
Before this module the transport was a literal, passed once at the bottom of
``gitea_mcp_server`` as ``bind_native_mcp_transport(transport="stdio")``. Every
guard that later asks "is this a trusted native session" resolves that question
through the value bound there, so the literal was effectively a constant in the
authorization chain rather than configuration.
This module is the seam. It owns three things and nothing else:
- the permitted set of transport identifiers,
- the default used when deployment configuration says nothing,
- the resolution of the configured value into a validated identifier.
It deliberately holds no state. The *bound* transport is pinned once, at bind
time, into the process-local native-runtime record owned by
:mod:`mcp_daemon_guard`, and is read back through
``mcp_daemon_guard.bound_transport()``. That split matters: configuration is
read exactly once, before any tool can dispatch, so a later environment change
cannot move the value a guard observes the same pinning rule already applied
to the session-state root under #695 AC2.
Nothing here consumes tool arguments, request bodies, or provenance fields. The
only input is the deployment environment, read at bind time.
Standing up a listener for a non-stdio transport is #938; this module only
makes the identifier expressible and validated.
"""
from __future__ import annotations
import os
from typing import Any, Mapping
# Deployment configuration key. Read once, at bind time, and never again.
TRANSPORT_ENV = "GITEA_MCP_TRANSPORT"
# The local, client-spawned transport. Unset configuration resolves to this,
# which is what keeps every existing stdio deployment byte-identical.
DEFAULT_TRANSPORT = "stdio"
# The sanctioned remote transport identifier. Accepting it here is what makes
# the bind pluggable; the endpoint that serves it belongs to #938. The name
# matches the MCP transport name so no second vocabulary has to be mapped.
REMOTE_TRANSPORT = "streamable-http"
# The permitted set. This is the only place transport identifiers are
# enumerated; guards consult it rather than restating any member.
#
# ``sse`` is a real MCP transport and is deliberately absent: it is the
# superseded remote transport, and admitting it would give the deployment two
# remote paths to reason about. An unregistered identifier must fail closed at
# bind time, and ``sse`` is held to that rule like any other.
SUPPORTED_TRANSPORTS = frozenset({DEFAULT_TRANSPORT, REMOTE_TRANSPORT})
# Recognition is not execution authorization (#931, review 635 B1/B2).
#
# SUPPORTED_TRANSPORTS answers "is this an identifier this system knows, and may
# it be bound, pinned and recorded?". It deliberately includes the remote
# identifier, because #931 requires the bind to become pluggable.
#
# EXECUTABLE_TRANSPORTS answers a strictly narrower question: "is this entrypoint
# commissioned to actually *serve* on that transport?". Only the local transport
# is. Handing ``streamable-http`` to ``mcp.run`` would start FastMCP's HTTP
# listener with no authentication, no TLS and no per-request principal — the
# endpoint #938 owns and gates. Recognition must therefore never imply execution.
#
# #938 commissions the remote listener by adding REMOTE_TRANSPORT here, together
# with the authentication and principal boundary its acceptance criteria require.
EXECUTABLE_TRANSPORTS = frozenset({DEFAULT_TRANSPORT})
# Which issue owns commissioning each recognized-but-not-executable transport.
# Used to make the refusal actionable rather than a generic denial.
TRANSPORT_EXECUTION_OWNER = {REMOTE_TRANSPORT: "#938"}
BLOCKER_TRANSPORT_NOT_BOUND = "transport_not_bound"
BLOCKER_TRANSPORT_NOT_RECOGNIZED = "transport_not_recognized"
BLOCKER_LISTENER_NOT_COMMISSIONED = "transport_listener_not_commissioned"
SOURCE_CONFIGURED = "deployment_configuration"
SOURCE_DEFAULT = "default"
class TransportConfigurationError(ValueError):
"""Raised when configured transport is outside :data:`SUPPORTED_TRANSPORTS`."""
def normalize_transport(value: Any) -> str:
"""Canonical form of a transport identifier; ``""`` when there is none.
Non-string values normalize to ``""`` rather than being coerced, so a
structured object smuggled in from a caller can never match a member of the
permitted set.
"""
if not isinstance(value, str):
return ""
return value.strip().lower()
def supported_transports() -> tuple[str, ...]:
"""Permitted identifiers, sorted, for messages and status payloads."""
return tuple(sorted(SUPPORTED_TRANSPORTS))
def is_supported_transport(value: Any) -> bool:
"""True when *value* normalizes to a member of the permitted set."""
return normalize_transport(value) in SUPPORTED_TRANSPORTS
def is_remote_transport(value: Any) -> bool:
"""True when *value* is a permitted transport that is not the local one."""
name = normalize_transport(value)
return name in SUPPORTED_TRANSPORTS and name != DEFAULT_TRANSPORT
def executable_transports() -> tuple[str, ...]:
"""Transports this entrypoint is commissioned to serve, sorted."""
return tuple(sorted(EXECUTABLE_TRANSPORTS))
def is_executable_transport(value: Any) -> bool:
"""True when *value* may actually be served by this entrypoint (#931).
Strictly narrower than :func:`is_supported_transport`. A recognized
identifier that is not executable is a correct, fully-bound configuration
whose listener has simply not been commissioned yet.
"""
return normalize_transport(value) in EXECUTABLE_TRANSPORTS
def assess_transport_execution(value: Any) -> dict[str, Any]:
"""Structured serve-authorization verdict for a bound transport (#931).
This is the decision that separates a *recognized* transport from one this
entrypoint may execute. It is deliberately a pure function of the bound
identifier so the serve path cannot reach a listener the deployment has not
commissioned.
Args:
value: The bound transport identifier, or ``None`` when unbound.
Returns:
dict with ``transport``, ``recognized``, ``executable``, ``allowed``,
``blocker_kind``, ``owner_issue``, ``reasons`` and
``exact_next_action``. ``allowed`` is true only for a bound, recognized,
commissioned transport.
"""
name = normalize_transport(value)
if not name:
return {
"transport": None,
"recognized": False,
"executable": False,
"allowed": False,
"blocker_kind": BLOCKER_TRANSPORT_NOT_BOUND,
"owner_issue": None,
"supported_transports": list(supported_transports()),
"executable_transports": list(executable_transports()),
"reasons": [
"no transport is bound; the serve path is fail-closed until "
"bind_native_mcp_transport succeeds (#695/#931)"
],
"exact_next_action": (
"Launch through the canonical entrypoint so "
"bind_native_mcp_transport runs before tool service."
),
}
recognized = name in SUPPORTED_TRANSPORTS
if not recognized:
return {
"transport": name,
"recognized": False,
"executable": False,
"allowed": False,
"blocker_kind": BLOCKER_TRANSPORT_NOT_RECOGNIZED,
"owner_issue": None,
"supported_transports": list(supported_transports()),
"executable_transports": list(executable_transports()),
"reasons": [
f"transport {name!r} is not a registered MCP transport (#931); "
"it should have been refused at bind time"
],
"exact_next_action": (
f"Set {TRANSPORT_ENV} to one of {list(supported_transports())}."
),
}
if name in EXECUTABLE_TRANSPORTS:
return {
"transport": name,
"recognized": True,
"executable": True,
"allowed": True,
"blocker_kind": None,
"owner_issue": None,
"supported_transports": list(supported_transports()),
"executable_transports": list(executable_transports()),
"reasons": [],
"exact_next_action": None,
}
owner = TRANSPORT_EXECUTION_OWNER.get(name)
owner_text = owner or "the issue that commissions this transport's listener"
return {
"transport": name,
"recognized": True,
"executable": False,
"allowed": False,
"blocker_kind": BLOCKER_LISTENER_NOT_COMMISSIONED,
"owner_issue": owner,
"supported_transports": list(supported_transports()),
"executable_transports": list(executable_transports()),
"reasons": [
f"transport {name!r} is registered and was bound and recorded, but "
f"this entrypoint is not commissioned to serve it (#931). Serving it "
f"would start a listener with no authentication, no transport "
f"security and no per-request principal; that endpoint is owned by "
f"{owner_text}."
],
"exact_next_action": (
f"Serve on {DEFAULT_TRANSPORT} until {owner_text} commissions the "
f"{name!r} listener with its authentication and principal boundary, "
f"which adds {name!r} to EXECUTABLE_TRANSPORTS."
),
}
def resolve_configured_transport(
env: Mapping[str, str] | None = None,
) -> dict[str, Any]:
"""Resolve the deployment-configured transport without raising.
Returns the resolution rather than a bare string so a caller can tell an
unset value (which legitimately yields :data:`DEFAULT_TRANSPORT`) from a
configured value that is not permitted (which must fail closed, never
silently degrade to the default).
Args:
env: Environment mapping to read; defaults to ``os.environ``.
Returns:
dict with ``transport`` (normalized; the default when unset),
``configured``, ``source``, ``raw``, ``supported``, and ``reasons``.
"""
source_env = os.environ if env is None else env
raw = source_env.get(TRANSPORT_ENV)
normalized = normalize_transport(raw)
configured = bool(normalized)
if not configured:
return {
"transport": DEFAULT_TRANSPORT,
"configured": False,
"source": SOURCE_DEFAULT,
"raw": raw,
"supported": True,
"supported_transports": list(supported_transports()),
"env_key": TRANSPORT_ENV,
"reasons": [],
}
supported = normalized in SUPPORTED_TRANSPORTS
reasons: list[str] = []
if not supported:
reasons.append(
f"{TRANSPORT_ENV}={normalized!r} is not a registered MCP transport "
f"(#931). Registered: {list(supported_transports())}."
)
return {
"transport": normalized,
"configured": True,
"source": SOURCE_CONFIGURED,
"raw": raw,
"supported": supported,
"supported_transports": list(supported_transports()),
"env_key": TRANSPORT_ENV,
"reasons": reasons,
}
def require_configured_transport(env: Mapping[str, str] | None = None) -> str:
"""Resolved transport identifier, or raise when it is not permitted.
Raises:
TransportConfigurationError: the configured identifier is unregistered.
"""
resolution = resolve_configured_transport(env)
if not resolution["supported"]:
raise TransportConfigurationError("; ".join(resolution["reasons"]))
return str(resolution["transport"])
+37
View File
@@ -74,6 +74,43 @@ def check_author_mutation_namespace(
return True, [] return True, []
def check_author_role_kind(
mutation_task: str,
profile: dict,
) -> tuple[bool, list[str]]:
"""Author-exclusive wall for durable-lock mutations (#953 F1).
``check_author_mutation_namespace`` walls off reviewer-bound sessions, which
is the whole gate for tasks whose required permission is itself author-only
(``gitea.pr.create``, ``gitea.repo.commit``). It is *not* sufficient for a
task gated on ``gitea.issue.comment``, which every configured role holds: a
merger, controller, or reconciler session would clear both the namespace
check and the permission gate and still reach the durable write.
Opt-in per call site and additive. It refuses any active profile whose
derived role kind is not exactly ``author`` for a task the router declares
author-required, and grants nothing to anyone a ``mixed`` profile is
refused rather than admitted.
"""
required_role = role_session_router.required_role_for_task(mutation_task)
if required_role != "author":
return True, []
allowed = profile.get("allowed_operations") or []
forbidden = profile.get("forbidden_operations") or []
active_role = derive_role_kind(allowed, forbidden)
if active_role == "author":
return True, []
profile_name = profile.get("profile_name") or ""
namespace = infer_mcp_namespace(profile_name)
return False, [
f"author mutation '{mutation_task}' blocked: active session role kind is "
f"'{active_role}', not 'author' ({profile_name} / {namespace}); this "
"operation writes a durable author issue lock and is author-exclusive",
]
def mutation_audit_context(mutation_task: str, profile: dict, *, def mutation_audit_context(mutation_task: str, profile: dict, *,
remote=None, repository=None) -> dict: remote=None, repository=None) -> dict:
"""Structured mutation metadata for audit records (#209).""" """Structured mutation metadata for audit records (#209)."""
+10
View File
@@ -75,6 +75,10 @@ AUTHOR_TASKS = frozenset({
"push_branch", "push_branch",
"bootstrap_author_issue_worktree", "bootstrap_author_issue_worktree",
"gitea_bootstrap_author_issue_worktree", "gitea_bootstrap_author_issue_worktree",
# #953: recovery of an incomplete bootstrap lock is an author-only durable
# state mutation and belongs to the same class as bootstrap itself.
"recover_incomplete_bootstrap_lock",
"gitea_recover_incomplete_bootstrap_lock",
"create_pr", "create_pr",
"comment_pr", "comment_pr",
"address_pr_change_requests", "address_pr_change_requests",
@@ -112,6 +116,12 @@ TASK_REQUIRED_ROLE = {
"claim_issue": "author", "claim_issue": "author",
"create_branch": "author", "create_branch": "author",
"push_branch": "author", "push_branch": "author",
# #953: without this entry ``required_role_for_task`` returns None and
# ``role_namespace_gate.check_author_mutation_namespace`` short-circuits to
# "allowed" — the namespace wall on the recovery tool would be inert. The
# capability map already records the same role; both tables must agree.
"recover_incomplete_bootstrap_lock": "author",
"gitea_recover_incomplete_bootstrap_lock": "author",
"create_pr": "author", "create_pr": "author",
"comment_pr": "author", "comment_pr": "author",
"address_pr_change_requests": "author", "address_pr_change_requests": "author",
-25
View File
@@ -252,31 +252,6 @@ Helpers: `scripts/worktree-start`, `scripts/worktree-review`,
- Never place raw tokens in LLM/MCP config. - Never place raw tokens in LLM/MCP config.
- Use `gitea_whoami` and `gitea_resolve_task_capability` before mutating. - Use `gitea_whoami` and `gitea_resolve_task_capability` before mutating.
## Fleet inventory
`gitea_whoami`, `gitea_get_runtime_context` and `gitea_assess_master_parity` each
describe only the server answering the call. Five namespaces independently
reporting the same revision never proved that five processes exist, that no sixth
exists, or that all five belong to one client cohort.
`gitea_assess_fleet_inventory` is the read-only capability that does prove it. It
takes no evidence parameters: it combines the control-plane runtime registry,
which each server writes about itself at native transport bind, with a process
observation the answering server performs. Classification is a pure function of
that snapshot, so `gitea-controller` and `gitea-reconciler` return the same
verdict for the same fleet.
Consume `mutation_gate_satisfied`. When it is false, report `blocked_reason`
verbatim and stop — `missing_members`, `duplicate_members`, `unexpected_members`,
`unregistered_processes`, `mixed_cohort` and `mixed_revision` are reported
separately because each needs a different operator action. Treat
`inventory_complete: false` and `single_cohort: null` as *unknown*, never as
healthy. The capability never terminates a duplicate process, restarts,
reconnects, or touches a lease.
Details, field meanings, and the gate-consumption sequence:
[`docs/mcp-fleet-inventory.md`](../../docs/mcp-fleet-inventory.md) (#949).
## Tool inventory ## Tool inventory
[`docs/mcp-tool-inventory.md`](../../docs/mcp-tool-inventory.md) is the canonical [`docs/mcp-tool-inventory.md`](../../docs/mcp-tool-inventory.md) is the canonical
+21 -13
View File
@@ -41,6 +41,27 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.issue.comment", "permission": "gitea.issue.comment",
"role": "author", "role": "author",
}, },
# #953: target-specific upgrade of an incomplete bootstrap lock (explicit
# operation, never a widening of lock_issue). Author-only, and the tool
# additionally proves exact-owner claimant match before writing.
"recover_incomplete_bootstrap_lock": {
"permission": "gitea.issue.comment",
"role": "author",
},
"gitea_recover_incomplete_bootstrap_lock": {
"permission": "gitea.issue.comment",
"role": "author",
},
# #953: read-only lock contract inspection. Read permission only — it must
# never be able to mutate.
"inspect_issue_lock_contract": {
"permission": "gitea.read",
"role": "author",
},
"gitea_inspect_issue_lock_contract": {
"permission": "gitea.read",
"role": "author",
},
# #860: dirty orphaned same-claimant worktree recovery (explicit operation). # #860: dirty orphaned same-claimant worktree recovery (explicit operation).
"recover_dirty_orphaned_issue_worktree": { "recover_dirty_orphaned_issue_worktree": {
"permission": "gitea.issue.comment", "permission": "gitea.issue.comment",
@@ -159,19 +180,6 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.read", "permission": "gitea.read",
"role": "reconciler", "role": "reconciler",
}, },
# #949: fleet inventory is strictly read-only evidence. It is deliberately
# not role-exclusive — the controller and reconciler gates are its named
# consumers, but any namespace holding gitea.read must be able to prove the
# exact-five-process/single-cohort invariant before it acts, and the verdict
# is a pure function of the snapshot so every namespace agrees.
"assess_fleet_inventory": {
"permission": "gitea.read",
"role": "reconciler",
},
"gitea_assess_fleet_inventory": {
"permission": "gitea.read",
"role": "reconciler",
},
# PR synchronization lifecycle: assess is read-only (any role with gitea.read); # PR synchronization lifecycle: assess is read-only (any role with gitea.read);
# update-by-merge is author-only and mutates the PR head via Gitea API. # update-by-merge is author-only and mutates the PR head via Gitea API.
"assess_pr_sync_status": { "assess_pr_sync_status": {
+2 -6
View File
@@ -14,7 +14,6 @@ from control_plane_db import (
ControlPlaneError, ControlPlaneError,
InvalidWorkKindError, InvalidWorkKindError,
LeaseRequiredError, LeaseRequiredError,
SCHEMA_VERSION,
WORK_KINDS, WORK_KINDS,
_ts, _ts,
_utc_now, _utc_now,
@@ -38,10 +37,7 @@ class ControlPlaneDBTest(unittest.TestCase):
rows = dict(conn.execute("SELECT key, value FROM schema_meta").fetchall()) rows = dict(conn.execute("SELECT key, value FROM schema_meta").fetchall())
finally: finally:
conn.close() conn.close()
# Pinned to the constant, not a literal: every additive migration bumps self.assertEqual(rows["schema_version"], "5")
# SCHEMA_VERSION, and the invariant under test is that the meta row
# records the version the code actually wrote (#949 added v6).
self.assertEqual(rows["schema_version"], str(SCHEMA_VERSION))
self.assertIn("DB coordinates", rows["architecture"]) self.assertIn("DB coordinates", rows["architecture"])
self.assertIn("bridge", rows["architecture"].lower()) self.assertIn("bridge", rows["architecture"].lower())
@@ -872,7 +868,7 @@ class SessionCheckpointTest(unittest.TestCase):
conn.close() conn.close()
self.assertIn("session_checkpoints", names) self.assertIn("session_checkpoints", names)
record = self._write() record = self._write()
self.assertEqual(record["checkpoint_schema_version"], SCHEMA_VERSION) self.assertEqual(record["checkpoint_schema_version"], 5)
# AC2 — checkpoints written for multi-role session fixtures. # AC2 — checkpoints written for multi-role session fixtures.
def test_multi_role_fixtures_each_get_a_row(self) -> None: def test_multi_role_fixtures_each_get_a_row(self) -> None:
@@ -1,282 +0,0 @@
"""Control-plane MCP server runtime registry (#949).
The registry is the half of the fleet evidence a caller cannot supply: each
server writes exactly one row about itself at native transport bind. These tests
pin the storage contract additive migration, deterministic ordering, and
PID-reuse/retention pruning that cannot hide a live duplicate.
"""
import os
import tempfile
import unittest
from datetime import datetime, timedelta, timezone
from unittest import mock
import control_plane_db
import mcp_fleet_inventory as mfi
def _stamp(delta_seconds=0):
moment = datetime.now(timezone.utc) + timedelta(seconds=delta_seconds)
return moment.replace(microsecond=0).strftime("%Y-%m-%dT%H:%M:%SZ")
class _DBCase(unittest.TestCase):
def setUp(self):
self._tmp = tempfile.TemporaryDirectory()
self.addCleanup(self._tmp.cleanup)
self.db_path = os.path.join(self._tmp.name, "control_plane.sqlite3")
self.db = control_plane_db.ControlPlaneDB(db_path=self.db_path)
def record(self, namespace, profile, role, pid, **overrides):
base = mfi.build_process_runtime_record(
namespace=namespace,
profile=profile,
role=role,
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
repository_root="/checkout/Gitea-Tools",
pid=pid,
startup_head="82d71b77028a7abd4f8ab4a4e4d89658a187f73d",
transport="stdio",
client_provenance="client_managed",
env={mfi.COHORT_ID_ENV: "ppid:40990"},
boot_id=f"boot{pid}",
)
base.update(overrides)
return base
class TestSchema(_DBCase):
def test_schema_version_is_bumped(self):
self.assertGreaterEqual(control_plane_db.SCHEMA_VERSION, 6)
def test_runtime_table_exists_on_a_fresh_database(self):
self.assertEqual(self.db.list_mcp_server_runtimes(), [])
def test_migration_is_additive_and_idempotent(self):
"""Re-opening an existing DB must not disturb the other tables."""
self.db.upsert_session(session_id="s-a", role="author", profile="prgs-author")
self.db.register_mcp_server_runtime(
self.record("gitea-author", "prgs-author", "author", 41000)
)
reopened = control_plane_db.ControlPlaneDB(db_path=self.db_path)
self.assertEqual(len(reopened.list_sessions()), 1)
self.assertEqual(len(reopened.list_mcp_server_runtimes()), 1)
class TestRegistration(_DBCase):
def test_registration_round_trips_every_field(self):
record = self.record("gitea-controller", "prgs-controller", "controller", 41001)
stored = self.db.register_mcp_server_runtime(record)
for key in (
"runtime_id",
"namespace",
"profile",
"role",
"remote",
"org",
"repo",
"repository_root",
"pid",
"cohort_id",
"cohort_source",
"client_provenance",
"boot_id",
"startup_head",
"transport",
"status",
):
self.assertEqual(stored[key], record[key], key)
def test_registration_requires_a_runtime_id(self):
record = self.record("gitea-author", "prgs-author", "author", 41002)
record["runtime_id"] = ""
with self.assertRaises(ValueError):
self.db.register_mcp_server_runtime(record)
def test_registration_requires_a_namespace(self):
record = self.record("gitea-author", "prgs-author", "author", 41003)
record["namespace"] = " "
with self.assertRaises(ValueError):
self.db.register_mcp_server_runtime(record)
def test_registration_requires_a_pid(self):
record = self.record("gitea-author", "prgs-author", "author", 41004)
record["pid"] = None
with self.assertRaises(ValueError):
self.db.register_mcp_server_runtime(record)
def test_five_members_register_independently(self):
for index, entry in enumerate(mfi.EXPECTED_PRGS_FLEET):
self.db.register_mcp_server_runtime(
self.record(
entry["namespace"], entry["profile"], entry["role"], 41100 + index
)
)
rows = self.db.list_mcp_server_runtimes()
self.assertEqual(len(rows), 5)
self.assertEqual(
sorted(r["profile"] for r in rows),
[
"prgs-author",
"prgs-controller",
"prgs-merger",
"prgs-reconciler",
"prgs-reviewer",
],
)
class TestPruning(_DBCase):
def test_reusing_a_pid_replaces_the_stale_row(self):
first = self.record("gitea-author", "prgs-author", "author", 41200)
self.db.register_mcp_server_runtime(first)
second = self.record("gitea-reviewer", "prgs-reviewer", "reviewer", 41200)
self.db.register_mcp_server_runtime(second)
rows = self.db.list_mcp_server_runtimes()
self.assertEqual(len(rows), 1)
self.assertEqual(rows[0]["runtime_id"], second["runtime_id"])
def test_pruning_never_removes_a_live_duplicate_on_another_pid(self):
"""The defect this must not have: hiding a second running server."""
first = self.record("gitea-author", "prgs-author", "author", 41300)
self.db.register_mcp_server_runtime(first)
second = self.record("gitea-author", "prgs-author", "author", 41301)
self.db.register_mcp_server_runtime(
second, retention_seconds=mfi.RUNTIME_RETENTION_SECONDS
)
rows = self.db.list_mcp_server_runtimes()
self.assertEqual(len(rows), 2)
self.assertEqual(sorted(r["pid"] for r in rows), [41300, 41301])
def test_retention_drops_rows_older_than_the_window(self):
old = self.record(
"gitea-merger",
"prgs-merger",
"merger",
41400,
registered_at=_stamp(-(mfi.RUNTIME_RETENTION_SECONDS + 3600)),
)
self.db.register_mcp_server_runtime(old)
fresh = self.record("gitea-author", "prgs-author", "author", 41401)
self.db.register_mcp_server_runtime(
fresh, retention_seconds=mfi.RUNTIME_RETENTION_SECONDS
)
rows = self.db.list_mcp_server_runtimes()
self.assertEqual([r["pid"] for r in rows], [41401])
def test_retention_is_skipped_when_not_requested(self):
old = self.record(
"gitea-merger",
"prgs-merger",
"merger",
41500,
registered_at=_stamp(-(mfi.RUNTIME_RETENTION_SECONDS + 3600)),
)
self.db.register_mcp_server_runtime(old)
self.db.register_mcp_server_runtime(
self.record("gitea-author", "prgs-author", "author", 41501)
)
self.assertEqual(len(self.db.list_mcp_server_runtimes()), 2)
class TestListing(_DBCase):
def test_listing_is_deterministically_ordered(self):
for index in range(5):
self.db.register_mcp_server_runtime(
self.record(
"gitea-author",
"prgs-author",
"author",
41600 + index,
registered_at="2026-07-28T02:00:00Z",
)
)
first = [r["runtime_id"] for r in self.db.list_mcp_server_runtimes()]
second = [r["runtime_id"] for r in self.db.list_mcp_server_runtimes()]
self.assertEqual(first, second)
self.assertEqual(first, sorted(first))
def test_status_filter_selects_running_rows(self):
running = self.record("gitea-author", "prgs-author", "author", 41700)
self.db.register_mcp_server_runtime(running)
stopped = self.record("gitea-merger", "prgs-merger", "merger", 41701)
self.db.register_mcp_server_runtime(stopped)
self.db.mark_mcp_server_runtime_stopped(stopped["runtime_id"])
rows = self.db.list_mcp_server_runtimes(statuses=("running",))
self.assertEqual([r["pid"] for r in rows], [41700])
def test_listing_without_a_filter_returns_stopped_rows_too(self):
stopped = self.record("gitea-merger", "prgs-merger", "merger", 41800)
self.db.register_mcp_server_runtime(stopped)
self.db.mark_mcp_server_runtime_stopped(stopped["runtime_id"])
rows = self.db.list_mcp_server_runtimes()
self.assertEqual(len(rows), 1)
self.assertEqual(rows[0]["status"], "stopped")
class TestHeartbeat(_DBCase):
def test_heartbeat_advances_only_the_timestamp(self):
record = self.record(
"gitea-author",
"prgs-author",
"author",
41900,
last_heartbeat_at="2026-07-28T02:00:00Z",
)
self.db.register_mcp_server_runtime(record)
self.db.heartbeat_mcp_server_runtime(record["runtime_id"])
row = self.db.list_mcp_server_runtimes()[0]
self.assertNotEqual(row["last_heartbeat_at"], "2026-07-28T02:00:00Z")
self.assertEqual(row["status"], "running")
self.assertEqual(row["pid"], 41900)
def test_heartbeat_for_an_unknown_runtime_is_a_no_op(self):
self.db.heartbeat_mcp_server_runtime("does-not-exist")
self.assertEqual(self.db.list_mcp_server_runtimes(), [])
class TestEndToEndSnapshot(_DBCase):
"""Registry rows feed the classifier without reshaping."""
def test_registered_fleet_classifies_as_healthy(self):
pids = []
for index, entry in enumerate(mfi.EXPECTED_PRGS_FLEET):
pid = 42000 + index
pids.append(pid)
self.db.register_mcp_server_runtime(
self.record(
entry["namespace"],
entry["profile"],
entry["role"],
pid,
registered_at="2026-07-28T02:00:00Z",
)
)
rows = self.db.list_mcp_server_runtimes(statuses=("running",))
scan = {
"available": True,
"processes": [
{"pid": pid, "started_at": "2026-07-28T01:00:00Z"} for pid in pids
],
"reason": None,
}
with mock.patch.object(mfi, "probe_pid_alive", return_value=True):
result = mfi.classify_fleet_inventory(
runtime_rows=rows,
process_scan=scan,
expected_binding={
"remote": "prgs",
"org": "Scaled-Tech-Consulting",
"repo": "Gitea-Tools",
},
)
self.assertTrue(result["inventory_complete"])
self.assertTrue(result["mutation_gate_satisfied"])
self.assertEqual(result["running_member_count"], 5)
if __name__ == "__main__":
unittest.main()
-337
View File
@@ -1,337 +0,0 @@
"""Native exposure of the fleet inventory capability (#949).
Covers the parts the pure classifier cannot: that the tool is registered and
documented, that it is gated on ``gitea.read`` rather than a role, that it
accepts no caller-supplied evidence, that it fails closed on an unreadable
registry, and that the workflow documentation explains gate consumption.
"""
import inspect
import os
import pathlib
import unittest
from unittest import mock
import gitea_mcp_server
import mcp_fleet_inventory as mfi
import task_capability_map
REPO_ROOT = pathlib.Path(__file__).resolve().parent.parent
TOOL_NAME = "gitea_assess_fleet_inventory"
class TestCapabilityMapExposure(unittest.TestCase):
"""AC14: the invariant is provable through sanctioned native calls."""
def test_task_resolves_to_a_read_permission(self):
self.assertEqual(
task_capability_map.required_permission("assess_fleet_inventory"),
"gitea.read",
)
self.assertEqual(
task_capability_map.required_permission(TOOL_NAME), "gitea.read"
)
def test_task_and_tool_keys_agree(self):
self.assertEqual(
task_capability_map.TASK_CAPABILITY_MAP["assess_fleet_inventory"],
task_capability_map.TASK_CAPABILITY_MAP[TOOL_NAME],
)
def test_task_is_not_role_exclusive(self):
"""Controller and reconciler both need it; so does any read-only gate."""
self.assertNotIn(
"assess_fleet_inventory", task_capability_map.ROLE_EXCLUSIVE_TASKS
)
self.assertNotIn(TOOL_NAME, task_capability_map.ROLE_EXCLUSIVE_TASKS)
def test_task_is_not_registered_as_an_issue_mutation(self):
self.assertNotIn(TOOL_NAME, task_capability_map.ISSUE_MUTATION_TOOL_TASKS)
def test_unknown_task_still_fails_closed(self):
with self.assertRaises(KeyError):
task_capability_map.required_permission("assess_fleet_inventory_typo")
class TestToolRegistration(unittest.TestCase):
def test_tool_is_registered_on_the_server(self):
self.assertTrue(hasattr(gitea_mcp_server, TOOL_NAME))
def test_tool_accepts_no_caller_supplied_evidence(self):
"""The defect #949 names: evidence parameters the caller controls."""
signature = inspect.signature(gitea_mcp_server.gitea_assess_fleet_inventory)
self.assertEqual(sorted(signature.parameters), ["host", "org", "remote", "repo"])
for forbidden in ("process", "probe_result", "registered_tools", "processes"):
self.assertNotIn(forbidden, signature.parameters)
def test_tool_is_documented_in_the_canonical_inventory(self):
doc = (REPO_ROOT / "docs" / "mcp-tool-inventory.md").read_text()
self.assertIn(f"`{TOOL_NAME}`", doc)
def test_docstring_states_the_read_only_guarantee(self):
doc = (gitea_mcp_server.gitea_assess_fleet_inventory.__doc__ or "").lower()
self.assertIn("read-only", doc)
self.assertIn("fails closed", doc)
class TestNamespaceResolution(unittest.TestCase):
"""Every configured profile must resolve to its real fleet namespace.
``role_namespace_gate.infer_mcp_namespace`` recognises only author and
reviewer and echoes the profile name for the rest, which would label the
controller, merger, and reconciler members with namespaces that do not
exist and then flag all three as role mismatches.
"""
def test_every_expected_profile_maps_to_its_namespace(self):
for entry in mfi.EXPECTED_PRGS_FLEET:
self.assertEqual(
mfi.namespace_for_profile(entry["profile"]),
entry["namespace"],
entry["profile"],
)
def test_server_helper_agrees_with_the_roster(self):
for entry in mfi.EXPECTED_PRGS_FLEET:
self.assertEqual(
gitea_mcp_server._fleet_namespace_for_profile(entry["profile"]),
entry["namespace"],
entry["profile"],
)
def test_unknown_profile_falls_back_to_the_supplied_default(self):
self.assertEqual(
mfi.namespace_for_profile("prgs-shadow", default="gitea-shadow"),
"gitea-shadow",
)
def test_unknown_profile_without_a_default_returns_the_profile(self):
self.assertEqual(mfi.namespace_for_profile("prgs-shadow"), "prgs-shadow")
def test_registered_row_uses_the_roster_namespace(self):
fake_db = mock.Mock()
with mock.patch.object(
gitea_mcp_server, "_control_plane_db_or_error", return_value=(fake_db, [])
), mock.patch.object(
gitea_mcp_server, "get_profile", return_value={"allowed_operations": []}
), mock.patch.object(
gitea_mcp_server.gitea_config,
"selected_profile_name",
return_value="prgs-reconciler",
):
record = gitea_mcp_server._register_fleet_runtime(transport="stdio")
self.assertEqual(record["namespace"], "gitea-reconciler")
class TestToolBehavior(unittest.TestCase):
"""The tool wires the registry and the process scan into the classifier."""
@staticmethod
def _healthy_rows():
return [
{
"runtime_id": f"{entry['namespace']}:{43000 + index}:boot",
"namespace": entry["namespace"],
"profile": entry["profile"],
"role": entry["role"],
"remote": "prgs",
"org": "Scaled-Tech-Consulting",
"repo": "Gitea-Tools",
"repository_root": "/checkout/Gitea-Tools",
"pid": 43000 + index,
"cohort_id": "ppid:40990",
"cohort_source": "parent_process",
"client_provenance": "client_managed",
"boot_id": "boot",
"startup_head": "82d71b77028a7abd4f8ab4a4e4d89658a187f73d",
"daemon_start_head": "82d71b77028a7abd4f8ab4a4e4d89658a187f73d",
"transport": "stdio",
"registered_at": "2026-07-28T02:00:00Z",
"last_heartbeat_at": "2026-07-28T02:00:00Z",
"status": "running",
}
for index, entry in enumerate(mfi.EXPECTED_PRGS_FLEET)
]
def _call(self, rows, *, scan=None, db=None, db_errors=None):
fake_db = db
if fake_db is None and db_errors is None:
fake_db = mock.Mock()
fake_db.list_mcp_server_runtimes.return_value = rows
scan = scan or {
"available": True,
"processes": [
{"pid": r["pid"], "started_at": "2026-07-28T01:00:00Z"} for r in rows
],
"reason": None,
}
with mock.patch.object(
gitea_mcp_server, "_profile_operation_gate", return_value=[]
), mock.patch.object(
gitea_mcp_server,
"_control_plane_db_or_error",
return_value=(fake_db, db_errors or []),
), mock.patch.object(
gitea_mcp_server, "_active_profile_name", return_value="prgs-controller"
), mock.patch.object(
mfi, "scan_mcp_server_processes", return_value=scan
), mock.patch.object(
mfi, "probe_pid_alive", return_value=True
):
return gitea_mcp_server.gitea_assess_fleet_inventory(
remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools"
)
def test_healthy_fleet_satisfies_the_gate_through_the_tool(self):
result = self._call(self._healthy_rows())
self.assertTrue(result["inventory_complete"])
self.assertTrue(result["mutation_gate_satisfied"])
self.assertEqual(result["running_member_count"], 5)
self.assertEqual(result["mutations_performed"], [])
def test_tool_reports_the_expected_repository_binding(self):
result = self._call(self._healthy_rows())
self.assertEqual(
result["expected_repository_binding"],
{"remote": "prgs", "org": "Scaled-Tech-Consulting", "repo": "Gitea-Tools"},
)
def test_tool_reads_only_running_rows(self):
rows = self._healthy_rows()
fake_db = mock.Mock()
fake_db.list_mcp_server_runtimes.return_value = rows
self._call(rows, db=fake_db)
fake_db.list_mcp_server_runtimes.assert_called_once_with(statuses=("running",))
def test_tool_never_writes_to_the_registry(self):
rows = self._healthy_rows()
fake_db = mock.Mock()
fake_db.list_mcp_server_runtimes.return_value = rows
self._call(rows, db=fake_db)
fake_db.register_mcp_server_runtime.assert_not_called()
fake_db.heartbeat_mcp_server_runtime.assert_not_called()
fake_db.mark_mcp_server_runtime_stopped.assert_not_called()
def test_unavailable_control_plane_fails_closed(self):
result = self._call([], db_errors=["control-plane DB substrate unavailable"])
self.assertFalse(result["inventory_complete"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertFalse(result["registry"]["available"])
self.assertIn("unavailable", result["blocked_reason"])
def test_registry_read_failure_fails_closed(self):
rows = self._healthy_rows()
fake_db = mock.Mock()
fake_db.list_mcp_server_runtimes.side_effect = RuntimeError("db locked")
result = self._call(rows, db=fake_db)
self.assertFalse(result["inventory_complete"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertIn("could not be read", result["registry"]["error"])
def test_missing_read_permission_blocks_without_touching_the_registry(self):
fake_db = mock.Mock()
with mock.patch.object(
gitea_mcp_server,
"_profile_operation_gate",
return_value=["profile may not read"],
), mock.patch.object(
gitea_mcp_server, "_control_plane_db_or_error", return_value=(fake_db, [])
):
result = gitea_mcp_server.gitea_assess_fleet_inventory(remote="prgs")
self.assertFalse(result["success"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertIn("permission_report", result)
self.assertEqual(result["mutations_performed"], [])
fake_db.list_mcp_server_runtimes.assert_not_called()
def test_answering_namespace_is_reported(self):
result = self._call(self._healthy_rows())
self.assertEqual(result["answering_namespace"], "gitea-controller")
def test_summary_is_present(self):
result = self._call(self._healthy_rows())
self.assertIn("5 of 5", result["summary"])
class TestStartupRegistration(unittest.TestCase):
"""The registry row is written by the process it describes."""
def test_registration_helper_exists_on_the_entrypoint_module(self):
self.assertTrue(hasattr(gitea_mcp_server, "_register_fleet_runtime"))
def test_registration_writes_one_row_for_this_process(self):
fake_db = mock.Mock()
with mock.patch.object(
gitea_mcp_server, "_control_plane_db_or_error", return_value=(fake_db, [])
), mock.patch.object(
gitea_mcp_server, "get_profile", return_value={"allowed_operations": []}
):
record = gitea_mcp_server._register_fleet_runtime(transport="stdio")
self.assertIsNotNone(record)
fake_db.register_mcp_server_runtime.assert_called_once()
written = fake_db.register_mcp_server_runtime.call_args.args[0]
self.assertEqual(written["pid"], os.getpid())
self.assertEqual(written["transport"], "stdio")
def test_registration_failure_never_blocks_startup(self):
with mock.patch.object(
gitea_mcp_server,
"_control_plane_db_or_error",
side_effect=RuntimeError("boom"),
):
self.assertIsNone(gitea_mcp_server._register_fleet_runtime())
def test_registration_is_skipped_when_the_control_plane_is_unavailable(self):
with mock.patch.object(
gitea_mcp_server,
"_control_plane_db_or_error",
return_value=(None, ["unavailable"]),
):
self.assertIsNone(gitea_mcp_server._register_fleet_runtime())
def test_entrypoint_registers_after_binding_native_transport(self):
"""Order matters: only a transport-bound process may claim a row."""
source = (REPO_ROOT / "gitea_mcp_server.py").read_text()
bind_at = source.index('bind_native_mcp_transport(transport="stdio")')
register_at = source.index('_register_fleet_runtime(transport="stdio")')
run_at = source.index('mcp.run(transport="stdio")')
self.assertLess(bind_at, register_at)
self.assertLess(register_at, run_at)
class TestWorkflowDocumentation(unittest.TestCase):
"""AC13: documentation explains how the gates consume the result."""
def setUp(self):
self.doc = (REPO_ROOT / "docs" / "mcp-fleet-inventory.md").read_text()
self.skill = (
REPO_ROOT / "skills" / "llm-project-workflow" / "SKILL.md"
).read_text()
def test_dedicated_document_exists(self):
self.assertIn("# Authoritative MCP fleet inventory", self.doc)
def test_document_names_both_consuming_namespaces(self):
self.assertIn("gitea-controller", self.doc)
self.assertIn("gitea-reconciler", self.doc)
def test_document_explains_gate_consumption(self):
self.assertIn("mutation_gate_satisfied", self.doc)
self.assertIn("blocked_reason", self.doc)
def test_document_states_the_non_inferences(self):
self.assertIn("Configuration is not existence", self.doc)
self.assertIn("Matching revisions are not a cohort", self.doc)
def test_document_preserves_the_neighbouring_issue_boundaries(self):
for issue in ("#950", "#951", "#952"):
self.assertIn(issue, self.doc)
def test_canonical_workflow_skill_links_the_document(self):
self.assertIn("docs/mcp-fleet-inventory.md", self.skill)
self.assertIn(TOOL_NAME, self.skill)
if __name__ == "__main__":
unittest.main()
+924
View File
@@ -0,0 +1,924 @@
"""Transport-neutral MCP bind seam (#931).
These tests drive the *real* bind boundary ``mark_sanctioned_daemon`` followed
by ``bind_native_mcp_transport`` from a canonical entrypoint path, with the
pytest allowance switched off rather than mocking the new accessor. The
distinction matters here for the same reason it mattered in #941: a suite that
only exercises the helper in isolation cannot observe a seam that the live path
never reaches.
Covered:
1. no configured transport defaults to the local transport
2. explicit local transport binds
3. the sanctioned remote identifier binds through the seam
4. an unregistered identifier is rejected at bind time
5. an invalid bind prevents the server reaching tool service
6. the unbound state fails closed where a bind is required
7. every transport-aware guard reads the same authoritative value
8. client-controlled input cannot alter the bound transport
9. the durable decision-lock record carries the selected transport
10. existing stdio behaviour is unchanged
11. repeated / conflicting bind attempts follow one fail-closed contract
12. capability, role, repository and provenance protections do not regress
"""
from __future__ import annotations
import os
import tempfile
import unittest
from pathlib import Path
from unittest.mock import patch
REPO_ROOT = Path(__file__).resolve().parent.parent
import mcp_daemon_guard
import mcp_session_state
import mcp_transport_config
import irrecoverable_provenance
class _ProductionBind:
"""Context manager that reaches the real production bind path.
Patches only the two things a unit test cannot otherwise satisfy: the
resolved canonical entrypoint frame, and the pytest allowance that would
short-circuit ``mark_sanctioned_daemon``. Everything downstream of those
validation, pinning, the rebind contract runs unmodified.
"""
def __init__(self, env: dict[str, str] | None = None):
self._env = env or {}
self._stack: list = []
def __enter__(self):
mcp_daemon_guard.clear_native_runtime_for_tests()
canonical = str((REPO_ROOT / "mcp_server.py").resolve())
self._stack = [
patch.object(
mcp_daemon_guard,
"_caller_official_entrypoint_path",
side_effect=lambda: canonical,
),
patch.object(mcp_daemon_guard, "is_pytest_runtime", return_value=False),
patch.dict(os.environ, self._env),
]
for ctx in self._stack:
ctx.__enter__()
# Start from a clean configuration unless the test set one.
if mcp_transport_config.TRANSPORT_ENV not in self._env:
os.environ.pop(mcp_transport_config.TRANSPORT_ENV, None)
mcp_daemon_guard.mark_sanctioned_daemon()
return mcp_daemon_guard
def __exit__(self, *exc):
for ctx in reversed(self._stack):
ctx.__exit__(*exc)
mcp_daemon_guard.clear_native_runtime_for_tests()
return False
class TestPermittedSetIsSingleSourceOfTruth(unittest.TestCase):
"""AC3: identifiers are enumerated once, in the seam."""
def test_guard_allowlist_is_the_seam_allowlist(self):
self.assertIs(
mcp_daemon_guard._PRODUCTION_TRANSPORTS,
mcp_transport_config.SUPPORTED_TRANSPORTS,
)
def test_default_is_a_member_of_the_permitted_set(self):
self.assertIn(
mcp_transport_config.DEFAULT_TRANSPORT,
mcp_transport_config.SUPPORTED_TRANSPORTS,
)
def test_remote_identifier_is_permitted_and_not_the_default(self):
self.assertIn(
mcp_transport_config.REMOTE_TRANSPORT,
mcp_transport_config.SUPPORTED_TRANSPORTS,
)
self.assertNotEqual(
mcp_transport_config.REMOTE_TRANSPORT,
mcp_transport_config.DEFAULT_TRANSPORT,
)
self.assertTrue(
mcp_transport_config.is_remote_transport(
mcp_transport_config.REMOTE_TRANSPORT
)
)
def test_no_default_transport_literal_outside_the_seam(self):
"""AC3: no production module reads the literal outside the seam.
Review 635 flagged that a fixed five-module list cannot catch a *new*
module reintroducing the literal. This globs every production module in
the repository root instead, so the guarantee holds for code that does
not exist yet.
"""
default = mcp_transport_config.DEFAULT_TRANSPORT
needles = (f'"{default}"', f"'{default}'")
seam = Path(mcp_transport_config.__file__).name
scanned: list[str] = []
offenders: list[str] = []
for path in sorted(REPO_ROOT.glob("*.py")):
if path.name == seam:
continue # the seam is the one place the literal may live
scanned.append(path.name)
for lineno, line in enumerate(
path.read_text(encoding="utf-8").splitlines(), start=1
):
code = line.split("#", 1)[0]
if any(needle in code for needle in needles):
offenders.append(f"{path.name}:{lineno}: {line.strip()}")
# Guard the guard: a glob that silently matched nothing would pass.
self.assertGreater(len(scanned), 20, "production glob matched too little")
self.assertIn("mcp_daemon_guard.py", scanned)
self.assertIn("gitea_mcp_server.py", scanned)
self.assertEqual(offenders, [], "\n".join(offenders))
class TestConfiguredTransportResolution(unittest.TestCase):
"""AC1: configuration supplies the identifier; unset still yields the default."""
def test_unset_yields_default(self):
res = mcp_transport_config.resolve_configured_transport(env={})
self.assertEqual(res["transport"], mcp_transport_config.DEFAULT_TRANSPORT)
self.assertFalse(res["configured"])
self.assertEqual(res["source"], mcp_transport_config.SOURCE_DEFAULT)
self.assertTrue(res["supported"])
self.assertEqual(res["reasons"], [])
def test_blank_and_whitespace_are_treated_as_unset(self):
for raw in ("", " ", "\t\n"):
res = mcp_transport_config.resolve_configured_transport(
env={mcp_transport_config.TRANSPORT_ENV: raw}
)
self.assertEqual(res["transport"], mcp_transport_config.DEFAULT_TRANSPORT)
self.assertFalse(res["configured"])
def test_explicit_default_is_reported_as_configured(self):
res = mcp_transport_config.resolve_configured_transport(
env={
mcp_transport_config.TRANSPORT_ENV: (
mcp_transport_config.DEFAULT_TRANSPORT
)
}
)
self.assertEqual(res["transport"], mcp_transport_config.DEFAULT_TRANSPORT)
self.assertTrue(res["configured"])
self.assertEqual(res["source"], mcp_transport_config.SOURCE_CONFIGURED)
def test_remote_identifier_resolves_and_is_supported(self):
res = mcp_transport_config.resolve_configured_transport(
env={
mcp_transport_config.TRANSPORT_ENV: (
mcp_transport_config.REMOTE_TRANSPORT
)
}
)
self.assertEqual(res["transport"], mcp_transport_config.REMOTE_TRANSPORT)
self.assertTrue(res["supported"])
def test_case_and_padding_are_normalized(self):
padded = f" {mcp_transport_config.REMOTE_TRANSPORT.upper()} "
res = mcp_transport_config.resolve_configured_transport(
env={mcp_transport_config.TRANSPORT_ENV: padded}
)
self.assertEqual(res["transport"], mcp_transport_config.REMOTE_TRANSPORT)
self.assertTrue(res["supported"])
def test_unregistered_identifier_is_not_silently_defaulted(self):
res = mcp_transport_config.resolve_configured_transport(
env={mcp_transport_config.TRANSPORT_ENV: "carrier-pigeon"}
)
self.assertFalse(res["supported"])
self.assertEqual(res["transport"], "carrier-pigeon")
self.assertNotEqual(res["transport"], mcp_transport_config.DEFAULT_TRANSPORT)
self.assertTrue(res["reasons"])
def test_superseded_sse_transport_is_not_registered(self):
"""A real MCP transport that this deployment does not sanction."""
self.assertFalse(mcp_transport_config.is_supported_transport("sse"))
res = mcp_transport_config.resolve_configured_transport(
env={mcp_transport_config.TRANSPORT_ENV: "sse"}
)
self.assertFalse(res["supported"])
def test_require_configured_transport_raises_on_unregistered(self):
with self.assertRaises(mcp_transport_config.TransportConfigurationError):
mcp_transport_config.require_configured_transport(
env={mcp_transport_config.TRANSPORT_ENV: "carrier-pigeon"}
)
def test_non_string_configuration_never_matches_permitted_set(self):
for value in (object(), 1, None, True, ["stdio"], {"t": "stdio"}):
self.assertEqual(mcp_transport_config.normalize_transport(value), "")
self.assertFalse(mcp_transport_config.is_supported_transport(value))
class TestBindSeam(unittest.TestCase):
"""AC1/AC2: the live bind path resolves, validates, and pins."""
def tearDown(self) -> None:
mcp_daemon_guard.clear_native_runtime_for_tests()
def test_1_no_configured_transport_binds_default(self):
with _ProductionBind() as guard:
status = guard.bind_native_mcp_transport()
self.assertEqual(
status["transport"], mcp_transport_config.DEFAULT_TRANSPORT
)
self.assertEqual(
guard.bound_transport(), mcp_transport_config.DEFAULT_TRANSPORT
)
self.assertTrue(status["production_native_mcp_transport"])
def test_2_explicit_default_transport_binds(self):
with _ProductionBind(
{
mcp_transport_config.TRANSPORT_ENV: (
mcp_transport_config.DEFAULT_TRANSPORT
)
}
) as guard:
status = guard.bind_native_mcp_transport()
self.assertEqual(
status["transport"], mcp_transport_config.DEFAULT_TRANSPORT
)
self.assertTrue(guard.is_production_native_mcp_transport())
def test_2b_explicit_argument_still_binds(self):
"""The pre-#931 call form keeps working for launchers and tests."""
with _ProductionBind() as guard:
status = guard.bind_native_mcp_transport(
transport=mcp_transport_config.DEFAULT_TRANSPORT
)
self.assertEqual(
status["transport"], mcp_transport_config.DEFAULT_TRANSPORT
)
def test_3_sanctioned_remote_identifier_binds_through_the_seam(self):
with _ProductionBind(
{
mcp_transport_config.TRANSPORT_ENV: (
mcp_transport_config.REMOTE_TRANSPORT
)
}
) as guard:
status = guard.bind_native_mcp_transport()
self.assertEqual(
status["transport"], mcp_transport_config.REMOTE_TRANSPORT
)
self.assertEqual(
guard.bound_transport(), mcp_transport_config.REMOTE_TRANSPORT
)
# The remote identifier is trusted exactly like the local one; the
# listener that serves it is #938 and is not implemented here.
self.assertTrue(guard.is_native_mcp_transport())
self.assertTrue(guard.is_production_native_mcp_transport())
guard.assert_production_mutation_runtime("remote-bind")
def test_4_unregistered_identifier_rejected_at_bind_time(self):
with _ProductionBind(
{mcp_transport_config.TRANSPORT_ENV: "carrier-pigeon"}
) as guard:
with self.assertRaises(guard.UnsanctionedRuntimeError) as ctx:
guard.bind_native_mcp_transport()
self.assertIn("carrier-pigeon", str(ctx.exception))
self.assertIn("#931", str(ctx.exception))
# Nothing was bound, so nothing may dispatch.
self.assertIsNone(guard.bound_transport())
self.assertFalse(guard.is_native_mcp_transport())
def test_4b_unregistered_explicit_argument_rejected(self):
with _ProductionBind() as guard:
with self.assertRaises(guard.UnsanctionedRuntimeError):
guard.bind_native_mcp_transport(transport="carrier-pigeon")
self.assertIsNone(guard.bound_transport())
def test_4c_superseded_sse_rejected_at_bind_time(self):
with _ProductionBind({mcp_transport_config.TRANSPORT_ENV: "sse"}) as guard:
with self.assertRaises(guard.UnsanctionedRuntimeError):
guard.bind_native_mcp_transport()
self.assertIsNone(guard.bound_transport())
def test_5_invalid_bind_prevents_tool_service(self):
"""A failed bind must stop the server before it serves tools."""
with _ProductionBind(
{mcp_transport_config.TRANSPORT_ENV: "carrier-pigeon"}
) as guard:
with self.assertRaises(guard.UnsanctionedRuntimeError):
guard.bind_native_mcp_transport()
# This is the exact expression the entrypoint passes to mcp.run.
with self.assertRaises(guard.UnsanctionedRuntimeError) as ctx:
guard.assert_transport_bound("tool service")
self.assertIn("No MCP transport is bound", str(ctx.exception))
def test_6_unbound_state_fails_closed(self):
"""Entrypoint claimed but never bound — the offline-import shape."""
with _ProductionBind() as guard:
self.assertIsNone(guard.bound_transport())
self.assertFalse(guard.is_native_mcp_transport())
with self.assertRaises(guard.UnsanctionedRuntimeError):
guard.assert_transport_bound("tool service")
with self.assertRaises(guard.UnsanctionedRuntimeError):
guard.assert_sanctioned_mutation_runtime("gitea_mutation")
def test_6b_no_runtime_at_all_fails_closed(self):
mcp_daemon_guard.clear_native_runtime_for_tests()
with patch.object(mcp_daemon_guard, "is_pytest_runtime", return_value=False):
self.assertIsNone(mcp_daemon_guard.bound_transport())
with self.assertRaises(mcp_daemon_guard.UnsanctionedRuntimeError):
mcp_daemon_guard.assert_transport_bound("tool service")
class TestOneAuthoritativeValue(unittest.TestCase):
"""AC: every transport-aware guard observes the same value."""
def tearDown(self) -> None:
mcp_daemon_guard.clear_native_runtime_for_tests()
def test_7_all_guards_read_the_same_bound_value(self):
with _ProductionBind(
{
mcp_transport_config.TRANSPORT_ENV: (
mcp_transport_config.REMOTE_TRANSPORT
)
}
) as guard:
guard.bind_native_mcp_transport()
expected = mcp_transport_config.REMOTE_TRANSPORT
self.assertEqual(guard.bound_transport(), expected)
self.assertEqual(guard.assert_transport_bound(), expected)
self.assertEqual(guard.native_runtime_status()["bound_transport"], expected)
self.assertEqual(guard.native_runtime_status()["transport"], expected)
self.assertEqual(
guard.mutation_provenance_fields()["bound_transport"], expected
)
self.assertEqual(
irrecoverable_provenance.assess_transport_for_auth_mint()[
"bound_transport"
],
expected,
)
def test_8_environment_change_after_bind_cannot_move_the_value(self):
"""Client- or environment-shaped input must not alter a bound transport."""
with _ProductionBind() as guard:
guard.bind_native_mcp_transport()
self.assertEqual(
guard.bound_transport(), mcp_transport_config.DEFAULT_TRANSPORT
)
# A stray launcher (or an attacker) rewrites config post-bind.
os.environ[mcp_transport_config.TRANSPORT_ENV] = (
mcp_transport_config.REMOTE_TRANSPORT
)
self.assertEqual(
guard.bound_transport(), mcp_transport_config.DEFAULT_TRANSPORT
)
self.assertEqual(
guard.mutation_provenance_fields()["bound_transport"],
mcp_transport_config.DEFAULT_TRANSPORT,
)
os.environ[mcp_transport_config.TRANSPORT_ENV] = "carrier-pigeon"
self.assertEqual(
guard.bound_transport(), mcp_transport_config.DEFAULT_TRANSPORT
)
def test_8b_bound_transport_takes_no_caller_argument(self):
"""The accessor cannot be steered by a tool parameter."""
import inspect
self.assertEqual(
list(inspect.signature(mcp_daemon_guard.bound_transport).parameters), []
)
def test_11_rebinding_the_same_transport_is_idempotent(self):
with _ProductionBind() as guard:
first = guard.bind_native_mcp_transport()
second = guard.bind_native_mcp_transport()
self.assertEqual(first["transport"], second["transport"])
self.assertEqual(
guard.bound_transport(), mcp_transport_config.DEFAULT_TRANSPORT
)
def test_11b_rebinding_a_different_transport_fails_closed(self):
with _ProductionBind() as guard:
guard.bind_native_mcp_transport()
with self.assertRaises(guard.UnsanctionedRuntimeError) as ctx:
guard.bind_native_mcp_transport(
transport=mcp_transport_config.REMOTE_TRANSPORT
)
self.assertIn("already bound", str(ctx.exception))
# The first value survives the attempt.
self.assertEqual(
guard.bound_transport(), mcp_transport_config.DEFAULT_TRANSPORT
)
def test_11c_rebinding_an_unregistered_transport_fails_closed(self):
with _ProductionBind() as guard:
guard.bind_native_mcp_transport()
with self.assertRaises(guard.UnsanctionedRuntimeError):
guard.bind_native_mcp_transport(transport="carrier-pigeon")
self.assertEqual(
guard.bound_transport(), mcp_transport_config.DEFAULT_TRANSPORT
)
class TestDurableRecordCarriesTransport(unittest.TestCase):
"""AC4: the identifier reaches a durable decision-lock record."""
def tearDown(self) -> None:
mcp_daemon_guard.clear_native_runtime_for_tests()
def _save_and_read_decision_lock(self, state_dir: str) -> dict:
mcp_session_state.save_state(
kind=mcp_session_state.KIND_DECISION_LOCK,
payload={"pr_number": 931, "action": "COMMENT"},
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
profile_identity="prgs-author",
state_dir=state_dir,
)
loaded = mcp_session_state.load_state(
kind=mcp_session_state.KIND_DECISION_LOCK,
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
profile_identity="prgs-author",
state_dir=state_dir,
)
self.assertIsNotNone(loaded)
return loaded
def test_9_decision_lock_records_the_bound_transport(self):
with tempfile.TemporaryDirectory() as tmp:
mcp_daemon_guard.clear_native_runtime_for_tests()
mcp_daemon_guard.install_test_native_runtime()
record = self._save_and_read_decision_lock(tmp)
self.assertIn("bound_transport", record)
self.assertEqual(
record["bound_transport"], mcp_daemon_guard.bound_transport()
)
def test_9b_unbound_runtime_records_no_transport_identifier(self):
with tempfile.TemporaryDirectory() as tmp:
mcp_daemon_guard.clear_native_runtime_for_tests()
record = self._save_and_read_decision_lock(tmp)
self.assertIn("bound_transport", record)
self.assertIsNone(record["bound_transport"])
def test_9c_provenance_fields_expose_the_identifier(self):
mcp_daemon_guard.clear_native_runtime_for_tests()
fields = mcp_daemon_guard.mutation_provenance_fields()
self.assertIn("bound_transport", fields)
self.assertIsNone(fields["bound_transport"])
class TestStdioBehaviourUnchanged(unittest.TestCase):
"""AC6 / prompt items 10 and 12: no regression on the existing path."""
def tearDown(self) -> None:
mcp_daemon_guard.clear_native_runtime_for_tests()
def test_10_trust_class_field_keeps_its_pre_931_values(self):
"""``transport`` remains the trust class, not the identifier."""
mcp_daemon_guard.clear_native_runtime_for_tests()
self.assertEqual(
mcp_daemon_guard.mutation_provenance_fields()["transport"], "untrusted"
)
mcp_daemon_guard.install_test_native_runtime()
self.assertEqual(
mcp_daemon_guard.mutation_provenance_fields()["transport"],
"test_native_mcp",
)
with _ProductionBind() as guard:
guard.bind_native_mcp_transport()
self.assertEqual(
guard.mutation_provenance_fields()["transport"], "native_mcp"
)
def test_10b_default_bind_reproduces_the_pre_931_runtime_record(self):
with _ProductionBind() as guard:
status = guard.bind_native_mcp_transport()
self.assertTrue(status["native_mcp_transport"])
self.assertTrue(status["production_native_mcp_transport"])
self.assertEqual(status["mode"], "production")
self.assertEqual(status["phase"], "transport_bound")
self.assertEqual(
status["transport"], mcp_transport_config.DEFAULT_TRANSPORT
)
guard.assert_sanctioned_mutation_runtime("native-ide")
guard.assert_production_mutation_runtime("native-ide")
def test_12_session_state_root_still_pinned_at_bind(self):
"""#695 AC2 must survive the seam."""
with tempfile.TemporaryDirectory() as legit:
with tempfile.TemporaryDirectory() as rogue:
with _ProductionBind(
{mcp_daemon_guard.SESSION_STATE_DIR_ENV: legit}
) as guard:
guard.bind_native_mcp_transport()
self.assertEqual(
guard.pinned_session_state_dir(), str(Path(legit).resolve())
)
os.environ[mcp_daemon_guard.SESSION_STATE_DIR_ENV] = rogue
self.assertEqual(
guard.pinned_session_state_dir(), str(Path(legit).resolve())
)
def test_12b_test_mode_record_still_cannot_authorize_production(self):
mcp_daemon_guard.clear_native_runtime_for_tests()
mcp_daemon_guard.install_test_native_runtime()
self.assertTrue(mcp_daemon_guard.is_native_mcp_transport())
self.assertFalse(mcp_daemon_guard.is_production_native_mcp_transport())
with self.assertRaises(mcp_daemon_guard.UnsanctionedRuntimeError):
mcp_daemon_guard.assert_production_mutation_runtime("prod-endpoint")
def test_12c_bind_still_requires_the_canonical_entrypoint(self):
"""A remote identifier does not relax provenance."""
mcp_daemon_guard.clear_native_runtime_for_tests()
with patch.object(mcp_daemon_guard, "is_pytest_runtime", return_value=False):
with patch.object(
mcp_daemon_guard,
"_caller_official_entrypoint_path",
return_value=None,
):
with self.assertRaises(mcp_daemon_guard.UnsanctionedRuntimeError):
mcp_daemon_guard.bind_native_mcp_transport(
transport=mcp_transport_config.REMOTE_TRANSPORT
)
self.assertIsNone(mcp_daemon_guard.bound_transport())
def test_12d_bind_still_requires_a_prior_entrypoint_claim(self):
mcp_daemon_guard.clear_native_runtime_for_tests()
canonical = str((REPO_ROOT / "mcp_server.py").resolve())
with patch.object(mcp_daemon_guard, "is_pytest_runtime", return_value=False):
with patch.object(
mcp_daemon_guard,
"_caller_official_entrypoint_path",
side_effect=lambda: canonical,
):
# No mark_sanctioned_daemon() first.
with self.assertRaises(
mcp_daemon_guard.UnsanctionedRuntimeError
) as ctx:
mcp_daemon_guard.bind_native_mcp_transport()
self.assertIn("no entrypoint claim", str(ctx.exception))
def test_12e_auth_mint_verdict_is_unchanged_for_the_default_transport(self):
with _ProductionBind() as guard:
guard.bind_native_mcp_transport()
verdict = irrecoverable_provenance.assess_transport_for_auth_mint()
self.assertTrue(verdict["allowed"])
self.assertTrue(verdict["native_mcp_transport"])
self.assertTrue(verdict["production_native_mcp_transport"])
self.assertEqual(verdict["reasons"], [])
def test_12f_auth_mint_still_refuses_an_unbound_runtime(self):
with _ProductionBind():
# Claimed but never bound.
verdict = irrecoverable_provenance.assess_transport_for_auth_mint()
self.assertFalse(verdict["allowed"])
self.assertTrue(verdict["reasons"])
self.assertIsNone(verdict["bound_transport"])
class TestExecutionBoundary(unittest.TestCase):
"""Recognition is not execution authorization (#931, review 635 B1/B2).
A registered remote identifier must still bind, pin and record the #931
seam while being refused at the serve boundary, because serving it would
start a listener with no authentication or per-request principal. That
listener belongs to #938.
"""
def tearDown(self) -> None:
mcp_daemon_guard.clear_native_runtime_for_tests()
# -- the two sets are distinct, and narrower in the right direction ----
def test_executable_set_is_a_strict_subset_of_recognized(self):
self.assertTrue(
mcp_transport_config.EXECUTABLE_TRANSPORTS
< mcp_transport_config.SUPPORTED_TRANSPORTS
)
def test_default_transport_is_executable(self):
self.assertTrue(
mcp_transport_config.is_executable_transport(
mcp_transport_config.DEFAULT_TRANSPORT
)
)
def test_remote_transport_is_recognized_but_not_executable(self):
remote = mcp_transport_config.REMOTE_TRANSPORT
self.assertTrue(mcp_transport_config.is_supported_transport(remote))
self.assertFalse(mcp_transport_config.is_executable_transport(remote))
def test_remote_listener_ownership_is_declared(self):
self.assertEqual(
mcp_transport_config.TRANSPORT_EXECUTION_OWNER[
mcp_transport_config.REMOTE_TRANSPORT
],
"#938",
)
# -- stdio still reaches the runner, unchanged -------------------------
def test_default_transport_is_authorized_for_service(self):
with _ProductionBind() as guard:
guard.bind_native_mcp_transport()
self.assertEqual(
guard.authorize_transport_execution("tool service"),
mcp_transport_config.DEFAULT_TRANSPORT,
)
self.assertTrue(guard.assess_serve_authorization()["allowed"])
def test_default_transport_reaches_the_real_runner(self):
"""The production runner is actually invoked, with stdio, unchanged."""
with _ProductionBind() as guard:
guard.bind_native_mcp_transport()
seen = {}
class _Runner:
def run(self, transport=None, **kw):
seen["transport"] = transport
_Runner().run(transport=guard.authorize_transport_execution("tool service"))
self.assertEqual(
seen["transport"], mcp_transport_config.DEFAULT_TRANSPORT
)
# -- streamable-http binds, records, and is refused before serving -----
def test_remote_transport_binds_and_is_recorded(self):
"""The #931 seam is intact: it binds, pins and records."""
with _ProductionBind(
{
mcp_transport_config.TRANSPORT_ENV: (
mcp_transport_config.REMOTE_TRANSPORT
)
}
) as guard:
guard.bind_native_mcp_transport()
self.assertEqual(
guard.bound_transport(), mcp_transport_config.REMOTE_TRANSPORT
)
self.assertEqual(
guard.mutation_provenance_fields()["bound_transport"],
mcp_transport_config.REMOTE_TRANSPORT,
)
def test_remote_transport_is_refused_at_the_serve_boundary(self):
with _ProductionBind(
{
mcp_transport_config.TRANSPORT_ENV: (
mcp_transport_config.REMOTE_TRANSPORT
)
}
) as guard:
guard.bind_native_mcp_transport()
with self.assertRaises(guard.TransportExecutionError) as ctx:
guard.authorize_transport_execution("tool service")
err = ctx.exception
self.assertEqual(
err.blocker_kind,
mcp_transport_config.BLOCKER_LISTENER_NOT_COMMISSIONED,
)
self.assertEqual(err.owner_issue, "#938")
self.assertEqual(err.transport, mcp_transport_config.REMOTE_TRANSPORT)
def test_refusal_names_transport_and_unmet_requirement_and_owner(self):
with _ProductionBind(
{
mcp_transport_config.TRANSPORT_ENV: (
mcp_transport_config.REMOTE_TRANSPORT
)
}
) as guard:
guard.bind_native_mcp_transport()
with self.assertRaises(guard.TransportExecutionError) as ctx:
guard.authorize_transport_execution("tool service")
text = str(ctx.exception)
self.assertIn(mcp_transport_config.REMOTE_TRANSPORT, text)
self.assertIn("#938", text)
self.assertIn("not commissioned", text)
self.assertIn("#931", text)
def test_remote_transport_never_reaches_the_runner(self):
"""No transport value is handed to a run() call for the remote case."""
with _ProductionBind(
{
mcp_transport_config.TRANSPORT_ENV: (
mcp_transport_config.REMOTE_TRANSPORT
)
}
) as guard:
guard.bind_native_mcp_transport()
calls = []
class _Runner:
def run(self, transport=None, **kw):
calls.append(transport)
with self.assertRaises(guard.TransportExecutionError):
_Runner().run(
transport=guard.authorize_transport_execution("tool service")
)
self.assertEqual(calls, [], "runner must never be invoked")
def test_no_http_listener_is_created_for_remote_transport(self):
"""Nothing in the refusal path touches uvicorn or a socket bind."""
with _ProductionBind(
{
mcp_transport_config.TRANSPORT_ENV: (
mcp_transport_config.REMOTE_TRANSPORT
)
}
) as guard:
guard.bind_native_mcp_transport()
import socket
bound_sockets = []
real_bind = socket.socket.bind
def _tripwire(self, addr): # pragma: no cover - must not run
bound_sockets.append(addr)
return real_bind(self, addr)
with patch.object(socket.socket, "bind", _tripwire):
with self.assertRaises(guard.TransportExecutionError):
guard.authorize_transport_execution("tool service")
self.assertEqual(bound_sockets, [], "no socket may be bound")
def test_no_mutation_is_authorized_after_the_denial(self):
"""The denial leaves no partial state that would let a tool dispatch."""
with _ProductionBind(
{
mcp_transport_config.TRANSPORT_ENV: (
mcp_transport_config.REMOTE_TRANSPORT
)
}
) as guard:
guard.bind_native_mcp_transport()
with self.assertRaises(guard.TransportExecutionError):
guard.authorize_transport_execution("tool service")
# Serve stays unauthorized on every subsequent query.
self.assertFalse(guard.assess_serve_authorization()["allowed"])
self.assertFalse(guard.native_runtime_status()["serve_authorized"])
with self.assertRaises(guard.TransportExecutionError):
guard.authorize_transport_execution("tool service")
# -- B2: the decision genuinely consumes bound_transport ---------------
def test_serve_decision_consumes_bound_transport(self):
"""Changing only bound_transport flips the verdict."""
with _ProductionBind() as guard:
guard.bind_native_mcp_transport()
self.assertTrue(guard.assess_serve_authorization()["allowed"])
with patch.object(
guard,
"bound_transport",
return_value=mcp_transport_config.REMOTE_TRANSPORT,
):
verdict = guard.assess_serve_authorization()
self.assertFalse(verdict["allowed"])
self.assertEqual(
verdict["transport"], mcp_transport_config.REMOTE_TRANSPORT
)
with self.assertRaises(guard.TransportExecutionError):
guard.authorize_transport_execution("tool service")
def test_serve_verdict_reports_the_bound_transport(self):
with _ProductionBind(
{
mcp_transport_config.TRANSPORT_ENV: (
mcp_transport_config.REMOTE_TRANSPORT
)
}
) as guard:
guard.bind_native_mcp_transport()
verdict = guard.assess_serve_authorization()
self.assertEqual(
verdict["transport"], mcp_transport_config.REMOTE_TRANSPORT
)
self.assertTrue(verdict["recognized"])
self.assertFalse(verdict["executable"])
# -- earlier and later boundaries are unchanged ------------------------
def test_unregistered_identifier_still_fails_at_bind_not_at_serve(self):
with _ProductionBind(
{mcp_transport_config.TRANSPORT_ENV: "carrier-pigeon"}
) as guard:
with self.assertRaises(guard.UnsanctionedRuntimeError) as ctx:
guard.bind_native_mcp_transport()
self.assertIn("not a registered MCP transport", str(ctx.exception))
self.assertIsNone(guard.bound_transport())
def test_unbound_execution_keeps_the_pre_existing_failure(self):
with _ProductionBind() as guard:
with self.assertRaises(guard.UnsanctionedRuntimeError) as ctx:
guard.authorize_transport_execution("tool service")
self.assertIn("No MCP transport is bound", str(ctx.exception))
self.assertNotIsInstance(ctx.exception, guard.TransportExecutionError)
def test_transport_execution_error_is_caught_by_existing_handlers(self):
"""Subclassing keeps every pre-existing fail-closed handler correct."""
self.assertTrue(
issubclass(
mcp_daemon_guard.TransportExecutionError,
mcp_daemon_guard.UnsanctionedRuntimeError,
)
)
def test_test_mode_runtime_is_not_servable(self):
"""The pytest-only record is outside the recognized and executable sets."""
mcp_daemon_guard.clear_native_runtime_for_tests()
mcp_daemon_guard.install_test_native_runtime()
self.assertNotIn(
mcp_daemon_guard.bound_transport(),
mcp_transport_config.SUPPORTED_TRANSPORTS,
)
self.assertFalse(mcp_daemon_guard.assess_serve_authorization()["allowed"])
def test_serve_authorization_does_not_leak_across_runtimes(self):
"""A later runtime's verdict never reflects an earlier one."""
with _ProductionBind(
{
mcp_transport_config.TRANSPORT_ENV: (
mcp_transport_config.REMOTE_TRANSPORT
)
}
) as guard:
guard.bind_native_mcp_transport()
self.assertFalse(guard.assess_serve_authorization()["allowed"])
with _ProductionBind() as guard:
guard.bind_native_mcp_transport()
self.assertTrue(guard.assess_serve_authorization()["allowed"])
mcp_daemon_guard.clear_native_runtime_for_tests()
self.assertEqual(
mcp_daemon_guard.assess_serve_authorization()["blocker_kind"],
mcp_transport_config.BLOCKER_TRANSPORT_NOT_BOUND,
)
def test_refusal_carries_no_credential_material(self):
with _ProductionBind(
{
mcp_transport_config.TRANSPORT_ENV: (
mcp_transport_config.REMOTE_TRANSPORT
),
"GITEA_TOKEN": "super-secret-value",
}
) as guard:
guard.bind_native_mcp_transport()
with self.assertRaises(guard.TransportExecutionError) as ctx:
guard.authorize_transport_execution("tool service")
blob = str(ctx.exception) + repr(ctx.exception.assessment)
self.assertNotIn("super-secret-value", blob)
self.assertNotIn("GITEA_TOKEN", blob)
class TestEntrypointWiring(unittest.TestCase):
"""The live entrypoint must use the seam and the execution guard."""
def test_entrypoint_binds_without_a_literal_transport(self):
text = (REPO_ROOT / "gitea_mcp_server.py").read_text(encoding="utf-8")
self.assertIn("mcp_daemon_guard.bind_native_mcp_transport()", text)
self.assertNotIn('bind_native_mcp_transport(transport="stdio")', text)
def test_entrypoint_serves_only_through_the_execution_guard(self):
text = (REPO_ROOT / "gitea_mcp_server.py").read_text(encoding="utf-8")
self.assertIn(
"mcp.run(transport=mcp_daemon_guard.authorize_transport_execution", text
)
self.assertNotIn('mcp.run(transport="stdio")', text)
self.assertNotIn(
"mcp.run(transport=mcp_daemon_guard.assert_transport_bound", text
)
def test_every_run_call_in_production_goes_through_the_guard(self):
"""Glob the production surface: no serve site may bypass the guard."""
offenders: list[str] = []
run_sites = 0
for path in sorted(REPO_ROOT.glob("*.py")):
for lineno, line in enumerate(
path.read_text(encoding="utf-8").splitlines(), start=1
):
code = line.split("#", 1)[0]
if "mcp.run(" not in code:
continue
run_sites += 1
if "authorize_transport_execution" not in code:
offenders.append(f"{path.name}:{lineno}: {line.strip()}")
self.assertEqual(run_sites, 1, "expected exactly one serve site")
self.assertEqual(offenders, [], "\n".join(offenders))
if __name__ == "__main__":
unittest.main()
File diff suppressed because it is too large Load Diff
+323
View File
@@ -0,0 +1,323 @@
"""Validation tooling for the remote-MCP threat model (#956).
#956 requires that "every boundary claim [is] traceable to a file and line
anchor that resolves at the reviewed commit". A prose document cannot enforce
that about itself, and #930 demonstrated the failure mode: its inventory cited
``gitea_mcp_server.py`` anchors generated at ``7bf4f125`` which no longer point
at the described code at ``aad5c8b4``. Nothing failed, because nothing checked.
These tests are that check. They enforce, in both directions:
* every ``file.py:NNN`` anchor cited in the prose is declared in the fixture;
* every declared anchor resolves the file exists, the line exists, and the
source line actually contains the substring the fixture claims for it;
* the document's structural obligations (assets, adversaries, boundaries,
credential rows, the co-residency ruling, and the child mapping) are present
and internally consistent.
A refactor that shifts a line number therefore breaks the suite instead of
silently rotting the security documentation.
"""
import json
import os
import re
import unittest
REPO_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
DOC_PATH = os.path.join(REPO_ROOT, "docs", "remote-mcp", "threat-model.md")
FIXTURE_PATH = os.path.join(
REPO_ROOT, "docs", "remote-mcp", "threat-model-anchors.json"
)
# ``module.py:123`` as it appears inside markdown inline code spans.
ANCHOR_RE = re.compile(r"`([A-Za-z0-9_./-]+\.py):(\d+)`")
# The epic children this document must map to a boundary (#929 children 2-10).
REQUIRED_CHILDREN = [931, 932, 933, 934, 935, 936, 937, 938, 939]
# The adversaries #956 names explicitly.
REQUIRED_ADVERSARIES = [
"compromised LLM client",
"prompt injection",
"malicious tool arguments",
"network attacker",
"curious operator",
]
def _read(path):
with open(path, "r", encoding="utf-8") as fh:
return fh.read()
def _heading_re(title):
"""Match a level-2 heading by title, with or without section numbering.
The document numbers its sections ('## 6. Decomposition ruling'), so an
exact-substring assertion would break on renumbering without the document
having actually lost anything.
"""
return re.compile(
r"^##\s+(?:\d+\.\s+)?" + re.escape(title), re.MULTILINE
)
def _section_body(doc, title):
"""Return the text of section *title*, bounded by the next level-2 heading.
Bounding matters: an unbounded slice runs to end-of-document, so the
walkthrough tables in a later section leak into the child-to-boundary
mapping and satisfy its coverage check with rows that assign no owner.
"""
match = _heading_re(title).search(doc)
if match is None:
return None
rest = doc[match.end():]
nxt = re.search(r"^##\s", rest, re.MULTILINE)
return rest[: nxt.start()] if nxt else rest
def _source_line(rel_path, lineno):
"""Return the 1-based *lineno* of *rel_path*, or None if out of range."""
abs_path = os.path.join(REPO_ROOT, rel_path)
if not os.path.exists(abs_path):
return None
with open(abs_path, "r", encoding="utf-8", errors="replace") as fh:
for idx, line in enumerate(fh, start=1):
if idx == lineno:
return line
return None
class ThreatModelFixtureTests(unittest.TestCase):
"""The fixture itself must be well-formed before it can prove anything."""
def setUp(self):
self.fixture = json.loads(_read(FIXTURE_PATH))
def test_fixture_declares_a_generation_commit(self):
sha = self.fixture.get("generated_against_commit") or ""
self.assertRegex(
sha,
r"^[0-9a-f]{40}$",
"the fixture must record the full commit its anchors were taken at",
)
def test_fixture_anchors_are_unique_and_well_formed(self):
seen = set()
for entry in self.fixture["anchors"]:
anchor = entry["anchor"]
self.assertNotIn(anchor, seen, f"duplicate anchor entry: {anchor}")
seen.add(anchor)
self.assertRegex(anchor, r"^[A-Za-z0-9_./-]+\.py:[1-9]\d*$", anchor)
self.assertTrue(
(entry.get("expect") or "").strip(),
f"anchor {anchor} declares no 'expect' substring, so it proves nothing",
)
class ThreatModelAnchorResolutionTests(unittest.TestCase):
"""#956 required positive test: every anchor resolves at the reviewed commit."""
def setUp(self):
self.fixture = json.loads(_read(FIXTURE_PATH))
self.doc = _read(DOC_PATH)
def test_every_declared_anchor_resolves_to_the_claimed_source_line(self):
failures = []
for entry in self.fixture["anchors"]:
rel_path, _, raw_lineno = entry["anchor"].partition(":")
lineno = int(raw_lineno)
line = _source_line(rel_path, lineno)
if line is None:
failures.append(f"{entry['anchor']}: file or line does not exist")
continue
if entry["expect"] not in line:
failures.append(
f"{entry['anchor']}: expected {entry['expect']!r}, "
f"found {line.strip()!r}"
)
self.assertEqual(
[], failures, "unresolved threat-model anchors:\n" + "\n".join(failures)
)
def test_every_anchor_cited_in_the_document_is_declared_in_the_fixture(self):
declared = {e["anchor"] for e in self.fixture["anchors"]}
cited = {f"{m.group(1)}:{m.group(2)}" for m in ANCHOR_RE.finditer(self.doc)}
undeclared = sorted(cited - declared)
self.assertEqual(
[],
undeclared,
"document cites anchors that no test verifies: " + ", ".join(undeclared),
)
def test_the_document_actually_cites_anchors(self):
cited = {f"{m.group(1)}:{m.group(2)}" for m in ANCHOR_RE.finditer(self.doc)}
self.assertGreaterEqual(
len(cited),
30,
"a boundary document with almost no anchors is not traceable",
)
def test_unresolvable_anchor_is_detected(self):
"""Negative control: the checker must fail on a deliberately bad anchor.
Without this, a checker that silently passed everything would look
identical to a correct one.
"""
self.assertIsNone(_source_line("gitea_config.py", 10**9))
self.assertIsNone(_source_line("no_such_module_for_956.py", 1))
real = _source_line("gitea_config.py", 54)
self.assertIsNotNone(real)
self.assertNotIn("this substring is not on that line", real)
class ThreatModelStructureTests(unittest.TestCase):
"""The document must contain what #956's acceptance criteria demand."""
def setUp(self):
self.doc = _read(DOC_PATH)
def test_records_the_commit_it_was_generated_against(self):
fixture = json.loads(_read(FIXTURE_PATH))
self.assertIn(
fixture["generated_against_commit"],
self.doc,
"the document must state the commit its anchors resolve at",
)
def test_names_every_required_adversary(self):
low = self.doc.lower()
for adversary in REQUIRED_ADVERSARIES:
self.assertIn(adversary.lower(), low, f"adversary not covered: {adversary}")
def test_maps_every_epic_child_from_two_through_ten(self):
for number in REQUIRED_CHILDREN:
self.assertIn(
f"#{number}",
self.doc,
f"epic child #{number} is not mapped to a boundary",
)
def test_credential_rows_declare_holder_boundary_and_blast_radius(self):
for column in ("Holder", "Boundary", "Blast radius"):
self.assertIn(
column,
self.doc,
f"the credential inventory must state each credential's {column.lower()}",
)
def test_states_an_explicit_co_residency_ruling(self):
"""AC3/AC5: an explicit ruling, not an implication."""
self.assertIsNotNone(
_heading_re("Decomposition ruling").search(self.doc),
"the document must contain an explicit decomposition-ruling section",
)
for service in ("Jenkins", "GlitchTip", "Sentry", "database"):
self.assertIn(service, self.doc, f"ruling does not address {service}")
self.assertRegex(
self.doc,
r"D1\b.*must not",
"the ruling must state the prohibition, not merely discuss it",
)
def test_contains_the_compromised_client_walkthrough(self):
"""#956 required negative/adversarial test."""
self.assertIsNotNone(
_heading_re("Adversarial walkthrough").search(self.doc),
"the required compromised-client walkthrough is missing",
)
self.assertIn("Before the migration", self.doc)
self.assertIn("After the migration", self.doc)
def test_every_boundary_states_what_it_protects_and_what_crossing_requires(self):
boundary_ids = set(re.findall(r"\bB(\d+)\b", self.doc))
self.assertGreaterEqual(
len(boundary_ids), 5, "too few trust boundaries to be a decomposition"
)
for column in (
"Protects",
"Crossing requires today",
"Crossing must require remotely",
):
self.assertIn(column, self.doc, f"boundary table is missing '{column}'")
def test_declares_itself_documentation_only(self):
self.assertIn("documentation only", self.doc.lower())
class ThreatModelConsistencyTests(unittest.TestCase):
"""Counts stated in prose must match the rows actually present."""
def setUp(self):
self.doc = _read(DOC_PATH)
def _declared_ids(self, prefix):
# Table rows begin '| CR1 |' / '| B3 |' / '| A2 |'.
return sorted(
{
int(m)
for m in re.findall(
r"^\|\s*%s(\d+)\s*\|" % prefix, self.doc, re.MULTILINE
)
}
)
def test_identifier_sequences_have_no_gaps(self):
for prefix, label in (
("A", "assets"),
("B", "boundaries"),
("CR", "credentials"),
):
ids = self._declared_ids(prefix)
self.assertTrue(ids, f"no {label} declared")
self.assertEqual(
list(range(1, len(ids) + 1)),
ids,
f"{label} identifiers must run 1..n with no gaps; got {ids}",
)
def test_stated_credential_count_matches_the_rows(self):
ids = self._declared_ids("CR")
match = re.search(r"(\d+)\s+credential(?:s)? in total", self.doc)
self.assertIsNotNone(match, "the credential inventory must state its own total")
self.assertEqual(
len(ids),
int(match.group(1)),
"stated credential total disagrees with the number of rows",
)
def test_every_boundary_is_owned_by_at_least_one_child(self):
"""Each boundary must be owned by a child *in the mapping table*.
Scanning the whole section would let a prose summary line ("Boundary
coverage: ... B5 (#936)") satisfy the assertion while the table row
that actually assigns the owner had been emptied verified by
deliberately blanking a row and watching a whole-section check still
pass. Only table rows count.
"""
mapping_section = _section_body(self.doc, "Child-to-boundary mapping")
self.assertIsNotNone(
mapping_section, "child-to-boundary mapping section is missing"
)
rows = [
line
for line in mapping_section.splitlines()
if line.lstrip().startswith("|") and re.search(r"#93\d", line)
]
self.assertGreaterEqual(
len(rows), len(REQUIRED_CHILDREN), "mapping table has too few child rows"
)
mapped = set(re.findall(r"\bB(\d+)\b", "\n".join(rows)))
declared = {str(i) for i in self._declared_ids("B")}
unmapped = sorted(declared - mapped, key=int)
self.assertEqual(
[],
unmapped,
"boundaries with no owning child: " + ", ".join("B" + u for u in unmapped),
)
if __name__ == "__main__":
unittest.main()
-788
View File
@@ -1,788 +0,0 @@
"""Authoritative fleet inventory classification (#949).
One test group per acceptance criterion. The classifier is pure, so every
scenario is expressed as a snapshot: registry rows plus a process observation.
PID liveness is the one impure input, so it is patched per test rather than
depending on whatever happens to be running on the machine.
"""
import os
import unittest
from datetime import datetime, timedelta, timezone
from unittest import mock
import mcp_fleet_inventory as mfi
NOW = datetime(2026, 7, 28, 3, 0, 0, tzinfo=timezone.utc)
HEAD_A = "82d71b77028a7abd4f8ab4a4e4d89658a187f73d"
HEAD_B = "35ed8a2fcb11134a37c862ca6eaca26e3028902a"
COHORT_A = "ppid:40990"
COHORT_B = "ppid:51022"
BINDING = {"remote": "prgs", "org": "Scaled-Tech-Consulting", "repo": "Gitea-Tools"}
_BASE_PID = 41000
def row(
namespace,
profile,
role,
pid,
*,
cohort_id=COHORT_A,
startup_head=HEAD_A,
registered_at="2026-07-28T02:00:00Z",
**overrides,
):
"""One control-plane runtime registry row."""
record = {
"runtime_id": f"{namespace}:{pid}:boot{pid}",
"namespace": namespace,
"profile": profile,
"role": role,
"remote": BINDING["remote"],
"org": BINDING["org"],
"repo": BINDING["repo"],
"repository_root": "/checkout/Gitea-Tools",
"pid": pid,
"cohort_id": cohort_id,
"cohort_source": "parent_process",
"client_provenance": "client_managed",
"boot_id": f"boot{pid}",
"startup_head": startup_head,
"daemon_start_head": startup_head,
"transport": "stdio",
"registered_at": registered_at,
"last_heartbeat_at": registered_at,
"status": "running",
}
record.update(overrides)
return record
def healthy_rows():
"""One live row per expected member, all one cohort, all one revision."""
return [
row(entry["namespace"], entry["profile"], entry["role"], _BASE_PID + index)
for index, entry in enumerate(mfi.EXPECTED_PRGS_FLEET)
]
def scan_for(rows, *, available=True, extra_pids=(), started_at="2026-07-28T01:00:00Z"):
"""A process observation that corroborates *rows* (plus any extra PIDs)."""
processes = [
{"pid": r["pid"], "started_at": started_at, "command": "python mcp_server.py"}
for r in rows
]
processes.extend(
{"pid": pid, "started_at": started_at, "command": "python mcp_server.py"}
for pid in extra_pids
)
return {"available": available, "processes": processes, "reason": None}
def classify(rows, scan=None, **kwargs):
kwargs.setdefault("expected_binding", BINDING)
kwargs.setdefault("now", NOW)
return mfi.classify_fleet_inventory(
runtime_rows=rows,
process_scan=scan if scan is not None else scan_for(rows),
**kwargs,
)
class _AliveMixin:
"""Treat every PID as alive unless a test declares it dead or unknown."""
def setUp(self):
super().setUp()
self.dead_pids = set()
self.unknown_pids = set()
def probe(pid):
if pid in self.unknown_pids:
return None
return pid not in self.dead_pids
patcher = mock.patch.object(mfi, "probe_pid_alive", side_effect=probe)
patcher.start()
self.addCleanup(patcher.stop)
class TestHealthyFleet(_AliveMixin, unittest.TestCase):
"""AC1: a healthy five-server fleet reports all five members exactly once."""
def test_five_members_each_reported_once(self):
result = classify(healthy_rows())
self.assertEqual(len(result["configured_members"]), 5)
self.assertEqual(len(result["running_members"]), 5)
self.assertEqual(result["running_member_count"], 5)
self.assertEqual(result["missing_members"], [])
self.assertEqual(result["duplicate_members"], [])
self.assertEqual(result["unexpected_members"], [])
self.assertTrue(result["exactly_one_per_profile"])
for member in result["configured_members"]:
self.assertEqual(member["instance_count"], 1)
self.assertEqual(member["health"], mfi.HEALTH_RUNNING)
def test_healthy_fleet_satisfies_the_mutation_gate(self):
result = classify(healthy_rows())
self.assertTrue(result["inventory_complete"])
self.assertTrue(result["mutation_gate_satisfied"])
self.assertIsNone(result["blocked_reason"])
self.assertTrue(result["single_cohort"])
self.assertFalse(result["mixed_cohort"])
self.assertFalse(result["mixed_revision"])
def test_every_expected_namespace_and_profile_is_present(self):
result = classify(healthy_rows())
self.assertEqual(
sorted(m["profile"] for m in result["configured_members"]),
[
"prgs-author",
"prgs-controller",
"prgs-merger",
"prgs-reconciler",
"prgs-reviewer",
],
)
class TestDuplicateMember(_AliveMixin, unittest.TestCase):
"""AC2: two processes on one profile are a duplicate and fail the gate."""
def test_duplicate_is_reported_with_every_pid(self):
rows = healthy_rows()
rows.append(row("gitea-author", "prgs-author", "author", 49999))
result = classify(rows)
self.assertEqual(len(result["duplicate_members"]), 1)
duplicate = result["duplicate_members"][0]
self.assertEqual(duplicate["profile"], "prgs-author")
self.assertEqual(duplicate["instance_count"], 2)
self.assertEqual(duplicate["pids"], [_BASE_PID, 49999])
def test_duplicate_fails_the_mutation_gate(self):
rows = healthy_rows()
rows.append(row("gitea-author", "prgs-author", "author", 49999))
result = classify(rows)
self.assertFalse(result["exactly_one_per_profile"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertIn("duplicate", result["blocked_reason"])
def test_duplicate_is_not_reported_as_missing_or_unexpected(self):
rows = healthy_rows()
rows.append(row("gitea-author", "prgs-author", "author", 49999))
result = classify(rows)
self.assertEqual(result["missing_members"], [])
self.assertEqual(result["unexpected_members"], [])
class TestMissingMember(_AliveMixin, unittest.TestCase):
"""AC3: a missing expected server is identified and fails the gate."""
def test_missing_member_is_named(self):
rows = [r for r in healthy_rows() if r["profile"] != "prgs-merger"]
result = classify(rows)
self.assertEqual(len(result["missing_members"]), 1)
self.assertEqual(result["missing_members"][0]["profile"], "prgs-merger")
self.assertEqual(result["missing_members"][0]["health"], mfi.HEALTH_MISSING)
def test_missing_member_fails_the_mutation_gate(self):
rows = [r for r in healthy_rows() if r["profile"] != "prgs-merger"]
result = classify(rows)
self.assertFalse(result["exactly_one_per_profile"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertIn("prgs-merger", result["blocked_reason"])
def test_missing_member_still_lists_all_configured_members(self):
rows = [r for r in healthy_rows() if r["profile"] != "prgs-merger"]
result = classify(rows)
self.assertEqual(len(result["configured_members"]), 5)
merger = [
m for m in result["configured_members"] if m["profile"] == "prgs-merger"
][0]
self.assertEqual(merger["instance_count"], 0)
self.assertEqual(merger["health"], mfi.HEALTH_MISSING)
class TestUnexpectedMember(_AliveMixin, unittest.TestCase):
"""AC4: an unexpected PRGS server is reported explicitly."""
def test_unexpected_member_is_its_own_category(self):
rows = healthy_rows()
rows.append(row("gitea-shadow", "prgs-shadow", "author", 47777))
result = classify(rows)
self.assertEqual(len(result["unexpected_members"]), 1)
self.assertEqual(result["unexpected_members"][0]["profile"], "prgs-shadow")
self.assertEqual(
result["unexpected_members"][0]["health"], mfi.HEALTH_UNEXPECTED
)
def test_unexpected_member_is_not_collapsed_into_duplicates_or_missing(self):
rows = healthy_rows()
rows.append(row("gitea-shadow", "prgs-shadow", "author", 47777))
result = classify(rows)
self.assertEqual(result["duplicate_members"], [])
self.assertEqual(result["missing_members"], [])
self.assertTrue(result["exactly_one_per_profile"])
self.assertFalse(result["no_unexpected_members"])
self.assertFalse(result["mutation_gate_satisfied"])
def test_unexpected_member_blocks_with_its_own_reason(self):
rows = healthy_rows()
rows.append(row("gitea-shadow", "prgs-shadow", "author", 47777))
result = classify(rows)
self.assertTrue(
any("unexpected" in reason for reason in result["blocked_reasons"])
)
class TestMixedCohort(_AliveMixin, unittest.TestCase):
"""AC5: members from different client cohorts are detected."""
def test_two_cohorts_are_detected(self):
rows = healthy_rows()
rows[0]["cohort_id"] = COHORT_B
result = classify(rows)
self.assertTrue(result["mixed_cohort"])
self.assertFalse(result["single_cohort"])
self.assertEqual(result["cohort_ids"], sorted([COHORT_A, COHORT_B]))
def test_mixed_cohort_fails_the_mutation_gate(self):
rows = healthy_rows()
rows[0]["cohort_id"] = COHORT_B
result = classify(rows)
self.assertFalse(result["mutation_gate_satisfied"])
self.assertTrue(any("cohort" in reason for reason in result["blocked_reasons"]))
def test_mixed_cohort_is_not_a_duplicate_or_missing_report(self):
rows = healthy_rows()
rows[0]["cohort_id"] = COHORT_B
result = classify(rows)
self.assertEqual(result["duplicate_members"], [])
self.assertEqual(result["missing_members"], [])
class TestMixedRevision(_AliveMixin, unittest.TestCase):
"""AC6: members running different startup revisions are detected."""
def test_two_revisions_are_detected(self):
rows = healthy_rows()
rows[0]["startup_head"] = HEAD_B
result = classify(rows)
self.assertTrue(result["mixed_revision"])
self.assertEqual(result["startup_revisions"], sorted([HEAD_A, HEAD_B]))
def test_mixed_revision_fails_the_mutation_gate(self):
rows = healthy_rows()
rows[0]["startup_head"] = HEAD_B
result = classify(rows)
self.assertFalse(result["mutation_gate_satisfied"])
self.assertTrue(
any("revision" in reason for reason in result["blocked_reasons"])
)
def test_mixed_revision_does_not_by_itself_imply_mixed_cohort(self):
rows = healthy_rows()
rows[0]["startup_head"] = HEAD_B
result = classify(rows)
self.assertTrue(result["single_cohort"])
self.assertFalse(result["mixed_cohort"])
class TestRevisionIsNotCohort(_AliveMixin, unittest.TestCase):
"""AC7: matching Git revisions alone do not establish a single cohort."""
def test_identical_revisions_with_unknown_cohort_stay_unknown(self):
rows = healthy_rows()
for r in rows:
r["cohort_id"] = None
result = classify(rows)
self.assertEqual(len({r["startup_head"] for r in rows}), 1)
self.assertIsNone(result["single_cohort"])
self.assertIsNone(result["mixed_cohort"])
def test_identical_revisions_with_unknown_cohort_fail_closed(self):
rows = healthy_rows()
for r in rows:
r["cohort_id"] = None
result = classify(rows)
self.assertFalse(result["inventory_complete"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertTrue(
any("cohort" in reason for reason in result["incomplete_reasons"])
)
def test_one_unknown_cohort_among_known_ones_still_fails_closed(self):
rows = healthy_rows()
rows[0]["cohort_id"] = None
result = classify(rows)
self.assertIsNone(result["single_cohort"])
self.assertFalse(result["mutation_gate_satisfied"])
def test_cohort_derivation_never_consults_revisions(self):
"""The cohort helper takes no revision input at all."""
identity = mfi.derive_cohort_identity({mfi.COHORT_ID_ENV: "cohort-x"})
self.assertEqual(identity["cohort_id"], "cohort-x")
self.assertEqual(identity["cohort_source"], "explicit_env")
def test_orphaned_process_reports_unknown_cohort(self):
with mock.patch("os.getppid", return_value=1):
identity = mfi.derive_cohort_identity({})
self.assertIsNone(identity["cohort_id"])
self.assertEqual(identity["cohort_source"], "unknown")
def test_parent_process_is_the_cohort_when_no_env_is_set(self):
with mock.patch("os.getppid", return_value=40990):
identity = mfi.derive_cohort_identity({})
self.assertEqual(identity["cohort_id"], COHORT_A)
self.assertEqual(identity["cohort_source"], "parent_process")
class TestConfigurationIsNotRunning(_AliveMixin, unittest.TestCase):
"""AC8: configuration without a live worker is not a running member."""
def test_no_registry_rows_means_every_member_is_missing(self):
result = classify([], scan=scan_for([]))
self.assertEqual(len(result["missing_members"]), 5)
self.assertEqual(result["running_members"], [])
self.assertFalse(result["mutation_gate_satisfied"])
def test_dead_pid_is_stale_not_running(self):
rows = healthy_rows()
self.dead_pids = {rows[0]["pid"]}
result = classify(rows)
self.assertEqual(len(result["running_members"]), 4)
self.assertEqual(len(result["stale_members"]), 1)
self.assertEqual(result["stale_members"][0]["liveness"], mfi.LIVENESS_DEAD)
self.assertEqual(result["stale_members"][0]["health"], mfi.HEALTH_STALE)
self.assertEqual(len(result["missing_members"]), 1)
def test_registry_row_without_a_matching_process_is_not_running(self):
rows = healthy_rows()
scan = scan_for(rows[1:]) # first member's process is absent
result = classify(rows, scan=scan)
self.assertEqual(len(result["running_members"]), 4)
self.assertEqual(
result["stale_members"][0]["liveness"], mfi.LIVENESS_UNOBSERVED
)
self.assertFalse(result["mutation_gate_satisfied"])
def test_recycled_pid_does_not_impersonate_a_dead_server(self):
rows = healthy_rows()
scan = scan_for(rows)
# The process now holding the first PID started *after* registration.
scan["processes"][0]["started_at"] = "2026-07-28T02:30:00Z"
result = classify(rows, scan=scan)
stale = [
m
for m in result["stale_members"]
if m["liveness"] == mfi.LIVENESS_PID_RECYCLED
]
self.assertEqual(len(stale), 1)
self.assertEqual(len(result["running_members"]), 4)
self.assertFalse(result["mutation_gate_satisfied"])
class TestIncompleteEvidenceFailsClosed(_AliveMixin, unittest.TestCase):
"""AC9: unknown or unavailable evidence produces a fail-closed result."""
def test_unavailable_process_listing_fails_closed(self):
rows = healthy_rows()
result = classify(
rows,
scan={"available": False, "processes": [], "reason": "ps unavailable"},
)
self.assertFalse(result["inventory_complete"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertIn("ps unavailable", result["incomplete_reasons"])
self.assertEqual(result["running_members"], [])
def test_unreadable_registry_fails_closed(self):
result = classify(
[],
scan=scan_for([]),
registry_available=False,
registry_error="registry unreadable",
)
self.assertFalse(result["inventory_complete"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertIn("registry unreadable", result["incomplete_reasons"])
def test_unregistered_running_process_fails_closed(self):
rows = healthy_rows()
result = classify(rows, scan=scan_for(rows, extra_pids=[59999]))
self.assertEqual([p["pid"] for p in result["unregistered_processes"]], [59999])
self.assertFalse(result["inventory_complete"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertTrue(
any("59999" in reason for reason in result["incomplete_reasons"])
)
def test_undeterminable_pid_liveness_is_unknown_not_healthy(self):
rows = healthy_rows()
self.unknown_pids = {rows[0]["pid"]}
result = classify(rows)
self.assertEqual(result["stale_members"][0]["liveness"], mfi.LIVENESS_UNKNOWN)
self.assertEqual(result["stale_members"][0]["health"], mfi.HEALTH_UNKNOWN)
self.assertFalse(result["inventory_complete"])
self.assertFalse(result["mutation_gate_satisfied"])
def test_unknown_startup_revision_fails_closed(self):
rows = healthy_rows()
rows[0]["startup_head"] = None
result = classify(rows)
self.assertFalse(result["inventory_complete"])
self.assertFalse(result["mutation_gate_satisfied"])
self.assertTrue(
any("startup revision" in r for r in result["incomplete_reasons"])
)
def test_unknown_is_distinguished_from_healthy(self):
rows = healthy_rows()
for r in rows:
r["cohort_id"] = None
result = classify(rows)
# Not "unhealthy" in the sense of a named defect: nothing is missing,
# duplicated or unexpected. It is *unknown*, and that still fails closed.
self.assertEqual(result["missing_members"], [])
self.assertEqual(result["duplicate_members"], [])
self.assertEqual(result["unexpected_members"], [])
self.assertIsNone(result["single_cohort"])
self.assertFalse(result["mutation_gate_satisfied"])
class TestBindingAndRoleConsistency(_AliveMixin, unittest.TestCase):
"""Repository-binding and role/profile mismatches fail closed."""
def test_repository_binding_mismatch_is_reported(self):
rows = healthy_rows()
rows[0]["repo"] = "Some-Other-Repo"
result = classify(rows)
self.assertEqual(len(result["repository_binding_mismatches"]), 1)
self.assertEqual(
result["repository_binding_mismatches"][0]["repo"], "Some-Other-Repo"
)
self.assertFalse(result["mutation_gate_satisfied"])
def test_incomplete_repository_binding_is_reported(self):
rows = healthy_rows()
rows[0]["org"] = None
result = classify(rows)
self.assertEqual(len(result["repository_binding_mismatches"]), 1)
self.assertIn("complete", result["repository_binding_mismatches"][0]["reason"])
self.assertFalse(result["mutation_gate_satisfied"])
def test_role_mismatch_is_reported(self):
rows = healthy_rows()
rows[0]["role"] = "merger" # prgs-author is configured as author
result = classify(rows)
self.assertEqual(len(result["role_mismatches"]), 1)
self.assertEqual(result["role_mismatches"][0]["expected_role"], "author")
self.assertEqual(result["role_mismatches"][0]["declared_role"], "merger")
self.assertFalse(result["mutation_gate_satisfied"])
def test_profile_served_from_the_wrong_namespace_is_reported(self):
rows = healthy_rows()
rows[0]["namespace"] = "gitea-reviewer"
result = classify(rows)
self.assertTrue(
any(
m.get("expected_namespace") == "gitea-author"
for m in result["role_mismatches"]
)
)
self.assertFalse(result["mutation_gate_satisfied"])
def test_matching_binding_produces_no_mismatch(self):
result = classify(healthy_rows())
self.assertEqual(result["repository_binding_mismatches"], [])
self.assertEqual(result["role_mismatches"], [])
class TestDeterminism(_AliveMixin, unittest.TestCase):
"""AC10 / stable ordering and deterministic structured output."""
def test_identical_snapshots_produce_identical_results(self):
rows = healthy_rows()
first = classify(rows, scan=scan_for(rows))
second = classify(healthy_rows(), scan=scan_for(healthy_rows()))
self.assertEqual(first, second)
def test_row_order_does_not_change_the_result(self):
rows = healthy_rows()
shuffled = list(reversed(healthy_rows()))
forward = classify(rows, scan=scan_for(rows))
backward = classify(shuffled, scan=scan_for(shuffled))
self.assertEqual(forward["running_members"], backward["running_members"])
self.assertEqual(forward["configured_members"], backward["configured_members"])
self.assertEqual(
forward["mutation_gate_satisfied"], backward["mutation_gate_satisfied"]
)
def test_members_are_sorted_by_namespace_then_profile_then_pid(self):
rows = healthy_rows()
rows.append(row("gitea-author", "prgs-author", "author", 40001))
result = classify(rows)
keys = [
(m["namespace"], m["profile"], m["pid"]) for m in result["running_members"]
]
self.assertEqual(keys, sorted(keys))
def test_answering_namespace_does_not_change_the_verdict(self):
"""AC10: controller and reconciler agree for one fleet snapshot."""
rows = healthy_rows()
controller = classify(
rows, scan=scan_for(rows), answering_namespace="gitea-controller"
)
reconciler = classify(
healthy_rows(),
scan=scan_for(healthy_rows()),
answering_namespace="gitea-reconciler",
)
self.assertEqual(controller["answering_namespace"], "gitea-controller")
self.assertEqual(reconciler["answering_namespace"], "gitea-reconciler")
for key in sorted(set(controller) - {"answering_namespace"}):
self.assertEqual(controller[key], reconciler[key], f"{key} disagreed")
def test_controller_and_reconciler_agree_on_an_unhealthy_fleet(self):
rows = [r for r in healthy_rows() if r["profile"] != "prgs-merger"]
controller = classify(
rows, scan=scan_for(rows), answering_namespace="gitea-controller"
)
reconciler = classify(
rows, scan=scan_for(rows), answering_namespace="gitea-reconciler"
)
self.assertEqual(controller["blocked_reason"], reconciler["blocked_reason"])
self.assertEqual(
controller["mutation_gate_satisfied"],
reconciler["mutation_gate_satisfied"],
)
class TestReadOnly(_AliveMixin, unittest.TestCase):
"""AC11: the capability mutates nothing."""
def test_classification_reports_no_mutations(self):
result = classify(healthy_rows())
self.assertEqual(result["mutations_performed"], [])
self.assertTrue(result["read_only"])
def test_classification_does_not_write_the_input_rows_back(self):
rows = healthy_rows()
snapshot = [dict(r) for r in rows]
classify(rows)
self.assertEqual(rows, snapshot)
class TestLivenessProbe(unittest.TestCase):
"""AC11: the real probe never sends a terminating signal.
Deliberately *not* using ``_AliveMixin`` these tests exercise
``probe_pid_alive`` itself, which the mixin replaces.
"""
def test_only_signal_zero_is_ever_sent(self):
with mock.patch("os.kill") as killer:
mfi.probe_pid_alive(4242)
killer.assert_called_once_with(4242, 0)
def test_classifying_a_duplicate_never_terminates_it(self):
rows = healthy_rows()
rows.append(row("gitea-author", "prgs-author", "author", 49999))
with mock.patch("os.kill") as killer:
result = mfi.classify_fleet_inventory(
runtime_rows=rows,
process_scan=scan_for(rows),
expected_binding=BINDING,
now=NOW,
)
self.assertTrue(killer.call_args_list, "liveness must actually be probed")
for call in killer.call_args_list:
self.assertEqual(call.args[1], 0, "only signal 0 may ever be sent")
self.assertEqual(result["mutations_performed"], [])
def test_liveness_probe_tolerates_a_missing_process(self):
with mock.patch("os.kill", side_effect=ProcessLookupError):
self.assertFalse(mfi.probe_pid_alive(4242))
def test_liveness_probe_treats_permission_error_as_alive(self):
with mock.patch("os.kill", side_effect=PermissionError):
self.assertTrue(mfi.probe_pid_alive(4242))
def test_liveness_probe_returns_unknown_on_other_os_errors(self):
with mock.patch("os.kill", side_effect=OSError):
self.assertIsNone(mfi.probe_pid_alive(4242))
def test_invalid_pid_is_unknown_and_probes_nothing(self):
with mock.patch("os.kill") as killer:
self.assertIsNone(mfi.probe_pid_alive(None))
self.assertIsNone(mfi.probe_pid_alive(0))
killer.assert_not_called()
class TestMultiClientRegression(_AliveMixin, unittest.TestCase):
"""The multi-LLM duplicate-server scenario that motivated #949."""
@staticmethod
def _two_client_rows():
first = healthy_rows()
second = [
row(
entry["namespace"],
entry["profile"],
entry["role"],
50000 + index,
cohort_id=COHORT_B,
)
for index, entry in enumerate(mfi.EXPECTED_PRGS_FLEET)
]
return first + second
def test_second_client_running_the_same_five_profiles_is_caught(self):
"""Two clients, ten servers, same profiles, same revision.
Every member self-reports the same parity, which is exactly why the old
per-process surfaces reported success. The fleet inventory must report
five duplicates, two cohorts, and a closed gate.
"""
rows = self._two_client_rows()
result = classify(rows, scan=scan_for(rows))
self.assertEqual(len(result["duplicate_members"]), 5)
self.assertFalse(result["exactly_one_per_profile"])
self.assertTrue(result["mixed_cohort"])
self.assertFalse(result["single_cohort"])
self.assertFalse(result["mixed_revision"], "both clients share a revision")
self.assertFalse(result["mutation_gate_satisfied"])
self.assertEqual(result["missing_members"], [])
def test_identical_parity_across_ten_servers_is_not_health(self):
rows = self._two_client_rows()
result = classify(rows, scan=scan_for(rows))
self.assertEqual(result["startup_revisions"], [HEAD_A])
self.assertFalse(result["mutation_gate_satisfied"])
def test_every_duplicate_pid_is_named_for_the_operator(self):
rows = self._two_client_rows()
result = classify(rows, scan=scan_for(rows))
for duplicate in result["duplicate_members"]:
self.assertEqual(len(duplicate["pids"]), 2)
class TestProcessScan(unittest.TestCase):
"""The process observation is server-side and fails closed."""
def test_scan_parses_mcp_server_processes(self):
stdout = (
" PID STARTED COMMAND\n"
" 41000 Mon Jul 27 20:00:00 2026 python /path/mcp_server.py\n"
" 41001 Mon Jul 27 20:00:01 2026 python /path/other_server.py\n"
)
result = mfi.scan_mcp_server_processes(
runner=lambda *a, **k: mock.Mock(stdout=stdout)
)
self.assertTrue(result["available"])
self.assertEqual([p["pid"] for p in result["processes"]], [41000])
def test_scan_failure_reports_unavailable_rather_than_empty(self):
def boom(*args, **kwargs):
raise OSError("ps missing")
result = mfi.scan_mcp_server_processes(runner=boom)
self.assertFalse(result["available"])
self.assertEqual(result["processes"], [])
self.assertIn("ps missing", result["reason"])
def test_scan_results_are_sorted_by_pid(self):
stdout = (
" PID STARTED COMMAND\n"
" 41005 Mon Jul 27 20:00:00 2026 python /path/mcp_server.py\n"
" 41001 Mon Jul 27 20:00:01 2026 python /path/mcp_server.py\n"
)
result = mfi.scan_mcp_server_processes(
runner=lambda *a, **k: mock.Mock(stdout=stdout)
)
self.assertEqual([p["pid"] for p in result["processes"]], [41001, 41005])
class TestRuntimeRecord(unittest.TestCase):
"""The row a server writes about itself."""
def test_record_describes_the_calling_process(self):
record = mfi.build_process_runtime_record(
namespace="gitea-controller",
profile="prgs-controller",
role="controller",
remote="prgs",
org=BINDING["org"],
repo=BINDING["repo"],
pid=41022,
startup_head=HEAD_A,
transport="stdio",
client_provenance="client_managed",
env={mfi.COHORT_ID_ENV: COHORT_A},
boot_id="deadbeefcafe0001",
)
self.assertEqual(
record["runtime_id"], "gitea-controller:41022:deadbeefcafe0001"
)
self.assertEqual(record["namespace"], "gitea-controller")
self.assertEqual(record["cohort_id"], COHORT_A)
self.assertEqual(record["cohort_source"], "explicit_env")
self.assertEqual(record["startup_head"], HEAD_A)
self.assertEqual(record["status"], "running")
def test_record_defaults_pid_to_the_current_process(self):
record = mfi.build_process_runtime_record(
namespace="gitea-author", profile="prgs-author", role="author"
)
self.assertEqual(record["pid"], os.getpid())
def test_each_boot_gets_a_distinct_runtime_id(self):
first = mfi.build_process_runtime_record(
namespace="gitea-author", profile="prgs-author", role="author", pid=1
)
second = mfi.build_process_runtime_record(
namespace="gitea-author", profile="prgs-author", role="author", pid=1
)
self.assertNotEqual(first["runtime_id"], second["runtime_id"])
def test_registered_at_uses_the_control_plane_timestamp_format(self):
record = mfi.build_process_runtime_record(
namespace="gitea-author", profile="prgs-author", role="author", pid=1
)
datetime.strptime(record["registered_at"], "%Y-%m-%dT%H:%M:%SZ")
class TestSummary(_AliveMixin, unittest.TestCase):
def test_healthy_summary_names_the_counts(self):
result = classify(healthy_rows())
self.assertIn("5 of 5", mfi.summarize(result))
def test_blocked_summary_repeats_the_blocked_reason(self):
rows = [r for r in healthy_rows() if r["profile"] != "prgs-merger"]
result = classify(rows)
self.assertIn(result["blocked_reason"], mfi.summarize(result))
class TestHeartbeatAge(_AliveMixin, unittest.TestCase):
def test_heartbeat_age_is_reported_for_diagnosis(self):
rows = healthy_rows()
stamp = (NOW - timedelta(minutes=30)).strftime("%Y-%m-%dT%H:%M:%SZ")
rows[0]["last_heartbeat_at"] = stamp
result = classify(rows)
member = [m for m in result["running_members"] if m["pid"] == rows[0]["pid"]][0]
self.assertEqual(member["heartbeat_age_seconds"], 1800)
def test_missing_heartbeat_is_reported_as_unknown_age(self):
rows = healthy_rows()
rows[0]["last_heartbeat_at"] = None
result = classify(rows)
member = [m for m in result["running_members"] if m["pid"] == rows[0]["pid"]][0]
self.assertIsNone(member["heartbeat_age_seconds"])
if __name__ == "__main__":
unittest.main()