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
22 changed files with 1698 additions and 6269 deletions
+11 -126
View File
@@ -37,48 +37,6 @@ CANONICAL_ROOT_ENV = "GITEA_CANONICAL_REPOSITORY_ROOT"
# Candidate git remote names probed when deriving repository identity. # Candidate git remote names probed when deriving repository identity.
_IDENTITY_REMOTE_CANDIDATES = ("prgs", "origin", "dadeschools", "mdcps") _IDENTITY_REMOTE_CANDIDATES = ("prgs", "origin", "dadeschools", "mdcps")
# Repository-authority modes (#973 B10). Exactly two values are supported.
# ``mode`` selects how repository authority is established, so an unrecognised
# value must never be normalised onto one of these: aliasing a trusted mode is
# precisely the defect. Omitting the argument keeps the documented safe default,
# ``validation``.
MODE_VALIDATION = "validation"
MODE_DERIVATION = "derivation"
SUPPORTED_MODES: tuple[str, ...] = (MODE_VALIDATION, MODE_DERIVATION)
# Reason code emitted when an unsupported mode is refused. Mirrors the existing
# ``webui.sanctioned_restart.DENY_UNKNOWN_MODE`` convention so callers and tests
# can assert the refusal cause rather than string-matching prose.
DENY_UNKNOWN_MODE = "unknown_mode"
def unsupported_mode_reason(mode: object) -> str | None:
"""Precise rejection reason for *mode*, or None when *mode* is supported.
Only the two documented string values are accepted, compared exactly — no
stripping, no case folding — so misspellings and whitespace variants are
refused rather than coerced. An empty string, ``None``, and any non-string
are all *explicitly supplied* unsupported values and are refused on the same
footing; none of them is normalised to a supported mode. Omitting the
argument entirely never reaches here with an unsupported value because the
parameter default is ``"validation"``.
"""
if isinstance(mode, str) and mode in SUPPORTED_MODES:
return None
supported = ", ".join(repr(m) for m in SUPPORTED_MODES)
if not isinstance(mode, str):
return (
f"unsupported repository-authority mode {mode!r} of type "
f"{type(mode).__name__}: only {supported} are supported; the mode "
"is refused before any repository assessment and no repository "
"identity was resolved through it (fail closed)"
)
return (
f"unsupported repository-authority mode {mode!r}: only {supported} are "
"supported; the mode is refused before any repository assessment and no "
"repository identity was resolved through it (fail closed)"
)
def configured_canonical_root( def configured_canonical_root(
profile: Mapping | None, profile: Mapping | None,
@@ -174,7 +132,6 @@ def assess_canonical_repository_root(
process_project_root: str, process_project_root: str,
remote: str | None = None, remote: str | None = None,
require_binding: bool = False, require_binding: bool = False,
mode: str = "validation",
) -> dict: ) -> dict:
"""Validate the canonical repository root binding, failing closed on forgery. """Validate the canonical repository root binding, failing closed on forgery.
@@ -183,41 +140,15 @@ def assess_canonical_repository_root(
``configured`` (whether a cross-repo binding was declared), ``configured`` (whether a cross-repo binding was declared),
``resolved_slug`` and ``source``. ``resolved_slug`` and ``source``.
*mode* accepts exactly the two values in :data:`SUPPORTED_MODES`: Without a configured binding the single-repo default is preserved: the
- ``"validation"`` (default): an independently trusted expected repository canonical root is derived from *process_project_root* and never blocks
slug is known (or required). Candidate root's observed identity must match. (unless *require_binding* explicitly demands one).
Unprovable or missing expected identity fails closed when *require_binding* is True.
- ``"derivation"``: caller is deriving the canonical repository identity.
No expected slug exists yet by design. Derivation succeeds if the configured
root exists, is a git repository, and carries a resolvable git remote identity.
Every other explicitly supplied value — unknown strings, misspellings, the With a configured binding the path must exist, be a git repository, and —
empty string, ``None``, and non-strings — is refused with ``proven`` False, when *expected_slug* is known — carry a matching repository identity. A
``block`` True, and ``reason_code`` :data:`DENY_UNKNOWN_MODE` (#973 B10). mismatched or (when *require_binding*) unprovable identity is a forged or
Omitting *mode* entirely keeps the documented ``"validation"`` default. conflicting binding and fails closed.
""" """
# #973 B10: refuse an unsupported repository-authority mode as the very first
# act, before any candidate-root existence check, path or symlink resolution,
# git top-level discovery, remote-URL or repository-identity discovery, and
# before any expected-versus-observed comparison or validation/derivation
# behaviour. An unsupported mode previously fell into the catch-all ``else``
# on the configured-root path (silently receiving validation semantics) and
# skipped the identity comparison entirely on the single-repository default
# path (strictly weaker than validation), so matching identities could return
# ``proven`` True. No repository identity may be resolved through a mode the
# contract does not define.
mode_reason = unsupported_mode_reason(mode)
if mode_reason is not None:
return _assessment(
proven=False,
reasons=[mode_reason],
configured=bool((configured_value or "").strip()),
canonical_repo_root=None,
resolved_slug=None,
source=source,
reason_code=DENY_UNKNOWN_MODE,
)
process_root = os.path.realpath(process_project_root) process_root = os.path.realpath(process_project_root)
declared = (configured_value or "").strip() declared = (configured_value or "").strip()
@@ -237,31 +168,12 @@ def assess_canonical_repository_root(
) )
# Single-repo default: canonical root follows the install checkout. # Single-repo default: canonical root follows the install checkout.
derived = resolve_repo_toplevel(process_root) or process_root derived = resolve_repo_toplevel(process_root) or process_root
resolved_slug = repository_identity_slug(derived, remote=remote)
reasons: list[str] = []
# #973 B10: the allowlist above guarantees *mode* is one of the two
# supported values here, so an unsupported value can no longer skip this
# identity comparison and end up strictly weaker than validation.
if mode == MODE_VALIDATION and expected_slug:
expected = expected_slug.strip()
if resolved_slug and resolved_slug.lower() != expected.lower():
reasons.append(
f"canonical repository root identity mismatch: '{derived}' resolves "
f"to repository '{resolved_slug}' but expected repository identity "
f"is '{expected}' (forged or conflicting binding, fail closed)"
)
elif not resolved_slug and require_binding:
reasons.append(
f"canonical repository root '{derived}' has no resolvable git "
f"remote identity to confirm authorization for '{expected}' "
"(fail closed)"
)
return _assessment( return _assessment(
proven=not reasons, proven=True,
reasons=reasons, reasons=[],
configured=False, configured=False,
canonical_repo_root=derived, canonical_repo_root=derived,
resolved_slug=resolved_slug, resolved_slug=None,
source=None, source=None,
) )
@@ -295,17 +207,6 @@ def assess_canonical_repository_root(
resolved_slug = repository_identity_slug(toplevel, remote=remote) resolved_slug = repository_identity_slug(toplevel, remote=remote)
reasons: list[str] = [] reasons: list[str] = []
if mode == MODE_DERIVATION:
if not resolved_slug:
reasons.append(
f"configured canonical repository root '{toplevel}' has no resolvable "
"git remote identity (fail closed)"
)
else:
# #973 B10: MODE_VALIDATION only. This arm is no longer a catch-all — the
# allowlist above admits no third value, so an unsupported mode can no
# longer silently receive validation semantics here.
expected = (expected_slug or "").strip() or None expected = (expected_slug or "").strip() or None
if expected: if expected:
if resolved_slug and resolved_slug.lower() != expected.lower(): if resolved_slug and resolved_slug.lower() != expected.lower():
@@ -320,12 +221,6 @@ def assess_canonical_repository_root(
f"remote identity to confirm authorization for '{expected}' " f"remote identity to confirm authorization for '{expected}' "
"(fail closed)" "(fail closed)"
) )
elif require_binding:
reasons.append(
f"canonical repository root '{toplevel}' has configured value "
f"'{configured_value}' but authoritative expected repository identity "
"is unprovable or missing (fail closed)"
)
return _assessment( return _assessment(
proven=not reasons, proven=not reasons,
@@ -357,19 +252,10 @@ def _assessment(
proven: bool, proven: bool,
reasons: list[str], reasons: list[str],
configured: bool, configured: bool,
canonical_repo_root: str | None, canonical_repo_root: str,
resolved_slug: str | None, resolved_slug: str | None,
source: str | None, source: str | None,
reason_code: str | None = None,
) -> dict: ) -> dict:
"""Build the assessment payload.
``canonical_repo_root`` is None only when the assessment refused to resolve
one at all — today exactly the unsupported-mode refusal (#973 B10), which
must not derive a trusted repository identity through an undefined mode.
``reason_code`` is a machine-checkable refusal cause; None for every
ordinary (non-coded) outcome.
"""
return { return {
"proven": proven, "proven": proven,
"block": not proven, "block": not proven,
@@ -378,5 +264,4 @@ def _assessment(
"canonical_repo_root": canonical_repo_root, "canonical_repo_root": canonical_repo_root,
"resolved_slug": resolved_slug, "resolved_slug": resolved_slug,
"source": source, "source": source,
"reason_code": reason_code,
} }
+226 -305
View File
@@ -27,12 +27,12 @@ import uuid
from contextlib import contextmanager from contextlib import contextmanager
from dataclasses import dataclass from dataclasses import dataclass
from datetime import datetime, timedelta, timezone from datetime import datetime, timedelta, timezone
from typing import Any, Iterator, Sequence from typing import Any, Iterator, Mapping, Sequence
import dependency_graph import dependency_graph
import gitea_audit import gitea_audit
SCHEMA_VERSION = 5 SCHEMA_VERSION = 6
# Assignable work kinds only — raw monitoring incidents are never work items. # Assignable work kinds only — raw monitoring incidents are never work items.
WORK_KINDS = frozenset({"issue", "pr"}) WORK_KINDS = frozenset({"issue", "pr"})
@@ -280,22 +280,6 @@ def _ts(dt: datetime | None = None) -> str:
return value.astimezone(timezone.utc).replace(microsecond=0).isoformat().replace("+00:00", "Z") return value.astimezone(timezone.utc).replace(microsecond=0).isoformat().replace("+00:00", "Z")
def _realpath_or_raw(value: str | None) -> str:
"""Normalize a filesystem path for compare-and-swap equality (#970).
Symlinks and ``..`` segments must not make two spellings of the same path
look different, but an unresolvable path must still compare as itself
rather than collapsing to empty — an empty result means "no path given".
"""
text = (value or "").strip()
if not text:
return ""
try:
return os.path.realpath(os.path.abspath(text))
except Exception:
return text
def _parse_ts(value: str | None) -> datetime | None: def _parse_ts(value: str | None) -> datetime | None:
if not value: if not value:
return None return None
@@ -435,6 +419,7 @@ class ControlPlaneDB:
self._migrate_incident_links_null_scope(conn) self._migrate_incident_links_null_scope(conn)
self._migrate_lease_lifecycle_columns(conn) self._migrate_lease_lifecycle_columns(conn)
self._migrate_session_ownership_columns(conn) self._migrate_session_ownership_columns(conn)
self._migrate_session_lifecycle_columns(conn)
self._migrate_usage_events_table(conn) self._migrate_usage_events_table(conn)
conn.execute( conn.execute(
"INSERT OR REPLACE INTO schema_meta(key, value) VALUES (?, ?)", "INSERT OR REPLACE INTO schema_meta(key, value) VALUES (?, ?)",
@@ -828,6 +813,7 @@ class ControlPlaneDB:
pid: int | None = None, pid: int | None = None,
status: str = "active", status: str = "active",
controller_instance_id: str | None = None, controller_instance_id: str | None = None,
owner_process_started_at: str | None = None,
) -> dict[str, Any]: ) -> dict[str, Any]:
"""Register/refresh a session row. """Register/refresh a session row.
@@ -837,16 +823,22 @@ class ControlPlaneDB:
controller instance can. It is never overwritten with ``None``, so a controller instance can. It is never overwritten with ``None``, so a
heartbeat from a caller that does not supply one cannot erase heartbeat from a caller that does not supply one cannot erase
ownership. 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() now = _ts()
instance = (controller_instance_id or "").strip() or None instance = (controller_instance_id or "").strip() or None
proc_start = (owner_process_started_at or "").strip() or None
with self._tx() as conn: with self._tx() as conn:
existing = conn.execute( existing = conn.execute(
"SELECT session_id FROM sessions WHERE session_id = ?", "SELECT session_id FROM sessions WHERE session_id = ?",
(session_id,), (session_id,),
).fetchone() ).fetchone()
if existing: if existing:
if instance is None: if instance is None and proc_start is None:
conn.execute( conn.execute(
""" """
UPDATE sessions UPDATE sessions
@@ -856,7 +848,21 @@ class ControlPlaneDB:
""", """,
(role, profile, namespace, pid, now, status, session_id), (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( conn.execute(
""" """
UPDATE sessions UPDATE sessions
@@ -870,18 +876,33 @@ class ControlPlaneDB:
instance, session_id, 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: else:
conn.execute( conn.execute(
""" """
INSERT INTO sessions( INSERT INTO sessions(
session_id, role, profile, namespace, pid, session_id, role, profile, namespace, pid,
started_at, last_heartbeat_at, status, started_at, last_heartbeat_at, status,
controller_instance_id controller_instance_id, owner_process_started_at
) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?) ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
""", """,
( (
session_id, role, profile, namespace, pid, now, now, session_id, role, profile, namespace, pid, now, now,
status, instance, status, instance, proc_start,
), ),
) )
row = conn.execute( row = conn.execute(
@@ -897,6 +918,165 @@ class ControlPlaneDB:
(_ts(), session_id), (_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( def list_sessions(
self, self,
*, *,
@@ -1655,6 +1835,29 @@ class ControlPlaneDB:
if name not in cols: if name not in cols:
conn.execute(f"ALTER TABLE sessions ADD COLUMN {name} {decl}") 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( def list_active_claims(
self, self,
*, *,
@@ -1929,168 +2132,6 @@ class ControlPlaneDB:
), ),
) )
def retire_lease_worktree_path(
self,
lease_id: str,
*,
expected_path: str | None = None,
expected_status: str | None = None,
expected_session_id: str | None = None,
expected_owner_pid: int | None = None,
reason: str = "missing_worktree_path_retired",
) -> dict[str, Any]:
"""Retire a missing worktree_path binding from a control-plane lease (#970).
Clears worktree_path on the lease row, updates provenance_json with
durable retirement audit proof, and writes a worktree_binding_retired
event.
The update is a compare-and-swap (#970 review 644 B1/B3): the caller
states the exact path it audited and, when known, the lease status,
owning session, and owner pid it classified against. Every stated value
must still match the stored row, and the ``UPDATE`` itself is keyed on
the stored ``worktree_path``, so a concurrent writer that moved or
replaced the binding between audit and apply loses the race instead of
having its value silently overwritten. A mismatch raises and mutates
nothing.
``expected_path`` is mandatory: a retirement that does not name the path
it intends to clear cannot be safe against concurrent recreation.
"""
now_s = _ts()
expected_norm = _realpath_or_raw(expected_path)
if not expected_norm:
raise ControlPlaneError(
f"cannot retire lease {lease_id} worktree_path: expected_path is "
"required for compare-and-swap retirement (fail closed)"
)
with self._tx() as conn:
cols = self._lease_columns(conn)
row = conn.execute(
"SELECT * FROM leases WHERE lease_id = ?",
(lease_id,),
).fetchone()
if not row:
raise ControlPlaneError(f"unknown lease_id {lease_id}")
record = dict(row)
current_wt = (record.get("worktree_path") or "").strip()
if not current_wt:
# Idempotent: the binding this caller audited is already gone.
return {
"lease_id": lease_id,
"retired": False,
"already_retired": True,
"prior_worktree_path": "",
"expected_worktree_path": expected_path,
"reason": reason,
"compare_and_swap": {
"matched": True,
"outcome": "already_retired",
},
}
if _realpath_or_raw(current_wt) != expected_norm:
raise ControlPlaneError(
f"cannot retire lease {lease_id} worktree_path: expected "
f"'{expected_path}' does not match current '{current_wt}' "
"(fail closed)"
)
for field, expected_value in (
("status", expected_status),
("session_id", expected_session_id),
):
if expected_value is None:
continue
current_value = record.get(field)
if str(current_value or "").strip() != str(expected_value).strip():
raise ControlPlaneError(
f"cannot retire lease {lease_id} worktree_path: lease "
f"{field} changed since audit (expected "
f"'{expected_value}', found '{current_value}'); fail closed"
)
if expected_owner_pid is not None:
current_pid = record.get("owner_pid")
if current_pid is not None and int(current_pid) != int(expected_owner_pid):
raise ControlPlaneError(
f"cannot retire lease {lease_id} worktree_path: lease "
f"owner_pid changed since audit (expected "
f"{expected_owner_pid}, found {current_pid}); fail closed"
)
# Parse and update provenance_json
raw_prov = record.get("provenance_json") or "{}"
try:
prov = json.loads(raw_prov) if isinstance(raw_prov, str) else dict(raw_prov)
except Exception:
prov = {}
if not isinstance(prov, dict):
prov = {}
prior_path = current_wt
prov.update({
"worktree_path_retired": True,
"retired_worktree_path": prior_path,
"retired_at": now_s,
"retirement_reason": reason,
"retired_from_status": record.get("status"),
"retired_from_session_id": record.get("session_id"),
"worktree_path": "",
})
prov_json = json.dumps(prov)
if "worktree_path" in cols:
# CAS: keyed on the exact stored path this caller audited.
cur = conn.execute(
"UPDATE leases SET worktree_path = '', provenance_json = ? "
"WHERE lease_id = ? AND worktree_path = ?",
(prov_json, lease_id, record.get("worktree_path")),
)
if cur.rowcount != 1:
raise ControlPlaneError(
f"cannot retire lease {lease_id} worktree_path: "
"compare-and-swap matched no row (concurrent change); "
"fail closed"
)
else:
conn.execute(
"UPDATE leases SET provenance_json = ? WHERE lease_id = ?",
(prov_json, lease_id),
)
conn.execute(
"""
INSERT INTO events(work_item_id, event_type, message, created_at)
VALUES (?, 'worktree_binding_retired', ?, ?)
""",
(
record["work_item_id"],
f"lease {lease_id} worktree_path '{prior_path}' retired: {reason}",
now_s,
),
)
return {
"lease_id": lease_id,
"retired": True,
"already_retired": False,
"prior_worktree_path": prior_path,
"expected_worktree_path": expected_path,
"retired_at": now_s,
"reason": reason,
"compare_and_swap": {
"matched": True,
"outcome": "retired",
"expected_status": expected_status,
"expected_session_id": expected_session_id,
"expected_owner_pid": expected_owner_pid,
},
}
def abandon_lease( def abandon_lease(
self, self,
*, *,
@@ -3229,123 +3270,3 @@ class ControlPlaneDB:
"live_lease_id": None if live_lease_id is None else str(live_lease_id), "live_lease_id": None if live_lease_id is None else str(live_lease_id),
"reconcile_action": "reconcile_required" if stale else "safe_to_resume", "reconcile_action": "reconcile_required" if stale else "safe_to_resume",
} }
def retire_session_checkpoint_worktree_path(
self,
session_id: str,
*,
checkpoint_id: str | None = None,
expected_path: str | None = None,
expected_status: str | None = None,
reason: str = "missing_worktree_path_retired",
) -> dict[str, Any]:
"""Retire a missing worktree_path from session_checkpoints (#970).
Compare-and-swap, mirroring :meth:`retire_lease_worktree_path` (#970
review 644 B3). ``expected_path`` names the exact stored path the caller
audited; the guarded ``UPDATE`` is keyed on that stored value, so a
checkpoint whose path was moved, replaced, or concurrently rewritten
after the audit is refused without mutation rather than blindly cleared.
Exactly one checkpoint row is targeted: by ``checkpoint_id`` when given,
otherwise by ``session_id``, which must identify a single row.
"""
now_s = _ts()
expected_norm = _realpath_or_raw(expected_path)
if not expected_norm:
raise ControlPlaneError(
"cannot retire session checkpoint worktree_path: expected_path "
"is required for compare-and-swap retirement (fail closed)"
)
if not checkpoint_id and not (session_id or "").strip():
raise ControlPlaneError(
"cannot retire session checkpoint worktree_path: checkpoint_id "
"or session_id is required (fail closed)"
)
with self._tx() as conn:
if checkpoint_id:
selector_sql = "SELECT * FROM session_checkpoints WHERE checkpoint_id = ?"
selector_params: tuple[Any, ...] = (checkpoint_id,)
selector_desc = f"checkpoint_id '{checkpoint_id}'"
else:
selector_sql = "SELECT * FROM session_checkpoints WHERE session_id = ?"
selector_params = (session_id,)
selector_desc = f"session_id '{session_id}'"
rows = [dict(r) for r in conn.execute(selector_sql, selector_params).fetchall()]
if not rows:
raise ControlPlaneError(
f"cannot retire session checkpoint worktree_path: no "
f"checkpoint matches {selector_desc} (fail closed)"
)
if len(rows) > 1:
raise ControlPlaneError(
f"cannot retire session checkpoint worktree_path: "
f"{selector_desc} matches {len(rows)} checkpoints; supply an "
"exact checkpoint_id (fail closed)"
)
record = rows[0]
target_checkpoint_id = record.get("checkpoint_id")
current_wt = (record.get("worktree_path") or "").strip()
if not current_wt:
# Idempotent: the binding this caller audited is already gone.
return {
"session_id": session_id,
"checkpoint_id": target_checkpoint_id,
"retired": False,
"already_retired": True,
"prior_worktree_path": "",
"expected_worktree_path": expected_path,
"reason": reason,
"compare_and_swap": {
"matched": True,
"outcome": "already_retired",
},
}
if _realpath_or_raw(current_wt) != expected_norm:
raise ControlPlaneError(
f"cannot retire session checkpoint worktree_path for "
f"{selector_desc}: expected '{expected_path}' does not match "
f"current '{current_wt}' (fail closed)"
)
if expected_status is not None:
current_status = record.get("status")
if str(current_status or "").strip() != str(expected_status).strip():
raise ControlPlaneError(
f"cannot retire session checkpoint worktree_path for "
f"{selector_desc}: status changed since audit (expected "
f"'{expected_status}', found '{current_status}'); fail closed"
)
cur = conn.execute(
"UPDATE session_checkpoints SET worktree_path = '', updated_at = ? "
"WHERE checkpoint_id = ? AND worktree_path = ?",
(now_s, target_checkpoint_id, record.get("worktree_path")),
)
if cur.rowcount != 1:
raise ControlPlaneError(
f"cannot retire session checkpoint worktree_path for "
f"{selector_desc}: compare-and-swap matched no row "
"(concurrent change); fail closed"
)
return {
"session_id": session_id,
"checkpoint_id": target_checkpoint_id,
"retired": True,
"already_retired": False,
"prior_worktree_path": current_wt,
"expected_worktree_path": expected_path,
"retired_at": now_s,
"reason": reason,
"compare_and_swap": {
"matched": True,
"outcome": "retired",
"expected_status": expected_status,
},
}
-5
View File
@@ -65,7 +65,6 @@ that gates each call, not which tools exist.
- `gitea_assess_work_issue_duplicate` - `gitea_assess_work_issue_duplicate`
- `gitea_assess_worktree_cleanup_integrity` - `gitea_assess_worktree_cleanup_integrity`
- `gitea_audit_config` - `gitea_audit_config`
- `gitea_audit_missing_worktree_bindings`
- `gitea_audit_runtime_recovery_contamination` - `gitea_audit_runtime_recovery_contamination`
- `gitea_audit_stable_branch_contamination` - `gitea_audit_stable_branch_contamination`
- `gitea_audit_worktree_cleanup` - `gitea_audit_worktree_cleanup`
@@ -127,20 +126,16 @@ that gates each call, not which tools exist.
- `gitea_post_heartbeat` - `gitea_post_heartbeat`
- `gitea_publish_unpublished_issue_branch` - `gitea_publish_unpublished_issue_branch`
- `gitea_quarantine_contaminated_review` - `gitea_quarantine_contaminated_review`
- `gitea_rebind_dirty_same_claimant_author_session`
- `gitea_reclaim_expired_workflow_lease` - `gitea_reclaim_expired_workflow_lease`
- `gitea_reconcile_after_restart`
- `gitea_reconcile_already_landed_pr` - `gitea_reconcile_already_landed_pr`
- `gitea_reconcile_issue_claims` - `gitea_reconcile_issue_claims`
- `gitea_reconcile_merged_cleanups` - `gitea_reconcile_merged_cleanups`
- `gitea_reconcile_missing_worktree_bindings`
- `gitea_reconcile_superseded_by_merged_pr` - `gitea_reconcile_superseded_by_merged_pr`
- `gitea_record_daemon_process_kill_attempt` - `gitea_record_daemon_process_kill_attempt`
- `gitea_record_irrecoverable_decision_lock_provenance` - `gitea_record_irrecoverable_decision_lock_provenance`
- `gitea_record_pre_review_command` - `gitea_record_pre_review_command`
- `gitea_record_shell_spawn_outcome` - `gitea_record_shell_spawn_outcome`
- `gitea_record_stable_branch_push_attempt` - `gitea_record_stable_branch_push_attempt`
- `gitea_recover_dirty_orphaned_issue_worktree`
- `gitea_recover_incomplete_bootstrap_lock` - `gitea_recover_incomplete_bootstrap_lock`
- `gitea_release_merger_pr_lease` - `gitea_release_merger_pr_lease`
- `gitea_release_reviewer_pr_lease` - `gitea_release_reviewer_pr_lease`
+5 -1
View File
@@ -19,7 +19,11 @@ The assessor classifies:
- **service_health** — process healthy / parity mutation-safe - **service_health** — process healthy / parity mutation-safe
- **clients** — connected client descriptors (optional inventory) - **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 - **checkpoints** — soft-depends on #660; skipped with reason when schema absent
- **leases** — live control-plane leases after restart - **leases** — live control-plane leases after restart
- **capabilities** — master-parity / stale-runtime (#610) - **capabilities** — master-parity / stale-runtime (#610)
-20
View File
@@ -1180,10 +1180,6 @@ RECOGNIZED_GITEA_ENV_KEYS = frozenset({
"GITEA_SERVER_PROVENANCE", "GITEA_SERVER_PROVENANCE",
"GITEA_AUTHOR_WORKTREE", "GITEA_AUTHOR_WORKTREE",
"GITEA_ACTIVE_WORKTREE", "GITEA_ACTIVE_WORKTREE",
"GITEA_REVIEWER_WORKTREE",
"GITEA_MERGER_WORKTREE",
"GITEA_CANONICAL_REPOSITORY_ROOT",
"GITEA_MCP_SESSION_STATE_TTL_HOURS",
"GITEA_DISABLE_KEYCHAIN", "GITEA_DISABLE_KEYCHAIN",
"GITEA_CONTROL_PLANE_DB", "GITEA_CONTROL_PLANE_DB",
"GITEA_DB_PATH", "GITEA_DB_PATH",
@@ -1193,22 +1189,6 @@ RECOGNIZED_GITEA_ENV_KEYS = frozenset({
"GITEA_IRRECOVERABLE_HMAC_SECRET", "GITEA_IRRECOVERABLE_HMAC_SECRET",
"GITEA_FORCE_MCP_RUNTIME_CHECK", "GITEA_FORCE_MCP_RUNTIME_CHECK",
"GITEA_FORCE_CLIENT_MANAGED", "GITEA_FORCE_CLIENT_MANAGED",
# #975: the client-identity inputs the server actually consumes at startup
# (CLIENT_NAME_ENV / CLIENT_INSTANCE_ENV / CLIENT_SESSION_ENV in
# gitea_mcp_server). Production read them while this allowlist omitted them,
# so the peer-env scan classified them as unsupported overrides and the
# capability resolver refused every mutation fleet-wide. Named individually
# on purpose: no prefix is added, so an unrecognised GITEA_* override is
# still refused exactly as it was before.
"GITEA_MCP_CLIENT",
"GITEA_MCP_CLIENT_INSTANCE",
"GITEA_MCP_CLIENT_SESSION",
# #975 review 652 B1: production also consumes HEARTBEAT_INTERVAL_ENV from
# mcp_worker_identity via gitea_mcp_server._start_worker_heartbeat. Omitting
# it reproduced the same unsupported-env → runtime_reconnect_required
# failure mode for the documented operator override. Named individually;
# no GITEA_* / GITEA_WORKER_* prefix is added.
"GITEA_WORKER_HEARTBEAT_INTERVAL_SECONDS",
}) })
RECOGNIZED_GITEA_ENV_PREFIXES = ( RECOGNIZED_GITEA_ENV_PREFIXES = (
+162 -356
View File
@@ -500,66 +500,10 @@ def _resolve_preflight_workspace_path(worktree_path: str | None = None) -> str:
return workspace return workspace
def _process_root_git_remote_url(remote_name: str) -> str | None: def _resolve_namespace_mutation_context(worktree_path: str | None = None) -> dict:
"""Best-effort local ``git remote get-url`` strictly inside ``PROJECT_ROOT``.
#973 (B8): Must derive expected repository identity from an authority
independent of the candidate configured canonical root.
"""
try:
proc = subprocess.run(
["git", "remote", "get-url", remote_name],
capture_output=True,
text=True,
cwd=PROJECT_ROOT,
)
if proc.returncode != 0:
return None
url = (proc.stdout or "").strip()
return url or None
except Exception:
return None
def _resolve_expected_repository_slug(
remote: str | None = None,
org: str | None = None,
repo: str | None = None,
) -> str | None:
"""Resolve expected repository slug from session context, parameters, or process root.
Must derive expected repository identity ONLY from trusted sources independent
of the candidate configured canonical root (#973 B8).
Bound session context takes precedence over request parameters so request-supplied
coordinates cannot replace a bound canonical root (#741 / #973 B9).
"""
bound = session_ctx.get_session_context() or {}
b_org = bound.get("org")
b_repo = bound.get("repository")
if b_org and b_repo:
return session_ctx.format_repository_slug(b_org, b_repo)
if org and repo:
return session_ctx.format_repository_slug(org, repo)
eff_remote = remote or bound.get("remote") or _effective_remote()
parsed = remote_repo_guard.parse_org_repo_from_remote_url(
_process_root_git_remote_url(eff_remote)
)
if parsed:
return session_ctx.format_repository_slug(parsed[0], parsed[1])
return None
def _resolve_namespace_mutation_context(
worktree_path: str | None = None,
remote: str | None = None,
) -> dict:
"""Canonical namespace workspace + repository root for guards (#460/#510/#706/#618).""" """Canonical namespace workspace + repository root for guards (#460/#510/#706/#618)."""
role = _effective_workspace_role() role = _effective_workspace_role()
configured_root, _source = _configured_canonical_root() configured_root, _source = _configured_canonical_root()
bound = session_ctx.get_session_context() or {}
eff_remote = remote or bound.get("remote") or _effective_remote()
expected_slug = _resolve_expected_repository_slug(eff_remote)
return nwb.resolve_namespace_mutation_context( return nwb.resolve_namespace_mutation_context(
role_kind=role, role_kind=role,
worktree_path=worktree_path, worktree_path=worktree_path,
@@ -572,12 +516,9 @@ def _resolve_namespace_mutation_context(
), ),
profile_name=get_profile().get("profile_name"), profile_name=get_profile().get("profile_name"),
configured_canonical_root=configured_root, configured_canonical_root=configured_root,
expected_slug=expected_slug,
remote=eff_remote,
) )
def _resolve_author_mutation_context(worktree_path: str | None = None) -> dict: def _resolve_author_mutation_context(worktree_path: str | None = None) -> dict:
"""Backward-compatible alias for namespace workspace context.""" """Backward-compatible alias for namespace workspace context."""
return _resolve_namespace_mutation_context(worktree_path) return _resolve_namespace_mutation_context(worktree_path)
@@ -679,9 +620,6 @@ def _preflight_workspace_details(worktree_path: str | None, dirty_files: list[st
inspected_root = _get_git_root(workspace) inspected_root = _get_git_root(workspace)
process_root = ctx["process_project_root"] process_root = ctx["process_project_root"]
canonical_root = ctx["canonical_repo_root"] canonical_root = ctx["canonical_repo_root"]
crr_assessment = ctx.get("canonical_root_assessment") or {}
resolved_slug = crr_assessment.get("resolved_slug")
expected_slug = ctx.get("expected_slug")
active_root = os.path.realpath(inspected_root or workspace) active_root = os.path.realpath(inspected_root or workspace)
if active_root == canonical_root: if active_root == canonical_root:
dirty_scope = "control checkout" dirty_scope = "control checkout"
@@ -711,10 +649,6 @@ def _preflight_workspace_details(worktree_path: str | None, dirty_files: list[st
"workspace_healthy": not bool( "workspace_healthy": not bool(
ctx.get("bound_worktree_missing") or ctx.get("author_worktree_block") ctx.get("bound_worktree_missing") or ctx.get("author_worktree_block")
), ),
"canonical_root_assessment": crr_assessment,
"expected_repository_slug": expected_slug,
"observed_repository_identity": resolved_slug,
"worktree_registration_result": ctx.get("in_git_worktree_list"),
} }
if not ctx["roots_aligned"]: if not ctx["roots_aligned"]:
details["workspace_root_mismatch"] = ( details["workspace_root_mismatch"] = (
@@ -724,7 +658,6 @@ def _preflight_workspace_details(worktree_path: str | None, dirty_files: list[st
return details return details
def _format_preflight_workspace_details(details: dict) -> str: def _format_preflight_workspace_details(details: dict) -> str:
parts = [ parts = [
f"MCP server process root: {details.get('mcp_server_process_root')}", f"MCP server process root: {details.get('mcp_server_process_root')}",
@@ -997,7 +930,10 @@ def _enforce_canonical_repository_root(
if not configured_value: if not configured_value:
return return
expected_slug = _resolve_expected_repository_slug(remote) bound = session_ctx.get_session_context() or {}
expected_slug = session_ctx.format_repository_slug(
bound.get("org"), bound.get("repository")
)
assessment = crr.assess_canonical_repository_root( assessment = crr.assess_canonical_repository_root(
configured_value=configured_value, configured_value=configured_value,
source=source, source=source,
@@ -1916,16 +1852,8 @@ def _verify_role_mutation_workspace(
if runtime_reasons: if runtime_reasons:
raise RuntimeError("; ".join(runtime_reasons)) raise RuntimeError("; ".join(runtime_reasons))
except Exception as exc: except Exception as exc:
# #975: ``_check_mcp_runtimes_diagnostics`` raises every reason it if "stale-runtime:" in str(exc):
# produces through this one RuntimeError, but only ``stale-runtime:`` raise RuntimeError(str(exc))
# was re-raised here — an ``unsupported-env:`` reason was swallowed
# while still failing the capability resolver, so this preflight and
# the resolver disagreed about the identical diagnostic. Both
# authoritative prefixes now propagate the same way. This can only ever
# widen what is refused, never widen what is permitted.
message = str(exc)
if any(prefix in message for prefix in RUNTIME_DIAGNOSTIC_HARD_PREFIXES):
raise RuntimeError(message)
pass pass
role = _effective_workspace_role() role = _effective_workspace_role()
@@ -1938,9 +1866,6 @@ def _verify_role_mutation_workspace(
# back to the install checkout and validated Gitea-Tools/branches/ instead # back to the install checkout and validated Gitea-Tools/branches/ instead
# of the target repository the namespace is actually bound to. # of the target repository the namespace is actually bound to.
_configured_root, _configured_source = _configured_canonical_root() _configured_root, _configured_source = _configured_canonical_root()
bound = session_ctx.get_session_context() or {}
eff_remote = remote or bound.get("remote") or _effective_remote()
expected_slug = _resolve_expected_repository_slug(eff_remote, org=org, repo=repo)
assessment = nwb.assess_namespace_mutation_workspace( assessment = nwb.assess_namespace_mutation_workspace(
role_kind=role, role_kind=role,
worktree_path=worktree_path, worktree_path=worktree_path,
@@ -1955,8 +1880,6 @@ def _verify_role_mutation_workspace(
profile_name=get_profile().get("profile_name"), profile_name=get_profile().get("profile_name"),
current_branch=git_state.get("current_branch"), current_branch=git_state.get("current_branch"),
configured_canonical_root=_configured_root, configured_canonical_root=_configured_root,
expected_slug=expected_slug,
remote=eff_remote,
) )
if assessment["block"]: if assessment["block"]:
raise RuntimeError( raise RuntimeError(
@@ -2261,7 +2184,6 @@ def _canonical_repository_slug(
process_project_root=PROJECT_ROOT, process_project_root=PROJECT_ROOT,
remote=remote, remote=remote,
require_binding=True, require_binding=True,
mode="derivation",
) )
slug = assessment.get("resolved_slug") slug = assessment.get("resolved_slug")
if assessment.get("block") or not slug: if assessment.get("block") or not slug:
@@ -3492,7 +3414,7 @@ def cleanup_in_progress_for_pr(
# ── Helpers ─────────────────────────────────────────────────────────────────── # ── Helpers ───────────────────────────────────────────────────────────────────
def _effective_remote(remote: str = "dadeschools") -> str: def _effective_remote(remote: str) -> str:
"""If remote is the default ('dadeschools') but the active profile base_url maps to a known remote, use that remote instead.""" """If remote is the default ('dadeschools') but the active profile base_url maps to a known remote, use that remote instead."""
try: try:
profile = get_profile() profile = get_profile()
@@ -13625,162 +13547,6 @@ def gitea_scan_already_landed_open_prs(
} }
@mcp.tool()
def gitea_audit_missing_worktree_bindings(
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
) -> dict:
"""Read-only: audit and classify worktree bindings whose paths are missing on disk (#970).
Correlates missing-path bindings from control-plane leases, session checkpoints,
and issue locks with repository, host, branch, issue/PR, session, and lease state.
Distinguishes deleted worktrees from moved paths, host/mount failures, and live
leases/sessions.
Args:
remote: Known instance 'dadeschools' or 'prgs'.
host: Override the Gitea host.
org: Override the owner/organization.
repo: Override the repository name.
Returns:
dict with audit counts, missing binding classifications, and resolution status.
"""
read_block = _profile_operation_gate("gitea.read")
if read_block:
return {
"success": False,
"performed": False,
"reasons": read_block,
"permission_report": _permission_block_report("gitea.read"),
}
import missing_worktree_reconcile
db, _ = _control_plane_db_or_error()
root = _canonical_local_git_root()
h, o, r = _resolve(remote, host, org, repo)
return missing_worktree_reconcile.audit_missing_worktree_bindings(
db,
project_root=root,
remote=remote,
org=o,
repo=r,
host=h,
)
@mcp.tool()
def gitea_reconcile_missing_worktree_bindings(
dry_run: bool = True,
# Deprecated: retained so callers that still pass it get an explicit deny.
operator_authorized: bool = False,
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
) -> dict:
"""Audit and safely resolve (retire) missing worktree path bindings (#970).
Identifies confirmed stale deleted worktree bindings, re-validates them
against a fresh authoritative read immediately before mutation to prevent
recreation and ownership races, and retires only the exact stale bindings
while preserving unrelated worktrees and Git metadata.
Apply mode (``dry_run=False``) is authorized **server-side** (#970 review
644 B2): it requires the reconciler-only ``gitea.branch.delete``
capability, the resolved ``reconcile_missing_worktree_bindings`` task
capability, and an active cleanup phase minted through
``gitea_authorize_reconciliation_cleanup_phase``. A caller-supplied
``operator_authorized`` is never authorization evidence (#709 F1 /
review 434) and is rejected outright. Dry-run stays available to any
``gitea.read`` profile and never mutates.
Args:
dry_run: If True (default), reports planned mutations without modifying state.
operator_authorized: Rejected. Authorization is a server-side artifact.
remote: Known instance 'dadeschools' or 'prgs'.
host: Override the Gitea host.
org: Override the owner/organization.
repo: Override the repository name.
Returns:
dict with before/after audit state, resolutions, and dimension status.
"""
import missing_worktree_reconcile
# Explicitly reject the self-assertable Boolean before anything else, so it
# can never combine with a legitimate gate to authorize a mutation.
if operator_authorized:
return {
"success": False,
"performed": False,
"reason": "operator_authorized_rejected",
"reasons": [missing_worktree_reconcile.OPERATOR_AUTHORIZED_REJECTION],
"exact_next_action": (
"authorize cleanup via gitea_authorize_reconciliation_cleanup_phase "
"from a reconciler profile, then re-run with dry_run=False"
),
}
read_block = _profile_operation_gate("gitea.read")
if read_block:
return {
"success": False,
"performed": False,
"reasons": read_block,
"permission_report": _permission_block_report("gitea.read"),
}
db, db_errs = _control_plane_db_or_error()
root = _canonical_local_git_root()
h, o, r = _resolve(remote, host, org, repo)
profile = get_profile()
role = (profile.get("role") or profile.get("role_kind") or "").strip().lower()
cleanup_authorization = None
if not dry_run:
cleanup_task = missing_worktree_reconcile.CLEANUP_TASK
required_permission = task_capability_map.required_permission(cleanup_task)
capability_blockers = _profile_operation_gate(required_permission)
role_matches = role == task_capability_map.required_role(cleanup_task)
cleanup_authorization = missing_worktree_reconcile.assess_cleanup_authorization(
role=role,
capability_blockers=capability_blockers,
operator_authorized=False,
task_capability_resolved=(not capability_blockers) and role_matches,
)
if not cleanup_authorization.get("authorized"):
return {
"success": False,
"performed": False,
"reason": "cleanup_authorization_required",
"reasons": cleanup_authorization.get("reasons") or [],
"cleanup_authorization": cleanup_authorization,
"permission_report": _permission_block_report(required_permission),
"exact_next_action": (
"resolve the reconcile_missing_worktree_bindings capability "
"from a reconciler profile and authorize cleanup via "
"gitea_authorize_reconciliation_cleanup_phase"
),
}
return missing_worktree_reconcile.reconcile_missing_worktree_bindings(
db,
project_root=root,
remote=remote,
org=o,
repo=r,
host=h,
dry_run=dry_run,
operator_authorized=False,
cleanup_authorization=cleanup_authorization,
)
@mcp.tool() @mcp.tool()
def gitea_audit_worktree_cleanup( def gitea_audit_worktree_cleanup(
remote: str = "dadeschools", remote: str = "dadeschools",
@@ -15355,17 +15121,22 @@ def _current_runtime_mode_report(refresh: bool = False) -> dict:
workspace_root = None workspace_root = None
aligned = None aligned = None
canonical_root = None canonical_root = None
resolved_slug = None
try: try:
ctx = _resolve_namespace_mutation_context(None) ctx = _resolve_namespace_mutation_context(None)
workspace_root = ctx.get("workspace_path") workspace_root = ctx.get("workspace_path")
canonical_root = ctx.get("canonical_repo_root") canonical_root = ctx.get("canonical_repo_root")
# Alignment keeps its established repository-level meaning (#615 F1):
# does this namespace target the repository the process is installed in,
# i.e. canonical_repo_root == process_project_root. It is deliberately
# NOT path equality between the task workspace and the process root --
# the global worktree rule requires task work to live in a branches/
# worktree, so that comparison would classify every correctly bound
# author, reviewer, and merger session as unsafe.
aligned = ctx.get("roots_aligned") aligned = ctx.get("roots_aligned")
crr_assessment = ctx.get("canonical_root_assessment") or {}
resolved_slug = crr_assessment.get("resolved_slug")
except Exception: except Exception:
# An unresolvable binding is reported as unknown alignment, never as
# proof of alignment.
aligned = None aligned = None
resolved_slug = None
try: try:
profile_name = get_profile()["profile_name"] profile_name = get_profile()["profile_name"]
except Exception: except Exception:
@@ -15378,7 +15149,6 @@ def _current_runtime_mode_report(refresh: bool = False) -> dict:
dirty_files=dirty_files, dirty_files=dirty_files,
active_task_workspace=workspace_root, active_task_workspace=workspace_root,
canonical_repository_root=canonical_root, canonical_repository_root=canonical_root,
repository_slug=resolved_slug,
workspace_roots_aligned=aligned, workspace_roots_aligned=aligned,
profile=profile_name, profile=profile_name,
declared_mode=stable_control_runtime.declared_runtime_mode(), declared_mode=stable_control_runtime.declared_runtime_mode(),
@@ -15722,10 +15492,6 @@ _WORKER_REGISTRY = None
_WORKER_IDENTITY: str | None = None _WORKER_IDENTITY: str | None = None
_WORKER_GENERATION: str | None = None _WORKER_GENERATION: str | None = None
_WORKER_REGISTRATION_ATTEMPTED = False _WORKER_REGISTRATION_ATTEMPTED = False
#: #975: the one heartbeat supervisor for this process's registration. One
#: process registers exactly one worker identity, so there is exactly one
#: supervisor and it is never replaced.
_WORKER_HEARTBEAT_SUPERVISOR = None
#: Env a client launcher may set to name itself and its session. Absent values #: 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 #: are reported as unknown; they are never guessed at, because guessing is what
@@ -15812,17 +15578,6 @@ def _active_worker_identity() -> str | None:
if outcome.get("registered"): if outcome.get("registered"):
_WORKER_IDENTITY = identity _WORKER_IDENTITY = identity
_WORKER_GENERATION = generation _WORKER_GENERATION = generation
# #975: registration is the only moment identity and fencing
# epoch are both known, so the heartbeat supervisor is started
# here. Without it ``last_heartbeat_at`` never left
# ``started_at`` and every healthy client lost ownership at the
# TTL.
_start_worker_heartbeat(
identity=identity,
fencing_epoch=outcome.get("fencing_epoch"),
hints=hints,
generation=generation,
)
return identity return identity
if not outcome.get("collision"): if not outcome.get("collision"):
return None return None
@@ -15833,78 +15588,6 @@ def _active_worker_identity() -> str | None:
return None return None
def _start_worker_heartbeat(
*,
identity: str,
fencing_epoch,
hints: dict,
generation: str | None,
) -> None:
"""Attach a heartbeat supervisor to the registration just created (#975).
Never raises: a supervisor that cannot start leaves the row exactly as
``register()`` wrote it, which is the pre-#975 behaviour, rather than
failing the tool call that happened to trigger lazy registration.
Under pytest the thread is deliberately not started. Tests drive
``WorkerHeartbeatSupervisor`` directly with an injected clock, so the suite
proves the lifecycle without leaving background sqlite writers behind.
"""
global _WORKER_HEARTBEAT_SUPERVISOR
if _WORKER_HEARTBEAT_SUPERVISOR is not None:
return
registry = _worker_registry()
if registry is None or fencing_epoch is None:
return
try:
supervisor = mcp_worker_identity.WorkerHeartbeatSupervisor(
registry,
worker_identity=identity,
fencing_epoch=int(fencing_epoch),
session_id=hints.get("session_id"),
generation_id=generation,
client_name=hints.get("client_name"),
pid=os.getpid(),
ttl_seconds=mcp_worker_identity.DEFAULT_HEARTBEAT_TTL_SECONDS,
interval_seconds=os.environ.get(
mcp_worker_identity.HEARTBEAT_INTERVAL_ENV
),
)
_WORKER_HEARTBEAT_SUPERVISOR = supervisor
if not mcp_daemon_guard.is_pytest_runtime():
supervisor.start()
except Exception:
return
def _worker_heartbeat_status() -> dict:
"""Read-only heartbeat observability for ``gitea_get_runtime_context`` (#975).
Reports the unsupervised case explicitly rather than omitting the key, so an
operator can tell "no supervisor" apart from "supervisor with no beats yet".
"""
supervisor = _WORKER_HEARTBEAT_SUPERVISOR
if supervisor is None:
return {
"supervised": False,
"running": False,
"reasons": [
"no worker heartbeat supervisor is attached to this process; the "
"registration is not being renewed and will go stale at its TTL"
],
}
try:
status = dict(supervisor.status())
except Exception as exc:
return {
"supervised": True,
"running": False,
"reasons": [f"heartbeat status unavailable: {type(exc).__name__}: {exc}"],
}
status.setdefault("reasons", [])
return status
def _active_role_kind_safe() -> str | None: def _active_role_kind_safe() -> str | None:
"""Best-effort role for the registry record; never raises into a tool call. """Best-effort role for the registry record; never raises into a tool call.
@@ -19546,11 +19229,6 @@ def gitea_get_runtime_context(
"fencing_epoch": provenance_assessment["fencing_epoch"], "fencing_epoch": provenance_assessment["fencing_epoch"],
"conflicting_live_sessions": provenance_assessment["conflicting_live_sessions"], "conflicting_live_sessions": provenance_assessment["conflicting_live_sessions"],
"provenance_assessment": provenance_assessment, "provenance_assessment": provenance_assessment,
# #975: whether this registration is actually being renewed. Read-only,
# and it grants nothing — ownership still comes from the attachment
# record above. It exists so "my heartbeat stopped" is diagnosable
# before the TTL turns it into session_attachment_missing.
"worker_heartbeat": _worker_heartbeat_status(),
"unconsumed_gitea_env": unconsumed_env, "unconsumed_gitea_env": unconsumed_env,
"preflight_ready": preflight["preflight_ready"], "preflight_ready": preflight["preflight_ready"],
"preflight_block_reasons": preflight["preflight_block_reasons"], "preflight_block_reasons": preflight["preflight_block_reasons"],
@@ -22059,17 +21737,6 @@ def gitea_route_task_session(
# self-recovery was removed from the read-only path (was _trigger_mcp_auto_restart). # self-recovery was removed from the read-only path (was _trigger_mcp_auto_restart).
# Recovery is owned exclusively by the IDE/client reconnect path. # Recovery is owned exclusively by the IDE/client reconnect path.
#: Every reason prefix ``_check_mcp_runtimes_diagnostics`` can emit. Callers
#: raise its reasons as one RuntimeError, and #975 found that the preflight
#: re-raise recognised only ``stale-runtime:``, silently dropping
#: ``unsupported-env:`` while the capability resolver still failed on it. This
#: lives beside the producer so a newly added reason family cannot be forgotten
#: by a distant re-raise predicate again.
RUNTIME_DIAGNOSTIC_HARD_PREFIXES: tuple[str, ...] = (
"stale-runtime:",
"unsupported-env:",
)
def _check_mcp_runtimes_diagnostics(task: str, matching_profiles: list[str]) -> list[str]: def _check_mcp_runtimes_diagnostics(task: str, matching_profiles: list[str]) -> list[str]:
"""Read-only: report missing or stale MCP runtimes (no config or process mutation). """Read-only: report missing or stale MCP runtimes (no config or process mutation).
@@ -24765,10 +24432,18 @@ def _run_post_restart_reconcile(
repo: str | None = None, repo: str | None = None,
mode: str | None = None, mode: str | None = None,
limit: int = 200, limit: int = 200,
apply_session_cleanup: bool = False,
) -> dict: ) -> 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 global _POST_RESTART_LAST_PROOF, _POST_RESTART_BOOT_RAN
import post_restart_reconcile as prr import post_restart_reconcile as prr
import session_lifecycle as sl
try: try:
_h, o, r = _resolve(remote, None, org, repo) _h, o, r = _resolve(remote, None, org, repo)
@@ -24782,13 +24457,48 @@ def _run_post_restart_reconcile(
inventory = _gather_post_restart_inventory( inventory = _gather_post_restart_inventory(
remote=remote, org=o, repo=r, limit=limit 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( proof = prr.reconcile_after_restart(
inventory, inventory,
mode=mode or _post_restart_reconcile_mode(), mode=mode or _post_restart_reconcile_mode(),
) )
payload = proof.as_dict() payload = proof.as_dict()
payload["success"] = True payload["success"] = True
payload["read_only"] = True payload["read_only"] = not apply_session_cleanup
payload["remote"] = remote payload["remote"] = remote
payload["org"] = o payload["org"] = o
payload["repo"] = r payload["repo"] = r
@@ -24797,6 +24507,9 @@ def _run_post_restart_reconcile(
"proposed_follow_ups lists durable issues the apply path may create; " "proposed_follow_ups lists durable issues the apply path may create; "
"this tool never creates them (log-only by default, #662 rollout)" "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_LAST_PROOF = payload
_POST_RESTART_BOOT_RAN = True _POST_RESTART_BOOT_RAN = True
return payload return payload
@@ -24823,8 +24536,9 @@ def gitea_reconcile_after_restart(
repo: str | None = None, repo: str | None = None,
mode: str | None = None, mode: str | None = None,
limit: int = 200, limit: int = 200,
apply_session_cleanup: bool = False,
) -> dict: ) -> 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 Gathers live control-plane sessions, leases, worktree bindings, and
master-parity evidence, then classifies them with the pure master-parity evidence, then classifies them with the pure
@@ -24832,12 +24546,17 @@ def gitea_reconcile_after_restart(
machine-readable completion proof listing resolved / unresolved dimensions machine-readable completion proof listing resolved / unresolved dimensions
and proposed durable follow-up issues. and proposed durable follow-up issues.
Read-only by design: never restarts MCP, never auto-resumes write Never restarts MCP, never auto-resumes write mutations, and never creates
mutations, and never creates Gitea issues (those are a separate apply Gitea issues. Default mode is ``log_only``; set
path). Default mode is ``log_only``; set
``GITEA_POST_RESTART_RECONCILE_MODE=enforce`` (or pass ``mode='enforce'``) ``GITEA_POST_RESTART_RECONCILE_MODE=enforce`` (or pass ``mode='enforce'``)
to set ``mutation_hold`` when anything remains unresolved. 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 Soft-depends on #660 for session checkpoints: when the checkpoint schema
module is absent the checkpoints dimension is ``skipped`` with an explicit module is absent the checkpoints dimension is ``skipped`` with an explicit
reason rather than inventing a schema. reason rather than inventing a schema.
@@ -24857,9 +24576,96 @@ def gitea_reconcile_after_restart(
repo=repo, repo=repo,
mode=mode, mode=mode,
limit=limit, 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() @mcp.tool()
def gitea_inspect_workflow_lease( def gitea_inspect_workflow_lease(
lease_id: str, lease_id: str,
-386
View File
@@ -29,7 +29,6 @@ implements:
from __future__ import annotations from __future__ import annotations
import atexit
import hashlib import hashlib
import os import os
import re import re
@@ -124,103 +123,6 @@ DEFAULT_REGISTRY_PATH = os.path.expanduser(
#: worker within a single operator coffee break. #: worker within a single operator coffee break.
DEFAULT_HEARTBEAT_TTL_SECONDS = 900.0 DEFAULT_HEARTBEAT_TTL_SECONDS = 900.0
#: Optional operator override for the beat interval, in seconds. Clamped by
#: :func:`heartbeat_interval_for` so it can never be set at or above the TTL.
HEARTBEAT_INTERVAL_ENV = "GITEA_WORKER_HEARTBEAT_INTERVAL_SECONDS"
#: Never beat faster than this, so a misconfigured interval cannot turn the
#: supervisor into a busy sqlite writer.
MIN_HEARTBEAT_INTERVAL_SECONDS = 1.0
def heartbeat_interval_for(
ttl_seconds: float = DEFAULT_HEARTBEAT_TTL_SECONDS,
override: str | float | None = None,
) -> float:
"""Beat interval for *ttl_seconds*: one third of the TTL, capped at a half.
#975. A third means two consecutive beats can be lost before the row is
allowed to look stale, and the half-TTL cap is a hard ceiling so no
override can produce an interval that expires the registration it is
supposed to renew. That is the whole liveness contract: a healthy worker
stays owned, and a worker that has genuinely stopped beating still expires
on the configured TTL — this function never touches the TTL itself.
"""
try:
ttl = float(ttl_seconds)
except (TypeError, ValueError):
ttl = DEFAULT_HEARTBEAT_TTL_SECONDS
if not ttl > 0:
ttl = DEFAULT_HEARTBEAT_TTL_SECONDS
interval = ttl / 3.0
if override is not None:
candidate = str(override).strip()
if candidate:
try:
parsed = float(candidate)
except (TypeError, ValueError):
parsed = None
if parsed is not None and parsed > 0:
interval = parsed
ceiling = ttl / 2.0
if interval > ceiling:
interval = ceiling
if interval < MIN_HEARTBEAT_INTERVAL_SECONDS:
interval = min(MIN_HEARTBEAT_INTERVAL_SECONDS, ceiling)
return interval
#: Recorded-vs-presented pairs compared by :func:`_heartbeat_expectation_drift`.
_HEARTBEAT_EXPECTATION_FIELDS = (
"session_id",
"generation_id",
"client_name",
"pid",
)
def _heartbeat_expectation_drift(
record: dict[str, Any],
*,
expected_session_id: str | None = None,
expected_generation_id: str | None = None,
expected_client_name: str | None = None,
expected_pid: int | None = None,
) -> list[tuple[str, Any, Any]]:
"""Return ``(field, recorded, presented)`` for every mismatched expectation.
Only supplied expectations are compared, so a caller that presents nothing
gets the pre-#975 behaviour. ``client_name`` is compared through
:func:`normalize_client_name` so one application's namespaces stay one
client while genuinely different clients stay distinct.
"""
presented = {
"session_id": expected_session_id,
"generation_id": expected_generation_id,
"client_name": (
normalize_client_name(expected_client_name)
if expected_client_name is not None
else None
),
"pid": expected_pid,
}
drift: list[tuple[str, Any, Any]] = []
for field in _HEARTBEAT_EXPECTATION_FIELDS:
want = presented[field]
if want is None:
continue
got = record.get(field)
if field == "pid":
match = got is not None and int(got) == int(want)
else:
match = got == want
if not match:
drift.append((field, got, want))
return drift
STATUS_ACTIVE = "active" STATUS_ACTIVE = "active"
STATUS_SUPERSEDED = "superseded" STATUS_SUPERSEDED = "superseded"
STATUS_RELEASED = "released" STATUS_RELEASED = "released"
@@ -884,25 +786,11 @@ class WorkerRegistry:
worker_identity: str, worker_identity: str,
fencing_epoch: int, fencing_epoch: int,
now: datetime | None = None, now: datetime | None = None,
expected_session_id: str | None = None,
expected_generation_id: str | None = None,
expected_client_name: str | None = None,
expected_pid: int | None = None,
) -> dict[str, Any]: ) -> dict[str, Any]:
"""Renew only the owning registration (#948 AC11). """Renew only the owning registration (#948 AC11).
A stale epoch is refused rather than silently renewed, so a superseded A stale epoch is refused rather than silently renewed, so a superseded
session that resumes cannot heartbeat its way back into ownership. session that resumes cannot heartbeat its way back into ownership.
#975 adds optional keyword-only *expectations*. Every one supplied must
match the recorded row or the renewal is refused. They exist because
identity plus epoch cannot express "renew the row I registered, and only
that row": a recycled PID, or a second session of the same client, would
otherwise be renewable by the wrong beater. They are fencing tokens,
never assertions — supplying one can only cause a refusal, never grant
anything, and omitting them preserves the pre-#975 behaviour exactly.
Drift reports the existing ``BLOCKER_FENCED`` rather than a new
``blocker_kind``, because consumers switch on that value.
""" """
stamp = _ts(now or _utc_now()) stamp = _ts(now or _utc_now())
with self._tx() as conn: with self._tx() as conn:
@@ -919,30 +807,6 @@ class WorkerRegistry:
"reasons": [f"no registration for {worker_identity!r}"], "reasons": [f"no registration for {worker_identity!r}"],
} }
record = self._row_to_record(row) record = self._row_to_record(row)
drift = _heartbeat_expectation_drift(
record,
expected_session_id=expected_session_id,
expected_generation_id=expected_generation_id,
expected_client_name=expected_client_name,
expected_pid=expected_pid,
)
if drift:
return {
"success": False,
"renewed": False,
"mutation_performed": False,
"blocker_kind": BLOCKER_FENCED,
"expectation_drift": drift,
"reasons": [
"presented worker expectations do not match the recorded "
"registration, so this beater does not own the row: "
+ "; ".join(
f"{field} recorded {recorded!r}, presented {presented!r}"
for field, recorded, presented in drift
)
+ " (#975)"
],
}
if record["status"] != STATUS_ACTIVE: if record["status"] != STATUS_ACTIVE:
return { return {
"success": False, "success": False,
@@ -1566,253 +1430,3 @@ def resolve_bound_remote(
"explicitly to avoid host drift (#948)." "explicitly to avoid host drift (#948)."
], ],
} }
# --- Production heartbeat lifecycle (#975) --------------------------------
#
# #948 delivered ``WorkerRegistry.heartbeat()`` and nothing ever called it, so
# ``last_heartbeat_at`` stayed pinned to ``started_at`` for every production
# worker and ``heartbeat_ttl_seconds`` stopped being a liveness window at all —
# it became a hard cap on how long any client could stay attached. This
# supervisor is the missing caller.
#
# It is a daemon thread rather than an asyncio task or a per-request refresh
# because the renewal has to survive an *idle* session: a client waiting on a
# lease with no tool call in flight must not lose ownership, which rules out
# request-driven refresh. Registration itself happens lazily inside a tool call
# on the server's event loop, and the registry performs blocking
# ``BEGIN IMMEDIATE`` sqlite writes, which must not run on that loop.
# ``daemon=True`` is deliberate: a hard kill takes the thread down with the
# process, so a dead worker still goes stale on the normal TTL.
#: Refusals that mean this worker no longer owns its row. Beating again could
#: only ever be an attempt to renew ownership it has already lost, so the
#: supervisor stops permanently instead of retrying.
TERMINAL_HEARTBEAT_BLOCKERS = frozenset({BLOCKER_FENCED, BLOCKER_NO_ATTACHMENT})
class WorkerHeartbeatSupervisor:
"""Periodically renew exactly one worker registration.
Contract:
* It never registers. A supervisor exists only for an already-registered
identity, so it cannot create a second registration or a second identity
system.
* Every beat presents the full expectation set, so it can renew only the row
matching this exact client name, worker identity, session, generation and
pid.
* A terminal refusal (fenced or missing) stops it permanently and records
why. A fenced session must never beat its way back into ownership.
* A transient failure (a locked database, say) is counted and the loop
continues, so one contended write does not silently end the heartbeat.
* Nothing here raises into a caller. ``beat_once`` returns its outcome and
the thread body swallows everything, because a heartbeat failure must
degrade to "not renewed" and never crash a tool call.
"""
def __init__(
self,
registry: WorkerRegistry,
*,
worker_identity: str,
fencing_epoch: int,
session_id: str | None = None,
generation_id: str | None = None,
client_name: str | None = None,
pid: int | None = None,
ttl_seconds: float = DEFAULT_HEARTBEAT_TTL_SECONDS,
interval_seconds: float | None = None,
clock=None,
) -> None:
self._registry = registry
self.worker_identity = worker_identity
self.fencing_epoch = int(fencing_epoch)
self.session_id = session_id
self.generation_id = generation_id
self.client_name = (
normalize_client_name(client_name) if client_name is not None else None
)
self.pid = int(pid) if pid is not None else None
self.ttl_seconds = float(ttl_seconds)
self.interval_seconds = (
heartbeat_interval_for(self.ttl_seconds)
if interval_seconds is None
else heartbeat_interval_for(self.ttl_seconds, interval_seconds)
)
#: Injectable so every TTL test uses controlled time and no test waits
#: for a real interval or a real TTL to elapse.
self._clock = clock or _utc_now
self._stop_event = threading.Event()
self._thread: threading.Thread | None = None
self._lock = threading.Lock()
self._atexit_registered = False
self.started = False
self.stopped_reason: str | None = None
self.beats_attempted = 0
self.beats_renewed = 0
self.transient_failures = 0
self.last_beat_at: str | None = None
self.last_result: dict[str, Any] | None = None
# -- one beat --
def beat_once(self, now: datetime | None = None) -> dict[str, Any]:
"""Renew once. Never raises; returns the registry outcome or a failure."""
if self.stopped_reason is not None:
return {
"success": False,
"renewed": False,
"beat_attempted": False,
"reasons": [
f"supervisor already stopped: {self.stopped_reason}"
],
}
stamp = now or self._clock()
self.beats_attempted += 1
try:
result = self._registry.heartbeat(
worker_identity=self.worker_identity,
fencing_epoch=self.fencing_epoch,
now=stamp,
expected_session_id=self.session_id,
expected_generation_id=self.generation_id,
expected_client_name=self.client_name,
expected_pid=self.pid,
)
except Exception as exc:
# Transient by assumption: an unexpected error is not proof this
# worker lost ownership, so it must not silently end the heartbeat.
self.transient_failures += 1
result = {
"success": False,
"renewed": False,
"transient": True,
"blocker_kind": None,
"reasons": [f"heartbeat raised {type(exc).__name__}: {exc}"],
}
self.last_result = result
return result
self.last_result = result
if result.get("renewed"):
self.beats_renewed += 1
self.last_beat_at = result.get("last_heartbeat_at") or _ts(stamp)
return result
blocker = result.get("blocker_kind")
if blocker in TERMINAL_HEARTBEAT_BLOCKERS:
self._stop_internal(
reason=(
f"refused with blocker_kind={blocker!r}: "
+ "; ".join(result.get("reasons") or [])
)
)
else:
self.transient_failures += 1
return result
# -- lifecycle --
def start(self) -> dict[str, Any]:
"""Start the beat thread. Idempotent; safe to call from a tool call."""
with self._lock:
if self.stopped_reason is not None:
return {
"started": False,
"reasons": [f"supervisor stopped: {self.stopped_reason}"],
}
if self._thread is not None and self._thread.is_alive():
return {"started": True, "already_running": True, "reasons": []}
self._stop_event.clear()
thread = threading.Thread(
target=self._run,
name=f"gitea-worker-heartbeat-{self.worker_identity}",
daemon=True,
)
self._thread = thread
self.started = True
if not self._atexit_registered:
# Orderly shutdown stops the heartbeat. A hard kill does not
# run this, which is correct: the row must then go stale.
atexit.register(self._atexit_stop)
self._atexit_registered = True
thread.start()
return {"started": True, "already_running": False, "reasons": []}
def _run(self) -> None:
while not self._stop_event.is_set():
# Wait first: registration already stamped a fresh heartbeat, so an
# immediate beat would be a redundant write on every server launch.
if self._stop_event.wait(self.interval_seconds):
return
try:
self.beat_once()
except Exception:
# beat_once is already total; this is the last-resort guard that
# keeps a supervisor thread from dying silently.
self.transient_failures += 1
if self.stopped_reason is not None:
return
def stop(self, reason: str = "stopped") -> dict[str, Any]:
"""Stop beating. Idempotent, and prompt because the loop waits on an Event."""
self._stop_internal(reason=reason)
thread = self._thread
if thread is not None and thread.is_alive():
thread.join(timeout=max(1.0, min(5.0, self.interval_seconds)))
return {
"stopped": True,
"reason": self.stopped_reason,
"thread_alive": bool(thread is not None and thread.is_alive()),
}
def _stop_internal(self, *, reason: str) -> None:
if self.stopped_reason is None:
self.stopped_reason = reason
self._stop_event.set()
if self._atexit_registered:
try:
atexit.unregister(self._atexit_stop)
except Exception:
pass
self._atexit_registered = False
def _atexit_stop(self) -> None:
try:
self.stop(reason="process exit")
except Exception:
pass
# -- observability --
def status(self) -> dict[str, Any]:
"""Read-only observability payload; safe to embed in a tool result."""
thread = self._thread
return {
"supervised": True,
"worker_identity": self.worker_identity,
"fencing_epoch": self.fencing_epoch,
"session_id": self.session_id,
"generation_id": self.generation_id,
"client_name": self.client_name,
"pid": self.pid,
"heartbeat_ttl_seconds": self.ttl_seconds,
"heartbeat_interval_seconds": self.interval_seconds,
"started": self.started,
"running": bool(
thread is not None
and thread.is_alive()
and self.stopped_reason is None
),
"stopped_reason": self.stopped_reason,
"beats_attempted": self.beats_attempted,
"beats_renewed": self.beats_renewed,
"transient_failures": self.transient_failures,
"last_heartbeat_at": self.last_beat_at,
"last_blocker_kind": (self.last_result or {}).get("blocker_kind"),
}
File diff suppressed because it is too large Load Diff
+17 -143
View File
@@ -8,10 +8,8 @@ poison workspace purity checks in another namespace.
from __future__ import annotations from __future__ import annotations
import os import os
import subprocess
import author_mutation_worktree as amw import author_mutation_worktree as amw
import canonical_repository_root as crr
ACTIVE_WORKTREE_ENV = amw.ACTIVE_WORKTREE_ENV ACTIVE_WORKTREE_ENV = amw.ACTIVE_WORKTREE_ENV
AUTHOR_WORKTREE_ENV = amw.AUTHOR_WORKTREE_ENV AUTHOR_WORKTREE_ENV = amw.AUTHOR_WORKTREE_ENV
@@ -154,60 +152,6 @@ def resolve_namespace_workspace(
return os.path.realpath(process_project_root), "MCP server process root (default)" return os.path.realpath(process_project_root), "MCP server process root (default)"
def verify_git_common_directory_membership(
workspace_path: str,
canonical_repo_root: str,
) -> tuple[bool, str | None]:
"""Verify that workspace_path belongs to canonical_repo_root via git common-dir or branches containment."""
ws = (workspace_path or "").strip()
root = (canonical_repo_root or "").strip()
if not ws or not root:
return False, "empty workspace or canonical root path"
try:
real_ws = os.path.realpath(os.path.abspath(ws))
real_root = os.path.realpath(os.path.abspath(root))
except Exception as exc:
return False, f"invalid workspace or root path: {exc}"
if real_ws == real_root:
return True, None
if not os.path.isdir(real_ws):
return False, f"workspace directory '{real_ws}' does not exist"
try:
res = subprocess.run(
["git", "-C", real_ws, "rev-parse", "--git-common-dir"],
capture_output=True,
text=True,
check=False,
)
if res.returncode == 0:
common_raw = (res.stdout or "").strip()
common_dir = amw._realpath_git_common_dir(real_ws, common_raw)
real_common = os.path.realpath(common_dir)
canonical_git = os.path.realpath(os.path.join(real_root, ".git"))
if real_common in (canonical_git, real_root):
return True, None
return (
False,
f"workspace '{real_ws}' git common directory '{real_common}' does not match "
f"canonical repository root '{real_root}' (.git at '{canonical_git}')"
)
else:
return (
False,
f"workspace '{real_ws}' is not a valid git repository or git rev-parse failed"
)
except Exception as exc:
return (
False,
f"failed to inspect git common directory for workspace '{real_ws}': {exc}"
)
def resolve_namespace_mutation_context( def resolve_namespace_mutation_context(
*, *,
role_kind: str, role_kind: str,
@@ -219,8 +163,6 @@ def resolve_namespace_mutation_context(
worktree: str | None = None, worktree: str | None = None,
profile_name: str | None = None, profile_name: str | None = None,
configured_canonical_root: str | None = None, configured_canonical_root: str | None = None,
expected_slug: str | None = None,
remote: str | None = None,
) -> dict: ) -> dict:
"""Shared workspace resolution for runtime_context and mutation guards. """Shared workspace resolution for runtime_context and mutation guards.
@@ -238,41 +180,11 @@ def resolve_namespace_mutation_context(
env_map = env if env is not None else os.environ env_map = env if env is not None else os.environ
process_root = os.path.realpath(process_project_root) process_root = os.path.realpath(process_project_root)
role = normalize_role_kind(role_kind, profile_name=profile_name) role = normalize_role_kind(role_kind, profile_name=profile_name)
configured = (configured_canonical_root or "").strip()
configured_val = (configured_canonical_root or "").strip() if configured:
if configured_val: canonical_root = os.path.realpath(configured)
crr_assessment = crr.assess_canonical_repository_root(
configured_value=configured_val,
source="configured_canonical_root",
expected_slug=expected_slug,
process_project_root=process_root,
remote=remote,
require_binding=True,
)
# #973 B10: a refused repository-authority mode deliberately resolves no
# canonical root, so downstream guards keep evaluating the install
# checkout rather than an identity derived through an undefined mode.
# ``roots_aligned`` still follows ``proven`` and stays False, and the
# refusal (with its reason_code) rides along in
# ``canonical_root_assessment`` — this is a fail-closed fallback, never a
# normalisation of the mode.
canonical_root = crr_assessment["canonical_repo_root"] or amw.resolve_canonical_repo_root(
process_root, process_root
)
roots_aligned = crr_assessment["proven"]
else: else:
crr_assessment = { canonical_root = amw.resolve_canonical_repo_root(process_root, process_root)
"proven": True,
"block": False,
"reasons": [],
"configured": False,
"canonical_repo_root": amw.resolve_canonical_repo_root(process_root, process_root),
"resolved_slug": None,
"source": None,
"reason_code": None,
}
canonical_root = crr_assessment["canonical_repo_root"]
roots_aligned = (canonical_root == process_root)
durable: dict | None = None durable: dict | None = None
if role == "author": if role == "author":
@@ -322,10 +234,7 @@ def resolve_namespace_mutation_context(
"ignored_bindings": demotions + (pollution.get("ignored_bindings") or []), "ignored_bindings": demotions + (pollution.get("ignored_bindings") or []),
"process_project_root": process_root, "process_project_root": process_root,
"canonical_repo_root": canonical_root, "canonical_repo_root": canonical_root,
"roots_aligned": roots_aligned, "roots_aligned": canonical_root == process_root,
"canonical_root_assessment": crr_assessment,
"expected_slug": expected_slug,
"remote": remote,
} }
if durable is not None: if durable is not None:
result["author_worktree_resolution"] = durable result["author_worktree_resolution"] = durable
@@ -337,15 +246,6 @@ def resolve_namespace_mutation_context(
result["author_worktree_reasons"] = list(durable.get("reasons") or []) result["author_worktree_reasons"] = list(durable.get("reasons") or [])
result["author_worktree_blocker_kind"] = durable.get("blocker_kind") result["author_worktree_blocker_kind"] = durable.get("blocker_kind")
result["operator_recovery"] = durable.get("operator_recovery") result["operator_recovery"] = durable.get("operator_recovery")
else:
path_exists = os.path.exists(workspace)
result["path_exists"] = path_exists
result["in_git_worktree_list"] = (
amw.path_in_git_worktree_list(workspace, canonical_root)
if path_exists
else False
)
result["bound_worktree_missing"] = not path_exists
return result return result
@@ -478,8 +378,8 @@ def format_namespace_workspace_binding_error(
def assess_namespace_mutation_workspace( def assess_namespace_mutation_workspace(
*, *,
role_kind: str, role_kind: str,
worktree_path: str | None = None, worktree_path: str | None,
worktree: str | None = None, worktree: str | None,
process_project_root: str, process_project_root: str,
env: dict[str, str] | os._Environ | None = None, env: dict[str, str] | os._Environ | None = None,
session_lease_worktree: str | None = None, session_lease_worktree: str | None = None,
@@ -487,8 +387,6 @@ def assess_namespace_mutation_workspace(
profile_name: str | None = None, profile_name: str | None = None,
current_branch: str | None = None, current_branch: str | None = None,
configured_canonical_root: str | None = None, configured_canonical_root: str | None = None,
expected_slug: str | None = None,
remote: str | None = None,
) -> dict: ) -> dict:
"""Evaluate namespace workspace binding before preflight/mutation.""" """Evaluate namespace workspace binding before preflight/mutation."""
ctx = resolve_namespace_mutation_context( ctx = resolve_namespace_mutation_context(
@@ -501,8 +399,6 @@ def assess_namespace_mutation_workspace(
session_lock_worktree=session_lock_worktree, session_lock_worktree=session_lock_worktree,
profile_name=profile_name, profile_name=profile_name,
configured_canonical_root=configured_canonical_root, configured_canonical_root=configured_canonical_root,
expected_slug=expected_slug,
remote=remote,
) )
mutation_workspace = ctx["workspace_path"] mutation_workspace = ctx["workspace_path"]
binding_source = ctx["workspace_binding_source"] binding_source = ctx["workspace_binding_source"]
@@ -526,27 +422,6 @@ def assess_namespace_mutation_workspace(
reasons = list(metadata.get("reasons") or []) reasons = list(metadata.get("reasons") or [])
operator_recovery = ctx.get("operator_recovery") operator_recovery = ctx.get("operator_recovery")
crr_reasons = list(ctx.get("canonical_root_assessment", {}).get("reasons") or [])
if crr_reasons:
reasons.extend(crr_reasons)
path_exists = ctx.get("path_exists")
if path_exists is None:
path_exists = os.path.exists(mutation_workspace)
if not path_exists:
if role != "author":
reasons.append(
f"{role} mutation blocked: configured workspace directory '{mutation_workspace}' does not exist (nonexistent worktree)"
)
else:
valid_common, common_err = verify_git_common_directory_membership(
mutation_workspace, ctx["canonical_repo_root"]
)
if not valid_common and common_err:
reasons.append(common_err)
if role == "author": if role == "author":
# #618 durable resolution already validated existence, membership, # #618 durable resolution already validated existence, membership,
# branches/, lock ownership, and traversal safety when present. # branches/, lock ownership, and traversal safety when present.
@@ -563,28 +438,27 @@ def assess_namespace_mutation_workspace(
) )
if branches["block"]: if branches["block"]:
reasons.extend(branches["reasons"]) reasons.extend(branches["reasons"])
elif role in {"reviewer", "merger"}: elif (
if mutation_workspace == process_root and not amw.is_path_under_branches(mutation_workspace, ctx["canonical_repo_root"]): role == "reviewer"
and mutation_workspace == process_root
and not amw.is_path_under_branches(mutation_workspace, ctx["canonical_repo_root"])
):
reasons.append( reasons.append(
f"{role} mutation blocked: workspace is the stable control checkout; " f"{role} mutation blocked: workspace is the stable control checkout; "
f"create or reconnect to a session-owned worktree under branches/ " f"create or reconnect to a session-owned worktree under branches/ "
f"or set {ROLE_WORKTREE_ENVS.get(role, ACTIVE_WORKTREE_ENV)} / " f"or set {ROLE_WORKTREE_ENVS.get(role, ACTIVE_WORKTREE_ENV)} / "
f"{ACTIVE_WORKTREE_ENV}" f"{ACTIVE_WORKTREE_ENV}"
) )
elif mutation_workspace != process_root and not amw.is_path_under_branches(mutation_workspace, ctx["canonical_repo_root"]): elif (
role in {"reviewer", "merger"}
and mutation_workspace != process_root
and not amw.is_path_under_branches(mutation_workspace, ctx["canonical_repo_root"])
):
reasons.append( reasons.append(
f"{role} mutation blocked: workspace '{mutation_workspace}' is not under " f"{role} mutation blocked: workspace '{mutation_workspace}' is not under "
f"'{ctx['canonical_repo_root']}/branches/'" f"'{ctx['canonical_repo_root']}/branches/'"
) )
if path_exists and amw.is_path_under_branches(mutation_workspace, ctx["canonical_repo_root"]):
in_list = ctx.get("in_git_worktree_list")
if in_list is False:
reasons.append(
f"{role} mutation blocked: workspace '{mutation_workspace}' is under branches/ "
f"but is not registered in git worktree list for '{ctx['canonical_repo_root']}'"
)
block = bool(reasons) block = bool(reasons)
return { return {
"block": block, "block": block,
+51 -6
View File
@@ -475,8 +475,37 @@ 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)] sessions = [s for s in (inventory.get("sessions") or []) if isinstance(s, Mapping)]
leases_for_sessions = [
L for L in (inventory.get("leases") or []) if isinstance(L, Mapping)
]
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 = [ orphan_sessions = [
s s
for s in sessions for s in sessions
@@ -484,15 +513,27 @@ def reconcile_after_restart(
and s.get("pid") is not None and s.get("pid") is not None
and not lease_lifecycle.is_process_alive(s.get("pid")) and not lease_lifecycle.is_process_alive(s.get("pid"))
] ]
if orphan_sessions: 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( items.append(
_item( _item(
DIM_SESSIONS, DIM_SESSIONS,
ITEM_UNRESOLVED, 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={ 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), "total_sessions": len(sessions),
"fleet": fleet_details,
}, },
follow_up=True, follow_up=True,
) )
@@ -502,8 +543,12 @@ def reconcile_after_restart(
_item( _item(
DIM_SESSIONS, DIM_SESSIONS,
ITEM_RESOLVED, ITEM_RESOLVED,
f"{len(sessions)} session row(s) reconciled (no dead-pid orphans)", f"{len(sessions)} session row(s) reconciled "
details={"total_sessions": len(sessions)}, 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
),
}
+12 -23
View File
@@ -153,8 +153,9 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.branch.push", "permission": "gitea.branch.push",
"role": "author", "role": "author",
}, },
# #662: post-restart reconcile is read-only inventory + pure classification. # #662: post-restart reconcile is inventory + pure classification.
# Durable follow-up issue creation is a separate apply path (not this task). # #969: optional apply_session_cleanup retires confirmed-stale session rows
# through the same tool; durable follow-up Gitea issues remain separate.
"reconcile_after_restart": { "reconcile_after_restart": {
"permission": "gitea.read", "permission": "gitea.read",
"role": "author", "role": "author",
@@ -163,6 +164,15 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.read", "permission": "gitea.read",
"role": "author", "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. # #644: Phase 2 Web Console recovery tasks.
"clear_stale_binding": { "clear_stale_binding": {
"permission": "gitea.read", "permission": "gitea.read",
@@ -386,25 +396,6 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.branch.delete", "permission": "gitea.branch.delete",
"role": "reconciler", "role": "reconciler",
}, },
# #970: auditing missing worktree bindings is read-only; retiring one is a
# control-plane cleanup mutation and carries the same reconciler-only
# authority as any other reconciliation cleanup (review 644 B2).
"audit_missing_worktree_bindings": {
"permission": "gitea.read",
"role": "reconciler",
},
"gitea_audit_missing_worktree_bindings": {
"permission": "gitea.read",
"role": "reconciler",
},
"reconcile_missing_worktree_bindings": {
"permission": "gitea.branch.delete",
"role": "reconciler",
},
"gitea_reconcile_missing_worktree_bindings": {
"permission": "gitea.branch.delete",
"role": "reconciler",
},
"work_issue": { "work_issue": {
"permission": "gitea.pr.create", "permission": "gitea.pr.create",
"role": "author", "role": "author",
@@ -672,8 +663,6 @@ ROLE_EXCLUSIVE_TASKS: frozenset[str] = frozenset(
"delete_branch", "delete_branch",
"cleanup_merged_pr_branch", "cleanup_merged_pr_branch",
"reconciliation_cleanup", "reconciliation_cleanup",
"reconcile_missing_worktree_bindings",
"gitea_reconcile_missing_worktree_bindings",
"work_issue", "work_issue",
"work-issue", "work-issue",
} }
+2 -2
View File
@@ -37,7 +37,7 @@ class ControlPlaneDBTest(unittest.TestCase):
rows = dict(conn.execute("SELECT key, value FROM schema_meta").fetchall()) rows = dict(conn.execute("SELECT key, value FROM schema_meta").fetchall())
finally: finally:
conn.close() conn.close()
self.assertEqual(rows["schema_version"], "5") self.assertEqual(rows["schema_version"], "6")
self.assertIn("DB coordinates", rows["architecture"]) self.assertIn("DB coordinates", rows["architecture"])
self.assertIn("bridge", rows["architecture"].lower()) self.assertIn("bridge", rows["architecture"].lower())
@@ -868,7 +868,7 @@ class SessionCheckpointTest(unittest.TestCase):
conn.close() conn.close()
self.assertIn("session_checkpoints", names) self.assertIn("session_checkpoints", names)
record = self._write() record = self._write()
self.assertEqual(record["checkpoint_schema_version"], 5) self.assertEqual(record["checkpoint_schema_version"], 6)
# AC2 — checkpoints written for multi-role session fixtures. # AC2 — checkpoints written for multi-role session fixtures.
def test_multi_role_fixtures_each_get_a_row(self) -> None: def test_multi_role_fixtures_each_get_a_row(self) -> None:
@@ -313,7 +313,6 @@ class TestCanonicalRootGuardBinding(_ServerHarness):
process_project_root=self.install_root, process_project_root=self.install_root,
remote="prgs", remote="prgs",
require_binding=True, require_binding=True,
mode="derivation",
) )
self.assertFalse(got.get("block"), got.get("reasons")) self.assertFalse(got.get("block"), got.get("reasons"))
self.assertEqual(got["resolved_slug"], TARGET_SLUG) self.assertEqual(got["resolved_slug"], TARGET_SLUG)
@@ -241,10 +241,9 @@ class TestNamespaceContextUsesConfiguredRoot(unittest.TestCase):
process_project_root=self.install, process_project_root=self.install,
env={}, env={},
configured_canonical_root=self.target, configured_canonical_root=self.target,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
) )
self.assertEqual(ctx["canonical_repo_root"], self.target) self.assertEqual(ctx["canonical_repo_root"], self.target)
self.assertTrue(ctx["roots_aligned"]) self.assertFalse(ctx["roots_aligned"])
def test_target_worktree_is_member_of_target_root(self): def test_target_worktree_is_member_of_target_root(self):
got = nwb.amw.assess_workspace_repo_membership( got = nwb.amw.assess_workspace_repo_membership(
@@ -269,7 +268,6 @@ class TestNamespaceContextUsesConfiguredRoot(unittest.TestCase):
env={}, env={},
current_branch="feat/issue-1", current_branch="feat/issue-1",
configured_canonical_root=self.target, configured_canonical_root=self.target,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
) )
self.assertFalse(assessment["block"], assessment.get("reasons")) self.assertFalse(assessment["block"], assessment.get("reasons"))
self.assertEqual(assessment["canonical_repo_root"], self.target) self.assertEqual(assessment["canonical_repo_root"], self.target)
+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()
File diff suppressed because it is too large Load Diff
-739
View File
@@ -1,739 +0,0 @@
"""Regression tests for Issue #973 blocker B10: repository-authority mode contract.
Before this repair ``assess_canonical_repository_root`` never checked its
``mode`` argument against an allowlist. Both dispatch points were permissive:
* the configured-root path tested ``mode == "derivation"`` and sent every other
value into a catch-all ``else``, so an unsupported mode silently received
*validation* semantics, and
* the single-repository default path tested ``mode == "validation"``, so an
unsupported mode skipped the identity comparison entirely and was strictly
*weaker* than validation.
The measured consequence was that ``mode="invalid_mode"`` with matching expected
and observed identities returned ``proven: True`` / ``block: False`` with no
reasons, and that an unsupported mode passed on the default path where
``"validation"`` correctly blocked.
These tests exercise the production module directly with real git repositories —
no patched stand-in for the function under test — and drive the production
enforcement and mutation-context consumers rather than only the intermediate
assessment.
"""
from __future__ import annotations
import ast
import inspect
import os
import subprocess
import tempfile
import unittest
from unittest.mock import patch
import canonical_repository_root as crr
import gitea_config
import gitea_mcp_server as mcp_server
import namespace_workspace_binding as nwb
import stable_control_runtime
INSTALL_SLUG = "Scaled-Tech-Consulting/Gitea-Tools"
TARGET_SLUG = "Scaled-Tech-Consulting/mcp-control-plane"
FOREIGN_SLUG = "Someone-Else/Evil-Repo"
# Explicitly supplied values that must all be refused. Omission is *not* in this
# list: omitting the argument keeps the documented ``"validation"`` default.
UNSUPPORTED_STRING_MODES = (
"invalid_mode",
"",
"validaton", # misspelling
"derivaton", # misspelling
"Validation", # case variant
"DERIVATION", # case variant
" validation", # leading whitespace
"validation ", # trailing whitespace
"derivation\n", # trailing newline
"validation,derivation",
)
UNSUPPORTED_NON_STRING_MODES = (
None,
0,
1,
True,
False,
3.14,
[],
["validation"],
{},
{"mode": "validation"},
("validation",),
object(),
)
def _init_repo(path: str, remote_url: str, *, user: str = "Test User") -> None:
os.makedirs(path, exist_ok=True)
subprocess.run(["git", "init", "-b", "master"], cwd=path, check=True,
capture_output=True)
subprocess.run(["git", "config", "user.email", "[email protected]"], cwd=path,
check=True, capture_output=True)
subprocess.run(["git", "config", "user.name", user], cwd=path, check=True,
capture_output=True)
with open(os.path.join(path, "README.md"), "w") as handle:
handle.write(f"{os.path.basename(path)}\n")
subprocess.run(["git", "add", "README.md"], cwd=path, check=True,
capture_output=True)
subprocess.run(["git", "commit", "-m", "initial"], cwd=path, check=True,
capture_output=True)
subprocess.run(["git", "remote", "add", "prgs", remote_url], cwd=path,
check=True, capture_output=True)
class _CanonicalRootFixture(unittest.TestCase):
"""Real install / target / foreign git repositories, as in the #973 suite."""
def setUp(self):
self._tmp = tempfile.TemporaryDirectory()
self.tmp_dir = os.path.realpath(self._tmp.name)
self.install_root = os.path.join(self.tmp_dir, "Gitea-Tools")
_init_repo(
self.install_root,
"https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools.git",
)
self.target_root = os.path.join(self.tmp_dir, "mcp-control-plane")
_init_repo(
self.target_root,
"https://gitea.prgs.cc/Scaled-Tech-Consulting/mcp-control-plane.git",
)
self.evil_root = os.path.join(self.tmp_dir, "Evil-Repo")
_init_repo(
self.evil_root,
"https://gitea.prgs.cc/Someone-Else/Evil-Repo.git",
user="Evil User",
)
self.target_branches = os.path.join(self.target_root, "branches")
self.target_worktree = os.path.join(self.target_branches, "rev-pr-99")
subprocess.run(
["git", "worktree", "add", "-b", "rev-pr-99", self.target_worktree],
cwd=self.target_root, check=True, capture_output=True,
)
def tearDown(self):
self._tmp.cleanup()
def assertRefusedForMode(self, assessment: dict, mode) -> None:
"""Assert a fail-closed refusal attributable to *mode* and nothing else."""
self.assertFalse(assessment["proven"], assessment)
self.assertTrue(assessment["block"], assessment)
self.assertEqual(assessment["reason_code"], crr.DENY_UNKNOWN_MODE, assessment)
self.assertEqual(len(assessment["reasons"]), 1, assessment)
reason = assessment["reasons"][0]
self.assertIn("unsupported repository-authority mode", reason)
self.assertIn(repr(mode), reason)
# No trusted repository identity may be derived through an invalid mode.
self.assertIsNone(assessment["resolved_slug"], assessment)
self.assertIsNone(assessment["canonical_repo_root"], assessment)
class TestB10UnsupportedModeIsRejected(_CanonicalRootFixture):
"""Direct assessment tests for invalid-mode parsing."""
def test_invalid_string_with_missing_expected_identity(self):
"""G1: blocks for the mode, not incidentally for a missing identity."""
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=None,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="invalid_mode",
)
self.assertRefusedForMode(assessment, "invalid_mode")
self.assertFalse(
any("unprovable or missing" in r for r in assessment["reasons"]),
"must block because the mode is unsupported, not because the expected "
"identity happened to be missing",
)
def test_invalid_string_with_matching_identities(self):
"""G2: the contract violation — matching identities used to return proven."""
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="invalid_mode",
)
self.assertRefusedForMode(assessment, "invalid_mode")
def test_invalid_string_with_conflicting_identities(self):
"""G3: refused for the mode, not for the incidental identity mismatch."""
assessment = crr.assess_canonical_repository_root(
configured_value=self.evil_root,
source="env",
expected_slug=INSTALL_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="invalid_mode",
)
self.assertRefusedForMode(assessment, "invalid_mode")
self.assertFalse(
any("identity mismatch" in r for r in assessment["reasons"]),
"the identity comparison must not have run at all",
)
def test_empty_string_mode(self):
"""G4a: an explicitly supplied empty string is an unsupported value."""
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="",
)
self.assertRefusedForMode(assessment, "")
def test_explicit_none_mode(self):
"""G4b: explicit None is refused; it is not treated as omission."""
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode=None,
)
self.assertRefusedForMode(assessment, None)
self.assertIn("of type NoneType", assessment["reasons"][0])
def test_representative_non_string_modes(self):
for mode in UNSUPPORTED_NON_STRING_MODES:
with self.subTest(mode=repr(mode)):
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode=mode,
)
self.assertRefusedForMode(assessment, mode)
self.assertIn(
f"of type {type(mode).__name__}", assessment["reasons"][0]
)
def test_unknown_strings_misspellings_and_whitespace_variants(self):
for mode in UNSUPPORTED_STRING_MODES:
with self.subTest(mode=repr(mode)):
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode=mode,
)
self.assertRefusedForMode(assessment, mode)
def test_supported_modes_are_exactly_two(self):
self.assertEqual(
crr.SUPPORTED_MODES, ("validation", "derivation")
)
self.assertIsNone(crr.unsupported_mode_reason("validation"))
self.assertIsNone(crr.unsupported_mode_reason("derivation"))
self.assertIsNotNone(crr.unsupported_mode_reason("invalid_mode"))
class TestB10SupportedModesUnchanged(_CanonicalRootFixture):
"""The repair must not disturb the two documented modes."""
def test_omitted_mode_defaults_to_validation(self):
"""G4c/G4d: omission still selects validation, proven by both outcomes."""
matching = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
)
self.assertTrue(matching["proven"], matching)
self.assertFalse(matching["block"], matching)
self.assertIsNone(matching["reason_code"], matching)
self.assertEqual(matching["resolved_slug"], TARGET_SLUG)
conflicting = crr.assess_canonical_repository_root(
configured_value=self.evil_root,
source="env",
expected_slug=INSTALL_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
)
self.assertFalse(conflicting["proven"], conflicting)
self.assertTrue(conflicting["block"], conflicting)
self.assertIsNone(conflicting["reason_code"], conflicting)
self.assertTrue(
any("identity mismatch" in r for r in conflicting["reasons"]),
"omission must behave exactly like explicit validation",
)
def test_explicit_validation_retains_strict_behavior(self):
"""G5a/G5b/G5c."""
ok = crr.assess_canonical_repository_root(
configured_value=self.target_root, source="env",
expected_slug=TARGET_SLUG, process_project_root=self.install_root,
remote="prgs", require_binding=True, mode="validation",
)
self.assertTrue(ok["proven"], ok)
mismatch = crr.assess_canonical_repository_root(
configured_value=self.evil_root, source="env",
expected_slug=INSTALL_SLUG, process_project_root=self.install_root,
remote="prgs", require_binding=True, mode="validation",
)
self.assertTrue(mismatch["block"], mismatch)
self.assertTrue(any("identity mismatch" in r for r in mismatch["reasons"]))
unprovable = crr.assess_canonical_repository_root(
configured_value=self.target_root, source="env",
expected_slug=None, process_project_root=self.install_root,
remote="prgs", require_binding=True, mode="validation",
)
self.assertTrue(unprovable["block"], unprovable)
self.assertTrue(
any("unprovable or missing" in r for r in unprovable["reasons"])
)
def test_explicit_derivation_retains_trusted_derivation(self):
"""G5d: derivation still resolves identity with no expected slug."""
derived = crr.assess_canonical_repository_root(
configured_value=self.target_root, source="env",
expected_slug=None, process_project_root=self.install_root,
remote="prgs", require_binding=True, mode="derivation",
)
self.assertTrue(derived["proven"], derived)
self.assertFalse(derived["block"], derived)
self.assertIsNone(derived["reason_code"], derived)
self.assertEqual(derived["resolved_slug"], TARGET_SLUG)
def test_derivation_without_resolvable_remote_still_fails_closed(self):
no_remote = os.path.join(self.tmp_dir, "no-remote-target")
os.makedirs(no_remote)
subprocess.run(["git", "init", "-b", "master"], cwd=no_remote, check=True,
capture_output=True)
assessment = crr.assess_canonical_repository_root(
configured_value=no_remote, source="env", expected_slug=None,
process_project_root=self.install_root, remote="prgs",
require_binding=True, mode="derivation",
)
self.assertTrue(assessment["block"], assessment)
self.assertTrue(
any("no resolvable" in r for r in assessment["reasons"]), assessment
)
class TestB10SingleRepositoryDefaultPath(_CanonicalRootFixture):
"""G6: on the unconfigured path an invalid mode used to be weaker than validation."""
def _assess(self, mode_kwargs: dict) -> dict:
return crr.assess_canonical_repository_root(
configured_value=None,
source=None,
expected_slug=FOREIGN_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=False,
**mode_kwargs,
)
def test_validation_blocks_a_foreign_expected_identity(self):
"""G6b: the reference behaviour the invalid mode must not undercut."""
got = self._assess({"mode": "validation"})
self.assertFalse(got["proven"], got)
self.assertTrue(got["block"], got)
self.assertTrue(any("identity mismatch" in r for r in got["reasons"]))
def test_invalid_mode_no_longer_passes_where_validation_blocks(self):
"""G6a: identical inputs, only the mode differs — must not fail open."""
invalid = self._assess({"mode": "invalid_mode"})
self.assertRefusedForMode(invalid, "invalid_mode")
validation = self._assess({"mode": "validation"})
self.assertEqual(
invalid["proven"], validation["proven"],
"an unsupported mode must never be more permissive than validation",
)
self.assertTrue(invalid["block"] and validation["block"])
def test_empty_string_mode_on_default_path(self):
"""G6c."""
self.assertRefusedForMode(self._assess({"mode": ""}), "")
def test_omitted_mode_on_default_path_still_validates(self):
got = self._assess({})
self.assertFalse(got["proven"], got)
self.assertTrue(any("identity mismatch" in r for r in got["reasons"]), got)
class TestB10RejectionOrdering(_CanonicalRootFixture):
"""Refusal must precede every form of candidate-root or Git inspection."""
def test_no_git_or_identity_discovery_runs_for_an_unsupported_mode(self):
with patch.object(crr, "resolve_repo_toplevel") as toplevel, \
patch.object(crr, "repository_identity_slug") as identity:
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="invalid_mode",
)
self.assertRefusedForMode(assessment, "invalid_mode")
toplevel.assert_not_called()
identity.assert_not_called()
def test_the_same_spies_do_fire_for_a_supported_mode(self):
"""Control: proves the previous test's assertions are not vacuous."""
with patch.object(crr, "resolve_repo_toplevel",
wraps=crr.resolve_repo_toplevel) as toplevel, \
patch.object(crr, "repository_identity_slug",
wraps=crr.repository_identity_slug) as identity:
crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="validation",
)
toplevel.assert_called()
identity.assert_called()
def test_nonexistent_candidate_root_still_reports_the_mode_refusal(self):
"""No mocks: the existence check cannot have run before the refusal."""
nonexistent = os.path.join(self.tmp_dir, "no-such-repository")
assessment = crr.assess_canonical_repository_root(
configured_value=nonexistent,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="invalid_mode",
)
self.assertRefusedForMode(assessment, "invalid_mode")
self.assertFalse(
any("does not exist" in r for r in assessment["reasons"]), assessment
)
def test_symlinked_candidate_root_is_not_resolved_for_an_unsupported_mode(self):
link = os.path.join(self.tmp_dir, "target-alias")
os.symlink(self.target_root, link)
assessment = crr.assess_canonical_repository_root(
configured_value=link,
source="env",
expected_slug=TARGET_SLUG,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="invalid_mode",
)
self.assertRefusedForMode(assessment, "invalid_mode")
# The refusal payload must not leak a resolved path for the candidate.
self.assertIsNone(assessment["canonical_repo_root"], assessment)
class TestB10NoModeInjectionSurface(unittest.TestCase):
"""``mode`` must not be reachable from requests, environment, or config."""
PRODUCTION_MODULES = ("gitea_mcp_server.py", "namespace_workspace_binding.py")
def _repo_root(self) -> str:
return os.path.dirname(os.path.dirname(os.path.abspath(__file__)))
def test_production_call_sites_only_pass_allowlisted_literal_modes(self):
"""Valid hardcoded modes reach the correct path; nothing else is passed."""
seen: list[tuple[str, str | None]] = []
for name in self.PRODUCTION_MODULES:
path = os.path.join(self._repo_root(), name)
with open(path) as handle:
tree = ast.parse(handle.read())
for node in ast.walk(tree):
if not isinstance(node, ast.Call):
continue
func = node.func
called = (
func.attr if isinstance(func, ast.Attribute)
else getattr(func, "id", None)
)
if called != "assess_canonical_repository_root":
continue
supplied = [k for k in node.keywords if k.arg == "mode"]
if not supplied:
seen.append((name, None))
continue
value = supplied[0].value
self.assertIsInstance(
value, ast.Constant,
f"{name}: mode must be a literal, never a variable or expression",
)
self.assertIn(
value.value, crr.SUPPORTED_MODES,
f"{name}: unsupported mode literal {value.value!r}",
)
seen.append((name, value.value))
self.assertTrue(seen, "expected production call sites to be found")
# Both documented modes are exercised by production, and omission is used.
self.assertIn(None, [mode for _, mode in seen])
self.assertIn("derivation", [mode for _, mode in seen])
def test_no_public_entry_point_exposes_a_mode_parameter(self):
for func in (
nwb.resolve_namespace_mutation_context,
nwb.assess_namespace_mutation_workspace,
mcp_server._resolve_namespace_mutation_context,
mcp_server._enforce_canonical_repository_root,
mcp_server._canonical_repository_slug,
mcp_server._trusted_session_repository,
mcp_server._resolve_expected_repository_slug,
):
with self.subTest(func=func.__name__):
self.assertNotIn("mode", inspect.signature(func).parameters)
def test_no_environment_key_selects_a_repository_authority_mode(self):
for key in gitea_config.RECOGNIZED_GITEA_ENV_KEYS:
self.assertNotIn(
"CANONICAL_REPOSITORY_MODE", key.upper(),
f"{key} would expose a repository-authority mode selector",
)
# An invented mode-ish variable is simply not consumed by crr.
env = {
"GITEA_CANONICAL_REPOSITORY_ROOT": "/some/path",
"GITEA_CANONICAL_REPOSITORY_MODE": "invalid_mode",
}
value, source = crr.configured_canonical_root(None, env)
self.assertEqual(value, "/some/path")
self.assertNotIn("mode", (source or "").lower())
self.assertIn(
"GITEA_CANONICAL_REPOSITORY_MODE",
gitea_config.get_unconsumed_gitea_env_overrides(env),
"an unknown GITEA_* key must still be rejected as unrecognised",
)
def test_repository_configuration_carries_no_mode_field(self):
profile = {
"canonical_repository_root": "/some/path",
"mode": "invalid_mode",
}
value, source = crr.configured_canonical_root(profile, {})
self.assertEqual(value, "/some/path")
self.assertEqual(source, "profile canonical_repository_root")
class TestB10ProductionEnforcementPaths(_CanonicalRootFixture):
"""Production enforcement, mutation-context, reviewer and merger consumers."""
def _force_invalid_mode(self):
"""Simulate a future call site threading an unsupported mode.
The real ``assess_canonical_repository_root`` still executes — only the
caller-side argument is substituted — so the refusal under test is
produced by production code, not by a stand-in. No caller-controlled
``mode`` parameter is added to any production signature to achieve this.
"""
real = crr.assess_canonical_repository_root
def _wrapper(**kwargs):
kwargs["mode"] = "invalid_mode"
return real(**kwargs)
return patch.object(crr, "assess_canonical_repository_root", _wrapper)
def test_mutation_context_fails_closed_under_a_refused_mode(self):
with self._force_invalid_mode():
ctx = nwb.resolve_namespace_mutation_context(
role_kind="reviewer",
worktree_path=self.target_worktree,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug=TARGET_SLUG,
remote="prgs",
)
assessment = ctx["canonical_root_assessment"]
self.assertFalse(ctx["roots_aligned"], ctx)
self.assertTrue(assessment["block"], assessment)
self.assertEqual(assessment["reason_code"], crr.DENY_UNKNOWN_MODE)
self.assertIsNone(assessment["resolved_slug"])
# The refused mode must not yield the candidate root as canonical.
self.assertNotEqual(ctx["canonical_repo_root"], self.target_root)
self.assertEqual(ctx["canonical_repo_root"], self.install_root)
def test_mutation_context_unchanged_for_the_supported_default(self):
ctx = nwb.resolve_namespace_mutation_context(
role_kind="reviewer",
worktree_path=self.target_worktree,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug=TARGET_SLUG,
remote="prgs",
)
self.assertTrue(ctx["roots_aligned"], ctx)
self.assertEqual(ctx["canonical_repo_root"], self.target_root)
self.assertIsNone(ctx["canonical_root_assessment"]["reason_code"])
def test_reviewer_and_merger_authorization_fails_closed_under_a_refused_mode(self):
for role in ("reviewer", "merger"):
with self.subTest(role=role), self._force_invalid_mode():
assessment = nwb.assess_namespace_mutation_workspace(
role_kind=role,
worktree_path=self.target_worktree,
worktree=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug=TARGET_SLUG,
remote="prgs",
)
self.assertTrue(assessment["block"], assessment)
self.assertTrue(
any("unsupported repository-authority mode" in r
for r in assessment["reasons"]),
assessment,
)
def test_reviewer_and_merger_authorization_unchanged_for_supported_modes(self):
for role in ("reviewer", "merger"):
with self.subTest(role=role):
assessment = nwb.assess_namespace_mutation_workspace(
role_kind=role,
worktree_path=self.target_worktree,
worktree=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug=TARGET_SLUG,
remote="prgs",
)
self.assertFalse(assessment["block"], assessment)
def test_final_mutation_gate_blocks_when_the_mode_refusal_unaligns_roots(self):
with self._force_invalid_mode():
ctx = nwb.resolve_namespace_mutation_context(
role_kind="reviewer",
worktree_path=self.target_worktree,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug=TARGET_SLUG,
remote="prgs",
)
report = stable_control_runtime.build_runtime_report(
process_root=self.install_root,
checkout_branch="master",
runtime_head="abcdef123456",
active_task_workspace=ctx["workspace_path"],
canonical_repository_root=ctx["canonical_repo_root"],
workspace_roots_aligned=ctx["roots_aligned"],
)
gate = stable_control_runtime.assess_runtime_mutation_gate(report)
self.assertTrue(gate["block"], gate)
def test_enforce_canonical_repository_root_uses_the_validation_default(self):
"""B8 boundary intact: a foreign configured root raises through production."""
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server, "_configured_canonical_root",
return_value=(self.evil_root, "env")), \
patch.object(mcp_server.session_ctx, "get_session_context",
return_value=None):
with self.assertRaises(RuntimeError) as raised:
mcp_server._enforce_canonical_repository_root(remote="prgs")
self.assertIn("identity mismatch", str(raised.exception))
def test_enforce_canonical_repository_root_refuses_a_threaded_invalid_mode(self):
bound = {"org": "Scaled-Tech-Consulting",
"repository": "mcp-control-plane", "remote": "prgs"}
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server, "_configured_canonical_root",
return_value=(self.target_root, "env")), \
patch.object(mcp_server.session_ctx, "get_session_context",
return_value=bound), \
self._force_invalid_mode():
with self.assertRaises(RuntimeError) as raised:
mcp_server._enforce_canonical_repository_root(remote="prgs")
self.assertIn("unsupported repository-authority mode", str(raised.exception))
def test_enforce_canonical_repository_root_passes_for_a_valid_binding(self):
bound = {"org": "Scaled-Tech-Consulting",
"repository": "mcp-control-plane", "remote": "prgs"}
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server, "_configured_canonical_root",
return_value=(self.target_root, "env")), \
patch.object(mcp_server.session_ctx, "get_session_context",
return_value=bound), \
patch.object(mcp_server.session_ctx, "assess_session_context",
return_value={"block": False, "reasons": []}), \
patch.object(mcp_server, "get_profile",
return_value={"profile_name": "prgs-reviewer"}):
mcp_server._enforce_canonical_repository_root(remote="prgs")
def test_canonical_repository_slug_derivation_unbroken(self):
"""B9 boundary intact: legitimate cross-repository derivation still works."""
profile = {
"profile_name": "prgs-author",
"allowed_repositories": [TARGET_SLUG],
"canonical_repository_root": self.target_root,
}
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server.session_ctx, "get_session_context",
return_value=None):
slug, reasons = mcp_server._canonical_repository_slug(profile, "prgs")
self.assertEqual(slug, TARGET_SLUG, reasons)
self.assertEqual(reasons, [])
def test_canonical_repository_slug_fails_closed_under_a_refused_mode(self):
profile = {
"profile_name": "prgs-author",
"allowed_repositories": [TARGET_SLUG],
"canonical_repository_root": self.target_root,
}
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server.session_ctx, "get_session_context",
return_value=None), \
self._force_invalid_mode():
slug, reasons = mcp_server._canonical_repository_slug(profile, "prgs")
self.assertIsNone(slug)
self.assertTrue(
any("unsupported repository-authority mode" in r for r in reasons),
reasons,
)
result = mcp_server._trusted_session_repository(
profile, "prgs", for_mutation=True
)
self.assertIsNone(result["org"])
self.assertIsNone(result["repository"])
self.assertTrue(result["reasons"])
if __name__ == "__main__":
unittest.main()
@@ -1,546 +0,0 @@
"""Regression tests for Issue #973: validated cross-repository canonical roots."""
from __future__ import annotations
import os
import shutil
import tempfile
import unittest
from unittest.mock import patch, MagicMock
from pathlib import Path
import subprocess
import gitea_config
import namespace_workspace_binding as nwb
import canonical_repository_root as crr
import stable_control_runtime
import gitea_mcp_server as mcp_server
class TestIssue973RecognizedEnvKeys(unittest.TestCase):
"""Test recognized environment variable keys under #973."""
def test_recognized_gitea_env_keys(self):
for key in (
"GITEA_CANONICAL_REPOSITORY_ROOT",
"GITEA_REVIEWER_WORKTREE",
"GITEA_MERGER_WORKTREE",
"GITEA_MCP_SESSION_STATE_TTL_HOURS",
):
self.assertIn(key, gitea_config.RECOGNIZED_GITEA_ENV_KEYS)
def test_get_unconsumed_gitea_env_overrides_ignores_recognized(self):
env = {
"GITEA_CANONICAL_REPOSITORY_ROOT": "/some/path",
"GITEA_REVIEWER_WORKTREE": "/some/reviewer/path",
"GITEA_MERGER_WORKTREE": "/some/merger/path",
"GITEA_MCP_SESSION_STATE_TTL_HOURS": "24",
"GITEA_UNRECOGNIZED_FOO_VAR": "bar",
}
unconsumed = gitea_config.get_unconsumed_gitea_env_overrides(env)
self.assertNotIn("GITEA_CANONICAL_REPOSITORY_ROOT", unconsumed)
self.assertNotIn("GITEA_REVIEWER_WORKTREE", unconsumed)
self.assertNotIn("GITEA_MERGER_WORKTREE", unconsumed)
self.assertNotIn("GITEA_MCP_SESSION_STATE_TTL_HOURS", unconsumed)
self.assertIn("GITEA_UNRECOGNIZED_FOO_VAR", unconsumed)
class TestIssue973CrossRepoCanonicalRoots(unittest.TestCase):
"""Test workspace binding and canonical root validation for cross-repo namespaces."""
def setUp(self):
self._tmp = tempfile.TemporaryDirectory()
self.tmp_dir = os.path.realpath(self._tmp.name)
# Create simulated installation root
self.install_root = os.path.join(self.tmp_dir, "Gitea-Tools")
os.makedirs(self.install_root)
subprocess.run(["git", "init", "-b", "master"], cwd=self.install_root, check=True)
subprocess.run(["git", "config", "user.email", "[email protected]"], cwd=self.install_root, check=True)
subprocess.run(["git", "config", "user.name", "Test User"], cwd=self.install_root, check=True)
with open(os.path.join(self.install_root, "README.md"), "w") as f:
f.write("install\n")
subprocess.run(["git", "add", "README.md"], cwd=self.install_root, check=True)
subprocess.run(["git", "commit", "-m", "initial"], cwd=self.install_root, check=True)
# Create simulated target repository root
self.target_root = os.path.join(self.tmp_dir, "mcp-control-plane")
os.makedirs(self.target_root)
subprocess.run(["git", "init", "-b", "master"], cwd=self.target_root, check=True)
subprocess.run(["git", "config", "user.email", "[email protected]"], cwd=self.target_root, check=True)
subprocess.run(["git", "config", "user.name", "Test User"], cwd=self.target_root, check=True)
with open(os.path.join(self.target_root, "README.md"), "w") as f:
f.write("target\n")
subprocess.run(["git", "add", "README.md"], cwd=self.target_root, check=True)
subprocess.run(["git", "commit", "-m", "initial"], cwd=self.target_root, check=True)
# Add remotes to simulate real git repositories with identities
subprocess.run(["git", "remote", "add", "prgs", "https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools.git"], cwd=self.install_root, check=True)
subprocess.run(["git", "remote", "add", "prgs", "https://gitea.prgs.cc/Scaled-Tech-Consulting/mcp-control-plane.git"], cwd=self.target_root, check=True)
# Create simulated foreign repository root
self.evil_root = os.path.join(self.tmp_dir, "Evil-Repo")
os.makedirs(self.evil_root)
subprocess.run(["git", "init", "-b", "master"], cwd=self.evil_root, check=True)
subprocess.run(["git", "config", "user.email", "[email protected]"], cwd=self.evil_root, check=True)
subprocess.run(["git", "config", "user.name", "Evil User"], cwd=self.evil_root, check=True)
with open(os.path.join(self.evil_root, "README.md"), "w") as f:
f.write("evil\n")
subprocess.run(["git", "add", "README.md"], cwd=self.evil_root, check=True)
subprocess.run(["git", "commit", "-m", "initial"], cwd=self.evil_root, check=True)
subprocess.run(["git", "remote", "add", "prgs", "https://gitea.prgs.cc/Someone-Else/Evil-Repo.git"], cwd=self.evil_root, check=True)
# Create branches/ directory and a valid registered worktree in target repository
self.target_branches = os.path.join(self.target_root, "branches")
self.target_worktree = os.path.join(self.target_branches, "rev-pr-99")
subprocess.run(["git", "worktree", "add", "-b", "rev-pr-99", self.target_worktree], cwd=self.target_root, check=True)
def tearDown(self):
self._tmp.cleanup()
def test_valid_same_repository_configuration(self):
ctx = nwb.resolve_namespace_mutation_context(
role_kind="reviewer",
worktree_path=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=None,
)
self.assertEqual(ctx["canonical_repo_root"], self.install_root)
self.assertTrue(ctx["roots_aligned"])
self.assertTrue(ctx["canonical_root_assessment"]["proven"])
def test_valid_cross_repo_canonical_root(self):
ctx = nwb.resolve_namespace_mutation_context(
role_kind="reviewer",
worktree_path=self.target_worktree,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
)
self.assertEqual(ctx["canonical_repo_root"], self.target_root)
self.assertTrue(ctx["roots_aligned"])
self.assertTrue(ctx["canonical_root_assessment"]["proven"])
def test_expected_repository_identity_match(self):
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="test",
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
process_project_root=self.install_root,
remote="prgs",
)
self.assertTrue(assessment["proven"])
self.assertFalse(assessment["block"])
def test_foreign_repository_identity_mismatch(self):
assessment = crr.assess_canonical_repository_root(
configured_value=self.evil_root,
source="test",
expected_slug="Scaled-Tech-Consulting/Gitea-Tools",
process_project_root=self.install_root,
remote="prgs",
)
self.assertFalse(assessment["proven"])
self.assertTrue(assessment["block"])
self.assertTrue(any("identity mismatch" in r for r in assessment["reasons"]))
def test_native_repository_binding_mismatch(self):
ctx = nwb.resolve_namespace_mutation_context(
role_kind="reviewer",
worktree_path=self.target_worktree,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.evil_root,
expected_slug="Scaled-Tech-Consulting/Gitea-Tools",
remote="prgs",
)
self.assertFalse(ctx["roots_aligned"])
self.assertFalse(ctx["canonical_root_assessment"]["proven"])
self.assertTrue(any("identity mismatch" in r for r in ctx["canonical_root_assessment"]["reasons"]))
def test_unpatched_foreign_configured_root_derives_expected_from_process_root_and_blocks(self):
"""B8: Production path test where foreign configured root cannot self-authorize."""
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server, "_configured_canonical_root", return_value=(self.evil_root, "env")):
expected_slug = mcp_server._resolve_expected_repository_slug("prgs")
self.assertEqual(expected_slug, "Scaled-Tech-Consulting/Gitea-Tools")
assessment = crr.assess_canonical_repository_root(
configured_value=self.evil_root,
source="env",
expected_slug=expected_slug,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
)
self.assertFalse(assessment["proven"])
self.assertTrue(assessment["block"])
self.assertEqual(assessment["resolved_slug"], "Someone-Else/Evil-Repo")
self.assertTrue(any("identity mismatch" in r for r in assessment["reasons"]))
def test_unpatched_valid_cross_repo_matching_session_context(self):
"""B8: Valid cross-repo namespace matches when session context is bound to target repo."""
bound_ctx = {"org": "Scaled-Tech-Consulting", "repository": "mcp-control-plane", "remote": "prgs"}
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server.session_ctx, "get_session_context", return_value=bound_ctx):
expected_slug = mcp_server._resolve_expected_repository_slug("prgs")
self.assertEqual(expected_slug, "Scaled-Tech-Consulting/mcp-control-plane")
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=expected_slug,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
)
self.assertTrue(assessment["proven"])
self.assertFalse(assessment["block"])
self.assertEqual(assessment["resolved_slug"], "Scaled-Tech-Consulting/mcp-control-plane")
def test_unprovable_expected_identity_fails_closed(self):
"""B8: If expected repository identity is unprovable for a configured root, fail closed."""
no_remote_root = os.path.join(self.tmp_dir, "no-remote-process-root")
os.makedirs(no_remote_root)
subprocess.run(["git", "init", "-b", "master"], cwd=no_remote_root, check=True)
with patch.object(mcp_server, "PROJECT_ROOT", no_remote_root), \
patch.object(mcp_server.session_ctx, "get_session_context", return_value=None):
expected_slug = mcp_server._resolve_expected_repository_slug("prgs")
self.assertIsNone(expected_slug)
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=expected_slug,
process_project_root=no_remote_root,
remote="prgs",
require_binding=True,
)
self.assertFalse(assessment["proven"])
self.assertTrue(assessment["block"])
self.assertTrue(any("unprovable or missing" in r for r in assessment["reasons"]))
def test_missing_canonical_root(self):
ctx = nwb.resolve_namespace_mutation_context(
role_kind="author",
worktree_path=None,
process_project_root=self.install_root,
env={},
configured_canonical_root="",
)
self.assertEqual(ctx["canonical_repo_root"], self.install_root)
self.assertTrue(ctx["roots_aligned"])
def test_nonexistent_configured_canonical_root(self):
nonexistent = os.path.join(self.tmp_dir, "nonexistent-repo")
ctx = nwb.resolve_namespace_mutation_context(
role_kind="reviewer",
worktree_path=self.target_worktree,
process_project_root=self.install_root,
env={},
configured_canonical_root=nonexistent,
)
self.assertFalse(ctx["roots_aligned"])
self.assertFalse(ctx["canonical_root_assessment"]["proven"])
self.assertTrue(any("does not exist" in r for r in ctx["canonical_root_assessment"]["reasons"]))
def test_non_git_configured_canonical_root(self):
non_git = os.path.join(self.tmp_dir, "non-git-dir")
os.makedirs(non_git)
ctx = nwb.resolve_namespace_mutation_context(
role_kind="reviewer",
worktree_path=self.target_worktree,
process_project_root=self.install_root,
env={},
configured_canonical_root=non_git,
)
self.assertFalse(ctx["roots_aligned"])
self.assertFalse(ctx["canonical_root_assessment"]["proven"])
self.assertTrue(any("not a git repository" in r for r in ctx["canonical_root_assessment"]["reasons"]))
def test_git_common_directory_membership_matching(self):
valid, err = nwb.verify_git_common_directory_membership(
self.target_worktree, self.target_root
)
self.assertTrue(valid, err)
self.assertIsNone(err)
def test_foreign_git_common_directory(self):
# Foreign worktree created under install_root
install_branches = os.path.join(self.install_root, "branches")
foreign_wt = os.path.join(install_branches, "foreign-wt")
subprocess.run(["git", "worktree", "add", "-b", "foreign-wt", foreign_wt], cwd=self.install_root, check=True)
valid, err = nwb.verify_git_common_directory_membership(
foreign_wt, self.target_root
)
self.assertFalse(valid)
self.assertIn("does not match", err)
def test_normalized_path_aliases(self):
alias_path = self.target_worktree + "/../rev-pr-99/./"
valid, err = nwb.verify_git_common_directory_membership(
alias_path, self.target_root
)
self.assertTrue(valid, err)
def test_safe_symlink_identity(self):
link_path = os.path.join(self.target_branches, "symlink-rev-99")
try:
os.symlink(self.target_worktree, link_path)
valid, err = nwb.verify_git_common_directory_membership(
link_path, self.target_root
)
self.assertTrue(valid, err)
finally:
if os.path.exists(link_path):
os.unlink(link_path)
def test_symlink_escape_or_foreign_alias(self):
outside_dir = os.path.join(self.tmp_dir, "outside-target")
os.makedirs(outside_dir)
link_escape = os.path.join(self.target_branches, "escape-link")
try:
os.symlink(outside_dir, link_escape)
assessment = nwb.assess_namespace_mutation_workspace(
role_kind="reviewer",
worktree_path=link_escape,
worktree=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
)
self.assertTrue(assessment["block"])
finally:
if os.path.exists(link_escape):
os.unlink(link_escape)
def test_reviewer_worktree_registered_and_valid(self):
assessment = nwb.assess_namespace_mutation_workspace(
role_kind="reviewer",
worktree_path=self.target_worktree,
worktree=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
)
self.assertFalse(assessment["block"])
def test_reviewer_worktree_unregistered_blocks(self):
unreg_wt = os.path.join(self.target_branches, "unregistered-reviewer")
os.makedirs(unreg_wt)
assessment = nwb.assess_namespace_mutation_workspace(
role_kind="reviewer",
worktree_path=unreg_wt,
worktree=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
)
self.assertTrue(assessment["block"])
self.assertTrue(any("is not registered in git worktree list" in r for r in assessment["reasons"]))
def test_merger_worktree_registered_and_valid(self):
assessment = nwb.assess_namespace_mutation_workspace(
role_kind="merger",
worktree_path=self.target_worktree,
worktree=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
)
self.assertFalse(assessment["block"])
def test_merger_worktree_unregistered_blocks(self):
unreg_wt = os.path.join(self.target_branches, "unregistered-merger")
os.makedirs(unreg_wt)
assessment = nwb.assess_namespace_mutation_workspace(
role_kind="merger",
worktree_path=unreg_wt,
worktree=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
)
self.assertTrue(assessment["block"])
self.assertTrue(any("is not registered in git worktree list" in r for r in assessment["reasons"]))
def test_reviewer_or_merger_worktree_outside_branches_blocks(self):
outside_wt = os.path.join(self.target_root, "outside_branches_wt")
os.makedirs(outside_wt)
for r_kind in ("reviewer", "merger"):
assessment = nwb.assess_namespace_mutation_workspace(
role_kind=r_kind,
worktree_path=outside_wt,
worktree=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
)
self.assertTrue(assessment["block"])
self.assertTrue(any("is not under" in r for r in assessment["reasons"]))
def test_nonexistent_reviewer_or_merger_worktree_blocks(self):
nonexistent_wt = os.path.join(self.target_branches, "nonexistent-wt")
for r_kind in ("reviewer", "merger"):
assessment = nwb.assess_namespace_mutation_workspace(
role_kind=r_kind,
worktree_path=nonexistent_wt,
worktree=None,
process_project_root=self.install_root,
env={},
configured_canonical_root=self.target_root,
expected_slug="Scaled-Tech-Consulting/mcp-control-plane",
)
self.assertTrue(assessment["block"])
self.assertTrue(any("does not exist" in r for r in assessment["reasons"]))
def test_safe_and_unsafe_mutation_alignment_outcomes(self):
# Safe alignment (same repo)
report_safe = stable_control_runtime.build_runtime_report(
process_root=self.install_root,
checkout_branch="master",
runtime_head="abcdef123456",
active_task_workspace=self.install_root,
canonical_repository_root=self.install_root,
workspace_roots_aligned=True,
)
gate_safe = stable_control_runtime.assess_runtime_mutation_gate(report_safe)
self.assertFalse(gate_safe["block"])
# Unsafe alignment
report_unsafe = stable_control_runtime.build_runtime_report(
process_root=self.install_root,
checkout_branch="master",
runtime_head="abcdef123456",
active_task_workspace=self.target_worktree,
canonical_repository_root=self.target_root,
workspace_roots_aligned=False,
)
gate_unsafe = stable_control_runtime.assess_runtime_mutation_gate(report_unsafe)
self.assertTrue(gate_unsafe["block"])
def test_reviewer_lease_lifecycle_production_path(self):
"""B5: Automated regression for reviewer lease acquire and release through production path."""
mock_whoami = {
"authenticated": True,
"username": "sysadmin",
"remote": "prgs",
"profile": {
"profile_name": "prgs-reviewer",
"role": "reviewer",
"role_kind": "reviewer",
"allowed_operations": ["gitea.read", "gitea.pr.comment", "gitea.pr.approve", "gitea.pr.request_changes"],
"forbidden_operations": [],
},
}
def mock_api_request(method, url, auth=None, json_data=None):
if method == "GET":
return {"number": 99, "head": {"sha": "abc1234"}, "state": "open", "merged": False, "merged_at": None}
elif method == "POST":
return {"id": 9999, "body": (json_data or {}).get("body", "")}
return {}
with patch.object(mcp_server, "gitea_whoami", return_value=mock_whoami), \
patch.object(mcp_server, "get_profile", return_value=mock_whoami["profile"]), \
patch.object(mcp_server, "_effective_workspace_role", return_value="reviewer"), \
patch.object(mcp_server, "_configured_canonical_root", return_value=(self.target_root, "env")), \
patch.object(mcp_server, "_reviewer_session_worktree", return_value=self.target_worktree), \
patch.object(mcp_server, "_auth", return_value="token mock-token"), \
patch.object(mcp_server, "_fetch_pr_comments", return_value=[]), \
patch.object(mcp_server, "api_request", side_effect=mock_api_request), \
patch("reviewer_pr_lease.assess_acquire_lease", return_value={"acquire_allowed": True, "reasons": [], "lease_body": "<!-- LEASE -->"}), \
patch("reviewer_pr_lease.find_active_reviewer_lease", return_value={"session_id": "sid-123", "reviewer": "sysadmin"}), \
patch("reviewer_pr_lease.get_session_lease", return_value={"session_id": "sid-123", "reviewer": "sysadmin"}), \
patch("reviewer_pr_lease.clear_session_lease") as mock_clear:
acq_res = mcp_server.gitea_acquire_reviewer_pr_lease(
pr_number=99,
remote="prgs",
worktree=self.target_worktree,
org="Scaled-Tech-Consulting",
repo="mcp-control-plane",
)
self.assertTrue(acq_res.get("success"), acq_res)
rel_res = mcp_server.gitea_release_reviewer_pr_lease(
pr_number=99,
worktree=self.target_worktree,
remote="prgs",
org="Scaled-Tech-Consulting",
repo="mcp-control-plane",
)
self.assertTrue(rel_res.get("success"), rel_res)
mock_clear.assert_called_once()
def test_unbound_session_derives_and_retains_configured_target_repository(self):
"""B9: Unbound session with a legitimate configured target root derives and retains target repository."""
prof = {
"profile_name": "prgs-author",
"allowed_repositories": ["Scaled-Tech-Consulting/mcp-control-plane"],
"canonical_repository_root": self.target_root,
}
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server.session_ctx, "get_session_context", return_value=None):
slug, reasons = mcp_server._canonical_repository_slug(prof, "prgs")
self.assertEqual(slug, "Scaled-Tech-Consulting/mcp-control-plane", f"reasons: {reasons}")
self.assertEqual(reasons, [])
res = mcp_server._trusted_session_repository(prof, "prgs")
self.assertEqual(res["org"], "Scaled-Tech-Consulting")
self.assertEqual(res["repository"], "mcp-control-plane")
def test_process_root_a_configured_target_b_retains_b(self):
"""B9: Process root A (Gitea-Tools) and configured target B (mcp-control-plane) keep B as expected identity when bound."""
bound_ctx = {"org": "Scaled-Tech-Consulting", "repository": "mcp-control-plane", "remote": "prgs"}
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server.session_ctx, "get_session_context", return_value=bound_ctx):
expected = mcp_server._resolve_expected_repository_slug("prgs")
self.assertEqual(expected, "Scaled-Tech-Consulting/mcp-control-plane")
def test_request_supplied_repository_c_cannot_replace_derived_or_bound_repository_b(self):
"""B9: Request-supplied repository C (Timesheet) cannot replace bound repository B (mcp-control-plane)."""
bound_ctx = {"org": "Scaled-Tech-Consulting", "repository": "mcp-control-plane", "remote": "prgs"}
with patch.object(mcp_server, "PROJECT_ROOT", self.install_root), \
patch.object(mcp_server.session_ctx, "get_session_context", return_value=bound_ctx):
expected = mcp_server._resolve_expected_repository_slug("prgs", org="Scaled-Tech-Consulting", repo="Timesheet")
self.assertEqual(expected, "Scaled-Tech-Consulting/mcp-control-plane")
def test_known_expected_repo_a_plus_candidate_root_b_fails_closed(self):
"""B9: Known expected repository A plus candidate root B fails closed in validation mode."""
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug="Scaled-Tech-Consulting/Gitea-Tools",
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="validation",
)
self.assertFalse(assessment["proven"])
self.assertTrue(assessment["block"])
self.assertTrue(any("identity mismatch" in r for r in assessment["reasons"]))
def test_validation_mode_with_unprovable_expected_identity_fails_closed(self):
"""B9: Validation mode with expected_slug=None and require_binding=True fails closed."""
assessment = crr.assess_canonical_repository_root(
configured_value=self.target_root,
source="env",
expected_slug=None,
process_project_root=self.install_root,
remote="prgs",
require_binding=True,
mode="validation",
)
self.assertFalse(assessment["proven"])
self.assertTrue(assessment["block"])
self.assertTrue(any("unprovable or missing" in r for r in assessment["reasons"]))
if __name__ == "__main__":
unittest.main()
File diff suppressed because it is too large Load Diff
+2 -7
View File
@@ -245,15 +245,10 @@ class TestNamespaceWorkspaceIntegration(unittest.TestCase):
def test_pr487_style_merge_binds_clean_merger_workspace( def test_pr487_style_merge_binds_clean_merger_workspace(
self, _exists, _isdir, mock_run self, _exists, _isdir, mock_run
): ):
def mock_git(cmd, *args, **kwargs): mock_run.return_value = MagicMock(returncode=0, stdout=f"{CONTROL_ROOT}/.git\n")
if "rev-parse" in cmd:
return MagicMock(returncode=0, stdout=f"{CONTROL_ROOT}/.git\n")
return MagicMock(returncode=0, stdout=f"worktree {CONTROL_ROOT}\nworktree {MERGER_CLEAN}\n")
mock_run.side_effect = mock_git
os.environ[nwb.AUTHOR_WORKTREE_ENV] = AUTHOR_DIRTY os.environ[nwb.AUTHOR_WORKTREE_ENV] = AUTHOR_DIRTY
os.environ[nwb.MERGER_WORKTREE_ENV] = MERGER_CLEAN
srv._preflight_resolved_role = "reviewer" srv._preflight_resolved_role = "reviewer"
with mock.patch.object(srv, "PROJECT_ROOT", MCP_PROCESS_ROOT): with mock.patch.object(srv, "PROJECT_ROOT", MCP_PROCESS_ROOT):
with mock.patch("gitea_mcp_server.get_profile", return_value=self._merger_profile()): with mock.patch("gitea_mcp_server.get_profile", return_value=self._merger_profile()):
resolved = srv._verify_role_mutation_workspace("prgs") resolved = srv._verify_role_mutation_workspace("prgs")
self.assertEqual(resolved, os.path.realpath(MERGER_CLEAN)) self.assertEqual(resolved, os.path.realpath(MCP_PROCESS_ROOT))
@@ -153,12 +153,6 @@ EXPECTED_ROLE_EXCLUSIVE_TASKS = frozenset(
"delete_branch", "delete_branch",
"cleanup_merged_pr_branch", "cleanup_merged_pr_branch",
"reconciliation_cleanup", "reconciliation_cleanup",
# #970 review 644 B2: retiring a missing worktree binding is a
# control-plane cleanup mutation, so it carries the same reconciler-only
# authority as every other reconciliation cleanup. Permission alone must
# not authorize it.
"reconcile_missing_worktree_bindings",
"gitea_reconcile_missing_worktree_bindings",
"work_issue", "work_issue",
"work-issue", "work-issue",
} }