Compare commits

..
Author SHA1 Message Date
sysadminandClaude Opus 4.8 e8bae606cb fix(reconcile): retire workflow session rows whose owners are no longer live
Implement a sanctioned session lifecycle for post-restart reconciliation so
active session rows with dead or reused owner PIDs can be terminalized
without deleting history. Protect live owners, live leases, and live
client-managed sessions; record durable session_retired audit events; keep
cleanup idempotent under concurrent reconciles.

Closes #969

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-29 05:27:42 -04:00
sysadmin 956fa15fe3 Merge pull request 'feat(mcp): client/session-aware runtime ownership and provenance (Closes #948)' (#968) from feat/issue-948-client-session-provenance into master 2026-07-29 03:29:47 -05:00
sysadminandClaude Opus 5 fa510dd28d docs(remote-mcp): restamp the #956 threat-model anchors onto the commit they resolve at
The anchors and their citations moved in the previous commit because
gitea_mcp_server.py gained the #948 worker-identity block. The fixture still
named the commit the old line numbers resolved at, so the recorded provenance
pointed at a tree where the new numbers do not hold.

No anchor target or expectation changes; only the recorded commit does.

Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_01F6Vomtndpq2gSBa88Tfcwy
2026-07-29 02:52:46 -04:00
sysadminandClaude Opus 5 1dd30ecb15 feat(mcp): client/session-aware runtime ownership and provenance (#948)
Two surfaces reported different provenance for one process.
`gitea_get_runtime_context` read the live environment and reported
`client_managed`; `mcp_namespace_health.classify_namespace_probe` derived
provenance from `_safe_env_summary()`, whose `SAFE_ENV_KEYS` allowlist never
contained `GITEA_CLIENT_MANAGED`, `GITEA_MCP_CLIENT_MANAGED`, or
`GITEA_SERVER_PROVENANCE`. That lookup could only ever miss, so the health
surface was structurally incapable of returning anything but `manual_launch`.

Neither model could name which client or which session owned a runtime, so a
healthy daemon serving a second client was indistinguishable from a duplicate,
and the profile-wide duplicate scan walled the whole fleet.

Introduce `mcp_worker_identity` as the one authority, splitting two claims the
old code ran together:

* launch provenance — was this hand-launched from a terminal? Answered from the
  environment, which is legitimate because the launcher sets it. Preserves the
  #686 wall unchanged.
* session ownership — which live client session owns this runtime now? Answered
  only from a live attachment record; no environment flag can establish it.

The module provides collision-resistant worker identities
(`<llm-name>-<UTC-timestamp>-<short-sha>`), an atomic SQLite registry with
fencing epochs, heartbeat-based liveness, generation takeover that supersedes
only a non-live claimant, cohort classification, and failure scoping.

Behaviour changes:

* Registering an existing worker identity fails closed; it is never replaced,
  adopted, or merged with. The caller mints a different identity instead.
* A generation held by a live session cannot be claimed by a second one. A
  generation whose claimant is not live is taken over with a higher fencing
  epoch, so stale ownership cannot permanently strand a healthy daemon.
* A superseded session presenting an old epoch is refused and performs no write.
* Liveness comes from heartbeat freshness; a live PID cannot resurrect an
  expired record, and a dead PID withdraws liveness.
* Workers sharing a role or profile no longer trigger a profile-wide duplicate
  block, provided each carries a distinct identity. Processes with no identity
  evidence remain classified as duplicates, so the #686 wall still holds.
* Runtime failures are scoped to a worker identity or generation, never to a
  profile or the fleet.
* Reconnect guidance no longer defaults to Codex. An unidentified client gets
  host-agnostic steps; `gitea_request_mcp_reconnect(client=...)` defaults to
  resolving the client from the live attachment record.
* `resolve_bound_remote` keeps a bound namespace on its remote instead of
  falling through to the `dadeschools` library default.

Absence of proof is now reported as `unproven` rather than asserted as
`manual_launch`. Both still fail closed — `is_client_managed` is unchanged, so
nothing previously refused is now permitted — but remediation names the proof
that is actually missing instead of describing a terminal launch it cannot
evidence. The #686 test is updated for that vocabulary and keeps every
wall-preserving assertion.

Threat-model anchors and their citations in docs/remote-mcp/threat-model.md are
restamped for the line movement in gitea_mcp_server.py.

Tests: tests/test_issue_948_client_session_provenance.py adds 43 cases covering
Codex/Gemini/Antigravity/Claude attachment, same-client new session, cross-client
takeover after a session ends, two live conflicting sessions, stale records,
missing attachment proof, environment flags without attachment, mixed
generations, duplicate cohorts, the hardcoded-client regression, explicit PRGS
selection, default-remote host drift, cross-surface agreement, and fail-closed
handling without false reconnect loops. Synthetic identifiers throughout.

Full suite from a branches/ worktree: 28F/5953P/6S at head vs 28F/5910P/6S at
merge base 8eada1fb, identical failing ID sets.

Co-Authored-By: Claude Opus 5 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_01F6Vomtndpq2gSBa88Tfcwy
2026-07-29 02:52:18 -04:00
sysadmin 8eada1fbe4 Merge pull request 'feat(mcp): gate Connected-but-unattached MCP namespaces (Closes #708)' (#967) from feat/issue-708-mcp-namespace-attachment into master 2026-07-28 18:31:13 -05:00
jcwalker3andClaude Opus 4.8 58bd852188 docs(remote-mcp): restamp the #956 threat model onto the commit it resolves at (review 637 B1)
Commit 126d76a set docs/remote-mcp/threat-model-anchors.json
generated_against_commit to e3fa3b26 and left docs/remote-mcp/threat-model.md
citing a143cd06, so the document and the fixture named different source
revisions. tests/test_issue_956_threat_model.py asserts the fixture value
appears in the document, and that assertion failed from 126d76a onward.

The mismatch was not cosmetic. Resolving all 58 anchors against each candidate
revision shows 0 unresolved at e3fa3b26 and 21 unresolved at a143cd06, so the
document's own claim about where its file:line citations resolve was false.
The fixture held the accurate value and the document held the stale one.

The B2 fix in ca5f078d inserted a net 16 lines into gitea_mcp_server.py at a
single point, shifting 21 anchors below it. All 21 are re-derived here by that
one uniform offset and each was verified against its recorded expect substring
rather than assumed, so no anchor is guessed and none became ambiguous. Both
artifacts now name ca5f078d, the commit those anchors were taken at, and the
document's inline citations are shifted to match. The boundary table's claim
about what the code enforces is restamped onto the same revision, because it
described that same source tree.

Verified: 58/58 anchors resolve at the declared generation commit, 58/58 at
this head, no two anchors share a file:line location, fixture and document
metadata agree, and tests/test_issue_956_threat_model.py passes 17/17.

Refs #708

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-28 17:25:24 -05:00
jcwalker3andClaude Opus 4.8 ca5f078d8a fix(mcp): separate not-connected from Connected-but-unattached (review 637 B2)
A required namespace absent from the connected-service inventory was never
counted as missing, because missing_namespaces was populated only while
walking connected_servers. The session therefore reported attachment_healthy
true, discovery_status namespaces_attached and error_type None while holding
no attachment proof for that namespace, and the gate text told the operator
"the host reports Connected" about a service nothing had reported Connected.

assess_connected_namespace_attachment now classifies the two conditions
apart. missing_namespaces keeps its meaning: required, Connected at the host,
absent from the session tool surface -- the actual #708 defect, still typed
mcp_connected_namespaces_missing. Required namespaces absent from the
connected inventory land in the new not_connected_namespaces list, typed
mcp_required_namespaces_not_connected, so a caller is never pointed at an
attachment recovery for a service that never connected. attachment_healthy
is false whenever any required namespace lacks attachment proof under either
condition, error_types reports every condition present so neither hides the
other, and namespace_conditions carries the per-namespace connected/attached
verdict. Evidence that disagrees with itself -- attached in the session yet
absent from the connected inventory -- is reported as contradictory and fails
closed rather than being read as proof.

attachment_gate_from_session states only what the recorded evidence supports:
the Connected wording appears solely when connected is true, the not-connected
wording when it is false, and a neutral refusal when connected status was
never recorded. Review and merge stay fail-closed in all three cases.

_record_live_namespace_attachment consumes the per-namespace verdicts instead
of pinning the session-wide error_type onto every unattached namespace and
instead of deriving health from a missing list that tracked only one of the
two conditions.

The two existing cases in test_issue_708_mcp_namespace_attachment.py declared
no required_namespaces, so they inherited a default set containing a namespace
their connected list omitted. Under the corrected rule that is a false-healthy
assertion, so each now declares the required set it actually means.

Closes #708

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-28 17:08:44 -05:00
sysadminandClaude Opus 4.8 126d76ad28 docs(remote-mcp): restamp the commit the #956 anchors resolve at
The anchors were re-derived against the #708 wiring commit; record that SHA so
the fixture states the tree its line numbers were taken from.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-28 16:46:37 -04:00
sysadminandClaude Opus 4.8 e3fa3b263d feat(mcp): gate Connected-but-unattached MCP namespaces (Closes #708)
The prior #708 slice added assess_connected_namespace_attachment() as a pure
decision function with no call site: nothing invoked it, so a session whose
role namespaces were Connected at the host but absent from the active session
tool surface still passed every mutation gate. Detection existed on paper only.

This makes it load-bearing.

Decision layer (mcp_namespace_health.py)
- assess_connected_namespace_attachment() gains secret-free telemetry
  (connected/attached/required/missing counts, discovery cache hit and age,
  reconnect_required, auto_attach_attempted, auto_recovered, error_type) and
  reports reconnect_required, auto_recovered and startup_ordering_race.
- Startup ordering: a namespace whose connect completed after the session tool
  snapshot cannot be in that snapshot, so parallel multi-role startup is
  identified as its own race with the affected namespaces listed.
- attachment_gate_from_session() is a fail-closed gate keyed by
  ATTACHMENT_GATED_TASKS. An unassessed namespace does not gate, matching #543
  semantics, so this never fabricates a block.
- SANCTIONED_ATTACH_RECOVERY_TOOL names gitea_request_mcp_reconnect (#678) as
  the only recovery.

Server wiring (gitea_mcp_server.py)
- New tool gitea_assess_mcp_namespace_attachment classifies the condition and
  records a per-namespace verdict in _LIVE_NAMESPACE_ATTACHMENT.
- gitea_submit_pr_review and gitea_merge_pr now consult
  _namespace_attachment_gate() alongside the existing #543 health gate, so both
  fail closed while a required namespace is unattached.
- Watchdog check-in emits status only, never namespace contents.

The typed condition mcp_connected_namespaces_missing stays distinct from config
drift (#672), transport-closed (#584) and resolver EOF (#685). Recovery never
routes through direct imports, CLI or raw API mutation, profile hopping,
session-state overrides, or process kills.

Docs
- docs/mcp-namespace-health.md documents the tool arguments, startup ordering,
  the fail-closed gate, and the telemetry contract.
- skills/llm-project-workflow/SKILL.md states Connected is not attached, and
  that preflight proof is live tool visibility plus gitea_whoami on the role
  namespace rather than host status alone.
- docs/mcp-tool-inventory.md lists the new tool.
- docs/remote-mcp/threat-model-anchors.json and threat-model.md: 21 #956 anchors
  restamped for the line shift these additions caused in gitea_mcp_server.py.
  Every anchor was re-derived from its recorded expect substring; none guessed.

Tests
- tests/test_issue_708_attachment_wiring.py (19 cases): typed detection, proof
  mapping, reconnect-only next action, auto-attach success and failure,
  reconnect rediscovery, multi-role startup ordering, telemetry including a
  no-secret-leak assertion, fail-closed gate per role, unassessed and unmapped
  tasks not gating, partial attachment gating only the affected role, and no
  healthy verdict without attachment proof.

Verification
- tests/test_issue_708_attachment_wiring.py + test_issue_708_mcp_namespace_attachment.py: 24 passed
- namespace/session/runtime/review sweep: 427 passed, 12 subtests
- full suite head: 31 failed, 5885 passed, 6 skipped, 1047 subtests
- full suite base 17ba1ff035: 30 failed, 5862 passed, 6 skipped, 1047 subtests
- failing identifier sets match, plus tests/test_mirror_refs.py DryRunBanner,
  which fails on the unmodified base in isolation and passes here: flaky, not a
  regression from this branch.
- Gate proven by execution, not inspection: registering the assessment blocks
  merge_pr and review_pr, and attaching the namespaces clears the block.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-28 16:46:09 -04:00
sysadmin 4533e86bcd Merge branch 'master' into feat/issue-708-mcp-namespace-attachment 2026-07-28 16:20:13 -04:00
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
sysadmin 08d9cf4cbd feat(mcp): detect Connected-but-namespaces-missing attachment failures (Closes #708) 2026-07-25 19:33:18 -04:00
28 changed files with 7679 additions and 100 deletions
+226 -7
View File
@@ -27,12 +27,12 @@ import uuid
from contextlib import contextmanager
from dataclasses import dataclass
from datetime import datetime, timedelta, timezone
from typing import Any, Iterator, Sequence
from typing import Any, Iterator, Mapping, Sequence
import dependency_graph
import gitea_audit
SCHEMA_VERSION = 5
SCHEMA_VERSION = 6
# Assignable work kinds only — raw monitoring incidents are never work items.
WORK_KINDS = frozenset({"issue", "pr"})
@@ -419,6 +419,7 @@ class ControlPlaneDB:
self._migrate_incident_links_null_scope(conn)
self._migrate_lease_lifecycle_columns(conn)
self._migrate_session_ownership_columns(conn)
self._migrate_session_lifecycle_columns(conn)
self._migrate_usage_events_table(conn)
conn.execute(
"INSERT OR REPLACE INTO schema_meta(key, value) VALUES (?, ?)",
@@ -812,6 +813,7 @@ class ControlPlaneDB:
pid: int | None = None,
status: str = "active",
controller_instance_id: str | None = None,
owner_process_started_at: str | None = None,
) -> dict[str, Any]:
"""Register/refresh a session row.
@@ -821,16 +823,22 @@ class ControlPlaneDB:
controller instance can. It is never overwritten with ``None``, so a
heartbeat from a caller that does not supply one cannot erase
ownership.
*owner_process_started_at* (#969) is the OS start time of the owner
process at registration time. When present it is retained across
heartbeats (never cleared by ``None``) so later PID-reuse checks do
not depend on a live ``ps`` probe of a long-dead process.
"""
now = _ts()
instance = (controller_instance_id or "").strip() or None
proc_start = (owner_process_started_at or "").strip() or None
with self._tx() as conn:
existing = conn.execute(
"SELECT session_id FROM sessions WHERE session_id = ?",
(session_id,),
).fetchone()
if existing:
if instance is None:
if instance is None and proc_start is None:
conn.execute(
"""
UPDATE sessions
@@ -840,7 +848,21 @@ class ControlPlaneDB:
""",
(role, profile, namespace, pid, now, status, session_id),
)
else:
elif instance is None:
conn.execute(
"""
UPDATE sessions
SET role = ?, profile = ?, namespace = ?, pid = ?,
last_heartbeat_at = ?, status = ?,
owner_process_started_at = COALESCE(?, owner_process_started_at)
WHERE session_id = ?
""",
(
role, profile, namespace, pid, now, status,
proc_start, session_id,
),
)
elif proc_start is None:
conn.execute(
"""
UPDATE sessions
@@ -854,18 +876,33 @@ class ControlPlaneDB:
instance, session_id,
),
)
else:
conn.execute(
"""
UPDATE sessions
SET role = ?, profile = ?, namespace = ?, pid = ?,
last_heartbeat_at = ?, status = ?,
controller_instance_id = ?,
owner_process_started_at = COALESCE(?, owner_process_started_at)
WHERE session_id = ?
""",
(
role, profile, namespace, pid, now, status,
instance, proc_start, session_id,
),
)
else:
conn.execute(
"""
INSERT INTO sessions(
session_id, role, profile, namespace, pid,
started_at, last_heartbeat_at, status,
controller_instance_id
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)
controller_instance_id, owner_process_started_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""",
(
session_id, role, profile, namespace, pid, now, now,
status, instance,
status, instance, proc_start,
),
)
row = conn.execute(
@@ -881,6 +918,165 @@ class ControlPlaneDB:
(_ts(), session_id),
)
def retire_session(
self,
*,
session_id: str,
reason: str,
actor_session_id: str | None = None,
details: Mapping[str, Any] | None = None,
now: datetime | None = None,
terminal_status: str = "retired",
) -> dict[str, Any]:
"""Terminalize one session row if it is still non-terminal (#969).
CAS on non-terminal status: concurrent retirements of the same row
yield exactly one ``retired`` outcome and subsequent
``already_terminal`` outcomes. Never deletes historical rows. Writes a
durable ``session_retired`` event for audit.
"""
moment = _ts(now)
reason_s = (reason or "").strip() or "unspecified"
term = (terminal_status or "retired").strip().lower() or "retired"
terminal_set = {
"retired",
"ended",
"terminal",
"dead",
"stale",
"orphaned",
}
with self._tx() as conn:
row = conn.execute(
"SELECT * FROM sessions WHERE session_id = ?",
(session_id,),
).fetchone()
if row is None:
return {
"outcome": "missing",
"reason": reason_s,
"prior_status": None,
"new_status": None,
"details": {"session_id": session_id},
}
prior = dict(row)
prior_status = str(prior.get("status") or "").strip().lower()
if prior_status in terminal_set:
return {
"outcome": "already_terminal",
"reason": reason_s,
"prior_status": prior_status,
"new_status": prior_status,
"details": {"session_id": session_id, "idempotent": True},
}
# Optional: refuse when an active lease still names this session.
active_lease = conn.execute(
"""
SELECT lease_id, status, expires_at FROM leases
WHERE session_id = ? AND status = 'active'
LIMIT 1
""",
(session_id,),
).fetchone()
if active_lease is not None:
return {
"outcome": "blocked",
"reason": "live_lease",
"prior_status": prior_status,
"new_status": prior_status,
"details": {
"session_id": session_id,
"lease_id": active_lease["lease_id"],
"blocker": "active_lease_row",
},
}
cols = {
r[1] for r in conn.execute("PRAGMA table_info(sessions)").fetchall()
}
if "retired_at" in cols and "retire_reason" in cols:
conn.execute(
"""
UPDATE sessions
SET status = ?, retired_at = ?, retire_reason = ?
WHERE session_id = ? AND status = ?
""",
(term, moment, reason_s, session_id, prior.get("status")),
)
else:
conn.execute(
"""
UPDATE sessions
SET status = ?
WHERE session_id = ? AND status = ?
""",
(term, session_id, prior.get("status")),
)
changed = conn.execute(
"SELECT changes()"
).fetchone()[0]
if not changed:
# Lost CAS race — re-read.
refreshed = conn.execute(
"SELECT status FROM sessions WHERE session_id = ?",
(session_id,),
).fetchone()
cur = (
str(refreshed["status"]).strip().lower()
if refreshed is not None
else None
)
return {
"outcome": "already_terminal"
if cur in terminal_set
else "blocked",
"reason": reason_s,
"prior_status": prior_status,
"new_status": cur,
"details": {
"session_id": session_id,
"cas_lost": True,
},
}
detail_payload: dict[str, Any] = {
"session_id": session_id,
"prior_status": prior_status,
"new_status": term,
"reason": reason_s,
"actor_session_id": actor_session_id,
"pid": prior.get("pid"),
"role": prior.get("role"),
"profile": prior.get("profile"),
}
if isinstance(details, Mapping):
for key, value in details.items():
if key not in detail_payload:
detail_payload[key] = value
message = (
f"session {session_id} retired reason={reason_s} "
f"prior_status={prior_status}"
)
# events.work_item_id is nullable; session retirement is not work-scoped.
conn.execute(
"""
INSERT INTO events(work_item_id, event_type, message, created_at)
VALUES (NULL, 'session_retired', ?, ?)
""",
(message[:2000], moment),
)
# Also persist a structured JSON line in the message when short enough
# by appending a compact summary (full detail stays in return value /
# audit log; events.message is human-readable).
return {
"outcome": "retired",
"reason": reason_s,
"prior_status": prior_status,
"new_status": term,
"details": detail_payload,
}
def list_sessions(
self,
*,
@@ -1639,6 +1835,29 @@ class ControlPlaneDB:
if name not in cols:
conn.execute(f"ALTER TABLE sessions ADD COLUMN {name} {decl}")
_SESSION_LIFECYCLE_COLUMNS: tuple[tuple[str, str], ...] = (
("owner_process_started_at", "TEXT"),
("retired_at", "TEXT"),
("retire_reason", "TEXT"),
)
def _migrate_session_lifecycle_columns(self, conn: sqlite3.Connection) -> None:
"""Add session retirement / PID-reuse provenance columns (#969).
Additive and idempotent. Pre-existing rows migrate with NULL; retirement
fills ``retired_at`` / ``retire_reason``, and new upserts may record
``owner_process_started_at`` for stronger identity checks.
"""
cols = {
row[1]
for row in conn.execute("PRAGMA table_info(sessions)").fetchall()
}
if not cols:
return
for name, decl in self._SESSION_LIFECYCLE_COLUMNS:
if name not in cols:
conn.execute(f"ALTER TABLE sessions ADD COLUMN {name} {decl}")
def list_active_claims(
self,
*,
+71
View File
@@ -86,3 +86,74 @@ When a namespace returns EOF, follow
When blocked, repair the IDE namespace and re-record a healthy
`client_namespace` assessment before retrying the mutation.
## Connected vs Attached Tool Surface (#708)
MCP servers can report **Connected** at the CLI / host inventory layer while the **active LLM session exposes none of their tool namespaces**.
### Core principle
* **Connected status at host layer ≠ attached tools in active session.**
* Required preflight proof is **live tool visibility + `gitea_whoami` call** through the target namespace, not host `Connected` status alone.
* When servers report Connected but namespaces are absent from attached tools, classify as `mcp_connected_namespaces_missing`.
### Forbidden unsafe fallbacks
When `mcp_connected_namespaces_missing` is detected, workflows must **fail closed** and must **never** encourage or perform:
* direct imports of MCP server Python modules
* CLI or raw Gitea API mutations as a substitute for native tools
* profile hopping to another MCP profile/namespace to bypass the empty session
* session-state overrides or hand-edited session/ledger files
* process kills (`pkill`), config mtime touches, or `.env` edits
Only sanctioned recovery: **client reconnect path**, followed by full preflight (`whoami` → capability resolve → task).
### Native detection tool
`gitea_assess_mcp_namespace_attachment` classifies the condition and records it in
the session. Pass the namespaces the host reports Connected and the namespaces
actually attached to the active session tool surface:
| Argument | Meaning |
|---|---|
| `connected_servers` | Namespaces the host/CLI reports Connected |
| `attached_session_namespaces` | Namespaces exposed in the active session tool surface |
| `required_namespaces` | Namespaces this workflow needs (defaults to the role namespaces) |
| `discovery_cache_hit` / `discovery_cache_age_seconds` | Client tool-discovery cache state |
| `auto_attach_attempted` / `auto_attach_succeeded` | Whether the runtime auto-attached |
| `session_tool_snapshot_at` / `namespace_connected_at` | Epoch seconds, to detect startup ordering races |
It returns `discovery_status`
(`namespaces_attached` | `connected_but_namespaces_missing` | `disconnected`),
`missing_namespaces`, `proof_of_connected_vs_attached` (per namespace
`{connected, attached}`), `error_type`, `reconnect_required`, `auto_recovered`,
`startup_ordering_race`, `late_attaching_namespaces`, `sanctioned_recovery_tool`,
and a reconnect-only `exact_next_action`.
### Startup ordering
`session_tool_snapshot_at` earlier than a namespace's `namespace_connected_at`
means that namespace could not have been in the session snapshot, however healthy
it looks now. That is reported as `startup_ordering_race` with the affected
namespaces listed — the multi-role parallel-connect case where Connected flips
true after the session tool list was already captured.
### Fail-closed gate
The recorded verdict gates mutations, mirroring the #543 health gate: a namespace
that has not been assessed does not gate, but one recorded Connected-but-unattached
blocks the mutation whose role namespace it is
(`review_pr` / `submit_review``gitea-reviewer`, `merge_pr``gitea-merger`,
`work_issue` / `create_pr``gitea-author`). `gitea_submit_pr_review` and
`gitea_merge_pr` return the block reason plus the hard-stop policy string, and the
only offered recovery is `gitea_request_mcp_reconnect`.
### Telemetry
The `telemetry` block carries `connected_count`, `attached_count`,
`required_count`, `missing_count`, `discovery_status`, `discovery_cache_hit`,
`discovery_cache_age_seconds`, `reconnect_required`, `auto_attach_attempted`,
`auto_recovered`, `startup_ordering_race`, and `error_type`. It contains namespace
names and counts only — never tokens, endpoints, env values, or filesystem paths.
+1
View File
@@ -56,6 +56,7 @@ that gates each call, not which tools exist.
- `gitea_assess_conflict_fix_push`
- `gitea_assess_gitea_operation_path`
- `gitea_assess_master_parity`
- `gitea_assess_mcp_namespace_attachment`
- `gitea_assess_mcp_namespace_health`
- `gitea_assess_pr_sync_status`
- `gitea_assess_review_merge_state_machine`
+5 -1
View File
@@ -19,7 +19,11 @@ The assessor classifies:
- **service_health** — process healthy / parity mutation-safe
- **clients** — connected client descriptors (optional inventory)
- **sessions** — active session rows with dead owner pids are unresolved
- **sessions** — active session rows with dead or reused owner pids are
unresolved until retired via `#969` (`session_lifecycle` /
`apply_session_cleanup=true` on `gitea_reconcile_after_restart`, or
`gitea_retire_stale_workflow_sessions`). Live owners, live leases, and live
client-managed sessions are never retired.
- **checkpoints** — soft-depends on #660; skipped with reason when schema absent
- **leases** — live control-plane leases after restart
- **capabilities** — master-parity / stale-runtime (#610)
+246
View File
@@ -0,0 +1,246 @@
{
"_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": "1dd30ecb1508b559868c2d5d94367bc055d5138e",
"anchors": [
{
"anchor": "gitea_mcp_server.py:25087",
"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:9192",
"expect": "assess_transport_for_auth_mint()"
},
{
"anchor": "gitea_mcp_server.py:9441",
"expect": "assess_transport_for_auth_mint()"
},
{
"anchor": "mcp_server.py:4",
"expect": "The transport is selected by deployment configuration"
},
{
"anchor": "gitea_mcp_server.py:15630",
"expect": "def _is_client_managed_process"
},
{
"anchor": "gitea_mcp_server.py:15644",
"expect": "def _provenance_mutation_block"
},
{
"anchor": "gitea_mcp_server.py:15652",
"expect": "unsupported_manual_launch"
},
{
"anchor": "gitea_mcp_server.py:19217",
"expect": "server_provenance"
},
{
"anchor": "gitea_mcp_server.py:21741",
"expect": "def _check_mcp_runtimes_diagnostics"
},
{
"anchor": "gitea_mcp_server.py:21761",
"expect": "\"ps\", \"-o\", \"pid,lstart,command\""
},
{
"anchor": "gitea_mcp_server.py:21805",
"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:19487",
"expect": "def gitea_list_profiles"
},
{
"anchor": "gitea_mcp_server.py:19538",
"expect": "gitea_config.resolve_token(p)"
},
{
"anchor": "gitea_mcp_server.py:19851",
"expect": "def gitea_audit_config"
},
{
"anchor": "gitea_mcp_server.py:19873",
"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:17909",
"expect": "\"jenkins-mcp\""
},
{
"anchor": "gitea_mcp_server.py:17915",
"expect": "external-mcp"
},
{
"anchor": "gitea_mcp_server.py:17936",
"expect": "\"glitchtip-mcp\""
},
{
"anchor": "gitea_mcp_server.py:17941",
"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:19331",
"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:2352",
"expect": "/tmp/gitea_issue_lock.json"
},
{
"anchor": "gitea_mcp_server.py:10957",
"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:12871",
"expect": "owner_pid_alive"
}
]
}
+409
View File
@@ -0,0 +1,409 @@
# 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:** `1dd30ecb1508b559868c2d5d94367bc055d5138e` (#708's
namespace-attachment gate). Originally generated against
`aad5c8b42361d380a8eeb07b94b90815e594c2c5` (`master`), re-anchored at
`a143cd065ba06e1a2bdc5143a19ec156e53650ef` when #931's transport bind seam shifted the
cited lines, and re-anchored again when #708 shifted them further.
- **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:25087`, 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
`1dd30ecb1508b559868c2d5d94367bc055d5138e`, 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:25087`) 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:15630`) or a refusal (`gitea_mcp_server.py:15652`); production transport before recovery-authorization mint (`irrecoverable_provenance.py:497`, consumed at `gitea_mcp_server.py:9192` and `gitea_mcp_server.py:9441`) | 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:10957`); 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:12871`). A legacy global slot still exists at `gitea_mcp_server.py:2352`, 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:19331` | 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:19487`) reports each profile's credential status by calling
`resolve_token` on it (`gitea_mcp_server.py:19538`), and `gitea_audit_config`
(`gitea_mcp_server.py:19851`) reports service credential status through
`service_summaries` (`gitea_mcp_server.py:19873`).
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:19487`) 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:19538`). 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:19851`) reports `MDCPS Jenkins: enabled, read-only, authenticated`.
That word `authenticated` is produced by `service_summaries` (`gitea_mcp_server.py:19873`,
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:17909`, `gitea_mcp_server.py:17915`, `gitea_mcp_server.py:17936`,
`gitea_mcp_server.py:17941`) 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:21761` and
`gitea_mcp_server.py:21805`, reached from `gitea_mcp_server.py:21741`.
**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:19217`),
derived from environment inspection (`gitea_mcp_server.py:15630`) 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:19538` and `gitea_mcp_server.py:19873` 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:15630`) 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:10957`), 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.
+555 -43
View File
@@ -1,7 +1,10 @@
#!/usr/bin/env python3
"""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):
python3 mcp_server.py
@@ -196,6 +199,7 @@ RECONCILER_WORKTREE_ENV = "GITEA_RECONCILER_WORKTREE"
import namespace_workspace_binding as nwb # noqa: E402
import canonical_repository_root as crr # noqa: E402 # #706 cross-repo canonical root
import mcp_namespace_health # noqa: E402
import mcp_worker_identity # noqa: E402 # #948 single client/session provenance authority
import stale_binding_recovery # noqa: E402
# Worktree env bindings inherited from the parent environment at daemon boot
@@ -6713,6 +6717,61 @@ def _live_namespace_health_gate(task: str) -> list[str]:
)
# Session-scoped MCP namespace *attachment* assessments (#708).
# Distinct from _LIVE_NAMESPACE_HEALTH: a namespace can probe healthy while never
# being attached to the active session tool surface ("Connected" != "tools available").
_LIVE_NAMESPACE_ATTACHMENT: dict[str, dict] = {}
def _record_live_namespace_attachment(assessment: dict | None) -> None:
"""Store a per-namespace attachment verdict for mutation gates (#708)."""
if not isinstance(assessment, dict):
return
proof = assessment.get("proof_of_connected_vs_attached")
if not isinstance(proof, dict):
return
missing = set(assessment.get("missing_namespaces") or [])
error_type = assessment.get("error_type")
# Per-namespace verdicts when the assessment supplies them (#708 B2): a session-wide
# error_type must never be pinned onto a namespace whose own condition differs, and a
# namespace that was never connected must not be recorded healthy just because the
# session-wide "missing" list only tracks Connected-but-unattached namespaces.
conditions = assessment.get("namespace_conditions")
conditions = conditions if isinstance(conditions, dict) else {}
for ns, state in proof.items():
ns_name = str(ns or "").strip()
if not ns_name or not isinstance(state, dict):
continue
condition = conditions.get(ns_name)
condition = condition if isinstance(condition, dict) else {}
connected = bool(condition.get("connected", state.get("connected")))
attached = bool(condition.get("attached", state.get("attached")))
if condition:
healthy = bool(condition.get("attachment_healthy"))
ns_error = condition.get("condition")
else:
healthy = attached and connected and ns_name not in missing
ns_error = None if healthy else error_type
_LIVE_NAMESPACE_ATTACHMENT[ns_name] = {
"namespace": ns_name,
"connected": connected,
"attached": attached,
"attachment_healthy": healthy,
"condition": ns_error,
"error_type": ns_error,
"discovery_status": assessment.get("discovery_status"),
"reconnect_required": bool(assessment.get("reconnect_required")),
"auto_recovered": bool(assessment.get("auto_recovered")),
}
def _namespace_attachment_gate(task: str) -> list[str]:
"""Fail closed on a recorded Connected-but-unattached namespace (#708)."""
return mcp_namespace_health.attachment_gate_from_session(
task, _LIVE_NAMESPACE_ATTACHMENT
)
def _decision_lock_binding(lock: dict | None = None) -> dict:
"""Resolve key fields for durable decision-lock storage."""
profile = get_profile()
@@ -7593,6 +7652,10 @@ def _evaluate_pr_review_submission(
if ns_gate:
reasons.extend(ns_gate)
return result
attach_gate = _namespace_attachment_gate("review_pr")
if attach_gate:
reasons.extend(attach_gate)
return result
if action not in _REVIEW_ACTIONS:
reasons.append(
@@ -11151,6 +11214,13 @@ def gitea_merge_pr(
reasons.extend(ns_gate)
return result
# Gate 0c — the merger namespace must be attached to the active session (#708).
# Connected at the host layer is not proof the tools are in this session.
attach_gate = _namespace_attachment_gate("merge_pr")
if attach_gate:
reasons.extend(attach_gate)
return result
# Gate 1 — valid merge method (no API call on a bad method).
if do not in _MERGE_METHODS:
reasons.append(
@@ -14847,7 +14917,7 @@ def _stale_runtime_reconnect_action() -> str:
return (
"blocker_kind=runtime_reconnect_required: call "
"gitea_request_mcp_reconnect(namespace=<active gitea-* namespace>, "
"reason='stale-runtime', client='codex') for a typed operator "
"reason='stale-runtime') for a typed operator "
"reconnect blocker with exact UI steps, then reconnect the IDE/client "
"MCP session so the server reloads at the current master head. Do not "
"call gitea_activate_profile, pkill, touch configs, or switch MCP role "
@@ -15409,34 +15479,166 @@ def _session_context_mutation_block(
return blocked
def _is_client_managed_process() -> bool:
"""Check whether the current MCP server process has client-managed launch provenance (#686)."""
val = (
os.environ.get("GITEA_CLIENT_MANAGED")
or os.environ.get("GITEA_MCP_CLIENT_MANAGED")
or os.environ.get("GITEA_SERVER_PROVENANCE")
or os.environ.get("GITEA_FORCE_CLIENT_MANAGED")
or ""
).strip().lower()
# --- #948 client/session runtime ownership -------------------------------
#
# One daemon launch is one *generation*. The session that owns it registers a
# unique worker identity against it, so a second healthy client attaching its
# own generation is no longer indistinguishable from a duplicate process.
# Everything here is process-local cache plus a local SQLite registry: no Gitea
# call, no network, no config write, and every failure degrades to "unproven"
# rather than raising into a tool call.
if val in ("0", "false", "no", "manual", "manual_launch"):
_WORKER_REGISTRY = None
_WORKER_IDENTITY: str | None = None
_WORKER_GENERATION: str | None = None
_WORKER_REGISTRATION_ATTEMPTED = False
#: Env a client launcher may set to name itself and its session. Absent values
#: are reported as unknown; they are never guessed at, because guessing is what
#: produced Codex reconnect steps for a Gemini operator.
CLIENT_NAME_ENV = "GITEA_MCP_CLIENT"
CLIENT_INSTANCE_ENV = "GITEA_MCP_CLIENT_INSTANCE"
CLIENT_SESSION_ENV = "GITEA_MCP_CLIENT_SESSION"
def _worker_registry():
"""Local worker registry, or ``None`` when it cannot be opened.
Under pytest the registry is only opened when a test has pinned a path, so
a test run never writes into the operator's real registry.
"""
global _WORKER_REGISTRY
if _WORKER_REGISTRY is not None:
return _WORKER_REGISTRY
if mcp_daemon_guard.is_pytest_runtime() and not (
os.environ.get(mcp_worker_identity.REGISTRY_PATH_ENV) or ""
).strip():
return None
try:
_WORKER_REGISTRY = mcp_worker_identity.WorkerRegistry()
except Exception:
return None
return _WORKER_REGISTRY
def _client_identity_hints() -> dict:
"""What the launcher told us about itself. Unset fields stay unset."""
return {
"client_name": (os.environ.get(CLIENT_NAME_ENV) or "").strip() or None,
"client_instance_id": (os.environ.get(CLIENT_INSTANCE_ENV) or "").strip()
or f"pid-{os.getpid()}",
"session_id": (os.environ.get(CLIENT_SESSION_ENV) or "").strip()
or f"proc-{os.getpid()}-{_process_boot_head_sha or 'nohead'}",
}
def _active_worker_identity() -> str | None:
"""This runtime's worker identity, registering it once per process.
A collision does not adopt the existing registration: a fresh identity is
minted and registered instead, which is the #948 AC32 requirement and also
the only safe response adopting would silently transfer another worker's
leases and fencing tokens.
"""
global _WORKER_IDENTITY, _WORKER_GENERATION, _WORKER_REGISTRATION_ATTEMPTED
if _WORKER_IDENTITY is not None:
return _WORKER_IDENTITY
if _WORKER_REGISTRATION_ATTEMPTED:
return None
_WORKER_REGISTRATION_ATTEMPTED = True
registry = _worker_registry()
if registry is None:
return None
hints = _client_identity_hints()
generation = _WORKER_GENERATION or mcp_worker_identity.new_generation_id(os.getpid())
native = mcp_daemon_guard.native_runtime_status()
try:
for _attempt in range(3):
identity = mcp_worker_identity.generate_worker_identity(
hints["client_name"], hints["session_id"]
)
outcome = registry.register(
worker_identity=identity,
client_name=hints["client_name"],
client_instance_id=hints["client_instance_id"],
session_id=hints["session_id"],
generation_id=generation,
role=_active_role_kind_safe(),
profile=(os.environ.get(gitea_config.ENV_PROFILE) or "").strip() or None,
remote=(os.environ.get("GITEA_MCP_REMOTE") or "").strip() or None,
repository_binding=PROJECT_ROOT,
pid=os.getpid(),
transport=native.get("bound_transport"),
token_fingerprint=native.get("token_fingerprint"),
pid_alive_probe=issue_lock_store.is_process_alive,
)
if outcome.get("registered"):
_WORKER_IDENTITY = identity
_WORKER_GENERATION = generation
return identity
if not outcome.get("collision"):
return None
# Collided: loop mints a different identity rather than reusing this
# one. The existing registration is left exactly as it was.
except Exception:
return None
return None
def _active_role_kind_safe() -> str | None:
"""Best-effort role for the registry record; never raises into a tool call.
The role is recorded for diagnostics only. It is deliberately not part of
the identity: role is a reusable capability definition that many live
workers may share (#948 AC37).
"""
try:
profile = get_profile() or {}
return _role_kind(
profile.get("allowed_operations") or [],
profile.get("forbidden_operations") or [],
)
except Exception:
return None
def _stdin_is_tty() -> bool:
"""Terminal evidence for launch provenance, tolerant of a detached stdin."""
try:
return bool(sys.stdin and sys.stdin.isatty())
except Exception:
return False
if val in ("1", "true", "yes", "client_managed"):
return True
# A terminal launch has an active TTY on stdin
try:
if sys.stdin and sys.stdin.isatty():
return False
except Exception:
pass
def _launch_provenance() -> dict:
"""This process's launch provenance, from the one shared authority (#948).
# Standard client launch or test runner with stdio pipe and profile env
if "GITEA_MCP_CONFIG" in os.environ or "GITEA_MCP_PROFILE" in os.environ or "GITEA_PROFILE_NAME" in os.environ:
return True
Previously each surface reimplemented this. ``gitea_get_runtime_context``
read the live environment while ``mcp_namespace_health`` read a filtered
summary that structurally could not see the provenance keys, so the two
reported different provenance for the same process. Both now call
``mcp_worker_identity``.
"""
return mcp_worker_identity.assess_launch_provenance(
dict(os.environ), stdin_is_tty=_stdin_is_tty()
)
return False
def _is_client_managed_process() -> bool:
"""Check whether the current MCP server process has client-managed launch provenance (#686).
#948: the decision logic moved to
``mcp_worker_identity.assess_launch_provenance`` unchanged same env keys,
same precedence, same TTY fallback so this keeps returning exactly what it
always returned while no longer being a second, divergent implementation.
Launch provenance answers "was this hand-launched from a terminal", which is
the #686 question. It does *not* answer which live session owns this
runtime; that is ``session_ownership`` and needs an attachment record.
"""
return bool(_launch_provenance()["client_managed"])
def _provenance_mutation_block(**extra_fields) -> dict | None:
@@ -18979,7 +19181,21 @@ def gitea_get_runtime_context(
source="gitea_get_runtime_context",
)
is_client_managed = _is_client_managed_process()
# #948: one assessment, shared with mcp_namespace_health. Reporting both
# dimensions from a single call is what makes the two surfaces agree — the
# contradiction they used to produce came from two implementations, not from
# two genuinely different observations.
provenance_assessment = mcp_worker_identity.assess_provenance(
registry=_worker_registry(),
worker_identity=_active_worker_identity(),
env=dict(os.environ),
native_transport_bound=mcp_daemon_guard.bound_transport() is not None,
profile=profile.get("profile_name"),
role=_role_kind(allowed, forbidden),
stdin_is_tty=_stdin_is_tty(),
pid_alive_probe=issue_lock_store.is_process_alive,
)
is_client_managed = provenance_assessment["is_client_managed"]
unconsumed_env = gitea_config.get_unconsumed_gitea_env_overrides()
result = {
@@ -18998,8 +19214,21 @@ def gitea_get_runtime_context(
"review_merge_blocked_reasons": blocked_reasons,
"suggested_fix": suggested_fix,
"safe_next_action": safe_next_action,
"server_provenance": "client_managed" if is_client_managed else "manual_launch",
"server_provenance": provenance_assessment["launch_provenance"],
"is_client_managed": is_client_managed,
# #948: the ownership dimension, which the environment cannot establish.
# A surface that needs "may this session mutate on behalf of its client"
# reads these, not server_provenance.
"session_ownership": provenance_assessment["session_ownership"],
"session_owned": provenance_assessment["session_owned"],
"worker_identity": provenance_assessment["worker_identity"],
"session_id": provenance_assessment["session_id"],
"client_instance_id": provenance_assessment["client_instance_id"],
"generation_id": provenance_assessment["generation_id"],
"client_name": provenance_assessment["client_name"],
"fencing_epoch": provenance_assessment["fencing_epoch"],
"conflicting_live_sessions": provenance_assessment["conflicting_live_sessions"],
"provenance_assessment": provenance_assessment,
"unconsumed_gitea_env": unconsumed_env,
"preflight_ready": preflight["preflight_ready"],
"preflight_block_reasons": preflight["preflight_block_reasons"],
@@ -19416,6 +19645,76 @@ def gitea_assess_mcp_namespace_health(
return result
@mcp.tool()
def gitea_assess_mcp_namespace_attachment(
connected_servers: list[str] | None = None,
attached_session_namespaces: list[str] | None = None,
required_namespaces: list[str] | None = None,
discovery_cache_age_seconds: float | None = None,
discovery_cache_hit: bool | None = None,
auto_attach_attempted: bool = False,
auto_attach_succeeded: bool = False,
session_tool_snapshot_at: float | None = None,
namespace_connected_at: dict | None = None,
) -> dict:
"""Detect Connected-but-namespaces-not-attached MCP sessions (#708).
MCP servers can report **Connected** at the CLI/host inventory layer while the
active LLM session exposes none of their tool namespaces. That is a *session
attachment* failure and is reported here as its own typed condition,
``mcp_connected_namespaces_missing`` deliberately distinct from config drift
(#672), transport-closed (#584), and resolver EOF (#685).
Connected is not proof that tools are available. The required preflight proof is
live tool visibility plus ``gitea_whoami`` on the role namespace, never host
Connected status alone.
The verdict is recorded in the session so ``gitea_submit_pr_review`` and
``gitea_merge_pr`` fail closed while a required namespace is unattached. The only
sanctioned recovery is the client attach/reconnect path followed by full preflight;
direct imports, CLI/API mutation, profile hopping, session-state overrides, and
process kills are forbidden and never suggested.
Args:
connected_servers: Namespaces the host/CLI reports as Connected.
attached_session_namespaces: Namespaces actually exposed in the active
session tool surface.
required_namespaces: Namespaces required for this workflow; defaults to the
canonical Gitea role namespaces.
discovery_cache_age_seconds: Age of the client tool-discovery cache entry.
discovery_cache_hit: Whether the tool list came from that cache.
auto_attach_attempted: Whether the runtime tried to auto-attach namespaces.
auto_attach_succeeded: Whether that auto-attach succeeded.
session_tool_snapshot_at: Epoch seconds the session tool snapshot was taken.
namespace_connected_at: Per-namespace epoch seconds that connect completed,
used to identify multi-role startup ordering races.
Returns:
dict with ``discovery_status``, ``missing_namespaces``,
``proof_of_connected_vs_attached``, ``error_type``, ``reconnect_required``,
``auto_recovered``, ``startup_ordering_race``, a secret-free ``telemetry``
block, and a reconnect-only ``exact_next_action``.
"""
result = mcp_namespace_health.assess_connected_namespace_attachment(
connected_servers=connected_servers,
attached_session_namespaces=attached_session_namespaces,
required_namespaces=required_namespaces,
discovery_cache_age_seconds=discovery_cache_age_seconds,
discovery_cache_hit=discovery_cache_hit,
auto_attach_attempted=auto_attach_attempted,
auto_attach_succeeded=auto_attach_succeeded,
session_tool_snapshot_at=session_tool_snapshot_at,
namespace_connected_at=namespace_connected_at,
)
_record_live_namespace_attachment(result)
# #606-style watchdog check-in (best-effort, fail open). Status only, no names.
sentry_observability.monitor_checkin(
"namespace_attachment",
"ok" if result.get("attachment_healthy") else "error",
)
return result
@mcp.tool()
def gitea_activate_profile(
profile_name: str,
@@ -21515,11 +21814,23 @@ def _check_mcp_runtimes_diagnostics(task: str, matching_profiles: list[str]) ->
if match:
profile = match.group(1)
# #948: provenance for a scanned peer comes from the shared authority,
# fed by that peer's own environment, so the fleet scan cannot disagree
# with what that peer reports about itself.
peer_env = {
m.group(1): m.group(2)
for m in re.finditer(r'\b(GITEA_[A-Z0-9_]+)=([^\s]+)', env_out)
}
# declared_only: this is a peer process. Its stdin is not ours to
# inspect, and it inherits GITEA_MCP_PROFILE from any shell that
# exported it, so only an explicit declaration counts as evidence here.
is_client_managed = bool(
re.search(r'\bGITEA_CLIENT_MANAGED=(1|true|yes|client_managed)\b', env_out, re.IGNORECASE)
or re.search(r'\bGITEA_MCP_CLIENT_MANAGED=(1|true|yes|client_managed)\b', env_out, re.IGNORECASE)
or re.search(r'\bGITEA_SERVER_PROVENANCE=client_managed\b', env_out, re.IGNORECASE)
mcp_worker_identity.assess_launch_provenance(
peer_env, declared_only=True
)["client_managed"]
)
peer_worker_identity = peer_env.get("GITEA_MCP_WORKER_IDENTITY")
peer_generation = peer_env.get("GITEA_MCP_GENERATION_ID")
for env_match in re.finditer(r'\b(GITEA_[A-Z0-9_]+)=([^\s]+)', env_out):
k, v = env_match.group(1), env_match.group(2)
@@ -21535,6 +21846,8 @@ def _check_mcp_runtimes_diagnostics(task: str, matching_profiles: list[str]) ->
"start_time": start_time,
"is_stale": is_stale,
"is_client_managed": is_client_managed,
"worker_identity": peer_worker_identity,
"generation_id": peer_generation,
}
if profile not in all_profile_procs:
all_profile_procs[profile] = []
@@ -21542,12 +21855,44 @@ def _check_mcp_runtimes_diagnostics(task: str, matching_profiles: list[str]) ->
running_profiles = {}
for profile, procs in all_profile_procs.items():
if len(procs) > 1:
# #948 AC40/AC43: sharing a profile is legitimate — profile is a
# reusable capability definition, not a worker identity. What is *not*
# legitimate is reusing one worker identity, or two live sessions
# claiming one generation. Distinctness has to be proven, though:
# processes carrying no identity evidence are indistinguishable, so
# they stay classified as duplicates and keep the #686 wall intact.
identified = [p for p in procs if p.get("worker_identity")]
distinct_identities = {p["worker_identity"] for p in identified}
all_identified = len(identified) == len(procs)
duplicate_identity = len(identified) != len(distinct_identities)
contested_generation = any(
len({p["worker_identity"] for p in identified if p.get("generation_id") == gen}) > 1
for gen in {p.get("generation_id") for p in identified if p.get("generation_id")}
)
if len(procs) > 1 and (
not all_identified or duplicate_identity or contested_generation
):
pids_str = ", ".join(str(p["pid"]) for p in procs)
if duplicate_identity or contested_generation:
detail = (
"The same worker identity or generation is claimed more than once, "
"so these are genuine duplicates rather than independent workers."
)
else:
detail = (
"Manual or duplicate launches defeat staleness detection and cannot "
"receive client stdio."
)
reasons.append(
f"stale-runtime: Duplicate MCP server process(es) detected for profile '{profile}' (PIDs: {pids_str}). "
"Manual or duplicate launches defeat staleness detection and cannot receive client stdio."
+ detail
)
# Otherwise: several independently identified workers share one profile.
# No reason is appended, deliberately. Every reason this function
# returns is raised as a hard RuntimeError by its callers, so recording
# legitimate concurrency here as "informational" would block exactly the
# case #948 exists to permit.
client_procs = [p for p in procs if p["is_client_managed"]]
if client_procs:
client_procs.sort(key=lambda p: p["start_time"], reverse=True)
@@ -21956,7 +22301,7 @@ def gitea_resolve_task_capability(
next_safe_action = (
"blocker_kind=runtime_reconnect_required: call "
"gitea_request_mcp_reconnect(namespace=<active gitea-* namespace>, "
"reason='stale-runtime', client='codex') for a typed operator "
"reason='stale-runtime') for a typed operator "
"blocker with exact UI steps, then reconnect/reload the IDE-managed "
"Gitea MCP server for this profile so it reloads current master. "
"Do not edit mcp_config.json by hand, pkill, or touch configs; the "
@@ -23524,7 +23869,7 @@ def gitea_workflow_dashboard(
def gitea_request_mcp_reconnect(
namespace: str | None = None,
reason: str | None = None,
client: str = "codex",
client: str | None = None,
remote: str = "dadeschools",
host: str | None = None,
session_id: str | None = None,
@@ -23554,8 +23899,12 @@ def gitea_request_mcp_reconnect(
to the active profile's inferred namespace.
reason: Why reconnect is requested: ``stale-runtime``, ``transport_eof``,
``missing_namespace``, ``not_required``, or free-form (normalized).
client: Operator UI surface ``codex`` (default), ``claude_code``, or
``generic``.
client: Operator UI surface ``codex``, ``claude_code``, or
``generic``. #948: omitting it resolves the client from this
runtime's live attachment record rather than assuming one vendor,
and falls back to ``generic`` when nothing identifies the client.
Emitting Codex panel steps to a Gemini/Antigravity operator left
them with no reachable recovery path.
remote: Known instance ``dadeschools`` or ``prgs`` (parity context).
host: Optional host override for parity context.
session_id: Optional session id to echo in the report.
@@ -23615,6 +23964,19 @@ def gitea_request_mcp_reconnect(
else:
effective_reason = mcp_client_reconnect.REASON_UNSPECIFIED
# #948: an omitted client is resolved from the live attachment record, so
# the steps describe the UI the operator is actually in front of.
if not (client or "").strip():
client = mcp_worker_identity.reconnect_client_for(
mcp_worker_identity.assess_provenance(
registry=_worker_registry(),
worker_identity=_active_worker_identity(),
env=dict(os.environ),
stdin_is_tty=_stdin_is_tty(),
pid_alive_probe=issue_lock_store.is_process_alive,
)
)
payload = mcp_client_reconnect.build_reconnect_request(
namespace=ns,
profile=profile_name,
@@ -24070,10 +24432,18 @@ def _run_post_restart_reconcile(
repo: str | None = None,
mode: str | None = None,
limit: int = 200,
apply_session_cleanup: bool = False,
) -> dict:
"""Gather + classify post-restart state; cache the latest proof (#662)."""
"""Gather + classify post-restart state; cache the latest proof (#662 / #969).
When *apply_session_cleanup* is true, confirmed-stale session rows (dead
owner / PID reuse, no live lease) are terminalized through
``session_lifecycle`` before re-classification so the sessions dimension
can resolve. Default remains dry (read-only) to preserve #662 rollout.
"""
global _POST_RESTART_LAST_PROOF, _POST_RESTART_BOOT_RAN
import post_restart_reconcile as prr
import session_lifecycle as sl
try:
_h, o, r = _resolve(remote, None, org, repo)
@@ -24087,13 +24457,48 @@ def _run_post_restart_reconcile(
inventory = _gather_post_restart_inventory(
remote=remote, org=o, repo=r, limit=limit
)
session_cleanup: dict | None = None
if apply_session_cleanup:
db, db_errs = _control_plane_db_or_error()
if db is None:
session_cleanup = {
"success": False,
"reasons": db_errs
or ["control-plane DB unavailable; cannot retire sessions"],
}
else:
profile = get_profile()
profile_name = (profile.get("profile_name") or "").strip() or "session"
actor = f"{profile_name}-{os.getpid()}"
session_cleanup = sl.retire_stale_sessions(
db,
sessions=inventory.get("sessions") or [],
leases=inventory.get("leases") or [],
dry_run=False,
actor_session_id=actor,
session_limit=max(1, int(limit)),
)
# Re-gather active sessions after mutation so classification sees
# the post-retirement fleet (retired rows drop out of active list).
try:
inventory["sessions"] = db.list_sessions(
statuses=("active",), limit=max(1, int(limit))
)
except Exception as exc: # noqa: BLE001
inventory["incomplete_reasons"] = list(
inventory.get("incomplete_reasons") or []
) + [f"post-retirement session re-list failed: {_redact(str(exc))}"]
fleet = (session_cleanup.get("fleet") or {}) if session_cleanup else {}
inventory["session_fleet"] = fleet
proof = prr.reconcile_after_restart(
inventory,
mode=mode or _post_restart_reconcile_mode(),
)
payload = proof.as_dict()
payload["success"] = True
payload["read_only"] = True
payload["read_only"] = not apply_session_cleanup
payload["remote"] = remote
payload["org"] = o
payload["repo"] = r
@@ -24102,6 +24507,9 @@ def _run_post_restart_reconcile(
"proposed_follow_ups lists durable issues the apply path may create; "
"this tool never creates them (log-only by default, #662 rollout)"
)
if session_cleanup is not None:
payload["session_cleanup"] = session_cleanup
payload["apply_session_cleanup"] = bool(apply_session_cleanup)
_POST_RESTART_LAST_PROOF = payload
_POST_RESTART_BOOT_RAN = True
return payload
@@ -24128,8 +24536,9 @@ def gitea_reconcile_after_restart(
repo: str | None = None,
mode: str | None = None,
limit: int = 200,
apply_session_cleanup: bool = False,
) -> dict:
"""Run post-restart MCP reconciliation and return a completion proof (#662).
"""Run post-restart MCP reconciliation and return a completion proof (#662 / #969).
Gathers live control-plane sessions, leases, worktree bindings, and
master-parity evidence, then classifies them with the pure
@@ -24137,12 +24546,17 @@ def gitea_reconcile_after_restart(
machine-readable completion proof listing resolved / unresolved dimensions
and proposed durable follow-up issues.
Read-only by design: never restarts MCP, never auto-resumes write
mutations, and never creates Gitea issues (those are a separate apply
path). Default mode is ``log_only``; set
Never restarts MCP, never auto-resumes write mutations, and never creates
Gitea issues. Default mode is ``log_only``; set
``GITEA_POST_RESTART_RECONCILE_MODE=enforce`` (or pass ``mode='enforce'``)
to set ``mutation_hold`` when anything remains unresolved.
*apply_session_cleanup* (#969): when true, confirmed-stale workflow session
rows (dead owner PID / PID reuse, no live lease, not a live client-managed
owner) are terminalized through the sanctioned ``session_lifecycle`` path
before re-classification. Default false preserves the historical read-only
gather+classify behaviour. Does not delete historical rows.
Soft-depends on #660 for session checkpoints: when the checkpoint schema
module is absent the checkpoints dimension is ``skipped`` with an explicit
reason rather than inventing a schema.
@@ -24162,9 +24576,96 @@ def gitea_reconcile_after_restart(
repo=repo,
mode=mode,
limit=limit,
apply_session_cleanup=bool(apply_session_cleanup),
)
@mcp.tool()
def gitea_retire_stale_workflow_sessions(
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
apply: bool = False,
limit: int = 500,
) -> dict:
"""Retire workflow session rows whose owners are no longer live (#969).
Plans (and optionally applies) terminalization of control-plane session
rows that are confirmed stale:
* owner PID absent / dead
* PID reused by an unrelated process (process start after session start)
* no live workflow lease
* not a live client-managed owner
Default ``apply=false`` is dry-run only. ``apply=true`` performs CAS
status updates to ``retired`` and writes durable ``session_retired``
events. Idempotent under concurrent reconciliation. Never deletes rows
and never touches sessions protected by a live lease or live owner.
"""
read_block = _profile_operation_gate("gitea.read")
if read_block:
return {
"success": False,
"read_only": not apply,
"reasons": read_block,
"permission_report": _permission_block_report("gitea.read"),
}
try:
_h, o, r = _resolve(remote, None, org, repo)
except ValueError as exc:
return {"success": False, "reasons": [str(exc)], "read_only": not apply}
db, db_errs = _control_plane_db_or_error()
if db is None:
return {
"success": False,
"read_only": not apply,
"reasons": db_errs or ["control-plane DB unavailable"],
}
import session_lifecycle as sl
profile = get_profile()
profile_name = (profile.get("profile_name") or "").strip() or "session"
actor = f"{profile_name}-{os.getpid()}"
# Prefer lease inventory scoped to the requested repo when available.
leases: list[dict] = []
try:
lease_result = lease_lifecycle.list_active_leases(
db,
remote=remote if remote in REMOTES else remote,
org=o,
repo=r,
role=None,
include_non_active=True,
limit=max(1, int(limit)),
)
leases = list(lease_result.get("leases") or [])
except Exception: # noqa: BLE001
try:
leases = db.list_leases(statuses=("active",), limit=max(1, int(limit)))
except Exception: # noqa: BLE001
leases = []
result = sl.retire_stale_sessions(
db,
leases=leases,
dry_run=not bool(apply),
actor_session_id=actor,
session_limit=max(1, int(limit)),
)
result["read_only"] = not bool(apply)
result["apply"] = bool(apply)
result["remote"] = remote
result["org"] = o
result["repo"] = r
return result
@mcp.tool()
def gitea_inspect_workflow_lease(
lease_id: str,
@@ -24718,7 +25219,11 @@ if __name__ == "__main__":
# basename-only stack frames, and import-only launch cannot reconstruct
# native transport; offline imports / standalone scripts fail closed.
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
# configuration (GITEA_MCP_TRANSPORT), defaults to stdio when unset, and is
# validated against the single permitted set in mcp_transport_config. An
# unregistered identifier raises here, before any tool can dispatch.
mcp_daemon_guard.bind_native_mcp_transport()
# Lock this session's launch profile into the environment so child CLI
# processes (e.g. review_pr.py) can detect and refuse profile
# side-channel overrides (#199).
@@ -24736,4 +25241,11 @@ if __name__ == "__main__":
sys.stderr.write(
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()
pytest = mcp_daemon_guard.is_pytest_runtime()
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:
reasons.append(
"irrecoverable provenance authorization requires production native "
@@ -510,6 +515,7 @@ def assess_transport_for_auth_mint() -> dict[str, Any]:
"allowed": not reasons,
"native_mcp_transport": native,
"production_native_mcp_transport": production,
"bound_transport": bound,
"pytest": pytest,
"reasons": reasons,
}
+21 -9
View File
@@ -92,7 +92,15 @@ OPERATOR_UI_STEPS: dict[str, tuple[str, ...]] = {
),
}
DEFAULT_CLIENT = "codex"
#: What an *unidentified* client gets. #948: this is deliberately the
#: host-agnostic step set rather than a specific product. Defaulting to one
#: vendor emitted Codex UI steps to a Gemini/Antigravity operator, who then had
#: no reachable recovery path — the guidance named a panel they do not have.
DEFAULT_CLIENT = "generic"
#: The historical default, kept addressable by name so Codex callers still get
#: Codex steps, without it silently becoming the fallback for unknown clients.
LEGACY_DEFAULT_CLIENT = "codex"
def normalize_reason(reason: str | None) -> str:
@@ -128,14 +136,18 @@ def normalize_reason(reason: str | None) -> str:
def normalize_client(client: str | None) -> str:
"""Return a known client key for operator UI steps."""
text = (client or "").strip().lower().replace(" ", "_").replace("-", "_")
if text in ("codex", "openai_codex", "openai"):
return "codex"
if text in ("claude", "claude_code", "claude_desktop", "anthropic"):
return "claude_code"
if text in OPERATOR_UI_STEPS:
return text
"""Return the UI-step key for a client.
#948: alias resolution is shared with ``mcp_worker_identity`` so a client
name means the same thing wherever it is read. A name we recognise but have
no bespoke steps for — Gemini, Antigravity, Grok — resolves to the generic
host-agnostic steps rather than to another vendor's panel.
"""
import mcp_worker_identity
canonical = mcp_worker_identity.normalize_client_name(client)
if canonical in OPERATOR_UI_STEPS:
return canonical
return DEFAULT_CLIENT
+176 -13
View File
@@ -32,6 +32,8 @@ import time
from pathlib import Path
from typing import Any
import mcp_transport_config
SANCTIONED_DAEMON_ENV = "GITEA_MCP_SANCTIONED_DAEMON"
ALLOW_DIRECT_IMPORT_ENV = "GITEA_ALLOW_DIRECT_MCP_IMPORT"
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
# 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_TEST = "test"
_PHASE_ENTRYPOINT_CLAIMED = "entrypoint_claimed"
@@ -58,6 +62,23 @@ class UnsanctionedRuntimeError(RuntimeError):
"""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:
if (os.environ.get(FORCE_PROVENANCE_FAIL_ENV) or "").strip() in {
"1",
@@ -171,23 +192,46 @@ def mark_sanctioned_daemon() -> dict[str, Any]:
return native_runtime_status()
def bind_native_mcp_transport(*, transport: str) -> dict[str, Any]:
"""Bind the live native MCP transport lifecycle (#695).
def bind_native_mcp_transport(*, transport: str | None = None) -> dict[str, Any]:
"""Bind the live native MCP transport lifecycle (#695 / #931).
Must be called from the resolved canonical entrypoint immediately before
the real MCP server transport loop (e.g. ``mcp.run(transport=\"stdio\")``).
Requires a prior successful :func:`mark_sanctioned_daemon` claim in this
process. Import-only or offline launch without this bind leaves
the real MCP server transport loop (``mcp.run``). Requires a prior
successful :func:`mark_sanctioned_daemon` claim in this process.
Import-only or offline launch without this bind leaves
: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
transport_name = (transport or "").strip().lower()
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)}."
)
if transport is None:
resolution = mcp_transport_config.resolve_configured_transport()
transport_name = str(resolution["transport"])
if not resolution["supported"]:
raise UnsanctionedRuntimeError(
"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()
if entrypoint_path is None:
@@ -216,6 +260,23 @@ def bind_native_mcp_transport(*, transport: str) -> dict[str, Any]:
"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).
# Changing GITEA_MCP_SESSION_STATE_DIR after bind must not manufacture a
# 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
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:
"""Backward-compatible name; #695 requires native transport, not env alone."""
if is_production_native_mcp_transport():
@@ -476,6 +619,21 @@ def native_runtime_status() -> dict[str, Any]:
"entrypoint_path": rt.get("entrypoint_path"),
"phase": rt.get("phase"),
"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"),
"session_state_dir": pinned_session_state_dir() or rt.get("session_state_dir"),
"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"]:
transport = "test_native_mcp"
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,
"bound_transport": st.get("bound_transport"),
"native_mcp_transport": bool(st["native_mcp_transport"]),
"production_native_mcp_transport": bool(
st.get("production_native_mcp_transport")
+397 -5
View File
@@ -51,6 +51,42 @@ EOF_PATTERNS = (
"eof",
)
ERROR_CONNECTED_NAMESPACES_MISSING = "mcp_connected_namespaces_missing"
# Distinct from the condition above. ``mcp_connected_namespaces_missing`` is reserved for
# the actual #708 defect: the host *does* report the service Connected, yet the namespace
# never entered the active session tool surface. A required namespace that is absent from
# the connected-service inventory has no Connected claim behind it at all, so reporting it
# under the #708 condition would assert something the evidence does not support and would
# point an operator at the wrong recovery.
ERROR_REQUIRED_NAMESPACES_NOT_CONNECTED = "mcp_required_namespaces_not_connected"
DISCOVERY_STATUS_ATTACHED = "namespaces_attached"
DISCOVERY_STATUS_CONNECTED_MISSING = "connected_but_namespaces_missing"
DISCOVERY_STATUS_NOT_CONNECTED = "required_namespaces_not_connected"
DISCOVERY_STATUS_DISCONNECTED = "disconnected"
# Namespaces that must be *attached to the active session* for a mutation task (#708).
# Connected-at-host is not attached-in-session; these are gated separately from the
# #543 health map because a namespace can be healthy on probe yet absent from the
# session tool surface.
ATTACHMENT_GATED_TASKS = {
"review_pr": "gitea-reviewer",
"submit_review": "gitea-reviewer",
"merge_pr": "gitea-merger",
"work_issue": "gitea-author",
"create_pr": "gitea-author",
}
# The only sanctioned recovery for an unattached namespace (#678 exposes it natively).
SANCTIONED_ATTACH_RECOVERY_TOOL = "gitea_request_mcp_reconnect"
UNSAFE_FALLBACK_WARNING = (
"Workflow Safety Hard Stop (#708): Connected-but-namespaces-missing recovery must "
"NEVER use direct imports, Gitea API mutations, profile hopping, session-state "
"overrides, PID kills, or config mtime touches. Use client reconnect only."
)
SAFE_ENV_KEYS = (
"GITEA_MCP_PROFILE",
"GITEA_PROFILE_NAME",
@@ -59,6 +95,318 @@ SAFE_ENV_KEYS = (
"GITEA_MCP_CONFIG",
)
# #948: provenance used to be derived from the summary this allowlist produces.
# The allowlist never carried a provenance key, so that derivation could only
# ever evaluate to ``manual_launch`` — whatever the process actually was — while
# ``gitea_get_runtime_context`` read the live environment and reported
# ``client_managed`` for the same process. Provenance is no longer derived here.
# It comes from ``mcp_worker_identity.assess_provenance``, the single authority
# every surface shares. This allowlist keeps its original and only job: deciding
# which env values are safe to echo back in diagnostics.
def assess_connected_namespace_attachment(
*,
connected_servers: list[str] | tuple[str, ...] | set[str] | None = None,
attached_session_namespaces: list[str] | tuple[str, ...] | set[str] | None = None,
required_namespaces: list[str] | tuple[str, ...] | set[str] | None = None,
discovery_cache_age_seconds: float | int | None = None,
discovery_cache_hit: bool | None = None,
auto_attach_attempted: bool = False,
auto_attach_succeeded: bool = False,
session_tool_snapshot_at: float | int | None = None,
namespace_connected_at: dict[str, float] | None = None,
) -> dict[str, Any]:
"""Assess whether host-connected MCP servers have attached tool namespaces in the active session (#708).
Addresses the Connected-but-namespaces-missing defect: CLI/host status may report Connected
while the active LLM session tool surface exposes 0 attached tool namespaces.
This is a *distinct* condition from config drift (#672), transport-closed (#584),
and resolver EOF (#685): the transport is up and the host reports Connected, yet the
namespace never entered the session tool surface.
Startup ordering (``session_tool_snapshot_at`` + ``namespace_connected_at``) identifies
the race where the session tool snapshot was taken before a role server finished
``initialize``/``list_tools``, which is why parallel multi-role startup can leave the
session with an empty namespace set while Connected later flips true.
Returns structured detection details plus secret-free telemetry.
"""
connected = [str(s).strip() for s in (connected_servers or []) if str(s).strip()]
attached = set(str(ns).strip() for ns in (attached_session_namespaces or []) if str(ns).strip())
req = [str(r).strip() for r in (required_namespaces or DEFAULT_NAMESPACES) if str(r).strip()]
required = list(dict.fromkeys(req))
connected_set = set(connected)
proof: dict[str, dict[str, bool]] = {}
for s in connected:
proof[s] = {"connected": True, "attached": s in attached}
for r in required:
if r not in proof:
proof[r] = {"connected": r in connected_set, "attached": r in attached}
# The genuine #708 condition: the host reports the service Connected and the namespace
# still never entered the active session tool surface.
missing = [r for r in required if r in connected_set and r not in attached]
# A separate condition: the service is required but absent from the connected-service
# inventory. Nothing here is Connected, so this must not borrow the #708 wording or its
# recovery — see ERROR_REQUIRED_NAMESPACES_NOT_CONNECTED.
not_connected = [r for r in required if r not in connected_set]
# Evidence that disagrees with itself: reported attached to the session while absent
# from the connected inventory. Neither statement is proof, so it fails closed.
contradictory = [r for r in not_connected if r in attached]
# No required namespace may lack attachment proof and still be called healthy.
attachment_healthy = (
not missing
and not not_connected
and len(connected) > 0
and len(required) > 0
)
if attachment_healthy:
discovery_status = DISCOVERY_STATUS_ATTACHED
elif not connected:
discovery_status = DISCOVERY_STATUS_DISCONNECTED
elif not_connected:
discovery_status = DISCOVERY_STATUS_NOT_CONNECTED
else:
discovery_status = DISCOVERY_STATUS_CONNECTED_MISSING
# Every condition actually present is reported; ``error_type`` names the primary one.
# Not-connected outranks Connected-but-unattached because a service that never
# connected cannot be recovered by attaching its namespace.
error_types: list[str] = []
if not_connected:
error_types.append(ERROR_REQUIRED_NAMESPACES_NOT_CONNECTED)
if missing:
error_types.append(ERROR_CONNECTED_NAMESPACES_MISSING)
if not attachment_healthy and not error_types:
error_types.append(ERROR_REQUIRED_NAMESPACES_NOT_CONNECTED)
error_type = None if attachment_healthy else error_types[0]
# Per-namespace verdict, so a mutation gate never has to infer one namespace's state
# from a whole-session summary. Each entry states only what its own evidence supports.
namespace_conditions: dict[str, dict[str, Any]] = {}
for ns_name, state in proof.items():
ns_connected = bool(state["connected"])
ns_attached = bool(state["attached"])
if ns_connected and ns_attached:
condition = None
elif ns_connected:
condition = ERROR_CONNECTED_NAMESPACES_MISSING
else:
condition = ERROR_REQUIRED_NAMESPACES_NOT_CONNECTED
namespace_conditions[ns_name] = {
"namespace": ns_name,
"required": ns_name in required,
"connected": ns_connected,
"attached": ns_attached,
"attachment_healthy": ns_connected and ns_attached,
"condition": condition,
"contradictory_evidence": ns_attached and not ns_connected,
}
# Startup-ordering race: a namespace that finished connecting *after* the session
# tool snapshot was taken cannot be in that snapshot, however healthy it looks now.
connected_at = {
str(k).strip(): v
for k, v in (namespace_connected_at or {}).items()
if str(k).strip() and isinstance(v, (int, float))
}
late_attaching: list[str] = []
if isinstance(session_tool_snapshot_at, (int, float)):
for ns_name, ts in connected_at.items():
if ts > session_tool_snapshot_at and ns_name not in attached:
late_attaching.append(ns_name)
late_attaching.sort()
startup_ordering_race = bool(late_attaching)
auto_recovered = bool(auto_attach_attempted and auto_attach_succeeded and attachment_healthy)
reconnect_required = not attachment_healthy
reasons: list[str] = []
remediation: list[str] = []
if not connected:
reasons.append("No MCP servers reported Connected.")
remediation.append("Start or reconnect Gitea MCP servers in client config.")
if missing:
reasons.append(
f"MCP server(s) {missing} report Connected at host/CLI layer but tool namespaces "
f"are missing from active session attached tools (Connected ≠ attached tools, #708)."
)
remediation.append(
"Reconnect the IDE/client MCP session to attach namespaces to the active session "
f"(sanctioned path: {SANCTIONED_ATTACH_RECOVERY_TOOL}), then re-run full preflight "
"(gitea_whoami -> gitea_resolve_task_capability -> task). "
"Do not use direct imports, CLI API mutations, profile hopping, or session file overrides."
)
if not_connected:
reasons.append(
f"required MCP namespace(s) {not_connected} are absent from the connected-service "
"inventory: nothing reports them Connected, so there is no attachment to claim and "
"the Connected-but-unattached condition does not apply to them (#708)."
)
remediation.append(
f"Connect the required MCP server(s) {not_connected} through the client, then "
"reconnect the IDE/client MCP session so their namespaces attach "
f"(sanctioned path: {SANCTIONED_ATTACH_RECOVERY_TOOL}), and re-run full preflight. "
"Do not use direct imports, CLI API mutations, profile hopping, or session file overrides."
)
if contradictory:
reasons.append(
f"contradictory evidence for namespace(s) {contradictory}: reported attached to the "
"active session while absent from the connected-service inventory; neither statement "
"is proof, so attachment is treated as unproven (fail closed, #708)."
)
if attachment_healthy:
reasons.append(
"All required MCP server namespaces are connected and attached to the active session."
)
if startup_ordering_race:
reasons.append(
f"startup ordering race: namespace(s) {late_attaching} finished connecting after the "
"active session tool snapshot was taken, so they cannot appear in that snapshot (#708)."
)
if auto_attach_attempted and not auto_attach_succeeded:
reasons.append(
"automatic namespace attachment was attempted and did not succeed; only the sanctioned "
"client reconnect path remains."
)
if auto_recovered:
reasons.append("namespaces were automatically attached; no operator reconnect was required.")
return {
"success": attachment_healthy,
"attachment_healthy": attachment_healthy,
"discovery_status": discovery_status,
"connected_servers": connected,
"attached_session_namespaces": list(attached),
"missing_namespaces": missing,
"not_connected_namespaces": not_connected,
"contradictory_namespaces": contradictory,
"proof_of_connected_vs_attached": proof,
"namespace_conditions": namespace_conditions,
"error_type": error_type,
"error_types": error_types,
"reasons": reasons,
"remediation": remediation,
"exact_next_action": (
"None; session tool namespaces attached."
if attachment_healthy
else (
f"Connect the required MCP server(s) {not_connected} through the client, then "
"reconnect the IDE/client MCP session so their tool namespaces attach to the "
"active session. Do not use direct imports, CLI API mutations, profile hopping, "
"or session-state overrides."
if not_connected
else (
"Reconnect the IDE/client MCP session so tool namespaces attach to the active "
"session. Do not use direct imports, CLI API mutations, profile hopping, or "
"session-state overrides."
)
)
),
"unsafe_fallback_policy": UNSAFE_FALLBACK_WARNING,
"sanctioned_recovery_tool": SANCTIONED_ATTACH_RECOVERY_TOOL,
"reconnect_required": reconnect_required,
"auto_attach_attempted": bool(auto_attach_attempted),
"auto_recovered": auto_recovered,
"startup_ordering_race": startup_ordering_race,
"late_attaching_namespaces": late_attaching,
# Secret-free structured signals (#708 AC5). Namespace names and counts only:
# never tokens, endpoints, env values, or filesystem paths.
"telemetry": {
"connected_count": len(connected),
"attached_count": len(attached),
"required_count": len(required),
"missing_count": len(missing),
"not_connected_count": len(not_connected),
"contradictory_count": len(contradictory),
"required_attached_count": sum(1 for r in required if r in attached),
"discovery_status": discovery_status,
"discovery_cache_hit": (
None if discovery_cache_hit is None else bool(discovery_cache_hit)
),
"discovery_cache_age_seconds": (
float(discovery_cache_age_seconds)
if isinstance(discovery_cache_age_seconds, (int, float))
else None
),
"reconnect_required": reconnect_required,
"auto_attach_attempted": bool(auto_attach_attempted),
"auto_recovered": auto_recovered,
"startup_ordering_race": startup_ordering_race,
"error_type": error_type,
"error_types": error_types,
},
}
def required_namespace_for_attachment(task: str) -> str | None:
"""Map a mutation task to the MCP namespace that must be *attached* (#708)."""
return ATTACHMENT_GATED_TASKS.get((task or "").strip())
def attachment_gate_from_session(
task: str,
session_attachment: dict[str, dict[str, Any]] | None,
) -> list[str]:
"""Fail-closed gate on recorded connected-but-unattached namespaces (#708).
Mirrors :func:`mutation_gate_from_session`: a namespace that has not been
assessed yet does not gate, so this never blocks a session that simply has
not run the assessment. Once an assessment records the namespace required
for *task* as Connected-but-unattached, the mutation fails closed and the
only offered recovery is the sanctioned client reconnect path.
"""
ns = required_namespace_for_attachment(task)
if not ns:
return []
store = session_attachment or {}
entry = store.get(ns)
if not entry:
return []
if entry.get("attached") and entry.get("attachment_healthy"):
return []
what = task or "mutation"
connected = entry.get("connected")
detail = entry.get("condition") or entry.get("error_type")
if connected is True:
# Only here has anything actually reported the service Connected, so only here may
# the block say so.
detail = detail or ERROR_CONNECTED_NAMESPACES_MISSING
blocked = (
f"live MCP namespace '{ns}' is recorded {detail}: the host reports Connected but the "
f"namespace is not attached to the active session tool surface; reconnect the "
f"IDE/client MCP session and re-run preflight before {what} "
"(fail closed, #708)"
)
elif connected is False:
detail = ERROR_REQUIRED_NAMESPACES_NOT_CONNECTED
blocked = (
f"live MCP namespace '{ns}' is recorded {detail}: it is absent from the "
f"connected-service inventory, so it is neither connected nor attached and no "
f"Connected status is claimed for it; connect the required MCP server, then "
f"reconnect the IDE/client MCP session and re-run preflight before {what} "
"(fail closed, #708)"
)
else:
# Connected status was never recorded. Refuse without asserting either condition.
detail = detail or ERROR_CONNECTED_NAMESPACES_MISSING
blocked = (
f"live MCP namespace '{ns}' is recorded {detail} with no connected-status evidence, "
f"so attachment to the active session tool surface is unproven; reconnect the "
f"IDE/client MCP session and re-run preflight before {what} "
"(fail closed, #708)"
)
return [blocked, UNSAFE_FALLBACK_WARNING]
def _as_list(value: Any) -> list[str] | None:
if value is None:
@@ -104,12 +452,23 @@ def classify_namespace_probe(
profile: str | None = None,
configured: bool = True,
probe_source: str | None = None,
worker_identity: str | None = None,
generation_id: str | None = None,
registry: Any | None = None,
pid_alive_probe: Any | None = None,
) -> dict[str, Any]:
"""Classify whether a required tool is callable through a live namespace.
``registered_tools`` is static/server-side evidence. ``probe_result`` is
live invocation evidence. Only ``probe_source=client_namespace`` proves the
IDE-managed path; ``offline_spawn`` is an offline subprocess check only.
#948: ``worker_identity``/``generation_id``/``registry`` carry the
client/session ownership evidence. Provenance is resolved by
``mcp_worker_identity.assess_provenance`` — the same call
``gitea_get_runtime_context`` makes — so the two surfaces cannot report
different provenance for one process. Omitting them yields the fail-closed
``unproven`` verdict, never a fabricated ``client_managed``.
"""
ns = (namespace or "").strip()
tool = required_tool or REQUIRED_NAMESPACE_TOOLS.get(ns) or "gitea_whoami"
@@ -226,14 +585,31 @@ def classify_namespace_probe(
blocks = namespace_health_blocks_task("merge_pr", healthy)
import gitea_config
import mcp_worker_identity
raw_env = process.get("env") if isinstance(process, dict) else None
unconsumed_env = gitea_config.get_unconsumed_gitea_env_overrides(raw_env)
is_client_managed = bool(
env_summary.get("GITEA_CLIENT_MANAGED") in ("1", "true", "yes", "client_managed")
or env_summary.get("GITEA_MCP_CLIENT_MANAGED") in ("1", "true", "yes", "client_managed")
or env_summary.get("GITEA_SERVER_PROVENANCE") == "client_managed"
# #948: one authority, shared with gitea_get_runtime_context. The env is
# passed whole rather than through SAFE_ENV_KEYS — the allowlist exists to
# decide what may be *echoed*, and using it to decide what may be *believed*
# is what made this surface structurally unable to report client_managed.
# ``declared_only``: ``process`` describes an observed peer, not this
# interpreter. Its stdin is unavailable and its launcher-config env is
# inherited from whatever shell started it, so only an explicit declaration
# is evidence. Absence of one is ``unproven``, not an asserted manual launch.
provenance_verdict = mcp_worker_identity.assess_provenance(
registry=registry,
worker_identity=worker_identity,
generation_id=generation_id,
env=raw_env if isinstance(raw_env, dict) else {},
namespace=ns,
profile=profile_name,
pid_alive_probe=pid_alive_probe,
declared_only=True,
)
provenance = "client_managed" if is_client_managed else "manual_launch"
provenance = provenance_verdict["provenance"]
is_client_managed = provenance_verdict["is_client_managed"]
return {
"success": healthy,
@@ -252,6 +628,15 @@ def classify_namespace_probe(
"remediation": remediation,
"provenance": provenance,
"is_client_managed": is_client_managed,
# Every non-client-session verdict fails closed. Consumers that only
# need "may this mutate?" read this and stay correct across the #948
# vocabulary split between ``manual_launch`` and ``unproven``.
"provenance_fail_closed": provenance_verdict["fail_closed"],
"provenance_assessment": provenance_verdict,
"worker_identity": provenance_verdict["worker_identity"],
"session_id": provenance_verdict["session_id"],
"generation_id": provenance_verdict["generation_id"],
"client_name": provenance_verdict["client_name"],
"unconsumed_gitea_env": unconsumed_env,
"diagnostics": {
"namespace": ns,
@@ -263,6 +648,13 @@ def classify_namespace_probe(
"probe_source": source,
"provenance": provenance,
"is_client_managed": is_client_managed,
"provenance_fail_closed": provenance_verdict["fail_closed"],
"provenance_blocker_kind": provenance_verdict["blocker_kind"],
"provenance_scope": provenance_verdict["scope"],
"worker_identity": provenance_verdict["worker_identity"],
"session_id": provenance_verdict["session_id"],
"generation_id": provenance_verdict["generation_id"],
"client_name": provenance_verdict["client_name"],
"unconsumed_gitea_env": unconsumed_env,
},
"blocks_merge_workflow": blocks,
+7 -3
View File
@@ -1,7 +1,10 @@
#!/usr/bin/env python3
"""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 sys
@@ -43,8 +46,9 @@ check_conflict_markers()
# #558 / #695: claim the official entrypoint before loading mutation modules.
# This alone does NOT authorize mutations — gitea_mcp_server binds the live
# native MCP transport (stdio) immediately before mcp.run. Import-only or
# offline launch without that bind fails closed on mutations.
# native MCP transport immediately before mcp.run, over the configured
# transport (#931). Import-only or offline launch without that bind fails
# closed on mutations.
try:
import mcp_daemon_guard
+4
View File
@@ -588,9 +588,13 @@ def save_state(
bool(prov.get("production_native_mcp_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:
body.setdefault("native_mcp_transport", False)
body.setdefault("transport", "untrusted")
body.setdefault("bound_transport", None)
envelope = {
"kind": kind,
+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"])
File diff suppressed because it is too large Load Diff
+57 -12
View File
@@ -475,24 +475,65 @@ def reconcile_after_restart(
)
)
# --- sessions -------------------------------------------------------
# --- sessions (#969: dead-owner / PID-reuse lifecycle) ---------------
# Prefer a precomputed fleet report from the gather/apply path when present;
# otherwise classify pure from inventory (injectable checkers stay default).
sessions = [s for s in (inventory.get("sessions") or []) if isinstance(s, Mapping)]
orphan_sessions = [
s
for s in sessions
if str(s.get("status") or "").lower() == "active"
and s.get("pid") is not None
and not lease_lifecycle.is_process_alive(s.get("pid"))
leases_for_sessions = [
L for L in (inventory.get("leases") or []) if isinstance(L, Mapping)
]
if orphan_sessions:
fleet_report = inventory.get("session_fleet")
if isinstance(fleet_report, Mapping) and "retireable_session_ids" in fleet_report:
fleet_details = dict(fleet_report)
retireable_ids = list(fleet_details.get("retireable_session_ids") or [])
resolved = bool(fleet_details.get("sessions_dimension_resolved", not retireable_ids))
else:
try:
import session_lifecycle as _sl
client_managed = inventory.get("client_managed_session_ids")
cm_set = None
if isinstance(client_managed, (list, tuple, set, frozenset)):
cm_set = {str(x) for x in client_managed}
fleet = _sl.classify_sessions(
sessions,
leases=leases_for_sessions,
now=started,
client_managed_sessions=cm_set,
)
fleet_details = fleet.as_dict()
retireable_ids = list(fleet_details.get("retireable_session_ids") or [])
resolved = bool(fleet_details.get("sessions_dimension_resolved"))
except Exception as exc: # noqa: BLE001 — fail closed to legacy signal
# Legacy fallback: dead-pid active rows only (pre-#969 behaviour).
orphan_sessions = [
s
for s in sessions
if str(s.get("status") or "").lower() == "active"
and s.get("pid") is not None
and not lease_lifecycle.is_process_alive(s.get("pid"))
]
retireable_ids = [s.get("session_id") for s in orphan_sessions]
resolved = not orphan_sessions
fleet_details = {
"total_sessions": len(sessions),
"retireable_session_ids": retireable_ids,
"legacy_fallback": True,
"fallback_error": str(exc),
}
if not resolved and retireable_ids:
items.append(
_item(
DIM_SESSIONS,
ITEM_UNRESOLVED,
f"{len(orphan_sessions)} active session row(s) with dead owner pid",
f"{len(retireable_ids)} session row(s) with dead/reused owner "
f"await retirement",
details={
"orphan_session_ids": [s.get("session_id") for s in orphan_sessions],
"orphan_session_ids": retireable_ids,
"retireable_session_ids": retireable_ids,
"total_sessions": len(sessions),
"fleet": fleet_details,
},
follow_up=True,
)
@@ -502,8 +543,12 @@ def reconcile_after_restart(
_item(
DIM_SESSIONS,
ITEM_RESOLVED,
f"{len(sessions)} session row(s) reconciled (no dead-pid orphans)",
details={"total_sessions": len(sessions)},
f"{len(sessions)} session row(s) reconciled "
f"(no retireable dead/reused owners)",
details={
"total_sessions": len(sessions),
"fleet": fleet_details,
},
)
)
+703
View File
@@ -0,0 +1,703 @@
"""Safe lifecycle for workflow session rows with dead or reused owners (#969).
Post-restart reconciliation previously left hundreds of ``active`` session rows
with dead owner PIDs permanently unresolved. A PID existence check alone is not
enough: operating systems reuse PIDs, so an unrelated live process can appear to
own a historical session.
This module is the pure classification + apply core for session retirement:
* distinguish live, disconnected, stale, protected, and terminal records
* refuse retirement when a live lease or live verified owner remains
* protect live client-managed sessions
* detect PID reuse via process start time vs session start / heartbeat
* terminalize confirmed-stale rows idempotently with durable audit events
* stay safe under concurrent reconciles (CAS on status)
Design mirrors ``lease_lifecycle`` / ``post_restart_reconcile``:
* Pure classification accepts injectable checkers so unit tests never touch
real processes.
* Apply mutations go only through :meth:`ControlPlaneDB.retire_session`.
* Historical rows are never deleted; status moves to a terminal value and an
events-row records the action + reason.
"""
from __future__ import annotations
import json
import os
import subprocess
from dataclasses import dataclass, field
from datetime import datetime, timedelta, timezone
from typing import Any, Callable, Mapping, Sequence
import control_plane_db as cpd
import lease_lifecycle
# Session status vocabulary.
SESSION_STATUS_ACTIVE = "active"
SESSION_STATUS_RETIRED = "retired"
SESSION_STATUS_ENDED = "ended"
SESSION_STATUS_TERMINAL = "terminal"
TERMINAL_SESSION_STATUSES = frozenset(
{
SESSION_STATUS_RETIRED,
SESSION_STATUS_ENDED,
SESSION_STATUS_TERMINAL,
"dead",
"stale",
"orphaned",
}
)
# Classification outcomes for one session row.
CLASS_LIVE = "live"
CLASS_DISCONNECTED = "disconnected"
CLASS_STALE = "stale"
CLASS_PROTECTED = "protected"
CLASS_TERMINAL = "terminal"
# Stable reason codes (audit + tests).
REASON_ALREADY_TERMINAL = "already_terminal"
REASON_LIVE_OWNER = "live_owner"
REASON_LIVE_LEASE = "live_lease"
REASON_CLIENT_MANAGED_LIVE = "client_managed_live"
REASON_DEAD_OWNER = "dead_owner"
REASON_PID_REUSE = "pid_reuse"
REASON_HEARTBEAT_STALE_DEAD = "heartbeat_stale_dead_owner"
REASON_MISSING_PID = "missing_pid_no_lease"
# Event type written to control-plane events table.
EVENT_SESSION_RETIRED = "session_retired"
# Default heartbeat window before a still-alive PID is treated as disconnected
# rather than live (does not alone authorize retirement).
DEFAULT_HEARTBEAT_STALE_SECONDS = 900
# PID reuse: process start must be strictly later than session started_at by
# more than this skew (ps lstart is second-resolution; clocks can lag).
PID_REUSE_SKEW = timedelta(seconds=2)
# Live lease freshness values that block retirement.
_LIVE_LEASE_FRESHNESS = frozenset({"active", "live"})
class SessionLifecycleError(RuntimeError):
"""Fail-closed session lifecycle policy error."""
def _utc_now() -> datetime:
return datetime.now(timezone.utc)
def _parse_ts(value: str | datetime | None) -> datetime | None:
if value is None:
return None
if isinstance(value, datetime):
if value.tzinfo is None:
return value.replace(tzinfo=timezone.utc)
return value.astimezone(timezone.utc)
return cpd._parse_ts(str(value))
def _ts(dt: datetime | None = None) -> str:
return cpd._ts(dt)
def process_start_time(pid: int | None) -> datetime | None:
"""Return the OS start time of *pid*, or None when it cannot be resolved.
Uses ``ps -o lstart=`` (POSIX). Failures return None rather than inventing
evidence missing start time never authorizes retirement of a live PID.
"""
if pid is None:
return None
try:
pid_i = int(pid)
except (TypeError, ValueError):
return None
if pid_i <= 0:
return None
try:
proc = subprocess.run(
["ps", "-o", "lstart=", "-p", str(pid_i)],
capture_output=True,
text=True,
check=False,
timeout=2,
)
except (OSError, subprocess.SubprocessError):
return None
if proc.returncode != 0:
return None
text = (proc.stdout or "").strip()
if not text:
return None
try:
# Example: "Wed Jul 29 09:14:36 2026"
naive = datetime.strptime(text, "%a %b %d %H:%M:%S %Y")
return naive.replace(tzinfo=timezone.utc)
except ValueError:
return None
def _lease_freshness_label(lease: Mapping[str, Any]) -> str:
fr = lease.get("freshness")
if isinstance(fr, Mapping):
return str(fr.get("freshness") or "").strip().lower()
if fr:
return str(fr).strip().lower()
# Fall back to classifying a raw lease row.
try:
return str(
lease_lifecycle.classify_lease_freshness(lease).get("freshness") or ""
).strip().lower()
except Exception: # noqa: BLE001 — pure classifier must not raise on bad rows
status = str(lease.get("status") or "").strip().lower()
return status or "unknown"
def live_lease_session_ids(
leases: Sequence[Mapping[str, Any]] | None,
*,
now: datetime | None = None,
pid_checker: Callable[[int | None], bool] = lease_lifecycle.is_process_alive,
) -> set[str]:
"""Session ids that still hold a live (or ambiguous-active) workflow lease."""
live: set[str] = set()
moment = now or _utc_now()
for lease in leases or ():
if not isinstance(lease, Mapping):
continue
status = str(lease.get("status") or "").strip().lower()
if status and status not in {
lease_lifecycle.LEASE_STATUS_ACTIVE,
"",
}:
# Explicit terminal lease statuses never protect a session.
if status in {
lease_lifecycle.LEASE_STATUS_RELEASED,
lease_lifecycle.LEASE_STATUS_EXPIRED,
lease_lifecycle.LEASE_STATUS_ABANDONED,
}:
continue
freshness = _lease_freshness_label(lease)
if freshness in _LIVE_LEASE_FRESHNESS or freshness in {"", "unknown"}:
# Ambiguous active rows: re-check with authoritative classifier.
try:
fr = lease_lifecycle.classify_lease_freshness(
lease, now=moment, pid_checker=pid_checker
)
freshness = str(fr.get("freshness") or "").strip().lower()
except Exception: # noqa: BLE001
freshness = "unknown"
if freshness in _LIVE_LEASE_FRESHNESS:
sid = str(lease.get("session_id") or "").strip()
if sid:
live.add(sid)
elif freshness == "unknown" and status in {
lease_lifecycle.LEASE_STATUS_ACTIVE,
"",
}:
# Fail closed: active lease with unknown freshness blocks retirement.
sid = str(lease.get("session_id") or "").strip()
if sid:
live.add(sid)
return live
@dataclass(frozen=True)
class SessionClassification:
"""Classification of one workflow session row."""
session_id: str
classification: str
reason: str
retireable: bool
status: str | None
pid: int | None
pid_alive: bool | None
pid_reused: bool
heartbeat_stale: bool
has_live_lease: bool
client_managed: bool
details: dict[str, Any] = field(default_factory=dict)
def as_dict(self) -> dict[str, Any]:
return {
"session_id": self.session_id,
"classification": self.classification,
"reason": self.reason,
"retireable": self.retireable,
"status": self.status,
"pid": self.pid,
"pid_alive": self.pid_alive,
"pid_reused": self.pid_reused,
"heartbeat_stale": self.heartbeat_stale,
"has_live_lease": self.has_live_lease,
"client_managed": self.client_managed,
"details": dict(self.details),
}
def classify_session(
row: Mapping[str, Any],
*,
now: datetime | None = None,
pid_checker: Callable[[int | None], bool] = lease_lifecycle.is_process_alive,
process_start_probe: Callable[[int | None], datetime | None] = process_start_time,
live_lease_sessions: set[str] | frozenset[str] | None = None,
client_managed_sessions: set[str] | frozenset[str] | None = None,
heartbeat_stale_seconds: int = DEFAULT_HEARTBEAT_STALE_SECONDS,
) -> SessionClassification:
"""Classify one session row for retirement decisions (#969).
Rules (first match wins where noted):
1. Non-active / already-terminal status ``terminal`` (not retireable).
2. Session holds a live lease ``protected`` (never retire).
3. PID missing and no live lease ``stale`` (retireable: missing owner).
4. PID alive + process start after session start ``stale`` (PID reuse).
5. PID alive + client-managed ``live`` protected (never retire).
6. PID alive + fresh heartbeat ``live``.
7. PID alive + stale heartbeat ``disconnected`` (not retireable alone).
8. PID dead ``stale`` (retireable).
"""
moment = now or _utc_now()
session_id = str(row.get("session_id") or "").strip()
status = str(row.get("status") or "").strip().lower() or None
raw_pid = row.get("pid")
try:
pid = int(raw_pid) if raw_pid is not None else None
except (TypeError, ValueError):
pid = None
client_managed = bool(
row.get("client_managed")
or row.get("is_client_managed")
or (
client_managed_sessions is not None
and session_id in client_managed_sessions
)
)
has_live_lease = bool(
live_lease_sessions is not None and session_id in live_lease_sessions
)
hb = _parse_ts(row.get("last_heartbeat_at"))
heartbeat_stale = bool(
hb is not None
and (moment - hb).total_seconds() > max(0, int(heartbeat_stale_seconds))
)
if hb is None and status == SESSION_STATUS_ACTIVE:
# No heartbeat evidence: treat as stale for liveness bookkeeping only.
heartbeat_stale = True
started = _parse_ts(row.get("started_at"))
recorded_proc_start = _parse_ts(row.get("owner_process_started_at"))
if status in TERMINAL_SESSION_STATUSES:
return SessionClassification(
session_id=session_id,
classification=CLASS_TERMINAL,
reason=REASON_ALREADY_TERMINAL,
retireable=False,
status=status,
pid=pid,
pid_alive=None,
pid_reused=False,
heartbeat_stale=heartbeat_stale,
has_live_lease=has_live_lease,
client_managed=client_managed,
)
if has_live_lease:
return SessionClassification(
session_id=session_id,
classification=CLASS_PROTECTED,
reason=REASON_LIVE_LEASE,
retireable=False,
status=status,
pid=pid,
pid_alive=pid_checker(pid) if pid is not None else None,
pid_reused=False,
heartbeat_stale=heartbeat_stale,
has_live_lease=True,
client_managed=client_managed,
details={"blocker": "live_workflow_lease"},
)
if pid is None:
return SessionClassification(
session_id=session_id,
classification=CLASS_STALE,
reason=REASON_MISSING_PID,
retireable=True,
status=status,
pid=None,
pid_alive=False,
pid_reused=False,
heartbeat_stale=heartbeat_stale,
has_live_lease=False,
client_managed=client_managed,
)
pid_alive = bool(pid_checker(pid))
pid_reused = False
proc_start: datetime | None = None
if pid_alive:
proc_start = recorded_proc_start or process_start_probe(pid)
# PID reuse: live process started after the session row itself was
# created. Anchor on started_at only — last_heartbeat alone is not a
# safe bound (synthetic inventories and long-lived processes would
# false-positive against a live ``ps`` probe).
if proc_start is not None and started is not None:
if proc_start > (started + PID_REUSE_SKEW):
pid_reused = True
if pid_reused:
return SessionClassification(
session_id=session_id,
classification=CLASS_STALE,
reason=REASON_PID_REUSE,
retireable=True,
status=status,
pid=pid,
pid_alive=True,
pid_reused=True,
heartbeat_stale=heartbeat_stale,
has_live_lease=False,
client_managed=client_managed,
details={
"process_started_at": _ts(proc_start) if proc_start else None,
"session_started_at": row.get("started_at"),
"last_heartbeat_at": row.get("last_heartbeat_at"),
},
)
if pid_alive and client_managed:
return SessionClassification(
session_id=session_id,
classification=CLASS_LIVE,
reason=REASON_CLIENT_MANAGED_LIVE,
retireable=False,
status=status,
pid=pid,
pid_alive=True,
pid_reused=False,
heartbeat_stale=heartbeat_stale,
has_live_lease=False,
client_managed=True,
)
if pid_alive and not heartbeat_stale:
return SessionClassification(
session_id=session_id,
classification=CLASS_LIVE,
reason=REASON_LIVE_OWNER,
retireable=False,
status=status,
pid=pid,
pid_alive=True,
pid_reused=False,
heartbeat_stale=False,
has_live_lease=False,
client_managed=client_managed,
)
if pid_alive and heartbeat_stale:
# Process still exists but has not heartbeated — disconnected, not
# confirmed stale. Do not retire; operator/reconnect owns next step.
return SessionClassification(
session_id=session_id,
classification=CLASS_DISCONNECTED,
reason=REASON_LIVE_OWNER,
retireable=False,
status=status,
pid=pid,
pid_alive=True,
pid_reused=False,
heartbeat_stale=True,
has_live_lease=False,
client_managed=client_managed,
details={"note": "alive_pid_stale_heartbeat_not_retired"},
)
# PID dead (or checker said not alive).
reason = REASON_DEAD_OWNER
if heartbeat_stale:
reason = REASON_HEARTBEAT_STALE_DEAD
return SessionClassification(
session_id=session_id,
classification=CLASS_STALE,
reason=reason,
retireable=True,
status=status,
pid=pid,
pid_alive=False,
pid_reused=False,
heartbeat_stale=heartbeat_stale,
has_live_lease=False,
client_managed=client_managed,
)
@dataclass(frozen=True)
class SessionFleetReport:
"""Fleet-wide classification summary for reconcile + apply."""
classifications: tuple[SessionClassification, ...]
live_count: int
disconnected_count: int
stale_count: int
protected_count: int
terminal_count: int
retireable: tuple[SessionClassification, ...]
counts_by_reason: dict[str, int]
def as_dict(self) -> dict[str, Any]:
return {
"live_count": self.live_count,
"disconnected_count": self.disconnected_count,
"stale_count": self.stale_count,
"protected_count": self.protected_count,
"terminal_count": self.terminal_count,
"retireable_count": len(self.retireable),
"retireable_session_ids": [c.session_id for c in self.retireable],
"counts_by_reason": dict(self.counts_by_reason),
"classifications": [c.as_dict() for c in self.classifications],
"sessions_dimension_resolved": len(self.retireable) == 0,
}
def classify_sessions(
sessions: Sequence[Mapping[str, Any]] | None,
*,
leases: Sequence[Mapping[str, Any]] | None = None,
now: datetime | None = None,
pid_checker: Callable[[int | None], bool] = lease_lifecycle.is_process_alive,
process_start_probe: Callable[[int | None], datetime | None] = process_start_time,
client_managed_sessions: set[str] | frozenset[str] | None = None,
heartbeat_stale_seconds: int = DEFAULT_HEARTBEAT_STALE_SECONDS,
) -> SessionFleetReport:
"""Classify a fleet of session rows against live leases (#969)."""
moment = now or _utc_now()
live_leases = live_lease_session_ids(
leases, now=moment, pid_checker=pid_checker
)
results: list[SessionClassification] = []
for row in sessions or ():
if not isinstance(row, Mapping):
continue
results.append(
classify_session(
row,
now=moment,
pid_checker=pid_checker,
process_start_probe=process_start_probe,
live_lease_sessions=live_leases,
client_managed_sessions=client_managed_sessions,
heartbeat_stale_seconds=heartbeat_stale_seconds,
)
)
live_count = sum(1 for c in results if c.classification == CLASS_LIVE)
disconnected_count = sum(
1 for c in results if c.classification == CLASS_DISCONNECTED
)
stale_count = sum(1 for c in results if c.classification == CLASS_STALE)
protected_count = sum(1 for c in results if c.classification == CLASS_PROTECTED)
terminal_count = sum(1 for c in results if c.classification == CLASS_TERMINAL)
retireable = tuple(c for c in results if c.retireable)
by_reason: dict[str, int] = {}
for c in results:
by_reason[c.reason] = by_reason.get(c.reason, 0) + 1
return SessionFleetReport(
classifications=tuple(results),
live_count=live_count,
disconnected_count=disconnected_count,
stale_count=stale_count,
protected_count=protected_count,
terminal_count=terminal_count,
retireable=retireable,
counts_by_reason=by_reason,
)
@dataclass(frozen=True)
class RetirementResult:
"""Outcome of one session retirement attempt."""
session_id: str
outcome: str # retired | already_terminal | skipped | blocked | missing
reason: str
prior_status: str | None = None
new_status: str | None = None
details: dict[str, Any] = field(default_factory=dict)
def as_dict(self) -> dict[str, Any]:
return {
"session_id": self.session_id,
"outcome": self.outcome,
"reason": self.reason,
"prior_status": self.prior_status,
"new_status": self.new_status,
"details": dict(self.details),
}
def apply_session_retirements(
db: cpd.ControlPlaneDB,
report: SessionFleetReport,
*,
dry_run: bool = False,
actor_session_id: str | None = None,
now: datetime | None = None,
) -> dict[str, Any]:
"""Terminalize every retireable session in *report* (idempotent).
Concurrent reconciles are safe: each retirement is a CAS on
``status='active'`` (or other non-terminal). A second pass that sees the
same session already retired records ``already_terminal`` rather than
duplicating audit noise beyond a single no-op outcome.
"""
moment = now or _utc_now()
results: list[RetirementResult] = []
retired = 0
already = 0
blocked = 0
missing = 0
for classification in report.retireable:
if not classification.retireable:
continue
if dry_run:
results.append(
RetirementResult(
session_id=classification.session_id,
outcome="skipped",
reason=classification.reason,
prior_status=classification.status,
new_status=SESSION_STATUS_RETIRED,
details={"dry_run": True, **classification.details},
)
)
continue
applied = db.retire_session(
session_id=classification.session_id,
reason=classification.reason,
actor_session_id=actor_session_id,
details={
"classification": classification.classification,
"pid": classification.pid,
"pid_alive": classification.pid_alive,
"pid_reused": classification.pid_reused,
"client_managed": classification.client_managed,
**classification.details,
},
now=moment,
)
outcome = str(applied.get("outcome") or "missing")
results.append(
RetirementResult(
session_id=classification.session_id,
outcome=outcome,
reason=str(applied.get("reason") or classification.reason),
prior_status=applied.get("prior_status"),
new_status=applied.get("new_status"),
details=dict(applied.get("details") or {}),
)
)
if outcome == "retired":
retired += 1
elif outcome == "already_terminal":
already += 1
elif outcome == "blocked":
blocked += 1
else:
missing += 1
# Non-retireable classifications are recorded for audit completeness when
# dry-run lists the fleet, but apply only mutates retireable rows.
return {
"success": True,
"dry_run": dry_run,
"retired_count": retired,
"already_terminal_count": already,
"blocked_count": blocked,
"missing_count": missing,
"planned_count": len(report.retireable),
"results": [r.as_dict() for r in results],
"audit_action": EVENT_SESSION_RETIRED,
"actor_session_id": actor_session_id,
"recorded_at": _ts(moment),
}
def retire_stale_sessions(
db: cpd.ControlPlaneDB,
*,
sessions: Sequence[Mapping[str, Any]] | None = None,
leases: Sequence[Mapping[str, Any]] | None = None,
dry_run: bool = False,
actor_session_id: str | None = None,
now: datetime | None = None,
pid_checker: Callable[[int | None], bool] = lease_lifecycle.is_process_alive,
process_start_probe: Callable[[int | None], datetime | None] = process_start_time,
client_managed_sessions: set[str] | frozenset[str] | None = None,
heartbeat_stale_seconds: int = DEFAULT_HEARTBEAT_STALE_SECONDS,
session_limit: int = 500,
) -> dict[str, Any]:
"""End-to-end plan + apply for stale session retirement (#969).
When *sessions* is omitted the control-plane DB is inventoried (active
rows only). Callers that already gathered inventory should pass it.
"""
moment = now or _utc_now()
if sessions is None:
sessions = db.list_sessions(
statuses=(SESSION_STATUS_ACTIVE,),
limit=max(1, int(session_limit)),
)
if leases is None:
try:
leases = db.list_leases(statuses=("active",), limit=max(1, int(session_limit)))
except Exception: # noqa: BLE001
leases = []
report = classify_sessions(
sessions,
leases=leases,
now=moment,
pid_checker=pid_checker,
process_start_probe=process_start_probe,
client_managed_sessions=client_managed_sessions,
heartbeat_stale_seconds=heartbeat_stale_seconds,
)
apply_result = apply_session_retirements(
db,
report,
dry_run=dry_run,
actor_session_id=actor_session_id,
now=moment,
)
return {
"success": True,
"fleet": report.as_dict(),
"apply": apply_result,
"sessions_dimension_resolved": (
report.as_dict()["sessions_dimension_resolved"]
if dry_run
else apply_result["planned_count"]
== (
apply_result["retired_count"]
+ apply_result["already_terminal_count"]
)
and apply_result["blocked_count"] == 0
),
}
+22
View File
@@ -204,6 +204,28 @@ proposed command before running it; `gitea_audit_runtime_recovery_contamination`
to inspect or (reconciler-only) clear the marker. Full contrast in
`docs/mcp-namespace-eof-recovery.md`.
## Connected is not attached (#708)
A host/CLI MCP inventory showing **Connected** is not proof the tools are usable.
The active session can expose **none** of a Connected server's tool namespaces —
a *session attachment* failure, distinct from config drift (#672),
transport-closed (#584), and resolver EOF (#685).
Required preflight proof is **live tool visibility plus `gitea_whoami` on the role
namespace**, never host Connected status alone. Call
`gitea_assess_mcp_namespace_attachment` with the Connected set and the namespaces
actually attached to the session; it returns the typed condition
**mcp_connected_namespaces_missing** with per-namespace connected-vs-attached proof,
and records the verdict so review and merge fail closed while a required namespace
is unattached.
Only sanctioned recovery: the client attach/reconnect path
(`gitea_request_mcp_reconnect`), then full preflight
(`gitea_whoami``gitea_resolve_task_capability` → task). Never recover by direct
module import, CLI/raw API mutation, profile hopping, session-state overrides,
process kills, or `.env`/mtime edits. A final report must not claim a healthy
session without attachment proof.
## Shell Spawn Hard-Stop Rule
`exit_code: -1` with empty stdout/stderr means the shell failed to spawn — not a
+12 -2
View File
@@ -153,8 +153,9 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.branch.push",
"role": "author",
},
# #662: post-restart reconcile is read-only inventory + pure classification.
# Durable follow-up issue creation is a separate apply path (not this task).
# #662: post-restart reconcile is inventory + pure classification.
# #969: optional apply_session_cleanup retires confirmed-stale session rows
# through the same tool; durable follow-up Gitea issues remain separate.
"reconcile_after_restart": {
"permission": "gitea.read",
"role": "author",
@@ -163,6 +164,15 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.read",
"role": "author",
},
# #969: explicit plan/apply path for dead-owner session retirement.
"retire_stale_workflow_sessions": {
"permission": "gitea.read",
"role": "author",
},
"gitea_retire_stale_workflow_sessions": {
"permission": "gitea.read",
"role": "author",
},
# #644: Phase 2 Web Console recovery tasks.
"clear_stale_binding": {
"permission": "gitea.read",
+2 -2
View File
@@ -37,7 +37,7 @@ class ControlPlaneDBTest(unittest.TestCase):
rows = dict(conn.execute("SELECT key, value FROM schema_meta").fetchall())
finally:
conn.close()
self.assertEqual(rows["schema_version"], "5")
self.assertEqual(rows["schema_version"], "6")
self.assertIn("DB coordinates", rows["architecture"])
self.assertIn("bridge", rows["architecture"].lower())
@@ -868,7 +868,7 @@ class SessionCheckpointTest(unittest.TestCase):
conn.close()
self.assertIn("session_checkpoints", names)
record = self._write()
self.assertEqual(record["checkpoint_schema_version"], 5)
self.assertEqual(record["checkpoint_schema_version"], 6)
# AC2 — checkpoints written for multi-role session fixtures.
def test_multi_role_fixtures_each_get_a_row(self) -> None:
+17 -3
View File
@@ -112,7 +112,19 @@ class TestIssue686ManualMcpProvenance(unittest.TestCase):
self.assertTrue(any("All matching profiles for task 'create_issue' (['prgs-author']) are running but stale" in r for r in reasons))
def test_namespace_health_classification_includes_provenance(self):
"""AC 1 & 4: mcp_namespace_health diagnostics include provenance and unconsumed_gitea_env."""
"""AC 1 & 4: mcp_namespace_health diagnostics include provenance and unconsumed_gitea_env.
#948 narrowed the vocabulary here. This process carries no client-managed
declaration, so the old code labelled it ``manual_launch`` asserting a
hand-launched terminal process it had no evidence for, and contradicting
``gitea_get_runtime_context``, which read the same process and reported
``client_managed``. Absence of proof is now reported as ``unproven``.
The #686 wall itself is unchanged and still asserted below:
``is_client_managed`` stays False, so nothing previously refused is now
permitted. Only the label on the *reason* changed, so remediation names
the proof that is actually missing.
"""
process = {
"pid": 5555,
"profile": "prgs-author",
@@ -129,10 +141,12 @@ class TestIssue686ManualMcpProvenance(unittest.TestCase):
process=process,
probe_source="client_namespace",
)
self.assertEqual(res["provenance"], "manual_launch")
self.assertEqual(res["provenance"], "unproven")
self.assertFalse(res["is_client_managed"])
# The wall is intact: no client-managed proof still fails closed.
self.assertTrue(res["provenance_fail_closed"])
self.assertEqual(res["unconsumed_gitea_env"], {"GITEA_DUMMY": "99"})
self.assertEqual(res["diagnostics"]["provenance"], "manual_launch")
self.assertEqual(res["diagnostics"]["provenance"], "unproven")
if __name__ == "__main__":
+269
View File
@@ -0,0 +1,269 @@
"""Regression tests for Issue #708 wiring: detection must reach a gate, not just exist.
The prior #708 slice added a pure decision function with no call site, so a session
whose namespaces were Connected-but-unattached still passed every mutation gate.
These tests pin the parts that make the detection load-bearing:
* the typed condition is distinct from #672 / #584 / #685,
* the session store records a per-namespace attachment verdict,
* review/merge mutations fail closed while a required namespace is unattached,
* recovery is reconnect-only and never suggests an unsafe fallback,
* startup ordering races and discovery-cache telemetry are reported,
* a healthy final report cannot be produced without attachment proof.
"""
import mcp_namespace_health
REQUIRED = ["gitea-author", "gitea-reviewer", "gitea-merger", "gitea-tools"]
def _connected_but_unattached():
"""The #708 signature: host says Connected, session tool surface is empty."""
return mcp_namespace_health.assess_connected_namespace_attachment(
connected_servers=REQUIRED,
attached_session_namespaces=[],
required_namespaces=REQUIRED,
)
# --- AC1: distinct typed detection -----------------------------------------
def test_connected_with_empty_session_tool_list_is_typed_distinctly():
res = _connected_but_unattached()
assert res["attachment_healthy"] is False
assert res["discovery_status"] == "connected_but_namespaces_missing"
assert res["error_type"] == "mcp_connected_namespaces_missing"
assert sorted(res["missing_namespaces"]) == sorted(REQUIRED)
# Not misclassified as config drift (#672), transport-closed (#584), or EOF (#685).
assert res["error_type"] not in {
"mcp_config_drift",
"transport_closed",
"mcp_client_eof",
}
def test_proof_carries_connected_and_attached_per_namespace():
res = _connected_but_unattached()
proof = res["proof_of_connected_vs_attached"]
for ns in REQUIRED:
assert proof[ns] == {"connected": True, "attached": False}
# --- AC2: recovery is auto-attach or reconnect-only -------------------------
def test_exact_next_action_is_reconnect_only():
res = _connected_but_unattached()
assert res["reconnect_required"] is True
action = res["exact_next_action"]
assert "Reconnect the IDE/client MCP session" in action
for forbidden in ("pkill", "chmod", "curl", ".env", "sys.path"):
assert forbidden not in action
def test_sanctioned_recovery_tool_is_named():
res = _connected_but_unattached()
assert res["sanctioned_recovery_tool"] == "gitea_request_mcp_reconnect"
assert "gitea_request_mcp_reconnect" in " ".join(res["remediation"])
def test_successful_auto_attach_reports_recovered_without_operator():
res = mcp_namespace_health.assess_connected_namespace_attachment(
connected_servers=REQUIRED,
attached_session_namespaces=REQUIRED,
required_namespaces=REQUIRED,
auto_attach_attempted=True,
auto_attach_succeeded=True,
)
assert res["attachment_healthy"] is True
assert res["auto_recovered"] is True
assert res["reconnect_required"] is False
def test_failed_auto_attach_still_requires_reconnect():
res = mcp_namespace_health.assess_connected_namespace_attachment(
connected_servers=REQUIRED,
attached_session_namespaces=[],
required_namespaces=REQUIRED,
auto_attach_attempted=True,
auto_attach_succeeded=False,
)
assert res["auto_recovered"] is False
assert res["reconnect_required"] is True
assert any("did not succeed" in r for r in res["reasons"])
def test_reconnect_rediscovery_clears_the_condition():
"""Attach state after a reconnect is healthy without any other change."""
before = _connected_but_unattached()
after = mcp_namespace_health.assess_connected_namespace_attachment(
connected_servers=REQUIRED,
attached_session_namespaces=REQUIRED,
required_namespaces=REQUIRED,
discovery_cache_hit=False,
discovery_cache_age_seconds=0.0,
)
assert before["attachment_healthy"] is False
assert after["attachment_healthy"] is True
assert after["discovery_status"] == "namespaces_attached"
assert after["missing_namespaces"] == []
# --- AC4: startup ordering across multiple role servers ---------------------
def test_multi_role_startup_ordering_race_is_reported():
res = mcp_namespace_health.assess_connected_namespace_attachment(
connected_servers=REQUIRED,
attached_session_namespaces=["gitea-author"],
required_namespaces=REQUIRED,
session_tool_snapshot_at=1000.0,
namespace_connected_at={
"gitea-author": 990.0,
"gitea-reviewer": 1005.0,
"gitea-merger": 1007.0,
},
)
assert res["startup_ordering_race"] is True
assert res["late_attaching_namespaces"] == ["gitea-merger", "gitea-reviewer"]
assert res["telemetry"]["startup_ordering_race"] is True
def test_no_ordering_race_when_snapshot_follows_connect():
res = mcp_namespace_health.assess_connected_namespace_attachment(
connected_servers=REQUIRED,
attached_session_namespaces=REQUIRED,
required_namespaces=REQUIRED,
session_tool_snapshot_at=2000.0,
namespace_connected_at={ns: 1000.0 for ns in REQUIRED},
)
assert res["startup_ordering_race"] is False
assert res["late_attaching_namespaces"] == []
# --- AC5: telemetry ---------------------------------------------------------
def test_telemetry_reports_cache_and_recovery_signals():
res = mcp_namespace_health.assess_connected_namespace_attachment(
connected_servers=REQUIRED,
attached_session_namespaces=["gitea-author"],
required_namespaces=REQUIRED,
discovery_cache_hit=True,
discovery_cache_age_seconds=42.5,
)
tel = res["telemetry"]
assert tel["connected_count"] == 4
assert tel["attached_count"] == 1
assert tel["missing_count"] == 3
assert tel["discovery_cache_hit"] is True
assert tel["discovery_cache_age_seconds"] == 42.5
assert tel["reconnect_required"] is True
assert tel["error_type"] == "mcp_connected_namespaces_missing"
def test_telemetry_leaks_no_secrets():
res = mcp_namespace_health.assess_connected_namespace_attachment(
connected_servers=REQUIRED,
attached_session_namespaces=[],
required_namespaces=REQUIRED,
)
blob = repr(res["telemetry"]).lower()
for leak in ("token", "authorization", "password", "secret", "/users/", "http"):
assert leak not in blob
# --- AC2/AC6: fail-closed mutation gate -------------------------------------
def _session_store_from(assessment):
"""Mirror the server-side recorder without importing the MCP server module."""
store = {}
missing = set(assessment["missing_namespaces"])
for ns, state in assessment["proof_of_connected_vs_attached"].items():
attached = bool(state["attached"])
store[ns] = {
"namespace": ns,
"connected": bool(state["connected"]),
"attached": attached,
"attachment_healthy": attached and ns not in missing,
"error_type": None if attached else assessment["error_type"],
}
return store
def test_review_and_merge_fail_closed_while_unattached():
store = _session_store_from(_connected_but_unattached())
for task in ("review_pr", "submit_review", "merge_pr", "create_pr", "work_issue"):
reasons = mcp_namespace_health.attachment_gate_from_session(task, store)
assert reasons, f"{task} must fail closed while its namespace is unattached"
assert "fail closed, #708" in reasons[0]
def test_gate_offers_only_sanctioned_recovery():
store = _session_store_from(_connected_but_unattached())
reasons = mcp_namespace_health.attachment_gate_from_session("merge_pr", store)
joined = " ".join(reasons)
assert "reconnect" in joined.lower()
assert "Workflow Safety Hard Stop (#708)" in joined
# Unsafe fallbacks appear only inside the prohibition, never as advice.
assert "NEVER use" in joined
def test_gate_passes_once_namespaces_are_attached():
healthy = mcp_namespace_health.assess_connected_namespace_attachment(
connected_servers=REQUIRED,
attached_session_namespaces=REQUIRED,
required_namespaces=REQUIRED,
)
store = _session_store_from(healthy)
for task in ("review_pr", "submit_review", "merge_pr", "create_pr", "work_issue"):
assert mcp_namespace_health.attachment_gate_from_session(task, store) == []
def test_unassessed_session_does_not_gate():
"""No recorded assessment must not fabricate a block (matches #543 semantics)."""
assert mcp_namespace_health.attachment_gate_from_session("merge_pr", {}) == []
assert mcp_namespace_health.attachment_gate_from_session("merge_pr", None) == []
def test_unmapped_task_is_not_gated():
store = _session_store_from(_connected_but_unattached())
assert mcp_namespace_health.attachment_gate_from_session("gitea_read", store) == []
def test_partial_attachment_gates_only_the_affected_role():
"""Author attached, reviewer not: author work proceeds, review fails closed."""
res = mcp_namespace_health.assess_connected_namespace_attachment(
connected_servers=REQUIRED,
attached_session_namespaces=["gitea-author", "gitea-tools"],
required_namespaces=REQUIRED,
)
store = _session_store_from(res)
assert mcp_namespace_health.attachment_gate_from_session("work_issue", store) == []
assert mcp_namespace_health.attachment_gate_from_session("review_pr", store)
assert mcp_namespace_health.attachment_gate_from_session("merge_pr", store)
# --- AC4: no false healthy report without attachment proof ------------------
def test_no_false_healthy_without_attachment_proof():
"""Connected alone never yields a healthy verdict."""
res = mcp_namespace_health.assess_connected_namespace_attachment(
connected_servers=REQUIRED,
attached_session_namespaces=None,
required_namespaces=REQUIRED,
)
assert res["success"] is False
assert res["attachment_healthy"] is False
assert res["telemetry"]["attached_count"] == 0
def test_attachment_gate_maps_each_role_namespace():
assert mcp_namespace_health.required_namespace_for_attachment("review_pr") == "gitea-reviewer"
assert mcp_namespace_health.required_namespace_for_attachment("merge_pr") == "gitea-merger"
assert mcp_namespace_health.required_namespace_for_attachment("create_pr") == "gitea-author"
assert mcp_namespace_health.required_namespace_for_attachment("nope") is None
@@ -0,0 +1,77 @@
"""Unit regression tests for Issue #708: Connected-but-namespaces-missing detection and attachment safety."""
import mcp_namespace_health
def test_assess_connected_namespace_attachment_success():
# required_namespaces is declared explicitly: these three are the namespaces this case
# is about. Relying on the default (which also requires gitea-tools) would ask for a
# healthy verdict covering a required namespace that was never connected — exactly the
# false-healthy classification these tests now forbid.
connected = ["gitea-author", "gitea-reviewer", "gitea-merger"]
attached = ["gitea-author", "gitea-reviewer", "gitea-merger"]
res = mcp_namespace_health.assess_connected_namespace_attachment(
connected_servers=connected,
attached_session_namespaces=attached,
required_namespaces=connected,
)
assert res["success"] is True
assert res["attachment_healthy"] is True
assert res["discovery_status"] == "namespaces_attached"
assert res["error_type"] is None
assert res["missing_namespaces"] == []
assert res["exact_next_action"] == "None; session tool namespaces attached."
def test_assess_connected_namespace_attachment_missing():
# Every required namespace here *is* connected, so the only condition present is the
# #708 one: Connected at the host, absent from the session tool surface.
connected = ["gitea-author", "gitea-reviewer", "gitea-merger"]
attached = ["gitea-author"]
res = mcp_namespace_health.assess_connected_namespace_attachment(
connected_servers=connected,
attached_session_namespaces=attached,
required_namespaces=connected,
)
assert res["success"] is False
assert res["attachment_healthy"] is False
assert res["discovery_status"] == "connected_but_namespaces_missing"
assert res["error_type"] == "mcp_connected_namespaces_missing"
assert "gitea-reviewer" in res["missing_namespaces"]
assert "gitea-merger" in res["missing_namespaces"]
assert "Reconnect the IDE/client MCP session" in res["exact_next_action"]
def test_assess_connected_namespace_attachment_disconnected():
res = mcp_namespace_health.assess_connected_namespace_attachment(
connected_servers=[],
attached_session_namespaces=[],
)
assert res["success"] is False
assert res["attachment_healthy"] is False
assert res["discovery_status"] == "disconnected"
def test_proof_of_connected_vs_attached_mapping():
connected = ["gitea-author", "gitea-reviewer"]
attached = ["gitea-author"]
res = mcp_namespace_health.assess_connected_namespace_attachment(
connected_servers=connected,
attached_session_namespaces=attached,
)
proof = res["proof_of_connected_vs_attached"]
assert proof["gitea-author"] == {"connected": True, "attached": True}
assert proof["gitea-reviewer"] == {"connected": True, "attached": False}
def test_unsafe_fallback_policy_enforcement():
res = mcp_namespace_health.assess_connected_namespace_attachment(
connected_servers=["gitea-author"],
attached_session_namespaces=[],
)
policy = res["unsafe_fallback_policy"]
assert "Workflow Safety Hard Stop (#708)" in policy
assert "direct imports" in policy
assert "API mutations" in policy
assert "profile hopping" in policy
assert "session-state overrides" in policy
@@ -0,0 +1,283 @@
"""Regression tests for Issue #708 B2: not-connected is not Connected-but-unattached.
The first #708 slice counted a namespace as ``missing`` only when it appeared in
``connected_servers``. A *required* namespace absent from that inventory was therefore
never counted at all, so the session reported ``attachment_healthy: true`` with
``error_type: None`` while holding no attachment proof for it and the gate text told the
operator "the host reports Connected" about a service nothing had reported Connected.
These tests pin the corrected distinction:
* required and not connected never yields a healthy verdict,
* it is typed ``mcp_required_namespaces_not_connected``, never the #708 condition,
* ``mcp_connected_namespaces_missing`` stays reserved for genuinely Connected namespaces,
* both categories still fail review and merge closed, with their own reason,
* per-namespace ``connected``/``attached`` evidence is reported accurately,
* one session's evidence cannot clear another session's block.
"""
import mcp_namespace_health
ROLES = ["gitea-author", "gitea-reviewer", "gitea-merger", "gitea-tools"]
NOT_CONNECTED = mcp_namespace_health.ERROR_REQUIRED_NAMESPACES_NOT_CONNECTED
CONNECTED_MISSING = mcp_namespace_health.ERROR_CONNECTED_NAMESPACES_MISSING
def _assess(connected, attached, required):
return mcp_namespace_health.assess_connected_namespace_attachment(
connected_servers=connected,
attached_session_namespaces=attached,
required_namespaces=required,
)
def _store(assessment):
"""The recorder contract: per-namespace verdicts feed the gate."""
return dict(assessment["namespace_conditions"])
# --- the four evidence combinations ----------------------------------------
def test_connected_and_attached_is_healthy():
res = _assess(ROLES, ROLES, ROLES)
assert res["attachment_healthy"] is True
assert res["error_type"] is None
assert res["missing_namespaces"] == []
assert res["not_connected_namespaces"] == []
assert res["discovery_status"] == "namespaces_attached"
def test_connected_but_unattached_keeps_the_708_condition():
res = _assess(ROLES, ["gitea-author"], ROLES)
assert res["attachment_healthy"] is False
assert res["error_type"] == CONNECTED_MISSING
assert res["not_connected_namespaces"] == []
assert sorted(res["missing_namespaces"]) == [
"gitea-merger",
"gitea-reviewer",
"gitea-tools",
]
assert res["discovery_status"] == "connected_but_namespaces_missing"
def test_required_but_not_connected_is_never_healthy():
"""The reviewer's exact reproduction from review 637."""
res = _assess(
["gitea-reviewer"], ["gitea-reviewer"], ["gitea-reviewer", "gitea-merger"]
)
assert res["attachment_healthy"] is False
assert res["success"] is False
assert res["error_type"] is not None
assert res["telemetry"]["not_connected_count"] == 1
def test_required_but_not_connected_is_typed_distinctly():
res = _assess(
["gitea-reviewer"], ["gitea-reviewer"], ["gitea-reviewer", "gitea-merger"]
)
assert res["error_type"] == NOT_CONNECTED
assert res["discovery_status"] == "required_namespaces_not_connected"
assert res["not_connected_namespaces"] == ["gitea-merger"]
# Not collapsed into the #708 condition, nor into config drift (#672) or
# transport-closed (#584).
assert CONNECTED_MISSING not in res["error_types"]
assert res["missing_namespaces"] == []
assert res["error_type"] not in {"mcp_config_drift", "transport_closed"}
def test_not_connected_reason_makes_no_connected_claim():
res = _assess(
["gitea-reviewer"], ["gitea-reviewer"], ["gitea-reviewer", "gitea-merger"]
)
about_merger = [r for r in res["reasons"] if "gitea-merger" in r]
assert about_merger, "the not-connected namespace must be named in the reasons"
for reason in about_merger:
assert "report Connected at host/CLI layer" not in reason
assert "gitea-merger" in res["exact_next_action"]
def test_neither_connected_nor_attached_reports_disconnected():
res = _assess([], [], ROLES)
assert res["attachment_healthy"] is False
assert res["discovery_status"] == "disconnected"
assert res["error_type"] == NOT_CONNECTED
assert sorted(res["not_connected_namespaces"]) == sorted(ROLES)
assert res["missing_namespaces"] == []
# --- mixed required namespaces ---------------------------------------------
def test_mixed_connected_attached_and_not_connected():
res = _assess(
["gitea-author"], ["gitea-author"], ["gitea-author", "gitea-merger"]
)
assert res["attachment_healthy"] is False
assert res["not_connected_namespaces"] == ["gitea-merger"]
assert res["missing_namespaces"] == []
conditions = res["namespace_conditions"]
assert conditions["gitea-author"]["attachment_healthy"] is True
assert conditions["gitea-author"]["condition"] is None
assert conditions["gitea-merger"]["attachment_healthy"] is False
assert conditions["gitea-merger"]["condition"] == NOT_CONNECTED
def test_both_conditions_present_are_both_reported():
"""One namespace Connected-but-unattached, another never connected."""
res = _assess(
["gitea-author", "gitea-reviewer"],
["gitea-author"],
["gitea-author", "gitea-reviewer", "gitea-merger"],
)
assert res["missing_namespaces"] == ["gitea-reviewer"]
assert res["not_connected_namespaces"] == ["gitea-merger"]
# Neither condition is hidden by the other; error_type names the primary one.
assert sorted(res["error_types"]) == sorted([NOT_CONNECTED, CONNECTED_MISSING])
assert res["error_type"] == NOT_CONNECTED
# --- unknown, partial, malformed, contradictory evidence --------------------
def test_unknown_attachment_evidence_is_not_healthy():
"""No session tool surface reported at all is unproven, not proven good."""
res = _assess(ROLES, None, ROLES)
assert res["attachment_healthy"] is False
assert res["telemetry"]["attached_count"] == 0
assert res["error_type"] == CONNECTED_MISSING
def test_partial_evidence_gates_only_the_unproven_roles():
res = _assess(ROLES, ["gitea-author", "gitea-tools"], ROLES)
store = _store(res)
assert mcp_namespace_health.attachment_gate_from_session("work_issue", store) == []
assert mcp_namespace_health.attachment_gate_from_session("review_pr", store)
assert mcp_namespace_health.attachment_gate_from_session("merge_pr", store)
def test_malformed_namespace_entries_are_discarded_not_trusted():
"""Blank and whitespace-only names must not become namespaces or proof."""
res = mcp_namespace_health.assess_connected_namespace_attachment(
connected_servers=["gitea-author", "", " "],
attached_session_namespaces=["gitea-author", ""],
required_namespaces=["gitea-author", " "],
)
assert "" not in res["proof_of_connected_vs_attached"]
assert " " not in res["proof_of_connected_vs_attached"]
assert res["attachment_healthy"] is True
assert res["telemetry"]["required_count"] == 1
def test_contradictory_attached_without_connected_fails_closed():
"""Attached in the session yet absent from the connected inventory."""
res = _assess(["gitea-author"], ["gitea-author", "gitea-merger"], ROLES)
assert res["attachment_healthy"] is False
assert "gitea-merger" in res["contradictory_namespaces"]
assert "gitea-merger" in res["not_connected_namespaces"]
assert res["namespace_conditions"]["gitea-merger"]["contradictory_evidence"] is True
assert res["namespace_conditions"]["gitea-merger"]["attachment_healthy"] is False
assert any("contradictory evidence" in r for r in res["reasons"])
assert res["telemetry"]["contradictory_count"] >= 1
def test_duplicate_required_entries_are_counted_once():
res = _assess(["gitea-author"], ["gitea-author"], ["gitea-author", "gitea-author"])
assert res["telemetry"]["required_count"] == 1
assert res["attachment_healthy"] is True
# --- review and merge behaviour for both failure categories -----------------
def test_merge_fails_closed_when_required_namespace_never_connected():
store = _store(_assess(["gitea-author"], ["gitea-author"], ROLES))
reasons = mcp_namespace_health.attachment_gate_from_session("merge_pr", store)
assert reasons
assert NOT_CONNECTED in reasons[0]
assert "fail closed, #708" in reasons[0]
def test_review_fails_closed_when_required_namespace_never_connected():
store = _store(_assess(["gitea-author"], ["gitea-author"], ROLES))
reasons = mcp_namespace_health.attachment_gate_from_session("review_pr", store)
assert reasons
assert NOT_CONNECTED in reasons[0]
assert "fail closed, #708" in reasons[0]
def test_not_connected_block_never_claims_the_host_reports_connected():
store = _store(_assess(["gitea-author"], ["gitea-author"], ROLES))
reasons = mcp_namespace_health.attachment_gate_from_session("merge_pr", store)
assert "the host reports Connected" not in reasons[0]
assert "absent from the connected-service inventory" in reasons[0]
def test_connected_but_unattached_block_still_says_connected():
store = _store(_assess(ROLES, ["gitea-author"], ROLES))
reasons = mcp_namespace_health.attachment_gate_from_session("merge_pr", store)
assert CONNECTED_MISSING in reasons[0]
assert "the host reports Connected" in reasons[0]
def test_both_categories_offer_only_the_sanctioned_recovery():
for store in (
_store(_assess(ROLES, [], ROLES)),
_store(_assess(["gitea-author"], ["gitea-author"], ROLES)),
):
reasons = mcp_namespace_health.attachment_gate_from_session("merge_pr", store)
joined = " ".join(reasons)
assert "reconnect" in joined.lower()
assert "Workflow Safety Hard Stop (#708)" in joined
for forbidden in ("pkill", "chmod", "curl ", ".env", "sys.path"):
assert forbidden not in reasons[0]
def test_entry_without_connected_evidence_blocks_without_asserting_either():
"""A legacy/partial store entry must fail closed and claim nothing it cannot prove."""
store = {"gitea-merger": {"namespace": "gitea-merger", "attached": False}}
reasons = mcp_namespace_health.attachment_gate_from_session("merge_pr", store)
assert reasons
assert "no connected-status evidence" in reasons[0]
assert "the host reports Connected" not in reasons[0]
assert "fail closed, #708" in reasons[0]
# --- session isolation ------------------------------------------------------
def test_one_session_evidence_does_not_clear_another_session_block():
"""Attachment evidence is per-session state; it must not travel between sessions."""
blocked_session = _store(_assess(["gitea-author"], ["gitea-author"], ROLES))
healthy_session = _store(_assess(ROLES, ROLES, ROLES))
assert (
mcp_namespace_health.attachment_gate_from_session("merge_pr", healthy_session)
== []
)
# The healthy session's verdict is not consulted for the blocked session.
assert mcp_namespace_health.attachment_gate_from_session("merge_pr", blocked_session)
# And the blocked session's store is unchanged by the healthy one existing.
assert blocked_session["gitea-merger"]["attachment_healthy"] is False
def test_gate_reads_only_the_store_it_is_given():
healthy_session = _store(_assess(ROLES, ROLES, ROLES))
assert mcp_namespace_health.attachment_gate_from_session("merge_pr", {}) == []
assert mcp_namespace_health.attachment_gate_from_session("merge_pr", None) == []
assert (
mcp_namespace_health.attachment_gate_from_session("merge_pr", healthy_session)
== []
)
# --- telemetry stays secret-free -------------------------------------------
def test_not_connected_telemetry_leaks_no_secrets():
res = _assess(["gitea-author"], ["gitea-author"], ROLES)
blob = repr(res["telemetry"]).lower()
for leak in ("token", "authorization", "password", "secret", "/users/", "http"):
assert leak not in blob
+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()
@@ -0,0 +1,662 @@
"""Client/session-aware runtime ownership and provenance (#948).
Covers the reproduced contradiction that motivated the issue: one surface
reporting ``client_managed`` while another reported ``manual_launch`` for the
same process, remediation hardcoded to one vendor, and a profile-wide duplicate
wall that could not tell two healthy clients apart.
All client and session identifiers here are synthetic.
"""
from __future__ import annotations
import os
import tempfile
import unittest
from datetime import datetime, timedelta, timezone
import mcp_client_reconnect
import mcp_namespace_health
import mcp_worker_identity as mwi
NOW = datetime(2026, 7, 29, 6, 0, 0, tzinfo=timezone.utc)
def _registry() -> mwi.WorkerRegistry:
"""A registry on a throwaway path; never the operator's real one."""
handle, path = tempfile.mkstemp(suffix=".sqlite3")
os.close(handle)
os.unlink(path)
return mwi.WorkerRegistry(path)
def _attach(
registry: mwi.WorkerRegistry,
*,
client: str,
session: str,
generation: str,
profile: str = "prgs-reviewer",
role: str = "reviewer",
pid: int = 4242,
now: datetime = NOW,
ttl: float = 900.0,
) -> dict:
"""Register one synthetic worker and return the outcome."""
identity = mwi.generate_worker_identity(client, session, now=now)
outcome = registry.register(
worker_identity=identity,
client_name=client,
client_instance_id=f"inst-{session}",
session_id=session,
generation_id=generation,
role=role,
profile=profile,
pid=pid,
heartbeat_ttl_seconds=ttl,
now=now,
)
outcome["identity"] = identity
return outcome
class IdentityFormatTests(unittest.TestCase):
"""AC27-29: collision-resistant `<llm-name>-<UTC-timestamp>-<short-sha>`."""
def test_identity_matches_required_format(self):
identity = mwi.generate_worker_identity("Gemini", "sess-0001", now=NOW)
parsed = mwi.parse_worker_identity(identity)
self.assertTrue(parsed["valid"], parsed["reasons"])
self.assertEqual(parsed["client_name"], "gemini")
self.assertEqual(parsed["minted_at"], "20260729T060000Z")
self.assertEqual(len(parsed["digest"]), 12)
def test_digest_varies_with_session_and_nonce(self):
base = dict(timestamp_ns=1, now=NOW)
a = mwi.generate_worker_identity("codex", "sess-A", nonce="n", **base)
b = mwi.generate_worker_identity("codex", "sess-B", nonce="n", **base)
c = mwi.generate_worker_identity("codex", "sess-A", nonce="m", **base)
self.assertNotEqual(a, b, "session must feed the digest")
self.assertNotEqual(a, c, "nonce must feed the digest")
def test_identity_is_not_role_or_profile(self):
"""AC26: identity is independent of role and profile."""
args = dict(timestamp_ns=7, nonce="fixed", now=NOW)
same = mwi.generate_worker_identity("claude", "sess-1", **args)
self.assertEqual(same, mwi.generate_worker_identity("claude", "sess-1", **args))
# Nothing role- or profile-derived appears in the identity.
self.assertNotIn("reviewer", same)
self.assertNotIn("prgs", same)
def test_malformed_identity_rejected(self):
self.assertFalse(mwi.parse_worker_identity("prgs-reviewer")["valid"])
self.assertFalse(mwi.parse_worker_identity("")["valid"])
self.assertFalse(mwi.parse_worker_identity(None)["valid"])
class PerClientAttachmentTests(unittest.TestCase):
"""Every supported client attaches and is reported as itself."""
def _assert_attached_as(self, client: str, expected_name: str):
registry = _registry()
outcome = _attach(
registry, client=client, session=f"sess-{client}", generation="gen-1"
)
self.assertTrue(outcome["registered"], outcome["reasons"])
verdict = mwi.assess_provenance(
registry=registry, worker_identity=outcome["identity"], env={}, now=NOW
)
self.assertEqual(verdict["session_ownership"], mwi.OWNERSHIP_OWNED)
self.assertEqual(verdict["provenance"], mwi.PROVENANCE_CLIENT_SESSION)
self.assertEqual(verdict["client_name"], expected_name)
self.assertTrue(verdict["session_owned"])
self.assertFalse(verdict["fail_closed"])
return verdict
def test_codex_attachment(self):
self._assert_attached_as("codex", "codex")
def test_gemini_attachment(self):
self._assert_attached_as("gemini", "gemini")
def test_antigravity_attachment(self):
self._assert_attached_as("antigravity", "antigravity")
def test_claude_attachment(self):
self._assert_attached_as("claude", "claude_code")
def test_unknown_client_is_named_not_guessed(self):
verdict = self._assert_attached_as("some_new_llm", "some_new_llm")
self.assertNotEqual(verdict["client_name"], "codex")
class SessionLifecycleTests(unittest.TestCase):
def test_same_client_new_session_gets_distinct_identity(self):
registry = _registry()
first = _attach(registry, client="codex", session="sess-1", generation="gen-1")
second = _attach(registry, client="codex", session="sess-2", generation="gen-2")
self.assertTrue(first["registered"])
self.assertTrue(second["registered"])
self.assertNotEqual(first["identity"], second["identity"])
# Both are live and neither blocks the other.
cohort = mwi.classify_cohort(registry.list_workers(), now=NOW)
self.assertEqual(cohort["live_worker_count"], 2)
self.assertFalse(cohort["blocked"], cohort["reasons"])
def test_different_client_attaches_after_previous_session_ends(self):
"""AC14: expiry then takeover with a higher fencing epoch."""
registry = _registry()
gone = _attach(
registry, client="codex", session="sess-old", generation="gen-shared", ttl=60
)
later = NOW + timedelta(hours=1)
self.assertFalse(
registry.is_live(registry.get(gone["identity"]), now=later)["live"]
)
arriving = _attach(
registry,
client="gemini",
session="sess-new",
generation="gen-other",
now=later,
)
claim = registry.claim_generation(
worker_identity=arriving["identity"],
generation_id="gen-shared",
now=later,
)
self.assertTrue(claim["claimed"], claim["reasons"])
self.assertIn(gone["identity"], claim["superseded_workers"])
self.assertGreater(claim["fencing_epoch"], gone["fencing_epoch"])
def test_superseded_session_is_fenced_on_resume(self):
"""AC15/AC16: the prior session cannot heartbeat its way back."""
registry = _registry()
old = _attach(
registry, client="codex", session="sess-old", generation="gen-shared", ttl=60
)
later = NOW + timedelta(hours=1)
new = _attach(
registry, client="gemini", session="sess-new", generation="gen-x", now=later
)
registry.claim_generation(
worker_identity=new["identity"], generation_id="gen-shared", now=later
)
resumed = registry.heartbeat(
worker_identity=old["identity"],
fencing_epoch=old["fencing_epoch"],
now=later,
)
self.assertFalse(resumed["renewed"])
self.assertFalse(resumed["mutation_performed"])
self.assertEqual(resumed["blocker_kind"], mwi.BLOCKER_FENCED)
def test_heartbeat_renews_only_the_owning_lease(self):
"""AC11: a wrong epoch never renews, and never mutates."""
registry = _registry()
worker = _attach(registry, client="codex", session="s", generation="g")
good = registry.heartbeat(
worker_identity=worker["identity"],
fencing_epoch=worker["fencing_epoch"],
now=NOW + timedelta(minutes=5),
)
self.assertTrue(good["renewed"])
bad = registry.heartbeat(
worker_identity=worker["identity"],
fencing_epoch=worker["fencing_epoch"] + 99,
now=NOW + timedelta(minutes=6),
)
self.assertFalse(bad["renewed"])
self.assertFalse(bad["mutation_performed"])
self.assertEqual(
registry.get(worker["identity"])["last_heartbeat_at"],
good["last_heartbeat_at"],
"a refused heartbeat must not advance the record",
)
class ConflictAndCollisionTests(unittest.TestCase):
def test_two_live_sessions_cannot_claim_one_generation(self):
registry = _registry()
first = _attach(registry, client="codex", session="s1", generation="gen-shared")
second = _attach(registry, client="gemini", session="s2", generation="gen-other")
claim = registry.claim_generation(
worker_identity=second["identity"],
generation_id="gen-shared",
now=NOW,
)
self.assertFalse(claim["claimed"])
self.assertFalse(claim["mutation_performed"])
self.assertEqual(claim["blocker_kind"], mwi.BLOCKER_CONFLICTING_SESSIONS)
self.assertEqual(
claim["conflicting_owners"][0]["worker_identity"], first["identity"]
)
# The sanctioned recovery must never be "kill the other process".
self.assertIn("Do not kill", claim["exact_next_action"])
def test_contested_generation_fails_closed_in_assessment(self):
registry = _registry()
first = _attach(registry, client="codex", session="s1", generation="gen-shared")
_attach(registry, client="gemini", session="s2", generation="gen-shared")
verdict = mwi.assess_provenance(
registry=registry, worker_identity=first["identity"], env={}, now=NOW
)
self.assertEqual(verdict["session_ownership"], mwi.OWNERSHIP_CONTESTED)
self.assertTrue(verdict["fail_closed"])
self.assertEqual(verdict["blocker_kind"], mwi.BLOCKER_CONTRADICTORY)
self.assertTrue(verdict["conflicting_live_sessions"])
def test_identity_collision_is_refused_without_corrupting_existing(self):
"""AC31: never replace, adopt, merge with, or corrupt the incumbent."""
registry = _registry()
incumbent = _attach(registry, client="codex", session="s1", generation="gen-1")
before = registry.get(incumbent["identity"])
collided = registry.register(
worker_identity=incumbent["identity"],
client_name="gemini",
client_instance_id="inst-other",
session_id="s2",
generation_id="gen-2",
pid=9999,
now=NOW,
)
self.assertFalse(collided["registered"])
self.assertTrue(collided["collision"])
self.assertFalse(collided["mutation_performed"])
self.assertEqual(collided["blocker_kind"], mwi.BLOCKER_IDENTITY_COLLISION)
self.assertEqual(collided["collision_kind"], "active_worker")
self.assertEqual(
registry.get(incumbent["identity"]), before, "incumbent must be untouched"
)
def test_after_collision_a_regenerated_identity_registers(self):
"""AC32/AC35: forced collision, safe regeneration, successful replacement."""
registry = _registry()
fixed = dict(timestamp_ns=99, nonce="deterministic", now=NOW)
forced = mwi.generate_worker_identity("codex", "sess-collide", **fixed)
first = registry.register(
worker_identity=forced,
client_name="codex",
client_instance_id="inst-1",
session_id="sess-collide",
generation_id="gen-1",
now=NOW,
)
self.assertTrue(first["registered"])
# A second worker deriving the same inputs collides deterministically.
again = mwi.generate_worker_identity("codex", "sess-collide", **fixed)
self.assertEqual(again, forced)
self.assertTrue(
registry.register(
worker_identity=again,
client_name="codex",
client_instance_id="inst-2",
session_id="sess-collide",
generation_id="gen-2",
now=NOW,
)["collision"]
)
replacement = mwi.generate_worker_identity(
"codex", "sess-collide", timestamp_ns=100, nonce="different", now=NOW
)
self.assertNotEqual(replacement, forced)
self.assertTrue(
registry.register(
worker_identity=replacement,
client_name="codex",
client_instance_id="inst-2",
session_id="sess-collide",
generation_id="gen-2",
now=NOW,
)["registered"]
)
def test_restarted_worker_inherits_nothing(self):
"""AC33/AC34: a restart mints a new identity and no prior epoch."""
registry = _registry()
before = _attach(
registry, client="codex", session="sess-before", generation="gen-1", ttl=60
)
later = NOW + timedelta(hours=2)
after = _attach(
registry, client="codex", session="sess-after", generation="gen-2", now=later
)
self.assertNotEqual(before["identity"], after["identity"])
self.assertNotEqual(
registry.get(after["identity"])["generation_id"],
registry.get(before["identity"])["generation_id"],
)
class LivenessTests(unittest.TestCase):
def test_stale_session_record_is_not_live(self):
registry = _registry()
worker = _attach(registry, client="codex", session="s", generation="g", ttl=300)
stale = registry.is_live(
registry.get(worker["identity"]), now=NOW + timedelta(hours=1)
)
self.assertFalse(stale["live"])
self.assertFalse(stale["heartbeat_fresh"])
def test_liveness_is_not_pid_comparison_alone(self):
"""AC7: a live PID does not resurrect an expired registration."""
registry = _registry()
worker = _attach(registry, client="codex", session="s", generation="g", ttl=60)
verdict = registry.is_live(
registry.get(worker["identity"]),
now=NOW + timedelta(hours=1),
pid_alive=True,
)
self.assertFalse(
verdict["live"], "a live PID must not override a dead heartbeat"
)
def test_dead_pid_withdraws_liveness_from_a_fresh_heartbeat(self):
registry = _registry()
worker = _attach(registry, client="codex", session="s", generation="g")
verdict = registry.is_live(
registry.get(worker["identity"]), now=NOW, pid_alive=False
)
self.assertFalse(verdict["live"])
def test_stale_ownership_does_not_permanently_strand_a_daemon(self):
registry = _registry()
stranded = _attach(
registry, client="codex", session="s-old", generation="gen-daemon", ttl=60
)
later = NOW + timedelta(hours=3)
rescuer = _attach(
registry, client="claude", session="s-new", generation="gen-tmp", now=later
)
claim = registry.claim_generation(
worker_identity=rescuer["identity"],
generation_id="gen-daemon",
now=later,
)
self.assertTrue(claim["claimed"], claim["reasons"])
self.assertIn(stranded["identity"], claim["superseded_workers"])
class EvidenceTests(unittest.TestCase):
def test_env_flag_alone_does_not_prove_session_ownership(self):
verdict = mwi.assess_provenance(
registry=None,
worker_identity=None,
env={"GITEA_CLIENT_MANAGED": "1", "GITEA_MCP_SANCTIONED_DAEMON": "1"},
now=NOW,
)
self.assertFalse(verdict["session_owned"])
self.assertEqual(verdict["session_ownership"], mwi.OWNERSHIP_UNOWNED)
self.assertTrue(verdict["env_flag_only"])
self.assertTrue(verdict["fail_closed"])
self.assertNotIn(mwi.EVIDENCE_ATTACHMENT_RECORD, verdict["evidence"])
self.assertFalse(verdict["env_signal"]["proves_session_ownership"])
def test_env_flag_still_answers_the_launch_question(self):
"""The #686 wall is preserved: env decides launch, not ownership."""
self.assertTrue(
mwi.assess_launch_provenance({"GITEA_CLIENT_MANAGED": "1"})["client_managed"]
)
self.assertFalse(
mwi.assess_launch_provenance({"GITEA_CLIENT_MANAGED": "0"})["client_managed"]
)
self.assertFalse(
mwi.assess_launch_provenance({}, stdin_is_tty=True)["client_managed"]
)
self.assertTrue(
mwi.assess_launch_provenance({"GITEA_MCP_PROFILE": "prgs-author"})[
"client_managed"
]
)
def test_missing_evidence_is_unproven_not_manual(self):
"""A missing proof must not be reported as a hand-launched process."""
verdict = mwi.assess_provenance(
registry=None, worker_identity=None, env={}, now=NOW
)
self.assertEqual(verdict["provenance"], mwi.PROVENANCE_UNPROVEN)
self.assertNotEqual(verdict["provenance"], mwi.PROVENANCE_MANUAL)
self.assertTrue(verdict["fail_closed"])
def test_declared_manual_launch_is_reported_as_manual(self):
verdict = mwi.assess_provenance(
registry=None,
worker_identity=None,
env={"GITEA_CLIENT_MANAGED": "0"},
now=NOW,
)
self.assertEqual(verdict["provenance"], mwi.PROVENANCE_MANUAL)
def test_fail_closed_refusal_names_its_scope_not_the_profile(self):
"""AC17/AC41: no refusal is profile-wide."""
verdict = mwi.assess_provenance(
registry=None,
worker_identity=None,
env={},
profile="prgs-reviewer",
role="reviewer",
now=NOW,
)
self.assertFalse(verdict["scope"]["profile_wide"])
self.assertEqual(verdict["blocker_kind"], mwi.BLOCKER_NO_ATTACHMENT)
class CohortScopingTests(unittest.TestCase):
def test_shared_profile_with_distinct_identities_does_not_block(self):
"""AC40: profile is not a singleton identity."""
registry = _registry()
_attach(
registry,
client="codex",
session="s1",
generation="g1",
profile="prgs-reviewer",
)
_attach(
registry,
client="gemini",
session="s2",
generation="g2",
profile="prgs-reviewer",
)
cohort = mwi.classify_cohort(registry.list_workers(), now=NOW)
self.assertFalse(cohort["blocked"], cohort["reasons"])
self.assertEqual(cohort["blocker_kind"], mwi.BLOCKER_NONE)
self.assertIn("prgs-reviewer", cohort["shared_profiles"])
self.assertTrue(cohort["profile_sharing_permitted"])
self.assertEqual(cohort["blocked_worker_identities"], [])
def test_duplicate_cohort_records_block_only_the_offenders(self):
registry = _registry()
_attach(registry, client="codex", session="s1", generation="gen-contested")
_attach(registry, client="gemini", session="s2", generation="gen-contested")
_attach(registry, client="claude", session="s3", generation="gen-fine")
cohort = mwi.classify_cohort(registry.list_workers(), now=NOW)
self.assertTrue(cohort["blocked"])
self.assertEqual(cohort["contested_generations"], ["gen-contested"])
self.assertEqual(len(cohort["blocked_worker_identities"]), 2)
def test_mixed_runtime_generations_are_scoped_independently(self):
"""AC17: one stale generation does not wall unrelated healthy ones."""
registry = _registry()
stale = _attach(
registry, client="codex", session="s1", generation="gen-stale", ttl=60
)
healthy_a = _attach(registry, client="gemini", session="s2", generation="gen-a")
healthy_b = _attach(registry, client="claude", session="s3", generation="gen-b")
scoped = mwi.scope_runtime_failure(
failure_kind="stale-runtime",
worker_identity=stale["identity"],
profile="prgs-reviewer",
all_live_workers=registry.list_workers(),
)
self.assertFalse(scoped["profile_wide"])
self.assertFalse(scoped["fleet_wide"])
self.assertEqual(len(scoped["affected_workers"]), 1)
self.assertEqual(scoped["unaffected_worker_count"], 2)
unaffected = {w["worker_identity"] for w in scoped["unaffected_workers"]}
self.assertEqual(unaffected, {healthy_a["identity"], healthy_b["identity"]})
class HardcodedClientRegressionTests(unittest.TestCase):
def test_unknown_client_does_not_resolve_to_codex(self):
for name in ("gemini", "antigravity", "grok", "some_new_llm", "", None):
with self.subTest(client=name):
self.assertNotEqual(
mcp_client_reconnect.normalize_client(name),
"codex",
"an unidentified client must never be handed Codex UI steps",
)
def test_known_clients_still_get_their_own_steps(self):
self.assertEqual(mcp_client_reconnect.normalize_client("codex"), "codex")
self.assertEqual(
mcp_client_reconnect.normalize_client("claude_code"), "claude_code"
)
def test_generic_steps_do_not_name_a_specific_vendor(self):
steps = " ".join(mcp_client_reconnect.operator_ui_steps("gemini"))
self.assertNotIn("Codex", steps)
def test_reconnect_client_is_derived_from_the_attachment_record(self):
registry = _registry()
worker = _attach(registry, client="antigravity", session="s", generation="g")
verdict = mwi.assess_provenance(
registry=registry, worker_identity=worker["identity"], env={}, now=NOW
)
self.assertEqual(mwi.reconnect_client_for(verdict), "antigravity")
def test_reconnect_client_is_unknown_rather_than_guessed(self):
self.assertEqual(mwi.reconnect_client_for({}), mwi.UNKNOWN_CLIENT)
class RemoteBindingTests(unittest.TestCase):
def test_explicit_prgs_selection_is_honoured(self):
resolved = mwi.resolve_bound_remote(
requested_remote="prgs", bound_remote="prgs", default_remote="dadeschools"
)
self.assertEqual(resolved["remote"], "prgs")
self.assertFalse(resolved["drifted"])
def test_omitted_remote_uses_the_binding_not_the_library_default(self):
"""The reported dadeschools host drift."""
resolved = mwi.resolve_bound_remote(
requested_remote=None, bound_remote="prgs", default_remote="dadeschools"
)
self.assertEqual(resolved["remote"], "prgs")
self.assertNotEqual(resolved["remote"], "dadeschools")
self.assertEqual(resolved["resolved_from"], "session_binding")
def test_contradicting_the_binding_is_refused(self):
resolved = mwi.resolve_bound_remote(
requested_remote="dadeschools",
bound_remote="prgs",
default_remote="dadeschools",
)
self.assertEqual(resolved["remote"], "prgs")
self.assertTrue(resolved["drifted"])
self.assertFalse(resolved["honoured_request"])
def test_unbound_session_falls_back_and_says_so(self):
resolved = mwi.resolve_bound_remote(
requested_remote=None, bound_remote=None, default_remote="dadeschools"
)
self.assertEqual(resolved["remote"], "dadeschools")
self.assertEqual(resolved["resolved_from"], "library_default")
self.assertTrue(resolved["reasons"])
class SurfaceAgreementTests(unittest.TestCase):
"""The reproduced contradiction: two surfaces, one process, two answers."""
def test_namespace_health_and_direct_assessment_agree(self):
registry = _registry()
worker = _attach(
registry,
client="gemini",
session="sess-agree",
generation="gen-agree",
profile="prgs-reviewer",
)
env = {"GITEA_MCP_PROFILE": "prgs-reviewer", "GITEA_CLIENT_MANAGED": "1"}
direct = mwi.assess_provenance(
registry=registry,
worker_identity=worker["identity"],
env=env,
profile="prgs-reviewer",
)
health = mcp_namespace_health.classify_namespace_probe(
"gitea-reviewer",
configured=True,
registered_tools=["gitea_whoami"],
probe_result={"success": True},
probe_source="client_namespace",
process={"pid": 4242, "profile": "prgs-reviewer", "env": env},
registry=registry,
worker_identity=worker["identity"],
)
self.assertEqual(health["provenance"], direct["provenance"])
self.assertEqual(health["is_client_managed"], direct["is_client_managed"])
self.assertEqual(health["worker_identity"], direct["worker_identity"])
self.assertEqual(health["session_id"], "sess-agree")
self.assertEqual(health["client_name"], "gemini")
def test_namespace_health_can_report_client_managed_at_all(self):
"""The old derivation was structurally incapable of this."""
env = {"GITEA_CLIENT_MANAGED": "1", "GITEA_MCP_PROFILE": "prgs-author"}
health = mcp_namespace_health.classify_namespace_probe(
"gitea-author",
configured=True,
registered_tools=["gitea_whoami"],
probe_result={"success": True},
probe_source="client_namespace",
process={"pid": 1234, "profile": "prgs-author", "env": env},
)
self.assertTrue(
health["is_client_managed"],
"a client-managed launch must be reportable as client-managed",
)
def test_namespace_health_without_attachment_fails_closed(self):
health = mcp_namespace_health.classify_namespace_probe(
"gitea-author",
configured=True,
registered_tools=["gitea_whoami"],
probe_result={"success": True},
probe_source="client_namespace",
process={"pid": 1234, "profile": "prgs-author", "env": {}},
)
self.assertTrue(health["provenance_fail_closed"])
self.assertEqual(health["provenance"], mwi.PROVENANCE_UNPROVEN)
self.assertIsNone(health["session_id"])
def test_no_false_reconnect_loop_for_an_owned_session(self):
"""A proven owner must not be told to reconnect."""
registry = _registry()
worker = _attach(registry, client="claude", session="s", generation="g")
verdict = mwi.assess_provenance(
registry=registry, worker_identity=worker["identity"], env={}, now=NOW
)
self.assertFalse(verdict["fail_closed"])
self.assertEqual(verdict["blocker_kind"], mwi.BLOCKER_NONE)
self.assertEqual(verdict["reasons"], [])
if __name__ == "__main__":
unittest.main()
+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()
+480
View File
@@ -0,0 +1,480 @@
"""Tests for dead-owner / PID-reuse session retirement (#969)."""
from __future__ import annotations
import os
import tempfile
import threading
import unittest
from concurrent.futures import ThreadPoolExecutor, as_completed
from datetime import datetime, timedelta, timezone
import control_plane_db as cpd
import post_restart_reconcile as prr
import session_lifecycle as sl
NOW = datetime(2026, 7, 29, 12, 0, 0, tzinfo=timezone.utc)
EARLIER = NOW - timedelta(hours=2)
LATER = NOW + timedelta(minutes=5)
def _session(
session_id: str,
*,
pid: int | None = 4242,
status: str = "active",
started_at: datetime = EARLIER,
last_heartbeat_at: datetime | None = None,
client_managed: bool = False,
owner_process_started_at: datetime | None = None,
role: str = "author",
) -> dict:
hb = last_heartbeat_at if last_heartbeat_at is not None else started_at
row = {
"session_id": session_id,
"role": role,
"profile": "prgs-author",
"pid": pid,
"status": status,
"started_at": cpd._ts(started_at),
"last_heartbeat_at": cpd._ts(hb),
"client_managed": client_managed,
}
if owner_process_started_at is not None:
row["owner_process_started_at"] = cpd._ts(owner_process_started_at)
return row
def _alive(pids: set[int]):
def _check(pid):
try:
return int(pid) in pids
except (TypeError, ValueError):
return False
return _check
def _starts(mapping: dict[int, datetime]):
def _probe(pid):
try:
return mapping.get(int(pid))
except (TypeError, ValueError):
return None
return _probe
class ClassifyDeadOwnerTests(unittest.TestCase):
def test_dead_owner_is_stale_and_retireable(self) -> None:
c = sl.classify_session(
_session("ghost", pid=2_000_000_000),
now=NOW,
pid_checker=_alive(set()),
process_start_probe=_starts({}),
)
self.assertEqual(c.classification, sl.CLASS_STALE)
self.assertTrue(c.retireable)
self.assertIn(c.reason, {sl.REASON_DEAD_OWNER, sl.REASON_HEARTBEAT_STALE_DEAD})
def test_missing_pid_is_stale(self) -> None:
c = sl.classify_session(
_session("no-pid", pid=None),
now=NOW,
pid_checker=_alive(set()),
)
self.assertEqual(c.classification, sl.CLASS_STALE)
self.assertEqual(c.reason, sl.REASON_MISSING_PID)
self.assertTrue(c.retireable)
class PidReuseTests(unittest.TestCase):
def test_pid_reuse_marks_stale_not_live(self) -> None:
# Process with same PID started AFTER the session was recorded.
c = sl.classify_session(
_session("reused", pid=77, started_at=EARLIER, last_heartbeat_at=EARLIER),
now=NOW,
pid_checker=_alive({77}),
process_start_probe=_starts({77: LATER}),
)
self.assertEqual(c.classification, sl.CLASS_STALE)
self.assertEqual(c.reason, sl.REASON_PID_REUSE)
self.assertTrue(c.pid_reused)
self.assertTrue(c.retireable)
def test_matching_process_start_is_live(self) -> None:
c = sl.classify_session(
_session(
"same-proc",
pid=88,
started_at=EARLIER,
last_heartbeat_at=NOW - timedelta(seconds=30),
owner_process_started_at=EARLIER - timedelta(seconds=5),
),
now=NOW,
pid_checker=_alive({88}),
process_start_probe=_starts({88: EARLIER - timedelta(seconds=5)}),
)
self.assertEqual(c.classification, sl.CLASS_LIVE)
self.assertFalse(c.retireable)
class LiveOwnerAndLeaseTests(unittest.TestCase):
def test_live_owner_not_retired(self) -> None:
c = sl.classify_session(
_session(
"live",
pid=os.getpid(),
last_heartbeat_at=NOW - timedelta(seconds=10),
),
now=NOW,
pid_checker=_alive({os.getpid()}),
process_start_probe=_starts({os.getpid(): EARLIER}),
)
self.assertEqual(c.classification, sl.CLASS_LIVE)
self.assertFalse(c.retireable)
def test_live_lease_blocks_retirement_even_if_pid_dead(self) -> None:
c = sl.classify_session(
_session("leased", pid=99999),
now=NOW,
pid_checker=_alive(set()),
live_lease_sessions={"leased"},
)
self.assertEqual(c.classification, sl.CLASS_PROTECTED)
self.assertEqual(c.reason, sl.REASON_LIVE_LEASE)
self.assertFalse(c.retireable)
def test_client_managed_live_never_retired(self) -> None:
c = sl.classify_session(
_session(
"client",
pid=55,
client_managed=True,
last_heartbeat_at=NOW - timedelta(seconds=5),
),
now=NOW,
pid_checker=_alive({55}),
process_start_probe=_starts({55: EARLIER}),
)
self.assertEqual(c.classification, sl.CLASS_LIVE)
self.assertEqual(c.reason, sl.REASON_CLIENT_MANAGED_LIVE)
self.assertFalse(c.retireable)
def test_client_managed_dead_pid_is_retireable(self) -> None:
# Dead client process is not a live client-managed session.
c = sl.classify_session(
_session("client-dead", pid=56, client_managed=True),
now=NOW,
pid_checker=_alive(set()),
)
self.assertEqual(c.classification, sl.CLASS_STALE)
self.assertTrue(c.retireable)
class TerminalAndDisconnectedTests(unittest.TestCase):
def test_already_terminal_not_retireable(self) -> None:
c = sl.classify_session(
_session("done", status="retired", pid=1),
now=NOW,
pid_checker=_alive(set()),
)
self.assertEqual(c.classification, sl.CLASS_TERMINAL)
self.assertFalse(c.retireable)
def test_alive_stale_heartbeat_is_disconnected_not_retired(self) -> None:
c = sl.classify_session(
_session(
"quiet",
pid=66,
last_heartbeat_at=NOW - timedelta(hours=5),
),
now=NOW,
pid_checker=_alive({66}),
process_start_probe=_starts({66: EARLIER}),
)
self.assertEqual(c.classification, sl.CLASS_DISCONNECTED)
self.assertFalse(c.retireable)
class FleetAndApplyTests(unittest.TestCase):
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.db_path = os.path.join(self._tmp.name, "cp.sqlite3")
self.db = cpd.ControlPlaneDB(self.db_path)
def tearDown(self) -> None:
self._tmp.cleanup()
def test_mixed_fleet_and_apply_retires_only_stale(self) -> None:
live_pid = os.getpid()
# Use wall-clock "now" so upsert timestamps align with classification.
moment = datetime.now(timezone.utc)
proc_start = moment - timedelta(hours=1)
self.db.upsert_session(
session_id="s-live",
role="author",
pid=live_pid,
status="active",
owner_process_started_at=cpd._ts(proc_start),
)
self.db.upsert_session(
session_id="s-dead", role="reviewer", pid=2_000_000_001, status="active"
)
self.db.upsert_session(
session_id="s-ended", role="merger", pid=3, status="ended"
)
sessions = self.db.list_sessions(limit=50)
report = sl.classify_sessions(
sessions,
leases=[],
now=moment,
pid_checker=_alive({live_pid}),
process_start_probe=_starts({live_pid: proc_start}),
)
self.assertGreaterEqual(report.stale_count, 1)
self.assertTrue(
any(c.session_id == "s-dead" and c.retireable for c in report.classifications)
)
self.assertTrue(
any(
c.session_id == "s-live" and not c.retireable
for c in report.classifications
)
)
first = sl.apply_session_retirements(
self.db, report, dry_run=False, actor_session_id="actor-1", now=moment
)
self.assertGreaterEqual(first["retired_count"], 1)
# After retirement, active list should exclude s-dead.
active = {
s["session_id"]
for s in self.db.list_sessions(statuses=("active",), limit=50)
}
self.assertNotIn("s-dead", active)
self.assertIn("s-live", active)
# Repeated cleanup is idempotent.
report2 = sl.classify_sessions(
self.db.list_sessions(limit=50),
leases=[],
now=moment,
pid_checker=_alive({live_pid}),
process_start_probe=_starts({live_pid: proc_start}),
)
second = sl.apply_session_retirements(
self.db, report2, dry_run=False, actor_session_id="actor-1", now=moment
)
# No double-retirement of the same row as a new mutation.
self.assertEqual(second["retired_count"], 0)
# Durable audit event present.
import sqlite3
conn = sqlite3.connect(self.db_path)
try:
events = conn.execute(
"SELECT event_type, message FROM events WHERE event_type = ?",
("session_retired",),
).fetchall()
finally:
conn.close()
self.assertTrue(events)
self.assertTrue(any("s-dead" in (m or "") for _, m in events))
def test_live_lease_blocks_db_retirement(self) -> None:
self.db.upsert_session(
session_id="s-leased", role="author", pid=2_000_000_002, status="active"
)
self.db.upsert_work_item(
remote="prgs",
org="org",
repo="repo",
kind="issue",
number=969,
)
result = self.db.assign_and_lease(
session_id="s-leased",
role="author",
remote="prgs",
org="org",
repo="repo",
kind="issue",
number=969,
)
self.assertEqual(result.outcome, "assigned")
# Inventory-style lease with explicit live freshness (authoritative for
# the pure classifier). DB apply also blocks on the active lease row.
leases = [
{
"lease_id": result.lease_id,
"session_id": "s-leased",
"status": "active",
"freshness": {"freshness": "active"},
}
]
report = sl.classify_sessions(
self.db.list_sessions(statuses=("active",), limit=20),
leases=leases,
now=NOW,
pid_checker=_alive(set()),
)
# Classifier protects via live lease set.
self.assertTrue(
any(
c.session_id == "s-leased" and c.classification == sl.CLASS_PROTECTED
for c in report.classifications
)
)
apply = sl.apply_session_retirements(
self.db, report, dry_run=False, actor_session_id="actor", now=NOW
)
self.assertEqual(apply["retired_count"], 0)
active = {
s["session_id"]
for s in self.db.list_sessions(statuses=("active",), limit=20)
}
self.assertIn("s-leased", active)
# Direct DB CAS also refuses while an active lease row remains.
blocked = self.db.retire_session(
session_id="s-leased",
reason=sl.REASON_DEAD_OWNER,
actor_session_id="actor",
now=NOW,
)
self.assertEqual(blocked["outcome"], "blocked")
self.assertEqual(blocked["reason"], "live_lease")
def test_concurrent_retirement_is_idempotent(self) -> None:
for i in range(20):
self.db.upsert_session(
session_id=f"ghost-{i}",
role="author",
pid=3_000_000 + i,
status="active",
)
def _worker() -> dict:
return sl.retire_stale_sessions(
self.db,
dry_run=False,
actor_session_id=f"actor-{threading.get_ident()}",
now=NOW,
pid_checker=_alive(set()),
process_start_probe=_starts({}),
session_limit=100,
)
outcomes = []
with ThreadPoolExecutor(max_workers=4) as pool:
futs = [pool.submit(_worker) for _ in range(4)]
for fut in as_completed(futs):
outcomes.append(fut.result())
total_retired = sum(o["apply"]["retired_count"] for o in outcomes)
# Exactly one successful retirement per ghost row across all workers.
self.assertEqual(total_retired, 20)
active = {
s["session_id"]
for s in self.db.list_sessions(statuses=("active",), limit=100)
}
for i in range(20):
self.assertNotIn(f"ghost-{i}", active)
class ReconcileIntegrationTests(unittest.TestCase):
def test_unresolved_until_retired_then_resolved(self) -> None:
inv = {
"inventory_complete": True,
"incomplete_reasons": [],
"service_health": {"healthy": True},
"clients": [{"session_id": "c1", "connected": True}],
"sessions": [
_session("ghost", pid=2_000_000_099, last_heartbeat_at=EARLIER),
],
"leases": [],
"checkpoints_available": False,
"worktree_bindings": [],
"pending_mutations": [],
"capabilities": {"stale": False},
"boot_head_sha": "a" * 40,
"current_head_sha": "a" * 40,
"queue_state": {"safe_to_resume": True},
}
proof = prr.reconcile_after_restart(inv, now=NOW, mode=prr.MODE_LOG_ONLY)
sess = next(i for i in proof.items if i.dimension == prr.DIM_SESSIONS)
self.assertEqual(sess.status, prr.ITEM_UNRESOLVED)
self.assertIn("ghost", sess.details.get("orphan_session_ids") or [])
# After retirement inventory (no active orphans) resolves.
inv2 = dict(inv)
inv2["sessions"] = []
inv2["session_fleet"] = {
"retireable_session_ids": [],
"sessions_dimension_resolved": True,
"live_count": 0,
"stale_count": 0,
}
proof2 = prr.reconcile_after_restart(inv2, now=NOW, mode=prr.MODE_LOG_ONLY)
sess2 = next(i for i in proof2.items if i.dimension == prr.DIM_SESSIONS)
self.assertEqual(sess2.status, prr.ITEM_RESOLVED)
def test_legacy_orphan_key_still_populated(self) -> None:
inv = {
"inventory_complete": True,
"service_health": {"healthy": True},
"clients": [],
"sessions": [_session("ghost", pid=2_000_000_100)],
"leases": [],
"checkpoints_available": False,
"worktree_bindings": [],
"pending_mutations": [],
"capabilities": {"stale": False},
"boot_head_sha": "a" * 40,
"current_head_sha": "a" * 40,
"queue_state": {"safe_to_resume": True},
}
proof = prr.reconcile_after_restart(inv, now=NOW)
sess = next(i for i in proof.items if i.dimension == prr.DIM_SESSIONS)
self.assertIn("orphan_session_ids", sess.details)
class SchemaMigrationTests(unittest.TestCase):
def test_lifecycle_columns_present(self) -> None:
with tempfile.TemporaryDirectory() as tmp:
path = os.path.join(tmp, "cp.sqlite3")
db = cpd.ControlPlaneDB(path)
db.upsert_session(session_id="s1", role="author", pid=1)
out = db.retire_session(
session_id="s1",
reason=sl.REASON_DEAD_OWNER,
actor_session_id="tester",
now=NOW,
)
self.assertEqual(out["outcome"], "retired")
rows = db.list_sessions(limit=5)
# May not appear under active filter
all_rows = db.list_sessions(limit=5)
# Re-open raw to check columns
import sqlite3
conn = sqlite3.connect(path)
try:
cols = {r[1] for r in conn.execute("PRAGMA table_info(sessions)")}
version = conn.execute(
"SELECT value FROM schema_meta WHERE key='schema_version'"
).fetchone()[0]
finally:
conn.close()
self.assertIn("retired_at", cols)
self.assertIn("retire_reason", cols)
self.assertIn("owner_process_started_at", cols)
self.assertEqual(version, "6")
if __name__ == "__main__":
unittest.main()