Files
Gitea-Tools/webui/inventory.py
sysadminandClaude Opus 4.8 1ca2b50406 fix(webui): authority-aware ownership + redacted contamination text (#641)
Addresses the two blockers from the PR #898 review at a81db754.

B1 - degraded ownership inventory was rendered as affirmative absence.
_build_session_rows read the sessions/leases/locks sections without
consulting their status, so a session row emitted lease_ids=() and
worktree_paths=() whether the session genuinely held nothing or the
lease store simply could not be read. The renderer printed both as
"none" and "unbound", contradicting the ownership_authority_complete
invariant documented on InventorySnapshot.

SessionRow now carries lease_authority and worktree_authority. A
worktree binding is correlated through lease work numbers, so it is
unproven when either the leases or the locks section fails to read --
this covers the narrow variant where locks hold real worktree paths but
a degraded leases section leaves work_numbers empty. The renderer emits
"unknown (inventory <status>)" with an authority-unproven badge instead
of none/unbound, the card names the unreadable sections, and an empty
session list from an unreadable sessions section no longer reads as
"no sessions recorded". snapshot_to_dict exports
ownership_authority_complete, ownership_section_status, and per-row
lease_authority / worktree_authority so /api/sessions consumers can
distinguish the two cases.

B2 - contamination payload strings bypassed redaction.
_inspect_contamination copied command_summary, session_id, role and
reason_class out of the marker payload with only str(), while every
inventory-sourced field on the same page arrives through
webui.inventory.scrub(). The write-time redactor
stable_branch_push_guard.redact_command is a narrow denylist that leaves
absolute $HOME paths, -H 'X-Api-Key: <value>', --password <value>, and
PRIVATE_KEY=<value> intact, and this is the first web surface to render
command_summary at all.

Adds webui.inventory.scrub_text(), which collapses $HOME and redacts
credential-shaped tokens and URL userinfo anywhere inside a string rather
than only at its start, and routes the marker payload through it. scrub()
and every existing caller are untouched. The command_summary field is
kept: it is the #630 evidence naming which daemon was killed. The module
docstring claiming absolute paths were already collapsed is corrected.

Also: removes the locks_by_session_hint dead loop and its discard (N1),
adds the missing trailing newline to webui/runtime_views.py (N4), drops
an unused dataclasses.field import, and documents both honesty rules in
docs/webui-local-dev.md.

Tests: tests/test_webui_sessions_view.py grows from 12 to 26 cases,
covering degraded and unavailable ownership sections in both the HTML and
JSON paths, the locks-readable/leases-degraded variant, a guard against
over-correcting clean inventory into "unknown", the previously untested
expired-lease flag, HTML escaping of hostile values in clean and degraded
renders, and each secret class the write-time denylist misses. The
STATUS_UNAVAILABLE import that was present but unused is now exercised.

Validation, from the issue worktree with venv/bin/python (Python 3.14.5,
pytest 9.1.1):

  pytest tests/test_webui_sessions_view.py tests/test_webui_*.py -q
    -> 512 passed, 376 subtests (was 498 / 372; +14 new tests)
  pytest tests/test_issue_854_semantic_container_exclusion.py -q
    -> 13 passed, 8 subtests
  13-file runtime/health/inventory/restart set
    -> 1 failed, 223 passed; the single failure is
       test_runtime_clarity.py::TestRuntimeClarity::
       test_activate_profile_succeeds_when_enabled, the identical test and
       assertion the reviewer recorded on master at 7af40fb5, so it is
       baseline-equivalent and not introduced here.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-25 01:25:28 -04:00

1018 lines
38 KiB
Python

"""Unified session/lease/lock/worktree inventory for the web console (#636).
Leases (#433), worktrees (#432), and runtime (#430) each ship their own MVP
view, each with its own shape and its own idea of what "owned" means. A
traffic-control or recovery operator has to read all three and correlate them
by hand, which is exactly the step that goes wrong under collision pressure.
This module aggregates them into one versioned, read-only snapshot so the
console, and any worker asking "what is safe to do next", read the same
inventory from the same authority.
Field authority is explicit and never blended. Every section declares where its
rows came from:
* ``control_plane_db`` — the #613 substrate: sessions, leases, assignments.
Authoritative for *exclusive ownership* (#600/#601).
* ``filesystem`` — durable per-issue lock files (:mod:`issue_lock_store`) and
registered git worktrees. Authoritative for *what exists on this machine*.
* ``gitea`` — remote issue/PR state, reached only through existing loaders.
Safety invariants:
* **Read-only.** The control-plane database is opened through a ``mode=ro``
URI. :class:`control_plane_db.ControlPlaneDB` creates directories and runs
migrations in its constructor, which an inventory read must never do, so this
module talks to sqlite directly rather than through that class.
* **Fail-soft, never fail-silent.** A subsystem that cannot be read degrades to
a section carrying ``status`` and ``reason``. It never raises, and it never
produces an empty list that reads like "nothing is there".
* **Never invent active ownership.** This is the invariant that matters most.
A degraded ownership source sets ``ownership_authority_complete`` false, and
while that flag is false no work item is reported unowned and no collision is
asserted. Absence of evidence is reported as absence of evidence.
* **Redaction at the boundary.** Absolute paths are collapsed against the home
directory, URLs lose userinfo and query strings, and no credential-shaped
value is emitted. No session token exists in these sources and none is read.
Phase 1 is read-only. Lease steal/release and worktree deletion are Phase 2+
and deliberately have no representation here, not even a disabled one.
"""
from __future__ import annotations
import os
import re
import sqlite3
import time
from dataclasses import dataclass, field
from datetime import datetime, timezone
from typing import Any, Callable
from urllib.parse import urlparse
import control_plane_db
import issue_lock_store
SCHEMA_VERSION = 1
API_VERSION = "v1"
#: Sections whose absence would make an ownership claim unprovable. If any of
#: these is not ``ok``, the snapshot refuses to describe anything as unowned.
OWNERSHIP_SECTIONS = ("sessions", "leases", "locks")
SECTION_NAMES = ("sessions", "leases", "locks", "worktrees", "namespaces")
STATUS_OK = "ok"
STATUS_DEGRADED = "degraded"
STATUS_UNAVAILABLE = "unavailable"
AUTHORITY_CONTROL_PLANE_DB = "control_plane_db"
AUTHORITY_FILESYSTEM = "filesystem"
AUTHORITY_GITEA = "gitea"
_CREDENTIAL_KEY_RE = re.compile(
r"(token|secret|password|passwd|api[_-]?key|authorization|bearer|credential)",
re.IGNORECASE,
)
_REDACTED = "[redacted]"
#: Credential-shaped *name* as it appears inside a free-form command line. This
#: is deliberately broader than :data:`_CREDENTIAL_KEY_RE` — it also matches a
#: bare ``key`` component, so ``PRIVATE_KEY=`` is caught. Over-redacting a
#: displayed string is safe; under-redacting one is not.
_TEXT_CREDENTIAL_NAME = (
r"[A-Za-z0-9_.\-]*"
r"(?:token|secret|password|passwd|key|authorization|bearer|credential)"
r"[A-Za-z0-9_.\-]*"
)
#: A value following such a name: single-quoted, double-quoted, or bare. The
#: bare form stops at a quote so an enclosing quote survives the redaction.
_TEXT_CREDENTIAL_VALUE = r"'[^']*'|\"[^\"]*\"|[^\s'\"]+"
_TEXT_CREDENTIAL_FLAG_RE = re.compile(
rf"(?P<key>(?<![\w\-])--?{_TEXT_CREDENTIAL_NAME})"
rf"(?P<sep>[=\s]+)"
rf"(?P<value>{_TEXT_CREDENTIAL_VALUE})",
re.IGNORECASE,
)
_TEXT_CREDENTIAL_ASSIGN_RE = re.compile(
rf"(?P<key>(?<![\w\-]){_TEXT_CREDENTIAL_NAME})"
rf"(?P<sep>\s*[:=]\s*)"
rf"(?P<value>{_TEXT_CREDENTIAL_VALUE})",
re.IGNORECASE,
)
_TEXT_URL_USERINFO_RE = re.compile(
r"(?P<scheme>\b[A-Za-z][A-Za-z0-9+.\-]*://)[^\s/@]+@"
)
@dataclass(frozen=True)
class InventorySection:
"""One subsystem's contribution, with its authority and health."""
name: str
authority: str
status: str
items: tuple[dict[str, Any], ...] = ()
reason: str | None = None
scan_ms: float | None = None
@property
def ok(self) -> bool:
return self.status == STATUS_OK
def to_dict(self) -> dict[str, Any]:
return {
"name": self.name,
"authority": self.authority,
"status": self.status,
"count": len(self.items),
"reason": self.reason,
"scan_ms": self.scan_ms,
"items": [dict(item) for item in self.items],
}
@dataclass(frozen=True)
class CollisionSignal:
"""A detected conflict between two ownership records."""
kind: str
message: str
severity: str = "warning"
issue_number: int | None = None
branch: str | None = None
worktree_path: str | None = None
session_ids: tuple[str, ...] = ()
def to_dict(self) -> dict[str, Any]:
return {
"kind": self.kind,
"severity": self.severity,
"message": self.message,
"issue_number": self.issue_number,
"branch": self.branch,
"worktree_path": self.worktree_path,
"session_ids": list(self.session_ids),
}
@dataclass(frozen=True)
class InventorySnapshot:
"""Versioned aggregate of every inventory section."""
generated_at: str
sections: tuple[InventorySection, ...]
collisions: tuple[CollisionSignal, ...] = ()
correlations: tuple[dict[str, Any], ...] = ()
schema_version: int = SCHEMA_VERSION
api_version: str = API_VERSION
scan_ms: float | None = None
_section_index: dict[str, InventorySection] = field(
default_factory=dict, repr=False, compare=False
)
def section(self, name: str) -> InventorySection | None:
return self._section_index.get(name)
@property
def degraded_sections(self) -> tuple[str, ...]:
return tuple(s.name for s in self.sections if not s.ok)
@property
def ownership_authority_complete(self) -> bool:
"""True only when every ownership-bearing section read cleanly.
While this is false the snapshot must not describe any work item as
unowned: a lease the reader could not load is not an absent lease.
"""
for name in OWNERSHIP_SECTIONS:
section = self._section_index.get(name)
if section is None or not section.ok:
return False
return True
@property
def status(self) -> str:
if all(s.ok for s in self.sections):
return STATUS_OK
return STATUS_DEGRADED
# ── redaction ────────────────────────────────────────────────────────────────
def redact_path(path: str | None) -> str | None:
"""Collapse an absolute path against ``$HOME`` for browser display."""
if not path:
return path
text = str(path)
home = os.path.expanduser("~")
if home and home != "/" and text.startswith(home):
return "~" + text[len(home) :]
return text
def redact_url(value: str | None) -> str | None:
"""Strip userinfo and query string from a URL."""
if not value:
return value
text = str(value)
try:
parsed = urlparse(text)
except ValueError:
return _REDACTED
if not parsed.scheme or not parsed.netloc:
return text
netloc = parsed.hostname or ""
if parsed.port:
netloc = f"{netloc}:{parsed.port}"
rebuilt = f"{parsed.scheme}://{netloc}{parsed.path}"
return rebuilt.rstrip("/") or rebuilt
def scrub(value: Any, *, key: str | None = None) -> Any:
"""Recursively drop credential-shaped values and redact paths/URLs.
Never raises: an unexpected object degrades to its ``repr`` rather than
propagating out of a read-only view.
"""
if key and _CREDENTIAL_KEY_RE.search(key):
return _REDACTED
if isinstance(value, dict):
return {str(k): scrub(v, key=str(k)) for k, v in value.items()}
if isinstance(value, (list, tuple)):
return [scrub(v, key=key) for v in value]
if isinstance(value, str):
if value.startswith(("http://", "https://")):
return redact_url(value)
if value.startswith("/") or value.startswith("~"):
return redact_path(value)
return value
if isinstance(value, (int, float, bool)) or value is None:
return value
return repr(value)
def collapse_home(text: str) -> str:
"""Collapse every ``$HOME`` occurrence *inside* a string, not just a prefix."""
home = os.path.expanduser("~")
if not home or home == "/":
return text
return text.replace(home, "~")
def scrub_text(value: Any) -> Any:
"""Redact a free-form text blob such as a recorded command line.
:func:`scrub` keys off structured field *names* and whole-value prefixes,
which is right for inventory records but blind to a secret embedded in the
middle of a sentence. This collapses ``$HOME`` and redacts credential-shaped
tokens and URL userinfo *anywhere* in the string, so operator-supplied text
rendered verbatim — contamination ``command_summary`` (#630) above all — is
held to the same standard as every other field on the page.
Returns ``None`` unchanged so callers can keep "absent" distinct from "".
"""
if value is None:
return None
text = value if isinstance(value, str) else str(value)
text = collapse_home(text)
text = _TEXT_URL_USERINFO_RE.sub(
lambda m: f"{m.group('scheme')}{_REDACTED}@", text
)
text = _TEXT_CREDENTIAL_FLAG_RE.sub(
lambda m: f"{m.group('key')}{m.group('sep')}{_REDACTED}", text
)
text = _TEXT_CREDENTIAL_ASSIGN_RE.sub(
lambda m: f"{m.group('key')}{m.group('sep')}{_REDACTED}", text
)
return text
# ── control-plane database (read-only) ───────────────────────────────────────
def _open_readonly(db_path: str) -> sqlite3.Connection:
"""Open the control-plane DB without creating or migrating anything."""
conn = sqlite3.connect(f"file:{db_path}?mode=ro", uri=True, timeout=5)
conn.row_factory = sqlite3.Row
return conn
def _table_names(conn: sqlite3.Connection) -> set[str]:
rows = conn.execute(
"SELECT name FROM sqlite_master WHERE type = 'table'"
).fetchall()
return {str(row[0]) for row in rows}
def _load_cp_db_sections(
*,
db_path: str | None = None,
limit: int = 200,
) -> tuple[InventorySection, InventorySection]:
"""Return the ``sessions`` and ``leases`` sections from the #613 DB."""
path = (db_path or control_plane_db.default_db_path()).strip()
def _both_unavailable(reason: str) -> tuple[InventorySection, InventorySection]:
return (
InventorySection(
name="sessions",
authority=AUTHORITY_CONTROL_PLANE_DB,
status=STATUS_UNAVAILABLE,
reason=reason,
),
InventorySection(
name="leases",
authority=AUTHORITY_CONTROL_PLANE_DB,
status=STATUS_UNAVAILABLE,
reason=reason,
),
)
if not path:
return _both_unavailable("control-plane database path is not configured")
if not os.path.exists(path):
return _both_unavailable(
f"control-plane database not present at {redact_path(path)}; "
"no session or lease authority available"
)
started = time.perf_counter()
try:
conn = _open_readonly(path)
except sqlite3.Error as exc:
return _both_unavailable(f"control-plane database could not be opened: {exc}")
try:
tables = _table_names(conn)
if "sessions" not in tables or "leases" not in tables:
missing = sorted({"sessions", "leases"} - tables)
return _both_unavailable(
"control-plane database is missing required tables: "
+ ", ".join(missing)
)
session_rows = [
dict(row)
for row in conn.execute(
"SELECT session_id, role, profile, namespace, pid, started_at,"
" last_heartbeat_at, status FROM sessions"
" ORDER BY last_heartbeat_at DESC LIMIT ?",
(max(1, int(limit)),),
).fetchall()
]
has_work_items = "work_items" in tables
if has_work_items:
lease_sql = (
"SELECT l.lease_id, l.session_id, l.role, l.phase, l.status,"
" l.expires_at, w.remote, w.org, w.repo, w.kind AS work_kind,"
" w.number AS work_number, w.state AS work_state,"
" s.pid AS session_pid, s.profile AS session_profile,"
" s.namespace AS session_namespace, s.status AS session_status"
" FROM leases l"
" JOIN work_items w ON w.work_item_id = l.work_item_id"
" LEFT JOIN sessions s ON s.session_id = l.session_id"
" ORDER BY l.expires_at DESC LIMIT ?"
)
else:
lease_sql = (
"SELECT l.lease_id, l.session_id, l.role, l.phase, l.status,"
" l.expires_at FROM leases l"
" ORDER BY l.expires_at DESC LIMIT ?"
)
lease_rows = [
dict(row)
for row in conn.execute(lease_sql, (max(1, int(limit)),)).fetchall()
]
except sqlite3.Error as exc:
return _both_unavailable(f"control-plane database read failed: {exc}")
finally:
conn.close()
elapsed = (time.perf_counter() - started) * 1000.0
now = datetime.now(timezone.utc)
sessions = tuple(
scrub(
{
"session_id": row.get("session_id"),
"role": row.get("role"),
"profile": row.get("profile"),
"namespace": row.get("namespace"),
"pid": row.get("pid"),
"pid_alive": issue_lock_store.is_process_alive(row.get("pid")),
"started_at": row.get("started_at"),
"last_heartbeat_at": row.get("last_heartbeat_at"),
"status": row.get("status"),
}
)
for row in session_rows
)
leases = tuple(
scrub(
{
"lease_id": row.get("lease_id"),
"session_id": row.get("session_id"),
"role": row.get("role"),
"phase": row.get("phase"),
"status": row.get("status"),
"expires_at": row.get("expires_at"),
"expired": _is_expired(row.get("expires_at"), now=now),
"remote": row.get("remote"),
"org": row.get("org"),
"repo": row.get("repo"),
"work_kind": row.get("work_kind"),
"work_number": row.get("work_number"),
"work_state": row.get("work_state"),
"session_pid": row.get("session_pid"),
"session_profile": row.get("session_profile"),
"session_namespace": row.get("session_namespace"),
"session_status": row.get("session_status"),
}
)
for row in lease_rows
)
degraded_reason = (
None
if has_work_items
else "work_items table absent; lease rows carry no work linkage"
)
lease_status = STATUS_OK if has_work_items else STATUS_DEGRADED
return (
InventorySection(
name="sessions",
authority=AUTHORITY_CONTROL_PLANE_DB,
status=STATUS_OK,
items=sessions,
scan_ms=round(elapsed, 3),
),
InventorySection(
name="leases",
authority=AUTHORITY_CONTROL_PLANE_DB,
status=lease_status,
items=leases,
reason=degraded_reason,
scan_ms=round(elapsed, 3),
),
)
def _is_expired(expires_at: str | None, *, now: datetime) -> bool | None:
if not expires_at:
return None
text = str(expires_at).strip().replace("Z", "+00:00")
try:
parsed = datetime.fromisoformat(text)
except ValueError:
return None
if parsed.tzinfo is None:
parsed = parsed.replace(tzinfo=timezone.utc)
return parsed <= now
# ── durable issue locks (filesystem) ─────────────────────────────────────────
def _load_locks_section(*, lock_dir: str | None = None) -> InventorySection:
started = time.perf_counter()
try:
paths = issue_lock_store.iter_lock_files(lock_dir)
except OSError as exc:
return InventorySection(
name="locks",
authority=AUTHORITY_FILESYSTEM,
status=STATUS_UNAVAILABLE,
reason=f"issue lock directory could not be listed: {exc}",
)
items: list[dict[str, Any]] = []
unreadable = 0
for path in paths:
try:
record = issue_lock_store.read_lock_file(path)
except (OSError, ValueError):
unreadable += 1
continue
if not record:
unreadable += 1
continue
try:
freshness = issue_lock_store.assess_lock_freshness(record)
except Exception: # noqa: BLE001 — a read-only view never raises
freshness = {"status": "unknown", "live": False, "stale": False}
claimant = record.get("claimant") or (
(record.get("work_lease") or {}).get("claimant") or {}
)
items.append(
scrub(
{
"issue_number": record.get("issue_number"),
"branch_name": record.get("branch_name"),
"remote": record.get("remote"),
"org": record.get("org"),
"repo": record.get("repo"),
"worktree_path": record.get("worktree_path"),
"pid": record.get("session_pid") or record.get("pid"),
"pid_alive": issue_lock_store.is_process_alive(
record.get("session_pid") or record.get("pid")
),
"claimant_username": (claimant or {}).get("username"),
"claimant_profile": (claimant or {}).get("profile"),
"lock_generation": record.get("lock_generation"),
"freshness_status": freshness.get("status"),
"live": bool(freshness.get("live")),
"stale": bool(freshness.get("stale")),
"freshness_reason": freshness.get("reason"),
"lock_path": record.get("lock_file_path") or path,
}
)
)
elapsed = (time.perf_counter() - started) * 1000.0
reason = (
f"{unreadable} lock file(s) were unreadable and are not represented"
if unreadable
else None
)
return InventorySection(
name="locks",
authority=AUTHORITY_FILESYSTEM,
status=STATUS_DEGRADED if unreadable else STATUS_OK,
items=tuple(items),
reason=reason,
scan_ms=round(elapsed, 3),
)
# ── worktrees (filesystem, via the #432 scanner) ─────────────────────────────
def _load_worktrees_section(
*, load_hygiene: Callable[[], Any] | None = None
) -> InventorySection:
started = time.perf_counter()
try:
loader = load_hygiene
if loader is None:
from webui.worktree_scanner import load_hygiene_snapshot
loader = load_hygiene_snapshot
snapshot = loader()
except Exception as exc: # noqa: BLE001 — fail soft, never fail the request
return InventorySection(
name="worktrees",
authority=AUTHORITY_FILESYSTEM,
status=STATUS_UNAVAILABLE,
reason=f"worktree scan failed: {exc}",
)
items = tuple(
scrub(
{
"rel_path": entry.rel_path,
"folder_name": entry.folder_name,
"classification": entry.classification,
"branch": entry.branch,
"head_sha": entry.head_sha,
"dirty_tracked": entry.dirty_tracked,
"dirty_untracked": entry.dirty_untracked,
"detached": entry.detached,
"registered_worktree": entry.registered_worktree,
"notes": entry.notes,
}
)
for entry in snapshot.entries
)
scan_error = getattr(snapshot, "scan_error", None)
elapsed = (time.perf_counter() - started) * 1000.0
return InventorySection(
name="worktrees",
authority=AUTHORITY_FILESYSTEM,
status=STATUS_DEGRADED if scan_error else STATUS_OK,
items=items,
reason=scan_error,
scan_ms=round(elapsed, 3),
)
# ── namespaces / capability summary ──────────────────────────────────────────
def _load_namespaces_section() -> InventorySection:
started = time.perf_counter()
try:
from gitea_auth import get_profile
profile = get_profile() or {}
except Exception as exc: # noqa: BLE001
return InventorySection(
name="namespaces",
authority=AUTHORITY_FILESYSTEM,
status=STATUS_UNAVAILABLE,
reason=f"active profile could not be resolved: {exc}",
)
allowed = list(profile.get("allowed_operations") or [])
forbidden = list(profile.get("forbidden_operations") or [])
profile_name = str(profile.get("profile_name") or "")
namespace = None
try:
import role_namespace_gate
namespace = role_namespace_gate.infer_mcp_namespace(profile_name)
except Exception: # noqa: BLE001 — namespace inference is advisory
namespace = None
item = scrub(
{
"profile_name": profile_name,
"role": profile.get("role"),
"mcp_namespace": namespace,
"allowed_operations": sorted(allowed),
"forbidden_operations": sorted(forbidden),
"capability_summary": {
"can_author": "gitea.pr.create" in allowed,
"can_review": "gitea.pr.approve" in allowed,
"can_merge": "gitea.pr.merge" in allowed,
"can_close_pr": "gitea.pr.close" in allowed,
},
"active": True,
}
)
elapsed = (time.perf_counter() - started) * 1000.0
return InventorySection(
name="namespaces",
authority=AUTHORITY_FILESYSTEM,
status=STATUS_OK,
items=(item,),
reason=(
"only the profile serving this web process is observable; other "
"namespaces are not enumerable from here"
),
scan_ms=round(elapsed, 3),
)
# ── correlation and collision detection ──────────────────────────────────────
def _issue_from_branch(branch: str | None) -> int | None:
match = re.search(r"issue-(\d+)", str(branch or ""), re.IGNORECASE)
return int(match.group(1)) if match else None
def correlate(
*,
leases: InventorySection,
locks: InventorySection,
worktrees: InventorySection,
sessions: InventorySection,
) -> tuple[tuple[dict[str, Any], ...], tuple[CollisionSignal, ...]]:
"""Join lease owner ↔ lock ↔ worktree ↔ namespace where evidence allows.
Correlation rows are emitted from whatever sections did load. Collision
signals are only emitted from sections that are ``ok``: a conflict inferred
from a partially-read source would be a false accusation.
"""
correlations: list[dict[str, Any]] = []
collisions: list[CollisionSignal] = []
worktree_by_branch: dict[str, dict[str, Any]] = {}
for entry in worktrees.items:
branch = (entry.get("branch") or "").strip()
if branch:
worktree_by_branch.setdefault(branch, entry)
session_by_id = {
str(s.get("session_id")): s for s in sessions.items if s.get("session_id")
}
# Lock-centred rows: a durable lock names an issue, a branch, and a worktree.
for lock in locks.items:
branch = (lock.get("branch_name") or "").strip()
worktree = worktree_by_branch.get(branch)
matching_leases = [
lease
for lease in leases.items
if lease.get("work_kind") == "issue"
and lease.get("work_number") == lock.get("issue_number")
]
correlations.append(
{
"issue_number": lock.get("issue_number"),
"branch": branch or None,
"lock_live": bool(lock.get("live")),
"lock_claimant": lock.get("claimant_profile"),
"lock_pid": lock.get("pid"),
"lock_pid_alive": lock.get("pid_alive"),
"worktree_rel_path": (worktree or {}).get("rel_path"),
"worktree_classification": (worktree or {}).get("classification"),
"worktree_registered": (worktree or {}).get("registered_worktree"),
"lease_ids": [
lease.get("lease_id")
for lease in matching_leases
if lease.get("lease_id")
],
"lease_sessions": [
lease.get("session_id")
for lease in matching_leases
if lease.get("session_id")
],
}
)
if locks.ok and worktrees.ok:
# A claim whose lease window is still open but has no registered
# worktree is an anomaly regardless of whether its pid is alive; a
# fully time-expired lease is on its way out and is not flagged.
if (
lock.get("freshness_status") != "expired"
and branch
and worktree is None
):
collisions.append(
CollisionSignal(
kind="lock-without-worktree",
severity="warning",
issue_number=lock.get("issue_number"),
branch=branch,
worktree_path=lock.get("worktree_path"),
message=(
f"Live lock on issue #{lock.get('issue_number')} names "
f"branch {branch!r} but no registered worktree carries "
"that branch (#404)"
),
)
)
if locks.ok:
# A lock whose recorded pid is gone is held by nobody: a clean #753
# dead-session recovery candidate. Subclassify by the lease window,
# because the two cases need different operator urgency. When the
# window is still open the lock would read as live to a naive
# timestamp check even though the owner is dead — the more dangerous
# case — so it is flagged distinctly from a fully time-expired lease.
if lock.get("pid_alive") is False and lock.get("stale"):
if lock.get("freshness_status") == "expired":
collisions.append(
CollisionSignal(
kind="stale-lock-dead-owner",
severity="warning",
issue_number=lock.get("issue_number"),
branch=branch or None,
message=(
f"Lock on issue #{lock.get('issue_number')} is stale "
f"and its recorded pid {lock.get('pid')} is not running "
"(#753 dead-session recovery candidate)"
),
)
)
else:
collisions.append(
CollisionSignal(
kind="live-lock-dead-owner",
severity="warning",
issue_number=lock.get("issue_number"),
branch=branch or None,
message=(
f"Lock on issue #{lock.get('issue_number')} has an "
"unexpired lease but its recorded pid "
f"{lock.get('pid')} is not running; it would read as "
"live to a timestamp check (#753 dead-session recovery "
"candidate)"
),
)
)
# The #635 trap: the lease has expired but the recorded pid is a
# still-running daemon, so neither dead-pid reclaim nor exact-owner
# renewal applies. This is the collision an operator must see.
elif (
lock.get("freshness_status") == "expired"
and lock.get("pid_alive") is True
):
collisions.append(
CollisionSignal(
kind="expired-lock-live-owner",
severity="error",
issue_number=lock.get("issue_number"),
branch=branch or None,
message=(
f"Lock on issue #{lock.get('issue_number')} has an expired "
f"lease but its recorded pid {lock.get('pid')} is still "
"running (daemon-pid deadlock; needs an operator decision, "
"#635/#760)"
),
)
)
# Two live locks on one branch, or two active leases on one work item.
if locks.ok:
by_branch: dict[str, list[dict[str, Any]]] = {}
for lock in locks.items:
if not lock.get("live"):
continue
branch = (lock.get("branch_name") or "").strip()
if branch:
by_branch.setdefault(branch, []).append(lock)
for branch, entries in sorted(by_branch.items()):
if len(entries) > 1:
collisions.append(
CollisionSignal(
kind="duplicate-live-lock",
severity="error",
branch=branch,
message=(
f"{len(entries)} live locks name branch {branch!r}: "
"issues "
+ ", ".join(
f"#{e.get('issue_number')}" for e in entries
)
),
)
)
if leases.ok:
by_work: dict[tuple[str, int], list[dict[str, Any]]] = {}
for lease in leases.items:
if str(lease.get("status") or "").lower() != "active":
continue
kind = str(lease.get("work_kind") or "").strip().lower()
number = lease.get("work_number")
if not kind or number is None:
continue
by_work.setdefault((kind, int(number)), []).append(lease)
for (kind, number), entries in sorted(by_work.items()):
sessions_held = {
str(e.get("session_id")) for e in entries if e.get("session_id")
}
if len(sessions_held) > 1:
collisions.append(
CollisionSignal(
kind="concurrent-active-lease",
severity="error",
issue_number=number if kind == "issue" else None,
session_ids=tuple(sorted(sessions_held)),
message=(
f"{len(sessions_held)} sessions hold an active lease on "
f"{kind} #{number}"
),
)
)
for entry in entries:
if entry.get("expired") is True:
collisions.append(
CollisionSignal(
kind="active-lease-past-expiry",
severity="warning",
issue_number=number if kind == "issue" else None,
session_ids=(
(str(entry.get("session_id")),)
if entry.get("session_id")
else ()
),
message=(
f"Lease {entry.get('lease_id')} on {kind} #{number} "
"is still marked active past its expiry"
),
)
)
# A lease whose owning session is gone is an orphan, not free work.
if leases.ok and sessions.ok:
for lease in leases.items:
if str(lease.get("status") or "").lower() != "active":
continue
session_id = str(lease.get("session_id") or "")
if session_id and session_id not in session_by_id:
collisions.append(
CollisionSignal(
kind="orphan-lease",
severity="error",
session_ids=(session_id,),
message=(
f"Active lease {lease.get('lease_id')} names session "
f"{session_id}, which has no session record"
),
)
)
return tuple(correlations), tuple(collisions)
# ── snapshot assembly ────────────────────────────────────────────────────────
def load_inventory_snapshot(
*,
db_path: str | None = None,
lock_dir: str | None = None,
load_hygiene: Callable[[], Any] | None = None,
include: tuple[str, ...] | None = None,
) -> InventorySnapshot:
"""Build the unified inventory snapshot.
Every section is loaded independently and fails soft. *include* restricts
the sections that are scanned; omitted sections are simply absent rather
than reported as empty, so a resource-split request cannot be mistaken for
a whole-inventory answer.
"""
started = time.perf_counter()
wanted = tuple(include) if include else SECTION_NAMES
sections: list[InventorySection] = []
sessions_section: InventorySection | None = None
leases_section: InventorySection | None = None
if "sessions" in wanted or "leases" in wanted:
sessions_section, leases_section = _load_cp_db_sections(db_path=db_path)
if "sessions" in wanted:
sections.append(sessions_section)
if "leases" in wanted:
sections.append(leases_section)
locks_section = (
_load_locks_section(lock_dir=lock_dir)
if "locks" in wanted
else _empty_section("locks", AUTHORITY_FILESYSTEM)
)
if "locks" in wanted:
sections.append(locks_section)
worktrees_section = (
_load_worktrees_section(load_hygiene=load_hygiene)
if "worktrees" in wanted
else _empty_section("worktrees", AUTHORITY_FILESYSTEM)
)
if "worktrees" in wanted:
sections.append(worktrees_section)
if "namespaces" in wanted:
sections.append(_load_namespaces_section())
correlations, collisions = correlate(
leases=leases_section or _empty_section("leases", AUTHORITY_CONTROL_PLANE_DB),
locks=locks_section,
worktrees=worktrees_section,
sessions=sessions_section
or _empty_section("sessions", AUTHORITY_CONTROL_PLANE_DB),
)
elapsed = (time.perf_counter() - started) * 1000.0
index = {section.name: section for section in sections}
return InventorySnapshot(
generated_at=datetime.now(timezone.utc).isoformat(),
sections=tuple(sections),
collisions=collisions,
correlations=correlations,
scan_ms=round(elapsed, 3),
_section_index=index,
)
def _empty_section(name: str, authority: str) -> InventorySection:
"""A section that was not requested — never a claim that it is empty."""
return InventorySection(
name=name,
authority=authority,
status=STATUS_UNAVAILABLE,
reason="section not requested in this scan",
)
def snapshot_to_dict(snapshot: InventorySnapshot) -> dict[str, Any]:
"""Serialize the snapshot for the versioned API."""
return {
"api_version": snapshot.api_version,
"schema_version": snapshot.schema_version,
"generated_at": snapshot.generated_at,
"status": snapshot.status,
"scan_ms": snapshot.scan_ms,
"ownership_authority_complete": snapshot.ownership_authority_complete,
"ownership_note": (
"Every ownership source read cleanly; an item absent from leases "
"and locks is genuinely unclaimed."
if snapshot.ownership_authority_complete
else "One or more ownership sources are degraded; nothing in this "
"snapshot may be treated as unowned. Collisions are reported only "
"from sections that read cleanly."
),
"degraded_sections": list(snapshot.degraded_sections),
"field_authority": {
"sessions": AUTHORITY_CONTROL_PLANE_DB,
"leases": AUTHORITY_CONTROL_PLANE_DB,
"locks": AUTHORITY_FILESYSTEM,
"worktrees": AUTHORITY_FILESYSTEM,
"namespaces": AUTHORITY_FILESYSTEM,
},
"sections": {section.name: section.to_dict() for section in snapshot.sections},
"correlations": [dict(row) for row in snapshot.correlations],
"collisions": [signal.to_dict() for signal in snapshot.collisions],
}