Compare commits

..
Author SHA1 Message Date
jcwalker3 f1e4809930 Merge branch 'master' into fix/issue-790-slice-a-heartbeat-policy 2026-07-23 12:31:44 -05:00
jcwalker3 dc1d0e045f Merge branch 'master' into fix/issue-790-slice-a-heartbeat-policy 2026-07-23 01:13:11 -05:00
jcwalker3 badc4e636b Merge branch 'master' into fix/issue-790-slice-a-heartbeat-policy 2026-07-23 00:06:18 -05:00
sysadminandClaude Opus 4.8 243f52dc79 feat(lease): make the author task heartbeat load-bearing (#790 Slice A)
Slice A of Issue #790, per the controller reassessment in comment 13958. Does
not close the issue: terminal retirement (Slice B) and the read-side generation
check plus the #760 renewal re-scope (Slice C) are deliberately not implemented.

The defect. `issue_lock_store.assess_lock_freshness` parsed `last_heartbeat_at`
and then never consulted it. Liveness was decided by an absolute four-hour
`expires_at` and by PID liveness, and the recorded PID is the long-lived MCP
daemon rather than the authoring task, so an abandoned claim stayed live for the
full four hours. A tree-wide search found the field written in exactly one place
and advanced by nothing. Issue #787 / PR #789 hit this; Issue #760 / PR #791 hit
it again, blocking reconciliation for over five hours after its work had landed.

A1 — central policy. New `lease_policy` declares every duration for every task
class in one place: author initial/sliding TTL 10 minutes, heartbeat cadence 2,
stale warning 5, missed-heartbeat grace 10, absolute cap 8 hours, recovery grace
10, terminal race-drain 2. It ships first so the first heartbeat and TTL
behavior to run reads from it (AC-N7). The duplicated four-hour literal is gone
from both `issue_lock_store` and `gitea_mcp_server`. Reviewer, merger, and
conflict-fix classes are declared but not rewired — Slice C moves those call
sites — and a test asserts the declaration still equals the constants #747 and
`pr_work_lease` own, so the two cannot drift apart unnoticed.

A2 — load-bearing freshness, with two deliberate asymmetries. An alive PID never
establishes freshness anywhere (AC-N2); it is recorded as evidence and no branch
returns live because of it. A dead PID still marks a lease stale, and that band
still precedes every heartbeat evaluation, so #753 dead-session recovery keys on
exactly the classification it always did. New bands `stale_missed_heartbeat` and
`stale_absolute_cap` are classified in `branch_cleanup_guard` rather than
falling through to unknown-status, and still block unless the ownership record
proves `reclaim_allowed is True`. A heartbeat lease carrying no heartbeat is
contradictory and fails closed. `assess_expired_lock_reclaim` accepts a lapsed
heartbeat as reclaim grounds for heartbeat-lifecycle leases only: under this
lifecycle the heartbeat is the liveness proof, and also requiring a dead PID
would reinstate the original defect.

A3/A4 — task-session identity and the writer. `mint_task_session_id` produces an
ownership key containing no process identifier, since the daemon PID is reused
by every task it serves and identifies none of them. `heartbeat_session_lock`
writes inside the existing per-issue flock under the #772 generation
compare-and-swap, verifying exact issue, branch, realpath-normalized worktree,
claimant username, claimant profile, and recorded session identifier. It cannot
acquire, take over, or revive: a lease past its grace is refused and must use
the reclaim path, so a session that stopped proving liveness cannot restore
ownership retroactively. New `gitea_heartbeat_issue_lock` gates on the same
authority as `lock_issue`, being strictly narrower.

A5 — legacy compatibility (AC-N8). The explicit `lifecycle_version` marker, never
a timestamp comparison, discriminates legacy from heartbeat leases: a legacy lock
has `last_heartbeat_at == created_at` forever precisely because nothing advanced
it, and a freshly minted heartbeat lease has them equal too, so the equality
carries no information in either direction. Legacy locks keep their recorded
absolute expiry and are never evaluated against the short grace, so deployment
cannot make an existing claim instantly reclaimable. They leave that state only
by terminal retirement (Slice B) or by `rebind_legacy_lock`, which re-verifies
the exact owner and mints a genuine identifier and first heartbeat while
preserving the original claim under `legacy_origin`. Rebinding a lapsed legacy
lease is refused; that belongs to #760 renewal or #601 reclaim.

A6 — native coverage. Review #499 proved assessor-level tests miss discard
points, so `tests/test_issue_790_heartbeat_mcp_path.py` drives the real tools
against a real git repository and a real durable lock: lock creation and
read-back, policy window, freshness, survival of `verify_lock_for_mutation`,
invariance of the duplicate-work and linked-open-PR gates, CAS rejection,
foreign-session and foreign-claimant refusal, alive-PID-only refusal, missed
heartbeat, legacy protection on deployment, and legacy rebinding.

Tests. New suites 55 passed. Lock and lease regression set (issue_lock_store,
lease_lifecycle, #753, #755, #760 x2, #768, #772, lock registration, worktree,
adoption, duplicate gate, branch cleanup guard, capability invariants, claim
heartbeat, worktrees) 383 passed with 98 subtests. Full suite 4295 passed, 11
failed, 6 skipped, 499 subtests passed, against a clean master baseline worktree
at 620ed6e9 that reports 11 failed and 4240 passed — the same eleven node IDs.
The 55-test delta is exactly the new suites; no new failures.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
Claude-Session: https://claude.ai/code/session_011u6GKSJwwrrYjguPjs1aK5
2026-07-22 04:19:48 -04:00
17 changed files with 2002 additions and 3173 deletions
+5 -176
View File
@@ -53,8 +53,6 @@ OUTCOME_CANDIDATE_SET_DRIFT = "candidate_set_drift"
SKIP_CLAIMED_BY_OTHER_SESSION = "claimed_by_other_session"
# #776: controller-supplied pre-rank exclusion.
SKIP_EXCLUDED_BY_CONTROLLER = "excluded_by_controller"
# #844: epic / child-only implementation container (pre-rank).
SKIP_EPIC_OR_CHILD_ONLY_CONTAINER = "epic_or_child_only_container"
# Ownership verdicts for a live claim on a candidate (#765).
OWNERSHIP_OWN = "own"
@@ -132,39 +130,6 @@ ROLE_ACTIONS: dict[str, tuple[tuple[str, ...], tuple[str, ...]]] = {
}
# Body phrases that prove an issue is an implementation container, not a
# unit of direct author work (#844). Matched case-insensitively against the
# issue body. Title alone is never sufficient (ordinary issues may mention
# "epic" incidentally).
_CHILD_ONLY_BODY_MARKERS: tuple[str, ...] = (
"implementation is delivered via child issues only",
"implementation is delivered through child issues only",
"implementation is delivered via child issues",
"implementation is delivered through child issues",
"do not implement product features in this epic",
"do not implement product features in this epic issue itself",
"no product feature implementation is claimed complete solely on this epic",
"implementable child issues remain independently eligible",
"owns the product roadmap and linkage",
"this epic owns the product roadmap",
"coordination container",
"child-only container",
"implementation is delegated to child",
)
# Explicit epic / umbrella labels (structured evidence preferred over title).
_EPIC_LABELS: frozenset[str] = frozenset(
{
"type:epic",
"epic",
"kind:epic",
"scope:epic",
"type:umbrella",
"umbrella",
}
)
@dataclass
class WorkCandidate:
"""One assignable Gitea issue or PR presented to the allocator."""
@@ -174,7 +139,6 @@ class WorkCandidate:
state: str = "open"
labels: tuple[str, ...] = ()
title: str = ""
body: str = ""
priority: int = 0
head_sha: str | None = None
# Routing signals (callers derive from Gitea / review feedback).
@@ -194,7 +158,6 @@ class WorkCandidate:
self.labels = tuple(
str(x).strip().lower() for x in (self.labels or ()) if str(x).strip()
)
self.body = str(self.body or "")
if self.kind not in WORK_KINDS:
raise InvalidWorkKindError(
f"candidate kind '{self.kind}' is not assignable; only "
@@ -208,7 +171,6 @@ class WorkCandidate:
"state": self.state,
"labels": list(self.labels),
"title": self.title,
"body": self.body,
"priority": self.priority,
"head_sha": self.head_sha,
"request_changes_current_head": self.request_changes_current_head,
@@ -222,51 +184,6 @@ class WorkCandidate:
}
def classify_epic_or_child_only_container(
c: WorkCandidate,
) -> tuple[bool, str | None]:
"""Return whether *c* is an epic / child-only implementation container (#844).
Exclusion uses structured evidence first (labels, body scope language).
A bare title containing the word "epic" is **not** enough — ordinary
implementable issues may mention epics incidentally. A title that is
explicitly prefixed ``Epic:`` only counts when the body also proves
child-only / no-direct-implementation scope (or an epic label is present).
PRs are never classified as containers here (they already have a head).
"""
if c.kind != "issue":
return False, None
labels = set(c.labels)
epic_label = sorted(labels & _EPIC_LABELS)
body_l = (c.body or "").lower()
title = (c.title or "").strip()
title_l = title.lower()
body_hits = [m for m in _CHILD_ONLY_BODY_MARKERS if m in body_l]
title_epic_prefix = title_l.startswith("epic:") or title_l.startswith("epic ")
if epic_label:
detail = f"label={epic_label[0]}"
if body_hits:
detail = f"{detail}; body_marker={body_hits[0]!r}"
return True, detail
if body_hits:
# Body proves child-only / umbrella scope. Title "Epic:" is corroborating
# but not required — containers without the word still exclude.
detail = f"body_marker={body_hits[0]!r}"
if title_epic_prefix:
detail = f"title_epic_prefix; {detail}"
return True, detail
# Title-only "Epic:" without body scope evidence is insufficient (#844 AC:
# eligibility does not rely solely on the word "Epic" in a title).
# Similarly, incidental "epic" mid-title without markers stays eligible.
return False, None
@dataclass
class SkipRecord:
kind: str
@@ -933,8 +850,7 @@ def allocate_next_work(
ownership_defects: list[dict[str, Any]] = []
controller_excluded: list[dict[str, Any]] = []
# #776 AC2 + #844: remove excluded numbers *and* epic/child-only containers
# *before* ranking / selection / lease so they never receive assignments.
# #776 AC2: remove excluded numbers *before* ranking / selection / lease.
rankable: list[WorkCandidate] = []
for c in candidates:
if int(c.number) in exclude_set:
@@ -1013,23 +929,6 @@ def allocate_next_work(
},
}
continue
# #844: epics / child-only containers are never direct implement targets.
is_container, container_detail = classify_epic_or_child_only_container(c)
if is_container:
detail = container_detail or "epic or child-only container"
reason = (
f"{c.kind}#{c.number} {SKIP_EPIC_OR_CHILD_ONLY_CONTAINER}: "
f"{detail}; implementation is delegated to child issues"
)
skipped.append(
SkipRecord(
c.kind,
c.number,
reason,
SKIP_EPIC_OR_CHILD_ONLY_CONTAINER,
)
)
continue
rankable.append(c)
ordered = sort_candidates(rankable)
@@ -1216,12 +1115,6 @@ def allocate_next_work(
"reasons": [
"dry-run only (apply=false); no assignment/lease created — "
"call again with apply=true to reserve via control-plane DB"
+ (
"; after apply, the required-role worker consumes via "
"gitea_adopt_workflow_lease (#843)"
if mode == ALLOCATION_MODE_CROSS_ROLE and expected_role != role_norm
else ""
)
],
"skipped": [s.as_dict() for s in skipped],
"terminal_pr": terminal_pr,
@@ -1253,9 +1146,6 @@ def allocate_next_work(
# Atomic reserve via #613 substrate.
ttl = lease_ttl_seconds if lease_ttl_seconds is not None else None
try:
cross_role_handoff = (
mode == ALLOCATION_MODE_CROSS_ROLE and lease_role != role_norm
)
kwargs: dict[str, Any] = {
"session_id": session_id,
"role": lease_role,
@@ -1267,8 +1157,7 @@ def allocate_next_work(
"expected_head_sha": selected.head_sha,
"allowed_actions": allowed,
"forbidden_actions": forbidden,
# #843: mark cross-role allocations as awaiting independent consume
"phase": "awaiting_handoff" if cross_role_handoff else "allocated",
"phase": "allocated",
}
if ttl is not None:
kwargs["lease_ttl_seconds"] = int(ttl)
@@ -1348,52 +1237,7 @@ def allocate_next_work(
"lease_role": lease_role,
"source": "control_plane_db.assign_and_lease",
}
consume_allocation = None
if cross_role_handoff and result.lease_id:
# Durable handoff marker so independent required-role workers can
# consume without sharing the controller session (#843).
handoff_prov = {
"cross_role_handoff": True,
"handoff_status": "pending",
"allocating_session_id": session_id,
"allocating_role": role_norm,
"required_role": expected_role,
"required_profile": selection["required_profile"],
"required_namespace": selection["required_namespace"],
"assignment_id": result.assignment_id,
"lease_id": result.lease_id,
"allocation_mode": mode,
"adopted_by_session_id": None,
}
try:
db.attach_lease_provenance(result.lease_id, handoff_prov)
except ControlPlaneError:
# Still return assignment evidence; consume path may be unavailable
handoff_prov["attach_failed"] = True
consume_allocation = {
"tool": "gitea_adopt_workflow_lease",
"lease_id": result.lease_id,
"assignment_id": result.assignment_id,
"required_role": expected_role,
"required_profile": selection["required_profile"],
"required_namespace": selection["required_namespace"],
"handoff_status": "pending",
"controller_session_required": False,
"instructions": (
f"From an independent {expected_role} session "
f"({selection['required_namespace']} / "
f"{selection['required_profile']}), call "
f"gitea_adopt_workflow_lease(lease_id={result.lease_id!r}) "
"to consume this controller allocation. The allocating "
"controller process does not need to remain alive. Wrong-role "
"and second-adoption attempts fail closed."
),
}
lease_proof["cross_role_handoff"] = True
lease_proof["handoff_status"] = "pending"
lease_proof["consume_tool"] = "gitea_adopt_workflow_lease"
out = {
return {
"success": True,
"outcome": OUTCOME_ASSIGNED,
"apply": True,
@@ -1427,17 +1271,8 @@ def allocate_next_work(
"lease_role": lease_role,
"lease_proof": lease_proof,
"selection_policy": SELECTION_POLICY,
"cross_role_handoff": bool(cross_role_handoff),
},
"next_valid_command": (
(
f"consume lease {result.lease_id} via gitea_adopt_workflow_lease "
f"as {expected_role}, then "
)
+ _next_command(lease_role, selected)
if cross_role_handoff
else _next_command(lease_role, selected)
),
"next_valid_command": _next_command(lease_role, selected),
"substrate": "control_plane_db",
"file_lock_only": False,
"comment_lease_only": False,
@@ -1452,14 +1287,9 @@ def allocate_next_work(
"downstream_note": (
"#612 incident bridge remains downstream of #600; "
"allocator never assigns raw monitoring incidents; "
"controller routes only under cross_role (#840); "
"cross-role assignments are consumable by independent "
"required-role workers via gitea_adopt_workflow_lease (#843)"
"controller routes only under cross_role (#840)"
),
}
if consume_allocation is not None:
out["consume_allocation"] = consume_allocation
return out
def _next_command(role: str, c: WorkCandidate) -> str:
@@ -1511,7 +1341,6 @@ def candidate_from_dict(data: dict[str, Any]) -> WorkCandidate:
state=str(data.get("state") or "open"),
labels=tuple(data.get("labels") or ()),
title=str(data.get("title") or ""),
body=str(data.get("body") or ""),
priority=priority,
head_sha=data.get("head_sha"),
request_changes_current_head=bool(data.get("request_changes_current_head")),
+13 -1
View File
@@ -163,7 +163,19 @@ _TERMINAL_OWNERSHIP_STATUSES = frozenset(
{"released", "abandoned", "done", "blocked", "terminal", "closed"}
)
_EXPIRED_STATUSES = frozenset({"expired"})
_STALE_STATUSES = frozenset({"stale", "stale_dead_process", "stale_missing_worktree"})
_STALE_STATUSES = frozenset(
{
"stale",
"stale_dead_process",
"stale_missing_worktree",
# #790 Slice A heartbeat-lifecycle bands. Listed here so they are
# *classified* rather than falling through to the unknown-status branch;
# they still block unless the ownership record proves
# ``reclaim_allowed is True``, so the O2 fail-closed rule is unchanged.
"stale_missed_heartbeat",
"stale_absolute_cap",
}
)
def _norm_str(value: Any) -> str:
+3 -169
View File
@@ -1637,13 +1637,11 @@ class ControlPlaneDB:
provenance: dict[str, Any] | None = None,
lease_ttl_seconds: int = DEFAULT_LEASE_TTL_SECONDS,
) -> dict[str, Any]:
"""Transfer or refresh a lease with provenance (#601 / #843).
"""Transfer or refresh a lease with provenance (#601).
* Same owner + active → refresh (owner-resume).
* Cross-role handoff pending + matching required role → atomic consume
(even while the allocating controller session still "owns" the lease).
* Expired/abandoned/released → create new assignment+lease with provenance.
* Active foreign (non-handoff) → raise ForeignLeaseError (never silent steal).
* Active foreign → raise ForeignLeaseError (never silent steal).
"""
now = _utc_now()
now_s = _ts(now)
@@ -1679,35 +1677,7 @@ class ControlPlaneDB:
status = "expired"
owner = lease["session_id"]
# Parse durable provenance for cross-role handoff consume (#843).
lease_prov: dict[str, Any] = {}
if "provenance_json" in lease.keys() and lease["provenance_json"]:
try:
loaded = json.loads(lease["provenance_json"])
if isinstance(loaded, dict):
lease_prov = loaded
except (TypeError, json.JSONDecodeError):
lease_prov = {}
handoff_pending = bool(lease_prov.get("cross_role_handoff")) and (
str(lease_prov.get("handoff_status") or "pending").strip().lower()
== "pending"
)
already_adopted = bool(
(lease["adopted_by_session_id"] if "adopted_by_session_id" in lease.keys() else None)
or lease_prov.get("adopted_by_session_id")
)
required_role = str(
lease_prov.get("required_role") or lease["role"] or ""
).strip().lower()
adopter_role = (role or "").strip().lower()
cross_role_consume = (
handoff_pending
and not already_adopted
and status == "active"
and owner != adopter_session_id
)
if status == "active" and owner != adopter_session_id and not cross_role_consume:
if status == "active" and owner != adopter_session_id:
raise ForeignLeaseError(
f"cannot adopt active foreign lease {lease_id} owned by {owner}"
)
@@ -1791,142 +1761,6 @@ class ControlPlaneDB:
"reasons": ["owner-resume: refreshed lease with provenance"],
}
# #843: controller→required-role handoff consume (atomic, same lease_id)
if cross_role_consume:
if not required_role:
raise ControlPlaneError(
f"cross-role handoff lease {lease_id} missing required_role"
)
if adopter_role != required_role:
raise ForeignLeaseError(
f"wrong role for cross-role handoff consume: "
f"required={required_role} adopter={adopter_role or 'none'} "
f"(fail closed)"
)
# CAS: only transfer if still owned by allocating session and unadopted
cols = self._lease_columns(conn)
adopted_col_null = (
"(adopted_by_session_id IS NULL OR adopted_by_session_id = '')"
if "adopted_by_session_id" in cols
else "1=1"
)
cas = conn.execute(
f"""
UPDATE leases
SET session_id = ?,
heartbeat_at = ?,
expires_at = ?,
phase = ?,
role = ?
WHERE lease_id = ?
AND status = 'active'
AND session_id = ?
AND {adopted_col_null}
""",
(
adopter_session_id,
now_s,
expires,
"adopted",
required_role,
lease_id,
owner,
),
)
if cas.rowcount != 1:
raise ForeignLeaseError(
f"cross-role handoff consume lost race for lease {lease_id} "
"(already adopted or no longer pending; fail closed)"
)
if "adopted_from_session_id" in cols:
conn.execute(
"""
UPDATE leases
SET adopted_from_session_id = ?, adopted_by_session_id = ?
WHERE lease_id = ?
""",
(owner, adopter_session_id, lease_id),
)
if "worktree_path" in cols and worktree_path:
conn.execute(
"UPDATE leases SET worktree_path = ? WHERE lease_id = ?",
(worktree_path, lease_id),
)
if "owner_pid" in cols and owner_pid is not None:
conn.execute(
"UPDATE leases SET owner_pid = ? WHERE lease_id = ?",
(owner_pid, lease_id),
)
if "expected_head_sha" in cols and expected_head_sha:
conn.execute(
"UPDATE leases SET expected_head_sha = ? WHERE lease_id = ?",
(expected_head_sha, lease_id),
)
# Merge handoff provenance + caller provenance
merged = dict(lease_prov)
merged.update(provenance or {})
merged["cross_role_handoff"] = True
merged["handoff_status"] = "adopted"
merged["adopted_from_session_id"] = owner
merged["adopted_by_session_id"] = adopter_session_id
merged["required_role"] = required_role
if "provenance_json" in cols:
conn.execute(
"UPDATE leases SET provenance_json = ? WHERE lease_id = ?",
(json.dumps(merged), lease_id),
)
# Transfer active assignment ownership atomically
asn_cas = conn.execute(
"""
UPDATE assignments
SET session_id = ?, role = ?
WHERE lease_id = ? AND status = 'active' AND session_id = ?
""",
(adopter_session_id, required_role, lease_id, owner),
)
if asn_cas.rowcount < 1:
# Fail closed: assignment must move with the lease
raise ControlPlaneError(
f"cross-role handoff: no active assignment for lease {lease_id} "
f"owned by {owner}"
)
lease2 = conn.execute(
"SELECT * FROM leases WHERE lease_id = ?", (lease_id,)
).fetchone()
asn = conn.execute(
"""
SELECT * FROM assignments
WHERE lease_id = ? AND status = 'active'
ORDER BY created_at DESC LIMIT 1
""",
(lease_id,),
).fetchone()
conn.execute(
"""
INSERT INTO events(work_item_id, event_type, message, created_at)
VALUES (?, 'lease_adopted', ?, ?)
""",
(
lease["work_item_id"],
f"cross-role handoff: {adopter_session_id} consumed "
f"{lease_id} from {owner} as {required_role}",
now_s,
),
)
return {
"outcome": "adopted_cross_role_handoff",
"lease": dict(lease2) if lease2 else dict(lease),
"assignment": dict(asn) if asn else None,
"reasons": [
"cross-role handoff: independent required-role worker consumed "
"controller allocation without abandonment"
],
"adopted_by_session_id": adopter_session_id,
"adopted_from_session_id": owner,
"required_role": required_role,
"handoff_status": "adopted",
}
# Non-active: create new lease + assignment (transfer)
new_lease_id = f"lease-{uuid.uuid4().hex[:16]}"
new_asn_id = f"asn-{uuid.uuid4().hex[:16]}"
+1
View File
@@ -100,6 +100,7 @@ that gates each call, not which tools exist.
- `gitea_get_profile`
- `gitea_get_runtime_context`
- `gitea_get_shell_health`
- `gitea_heartbeat_issue_lock`
- `gitea_heartbeat_reviewer_pr_lease`
- `gitea_inspect_workflow_lease`
- `gitea_issue_irrecoverable_provenance_authorization`
-102
View File
@@ -292,108 +292,6 @@ health, workflow/schema SHA-256 hashes, and stale-runtime warnings when the
checkout is behind merged safety-gate changes. Restart guidance links to #420;
no tokens or MCP restart actions are exposed.
## Workflow-event timeline (#637)
`GET /api/v1/timeline` is a read-only, versioned aggregation of workflow
events from every available source into one normalised, filterable stream. It
is the model layer for the Phase 1 timeline console view (a later child issue
of #631); this issue ships the schema, adapters, and read API only.
### Schema (versioned)
`webui/timeline.py` declares `TIMELINE_SCHEMA_VERSION` (currently `1`) and the
frozen `WorkflowEvent` record. Every response carries `schema_version` so a
consumer can branch on shape. One event:
```json
{
"source": "control_plane",
"event_type": "lease.renew",
"event_key": "cp:1421",
"timestamp": "2026-07-23T02:00:00Z",
"actor": null,
"role": null,
"issue_number": 637,
"pr_number": null,
"session_id": null,
"tool_name": null,
"decision": null,
"message": "lease renewed",
"correlation_id": "issue#637",
"evidence_refs": [],
"sensitive": true
}
```
`event_key` is stable and unique per source (`cp:<event_id>`,
`cth:<kind>:<number>:<comment_id>`), so pagination and dedup are deterministic.
### Sources and field authority
| Source | Adapter | Authority |
|---|---|---|
| Control-plane `events``work_items` | `adapt_cp_events` | `event_type`, `message`, `timestamp`, issue/PR scope come from the CP database, read through a `mode=ro` URI (never creates the DB or runs migrations) |
| Gitea Canonical Thread Handoff comments | `adapt_cth_comments` | `actor`, `role` (next owner), `decision`, `evidence_refs`, `timestamp` come from the parsed CTH comment body (`canonical_thread_handoff`) |
Handoff comments are thread-scoped: they are only read when the request filters
by a single `issue` or `pr`. Otherwise the handoff source reports `not run`
with a reason — it is never rendered as empty-and-healthy. Each source degrades
independently: an unavailable control-plane DB or a failed comment fetch is a
`sources[]` entry with `ok:false` and a `reason`, never a dropped timeline.
### Query parameters
`issue`, `pr`, `session` (conjunctive filters); `limit` (default 50, max 500)
and `offset` for pagination; `remote`, `org`, `repo` to override the default
registry-project scope. Events sort ascending by
`(timestamp, source_rank, event_key)`; missing timestamps sort last.
### Filter authority, and refusing what cannot be answered
A filter dimension is only meaningful for a source whose records carry it.
Each source declares its own support in `_SOURCE_FILTER_SUPPORT` and reports it
per response as `supported_filters` / `unsupported_filters`:
| Source | issue | pr | session |
|---|---|---|---|
| `control_plane` | yes | yes | **no** — the `events` table is `(event_id, work_item_id, event_type, message, created_at)` and records no session |
| `gitea_handoff` | yes | yes | yes — a CTH comment declares its own `Session:` field |
`session_id` is read only from that declared CTH field. It is never inferred
from a work item, an actor, or message text, and a value that is
redaction-altering or bare-secret-shaped is dropped rather than emitted.
When **no source that ran** can carry a requested dimension, the request is
refused rather than answered: the response is `422` with `ok:false` and a
structured `error` naming `unsupported_filters` and the per-source reason. A
`200` with zero events would tell an operator that no such activity exists,
which is a stronger — and false — claim than "this cannot be answered here".
A source that *can* answer the dimension and simply matched nothing still
returns `200` with `ok:true` and an empty page.
### Redaction
Every free-text field (event messages, decision/proof text, roles, actors) is
passed through the console redaction policy (`webui.console_redaction`, backed
by `gitea_audit.redact`) before it leaves the module, failing closed to the
placeholder. No unredacted tool arguments or secrets are ever emitted, and a
generation error never drops raw data to a caller or a log.
Redaction also runs *before* any structured value is derived from free text.
`evidence_refs` are extracted from already-redacted proof/decision text, and a
commit reference is recognised only where the text declares one (`commit`,
`head`, `base`, `sha`, …). An undeclared 40-character hex run has the exact
shape of a Gitea access token, so it is never lifted out of prose into a
structured field. Every reference is then independently revalidated against an
allowed shape and a second redaction pass immediately before serialization;
anything unproven is dropped and the event is flagged `sensitive`.
### Tests
```bash
pytest tests/test_webui_timeline.py -q
```
## Tests
```bash
+155 -60
View File
@@ -2020,6 +2020,7 @@ import allocator_dependencies # noqa: E402
import dependency_graph # noqa: E402 # #784 durable dependency edges
import control_plane_db # noqa: E402
import lease_lifecycle # noqa: E402
import lease_policy # noqa: E402
import workflow_dashboard # noqa: E402 # #605 live queue/lease dashboard
import incident_bridge # noqa: E402
import sentry_observability # noqa: E402 (#606 optional Sentry observability)
@@ -2262,7 +2263,6 @@ import canonical_comment_validator as ccv # noqa: E402
# GITEA_ISSUE_LOCK_DIR, bound to the current MCP session via a per-PID pointer.
# Legacy global path retained only for test/doc references — do not seed manually.
ISSUE_LOCK_FILE = "/tmp/gitea_issue_lock.json"
WORK_LEASE_TTL_HOURS = 4
AUTHOR_ISSUE_WORK_LEASE = "author_issue_work"
VALID_WORK_LEASE_OPERATIONS = frozenset({
AUTHOR_ISSUE_WORK_LEASE,
@@ -2562,7 +2562,12 @@ def _build_author_issue_work_lease(
host: str | None,
) -> dict:
created = _work_lease_now()
expires = created + timedelta(hours=WORK_LEASE_TTL_HOURS)
# #790 Slice A: the window comes from the central policy, not a literal here.
# It is also now a *sliding* window — the lease lives ``initial_ttl_minutes``
# past its last valid heartbeat rather than a fixed four hours past its
# creation, so an abandoned task stops holding the claim within one TTL.
policy = lease_policy.policy_for(lease_policy.TASK_CLASS_AUTHOR_ISSUE_WORK)
expires = created + timedelta(minutes=policy.initial_ttl_minutes)
return {
"operation_type": AUTHOR_ISSUE_WORK_LEASE,
"issue_number": issue_number,
@@ -2573,6 +2578,15 @@ def _build_author_issue_work_lease(
"created_at": _work_lease_timestamp(created),
"expires_at": _work_lease_timestamp(expires),
"last_heartbeat_at": _work_lease_timestamp(created),
# #790 AC-N1: the ownership key for this task. Distinct from the recorded
# PID, which is the shared daemon and identifies no individual task.
"task_session_id": issue_lock_store.mint_task_session_id(
AUTHOR_ISSUE_WORK_LEASE
),
# #790 AC-N8: the explicit lifecycle marker. Its absence — never a
# timestamp comparison — is what makes a lock legacy.
"lifecycle_version": lease_policy.LIFECYCLE_HEARTBEAT_V1,
"heartbeat_count": 1,
}
@@ -4342,6 +4356,138 @@ def gitea_lock_issue(
return result
@mcp.tool()
def gitea_heartbeat_issue_lock(
issue_number: int,
branch_name: str,
task_session_id: str | None = None,
remote: str = "dadeschools",
host: str | None = None,
org: str | None = None,
repo: str | None = None,
worktree_path: str | None = None,
expected_generation: int | None = None,
) -> dict:
"""Prove an owned author issue lease is still active (#790 Slice A).
The task-liveness signal the lifecycle was missing. Before this, an author
lease carried a fixed four-hour expiry that nothing could shorten, and the
only liveness evidence was the recorded PID the long-lived MCP daemon,
which stays alive across every task it serves and so proved nothing about
whether the authoring task still held the work.
Each successful call slides the lease ``initial_ttl_minutes`` past *now*
from the central policy, so an actively heartbeating session is never
evicted while an abandoned one releases its claim within one TTL.
What this tool cannot do, by construction:
* **Acquire.** It refuses when no durable lock exists.
* **Take over.** Exact issue, branch, realpath-normalized worktree,
claimant username, claimant profile, and recorded task-session identifier
must all match; a superseded session holding an older identifier is
refused.
* **Revive.** A lease already past its grace is not heartbeatable that
would let a session restore ownership it had stopped proving. It must use
the sanctioned reclaim path, which mints a new generation.
A lock predating the heartbeat lifecycle is rebound rather than heartbeated:
its exact owner is re-verified and a genuine task-session identifier and
first heartbeat are minted (#790 AC-N8). The rebind is decided server-side
from the durable lifecycle marker; there is no caller-facing switch.
Args:
issue_number: The locked issue number.
branch_name: The branch recorded on the lock.
task_session_id: The identifier this session received when it acquired
or rebound the lock. It is a fencing token, not an ownership
assertion: it is compared against durable state and can only ever
cause a refusal, never grant anything. Omitted only when rebinding a
legacy lock, which has no identifier yet and mints one.
remote: Known instance 'dadeschools' or 'prgs'.
host: Override the Gitea host.
org: Override the owner/organization.
repo: Override the repository name.
worktree_path: Author worktree recorded on the lock.
expected_generation: Optional fencing value. The per-issue flock already
serializes the read and the write, so this is for a caller that
wants to pin the generation it last observed across calls; a moved
generation fails closed.
Returns:
dict with 'success', 'performed', the sliding 'expires_at',
'last_heartbeat_at', 'lock_generation', 'task_session_id', the applied
'policy', and post-write 'freshness'; on refusal 'success'/'performed'
False with 'reasons' naming exactly what did not match.
"""
blocked = _profile_permission_block(
task_capability_map.required_permission("heartbeat_issue_lock"),
issue_number=issue_number,
remote=remote,
host=host,
org=org,
repo=repo,
org_explicit=org is not None,
repo_explicit=repo is not None,
)
if blocked:
return blocked
resolved_worktree = issue_lock_worktree.resolve_author_worktree_path(
worktree_path, _canonical_local_git_root()
)
h, o, r = _resolve(remote, host, org, repo)
claimant = _work_lease_claimant(h)
identity = claimant.get("username")
profile = claimant.get("profile")
existing = _load_existing_issue_lock(
remote=remote, org=o, repo=r, issue_number=issue_number
)
if not existing:
return {
"success": False,
"performed": False,
"issue_number": issue_number,
"reasons": [
f"no durable lock for issue #{issue_number}; heartbeat cannot "
"acquire a claim (fail closed)"
],
}
if issue_lock_store.is_legacy_lease(existing):
# AC-N8 exit route one: canonical exact-owner rebinding. The other exit
# is terminal retirement, which is Slice B.
outcome = issue_lock_store.rebind_legacy_lock(
remote=remote,
org=o,
repo=r,
issue_number=issue_number,
branch_name=branch_name,
worktree_path=resolved_worktree,
identity=identity,
profile=profile,
expected_generation=expected_generation,
)
outcome["operation"] = "legacy_rebind"
return outcome
outcome = issue_lock_store.heartbeat_session_lock(
remote=remote,
org=o,
repo=r,
issue_number=issue_number,
branch_name=branch_name,
worktree_path=resolved_worktree,
identity=identity,
profile=profile,
task_session_id=str(task_session_id or ""),
expected_generation=expected_generation,
)
outcome["operation"] = "heartbeat"
return outcome
@mcp.tool()
def gitea_assess_work_issue_duplicate(
issue_number: int,
@@ -19912,7 +20058,6 @@ def _allocator_candidates_from_gitea(
state="open",
labels=tuple(labels),
title=title,
body=body,
priority=20 if "status:ready" in labels else 1,
blocked=blocked,
dependency_unmet=dep_unmet,
@@ -21149,21 +21294,10 @@ def gitea_adopt_workflow_lease(
remote: str = "dadeschools",
host: str | None = None,
) -> dict:
"""Adopt a control-plane lease through the sanctioned path (#601 / #843).
"""Adopt a control-plane lease through the sanctioned path (#601).
Same-owner resume refreshes provenance. Foreign active leases are refused
unless the lease is a pending controller cross-role handoff and the caller
holds the required role (independent consume without sharing the
controller session). Expired leases may be reclaimed; provenance records
adopted_from/by. Terminal (abandoned/released) leases cannot be adopted.
#843 F1: the adopter role is derived authoritatively from the active
authenticated profile never from caller input. A supplied ``role`` that
does not exactly match the profile-derived role is rejected (no silent
accept or reinterpretation), and handoff provenance ``required_profile`` /
``required_namespace`` restrictions are validated against the same
authoritative caller context. Caller-supplied role/profile/namespace can
never grant authority.
Same-owner resume refreshes provenance. Foreign active leases are refused.
Expired leases may be reclaimed; provenance records adopted_from/by.
"""
read_block = _profile_operation_gate("gitea.read")
if read_block:
@@ -21172,64 +21306,25 @@ def gitea_adopt_workflow_lease(
"reasons": read_block,
"permission_report": _permission_block_report("gitea.read"),
}
profile = get_profile()
profile_name = (profile.get("profile_name") or "").strip() or "session"
active_role = (_profile_role_kind(profile) or "").strip().lower()
if not active_role:
return {
"success": False,
"outcome": "blocked",
"mutation_performed": False,
"reasons": [
"active profile role could not be derived authoritatively; "
"refusing lease adoption (fail closed, #843)"
],
"lease_id": lease_id,
"authoritative_source": "control_plane_db",
"file_lock_only": False,
"comment_lease_only": False,
}
if role is not None and str(role).strip():
supplied_role = str(role).strip().lower()
if supplied_role != active_role:
return {
"success": False,
"outcome": "blocked",
"mutation_performed": False,
"profile_role_kind": active_role,
"supplied_role": supplied_role,
"reasons": [
f"caller-supplied role '{supplied_role}' does not match "
f"the authenticated profile-derived role '{active_role}'; "
"caller-supplied role/profile/namespace can never grant "
"authority (fail closed, #843)"
],
"lease_id": lease_id,
"authoritative_source": "control_plane_db",
"file_lock_only": False,
"comment_lease_only": False,
}
db, errs = _control_plane_db_or_error()
if db is None:
return {"success": False, "reasons": errs}
profile = get_profile()
profile_name = (profile.get("profile_name") or "").strip() or "session"
active_role = _profile_role_kind(profile) or "author"
sid = (session_id or "").strip() or (
f"{profile_name}-{os.getpid()}-{uuid.uuid4().hex[:8]}"
)
adopter_namespace = allocator_service.DEFAULT_ROLE_NAMESPACES.get(
active_role, f"gitea-{active_role}"
)
try:
return lease_lifecycle.adopt_lease(
db,
lease_id=lease_id,
adopter_session_id=sid,
role=active_role,
role=(role or active_role).strip() or "author",
worktree_path=worktree_path,
expected_head_sha=expected_head_sha,
owner_pid=os.getpid(),
operator_authorized=bool(operator_authorized),
adopter_profile_name=profile_name,
adopter_namespace=adopter_namespace,
)
except (lease_lifecycle.LeaseLifecycleError, control_plane_db.ControlPlaneError) as exc:
return {
+553 -29
View File
@@ -15,15 +15,27 @@ import json
import os
import re
import tempfile
import uuid
from contextlib import contextmanager
from datetime import datetime, timedelta, timezone
from typing import Any
import lease_policy
LOCK_DIR_ENV = "GITEA_ISSUE_LOCK_DIR"
DEFAULT_LOCK_DIR = os.path.expanduser("~/.cache/gitea-tools/issue-locks")
WORK_LEASE_TTL_HOURS = 4
AUTHOR_ISSUE_WORK_LEASE = "author_issue_work"
# Freshness classifications. ``STATUS_STALE`` remains the dead-PID band that
# #753 recovery keys on; the two bands below are new in #790 Slice A and apply
# only to leases minted under the heartbeat lifecycle.
STATUS_LIVE = "live"
STATUS_EXPIRED = "expired"
STATUS_ABSENT = "absent"
STATUS_STALE = "stale"
STATUS_STALE_MISSED_HEARTBEAT = "stale_missed_heartbeat"
STATUS_STALE_ABSOLUTE_CAP = "stale_absolute_cap"
_SAFE_SEGMENT_RE = re.compile(r"[^A-Za-z0-9._+-]+")
@@ -253,6 +265,331 @@ def bind_session_lock(
return path
def _ownership_refusals(
lock: dict[str, Any],
*,
issue_number: int,
branch_name: str,
worktree_path: str,
identity: str | None,
profile: str | None,
) -> list[str]:
"""Exact-ownership mismatches between a durable lock and a live caller.
Shared by the heartbeat writer and the legacy rebind path so the two cannot
disagree about what "the same owner" means. Every field is compared against
durable state; nothing is taken on the caller's word beyond the identity the
server itself resolved.
"""
reasons: list[str] = []
if lock.get("issue_number") != issue_number:
reasons.append(
f"lock targets issue #{lock.get('issue_number')}, not #{issue_number}"
)
if str(lock.get("branch_name") or "") != str(branch_name or ""):
reasons.append(
f"lock branch '{lock.get('branch_name')}' does not match '{branch_name}'"
)
if not _same_realpath(str(lock.get("worktree_path") or ""), worktree_path):
reasons.append(
f"lock worktree '{lock.get('worktree_path')}' does not match "
f"'{worktree_path}'"
)
lease = lock.get("work_lease") if isinstance(lock, dict) else None
claimant = lease.get("claimant") if isinstance(lease, dict) else None
claimant = claimant if isinstance(claimant, dict) else {}
recorded_identity = str(claimant.get("username") or "").strip()
recorded_profile = str(claimant.get("profile") or "").strip()
if not recorded_identity or not recorded_profile:
reasons.append("lock does not record both a claimant username and profile")
if recorded_identity and recorded_identity != str(identity or "").strip():
reasons.append(
f"lock claimant '{recorded_identity}' does not match active identity "
f"'{str(identity or '').strip() or 'unknown'}'"
)
if recorded_profile and recorded_profile != str(profile or "").strip():
reasons.append(
f"lock profile '{recorded_profile}' does not match active profile "
f"'{str(profile or '').strip() or 'unknown'}'"
)
return reasons
def _refusal(reasons: list[str], **extra: Any) -> dict[str, Any]:
return {"success": False, "performed": False, "reasons": reasons, **extra}
def heartbeat_session_lock(
*,
remote: str,
org: str,
repo: str,
issue_number: int,
branch_name: str,
worktree_path: str,
identity: str | None,
profile: str | None,
task_session_id: str,
expected_generation: int | None = None,
lock_dir: str | None = None,
now: datetime | None = None,
) -> dict[str, Any]:
"""Slide a heartbeat-lifecycle lease forward (#790 Slice A, A4).
The write happens inside the same per-issue ``flock`` that serializes
acquisition, and under the #772 generation compare-and-swap, so a heartbeat
can never race a concurrent reclaim: whichever lands first moves the
generation and the other fails closed.
Refuses — never revives — in every ambiguous case. A lease that has already
lapsed past its grace is *not* heartbeatable: allowing that would let a
session that stopped proving liveness restore ownership retroactively, which
is precisely the revival AC-N5 forbids. Such a session must go through the
sanctioned reclaim path, which mints a fresh generation.
"""
current = _lease_now(now)
root = _ensure_lock_dir(lock_dir)
path = lock_file_path(
remote=remote, org=org, repo=repo, issue_number=issue_number, lock_dir=root
)
declared_session = str(task_session_id or "").strip()
if not declared_session:
return _refusal(["no task_session_id supplied (fail closed)"])
sentinel = flock_path(path)
try:
with _exclusive_file_lock(sentinel):
lock = read_lock_file(path)
if not lock:
return _refusal([f"no durable lock for issue #{issue_number}"])
if is_legacy_lease(lock):
return _refusal(
[
"lock predates the heartbeat lifecycle; it must be rebound "
"by its exact owner before it can be heartbeated"
],
lifecycle=lease_lifecycle_version(lock),
legacy_lease=True,
)
reasons = _ownership_refusals(
lock,
issue_number=issue_number,
branch_name=branch_name,
worktree_path=worktree_path,
identity=identity,
profile=profile,
)
recorded_session = lease_task_session_id(lock)
if not recorded_session:
reasons.append(
"lock declares the heartbeat lifecycle but records no "
"task_session_id (fail closed)"
)
elif recorded_session != declared_session:
# A superseded session holding an old identifier cannot heartbeat
# over the session that replaced it.
reasons.append(
"task_session_id does not match the session recorded on the lock"
)
if reasons:
return _refusal(reasons)
current_generation = lock_generation(lock)
if (
expected_generation is not None
and current_generation != expected_generation
):
return _refusal(
[
f"lock generation changed: expected {expected_generation}, "
f"found {current_generation}; another session reclaimed or "
"replaced this claim (fail closed)"
],
lock_generation=current_generation,
)
freshness = assess_lock_freshness(lock, now=current)
if not freshness.get("live"):
return _refusal(
[
f"lease is not live ({freshness.get('status')}): "
f"{freshness.get('reason')}; a lapsed lease must be "
"reclaimed, not heartbeated"
],
freshness=freshness,
)
policy = lease_policy.policy_for(lease_task_class(lock))
expires = current + timedelta(minutes=policy.initial_ttl_minutes)
record = dict(lock)
lease = dict(record.get("work_lease") or {})
prior_heartbeat = lease.get("last_heartbeat_at")
lease["last_heartbeat_at"] = _format_lease_timestamp(current)
lease["expires_at"] = _format_lease_timestamp(expires)
try:
lease["heartbeat_count"] = int(lease.get("heartbeat_count") or 0) + 1
except (TypeError, ValueError):
lease["heartbeat_count"] = 1
record["work_lease"] = lease
record["lock_generation"] = current_generation + 1
save_lock_file(path, record)
except LockContentionError as exc:
return _refusal([f"issue #{issue_number} lock contention: {exc} (fail closed)"])
return {
"success": True,
"performed": True,
"issue_number": issue_number,
"branch_name": branch_name,
"worktree_path": worktree_path,
"task_session_id": declared_session,
"lock_generation": record["lock_generation"],
"prior_generation": current_generation,
"prior_heartbeat_at": prior_heartbeat,
"last_heartbeat_at": lease["last_heartbeat_at"],
"expires_at": lease["expires_at"],
"heartbeat_count": lease["heartbeat_count"],
"lock_file_path": path,
"policy": lease_policy.describe(lease_task_class(record)),
"freshness": assess_lock_freshness(record, now=current),
}
def rebind_legacy_lock(
*,
remote: str,
org: str,
repo: str,
issue_number: int,
branch_name: str,
worktree_path: str,
identity: str | None,
profile: str | None,
expected_generation: int | None = None,
lock_dir: str | None = None,
now: datetime | None = None,
) -> dict[str, Any]:
"""Move a legacy lock into the heartbeat lifecycle (#790 AC-N8).
One of the two sanctioned exits from the preserved-expiry legacy state; the
other is terminal retirement, which is Slice B. Only the exact recorded
owner may rebind, and only while the legacy lock is still live under its
original absolute expiry — an already-expired legacy lease belongs to the
#760 renewal path or #601 reclaim, and this must not become a second, weaker
way to revive one.
The rebind mints a genuine task-session identifier and a genuine first
heartbeat. It does not fabricate history: the original creation and expiry
are preserved under ``legacy_origin`` for audit, and the new lifecycle's
absolute cap runs from the rebind, not from the legacy claim.
"""
current = _lease_now(now)
root = _ensure_lock_dir(lock_dir)
path = lock_file_path(
remote=remote, org=org, repo=repo, issue_number=issue_number, lock_dir=root
)
sentinel = flock_path(path)
try:
with _exclusive_file_lock(sentinel):
lock = read_lock_file(path)
if not lock:
return _refusal([f"no durable lock for issue #{issue_number}"])
if not is_legacy_lease(lock):
return _refusal(
[
"lock is already on the heartbeat lifecycle; use the "
"heartbeat path"
],
lifecycle=lease_lifecycle_version(lock),
legacy_lease=False,
)
reasons = _ownership_refusals(
lock,
issue_number=issue_number,
branch_name=branch_name,
worktree_path=worktree_path,
identity=identity,
profile=profile,
)
if reasons:
return _refusal(reasons)
current_generation = lock_generation(lock)
if (
expected_generation is not None
and current_generation != expected_generation
):
return _refusal(
[
f"lock generation changed: expected {expected_generation}, "
f"found {current_generation} (fail closed)"
],
lock_generation=current_generation,
)
freshness = assess_lock_freshness(lock, now=current)
if not freshness.get("live"):
return _refusal(
[
f"legacy lease is not live ({freshness.get('status')}): "
f"{freshness.get('reason')}; rebinding is not a recovery "
"path for a lapsed lease"
],
freshness=freshness,
)
policy = lease_policy.policy_for(lease_task_class(lock))
expires = current + timedelta(minutes=policy.initial_ttl_minutes)
session_id = mint_task_session_id(lease_task_class(lock))
record = dict(lock)
lease = dict(record.get("work_lease") or {})
legacy_origin = {
"created_at": lease.get("created_at"),
"expires_at": lease.get("expires_at"),
"last_heartbeat_at": lease.get("last_heartbeat_at"),
"lifecycle": lease_policy.LIFECYCLE_LEGACY,
}
lease["lifecycle_version"] = lease_policy.LIFECYCLE_HEARTBEAT_V1
lease["task_session_id"] = session_id
lease["created_at"] = _format_lease_timestamp(current)
lease["last_heartbeat_at"] = _format_lease_timestamp(current)
lease["expires_at"] = _format_lease_timestamp(expires)
lease["heartbeat_count"] = 1
record["work_lease"] = lease
record["legacy_rebind"] = {
"rebound_at": _format_lease_timestamp(current),
"task_session_id": session_id,
"prior_generation": current_generation,
"legacy_origin": legacy_origin,
"reason": (
"legacy lock rebound into the heartbeat lifecycle by its exact "
"recorded owner"
),
}
record["lock_generation"] = current_generation + 1
save_lock_file(path, record)
except LockContentionError as exc:
return _refusal([f"issue #{issue_number} lock contention: {exc} (fail closed)"])
return {
"success": True,
"performed": True,
"issue_number": issue_number,
"task_session_id": session_id,
"lock_generation": record["lock_generation"],
"prior_generation": current_generation,
"lifecycle": lease_policy.LIFECYCLE_HEARTBEAT_V1,
"legacy_rebind": record["legacy_rebind"],
"expires_at": lease["expires_at"],
"last_heartbeat_at": lease["last_heartbeat_at"],
"lock_file_path": path,
"freshness": assess_lock_freshness(record, now=current),
}
def read_session_issue_lock(lock_dir: str | None = None) -> dict[str, Any] | None:
root = (lock_dir or default_lock_dir()).strip()
pointer = read_lock_file(session_pointer_path(root))
@@ -336,6 +673,16 @@ def _parse_lease_timestamp(value: str | None) -> datetime | None:
return None
def _format_lease_timestamp(value: datetime) -> str:
"""Serialize a lease timestamp in the durable ``...Z`` form already on disk."""
return (
value.astimezone(timezone.utc)
.replace(microsecond=0)
.isoformat()
.replace("+00:00", "Z")
)
def lease_expires_at(lock: dict[str, Any] | None) -> datetime | None:
if not lock:
return None
@@ -356,60 +703,216 @@ def is_lease_live(lock: dict[str, Any] | None, *, now: datetime | None = None) -
return assess_lock_freshness(lock, now=now)["live"]
def lease_task_class(lock_data: dict[str, Any] | None) -> str:
"""Policy task class for a durable lock; author work when unrecorded."""
lease = lock_data.get("work_lease") if isinstance(lock_data, dict) else None
if isinstance(lease, dict):
recorded = str(lease.get("operation_type") or "").strip()
if recorded:
return recorded
return AUTHOR_ISSUE_WORK_LEASE
def lease_lifecycle_version(lock_data: dict[str, Any] | None) -> str:
"""Read the durable lifecycle marker (#790 AC-N8).
The marker is the *only* discriminator between a heartbeat-lifecycle lease
and a legacy one. Timestamps are deliberately not consulted: a lock minted
before this lifecycle existed has ``last_heartbeat_at == created_at``
forever, and reading that equality as "recently heartbeated" would treat
every never-heartbeated legacy lock as fresh — the precise inversion AC-N8
forbids. A newly minted heartbeat lease also has the two equal, so the
equality carries no information in either direction.
"""
lease = lock_data.get("work_lease") if isinstance(lock_data, dict) else None
if isinstance(lease, dict):
recorded = str(lease.get("lifecycle_version") or "").strip()
if recorded:
return recorded
return lease_policy.LIFECYCLE_LEGACY
def is_legacy_lease(lock_data: dict[str, Any] | None) -> bool:
"""True when a lock predates the shared heartbeat lifecycle."""
return lease_lifecycle_version(lock_data) != lease_policy.LIFECYCLE_HEARTBEAT_V1
def lease_task_session_id(lock_data: dict[str, Any] | None) -> str:
"""Recorded per-task session identifier, or empty for a legacy lock."""
lease = lock_data.get("work_lease") if isinstance(lock_data, dict) else None
if isinstance(lease, dict):
return str(lease.get("task_session_id") or "").strip()
return ""
def mint_task_session_id(task_class: str = AUTHOR_ISSUE_WORK_LEASE) -> str:
"""Mint an ownership key for one task (#790 AC-N1).
Deliberately contains no process identifier. The recorded PID belongs to the
long-lived MCP daemon, which outlives any individual task and is reused by
every task it serves, so PID digits cannot identify *which* task holds a
claim. The PID is still recorded alongside this value as evidence.
"""
prefix = _sanitize_segment(str(task_class or AUTHOR_ISSUE_WORK_LEASE))
return f"{prefix}-{uuid.uuid4().hex[:16]}"
def _lease_heartbeat_at(lock_data: dict[str, Any] | None) -> datetime | None:
lease = lock_data.get("work_lease") if isinstance(lock_data, dict) else None
heartbeat_at = None
if isinstance(lock_data, dict):
heartbeat_at = _parse_lease_timestamp(lock_data.get("last_heartbeat_at"))
if heartbeat_at is None and isinstance(lease, dict):
heartbeat_at = _parse_lease_timestamp(lease.get("last_heartbeat_at"))
return heartbeat_at
def assess_lock_freshness(
lock_data: dict[str, Any] | None,
*,
now: datetime | None = None,
) -> dict[str, Any]:
"""Classify a lock as live, expired, stale, or absent."""
"""Classify a lock as live, expired, stale, or absent.
#790 Slice A makes the heartbeat load-bearing. Before this change
``last_heartbeat_at`` was parsed and then never consulted: liveness was
decided entirely by the absolute ``expires_at`` and by PID liveness, and
since the recorded PID is the long-lived MCP daemon, an abandoned author
task stayed "live" for the full four-hour TTL.
Two rules govern the rewrite:
* **An alive PID never establishes freshness** (AC-N2). It proves the daemon
is up, nothing about the task. It is recorded as evidence and no branch
returns ``live`` because of it.
* **A dead PID still corroborates staleness.** The dead-PID band is
unchanged and still precedes every heartbeat evaluation, so #753
dead-session recovery keys on exactly the classification it always did.
Legacy leases (AC-N8) keep their recorded absolute expiry and are never
evaluated against the short heartbeat grace, so deploying this change cannot
make an existing claim instantly reclaimable.
"""
current = _lease_now(now)
if not lock_data:
return {
"status": "absent",
"status": STATUS_ABSENT,
"live": False,
"stale": False,
"reason": "no lock record",
}
expires_at = lease_expires_at(lock_data)
lease = lock_data.get("work_lease")
heartbeat_at = _parse_lease_timestamp(lock_data.get("last_heartbeat_at"))
if heartbeat_at is None and isinstance(lease, dict):
heartbeat_at = _parse_lease_timestamp(lease.get("last_heartbeat_at"))
expires_at = lease_expires_at(lock_data)
heartbeat_at = _lease_heartbeat_at(lock_data)
created_at = (
_parse_lease_timestamp(lease.get("created_at"))
if isinstance(lease, dict)
else None
)
pid = lock_data.get("session_pid")
if pid is None:
pid = lock_data.get("pid")
# Evidence only. Never consulted to grant liveness (AC-N2).
pid_alive = is_process_alive(pid) if pid is not None else False
if expires_at and expires_at <= current:
return {
"status": "expired",
"live": False,
"stale": True,
"reason": f"lease expired at {expires_at.isoformat()}",
"pid_alive": pid_alive,
}
lifecycle = lease_lifecycle_version(lock_data)
legacy = lifecycle != lease_policy.LIFECYCLE_HEARTBEAT_V1
policy = lease_policy.policy_for(lease_task_class(lock_data))
if pid is not None and not pid_alive:
return {
"status": "stale",
"live": False,
"stale": True,
"reason": f"owner pid {pid} is not alive",
"pid_alive": False,
}
return {
"status": "live",
"live": True,
"stale": False,
"reason": "lock heartbeat and lease are fresh",
evidence: dict[str, Any] = {
"pid_alive": pid_alive,
"lifecycle": lifecycle,
"legacy_lease": legacy,
"task_session_id": lease_task_session_id(lock_data) or None,
"heartbeat_at": heartbeat_at.isoformat() if heartbeat_at else None,
"expires_at": expires_at.isoformat() if expires_at else None,
}
def _result(status: str, *, live: bool, reason: str, **extra: Any) -> dict[str, Any]:
return {
"status": status,
"live": live,
"stale": not live and status != STATUS_ABSENT,
"reason": reason,
**evidence,
**extra,
}
if legacy:
# AC-N8: the preserved absolute expiry is the only clock for a lock
# written before task-session heartbeats existed.
if expires_at and expires_at <= current:
return _result(
STATUS_EXPIRED,
live=False,
reason=f"lease expired at {expires_at.isoformat()}",
)
if pid is not None and not pid_alive:
return _result(
STATUS_STALE, live=False, reason=f"owner pid {pid} is not alive"
)
return _result(
STATUS_LIVE,
live=True,
reason=(
"legacy lease is within its recorded absolute expiry; the "
"heartbeat grace does not apply retroactively"
),
legacy_expiry_preserved=True,
)
# ── Heartbeat lifecycle ──
if pid is not None and not pid_alive:
# Unchanged dead-PID band: #753 recovery depends on this exact status.
return _result(STATUS_STALE, live=False, reason=f"owner pid {pid} is not alive")
if heartbeat_at is None:
# Contradictory: a heartbeat lease must carry a heartbeat. Fail closed.
return _result(
STATUS_STALE_MISSED_HEARTBEAT,
live=False,
reason=(
f"lease declares lifecycle '{lifecycle}' but records no "
"last_heartbeat_at (fail closed)"
),
)
if policy.absolute_cap_hours and created_at is not None:
cap_at = created_at + timedelta(hours=policy.absolute_cap_hours)
if cap_at <= current:
return _result(
STATUS_STALE_ABSOLUTE_CAP,
live=False,
reason=(
f"lease exceeded its {policy.absolute_cap_hours}h absolute cap "
f"at {cap_at.isoformat()}; canonical re-adoption is required"
),
absolute_cap_at=cap_at.isoformat(),
)
grace_at = heartbeat_at + timedelta(minutes=policy.missed_heartbeat_grace_minutes)
if grace_at <= current or (expires_at is not None and expires_at <= current):
return _result(
STATUS_STALE_MISSED_HEARTBEAT,
live=False,
reason=(
f"no valid heartbeat since {heartbeat_at.isoformat()}; the "
f"{policy.missed_heartbeat_grace_minutes}min grace lapsed at "
f"{grace_at.isoformat()}"
),
missed_heartbeat_since=grace_at.isoformat(),
)
warning_at = heartbeat_at + timedelta(minutes=policy.stale_warning_minutes)
return _result(
STATUS_LIVE,
live=True,
reason="lease heartbeat is fresh within the configured grace",
heartbeat_warning=warning_at <= current,
)
def _same_realpath(left: str | None, right: str | None) -> bool:
if not left or not right:
@@ -446,6 +949,27 @@ def assess_expired_lock_reclaim(
"reasons": ["lock is still live; cannot reclaim (fail closed)"],
"freshness": freshness,
}
status = str(freshness.get("status") or "")
if status in (STATUS_STALE_MISSED_HEARTBEAT, STATUS_STALE_ABSOLUTE_CAP):
# #790: under the heartbeat lifecycle the heartbeat *is* the liveness
# proof, so a session that stopped heartbeating past its grace has
# released its claim by definition. Requiring a dead PID on top of that
# would reinstate the original defect — the recorded PID is the shared
# daemon, which stays alive across every abandoned task it ever served.
#
# This band is unreachable for a legacy lease (AC-N8), so no lock
# written before this lifecycle can be reclaimed by this path.
return {
"reclaim_allowed": True,
"reasons": [
f"heartbeat-lifecycle lease is {status}: {freshness.get('reason')}"
],
"freshness": freshness,
"prior_branch": existing_lock.get("branch_name"),
"prior_worktree": existing_lock.get("worktree_path"),
"prior_pid": existing_lock.get("session_pid") or existing_lock.get("pid"),
"prior_task_session_id": lease_task_session_id(existing_lock) or None,
}
pid = existing_lock.get("session_pid")
if pid is None:
pid = existing_lock.get("pid")
+13 -209
View File
@@ -39,7 +39,6 @@ SAFE_RELEASE_OWNED = "release_owned"
SAFE_STALE_PROMPT = "stale_prompt_lease"
SAFE_UNKNOWN = "inspect_only"
SAFE_NO_AUTHORITY = "file_or_comment_not_authoritative"
SAFE_CONSUME_CROSS_ROLE = "consume_cross_role_handoff"
LEASE_STATUS_ACTIVE = "active"
LEASE_STATUS_RELEASED = "released"
@@ -251,23 +250,6 @@ def decide_safe_next_action(
"same_owner": True,
"also_allowed": [SAFE_ABANDON_ALLOWED, SAFE_RELEASE_OWNED],
}
handoff = is_pending_cross_role_handoff({"lease": lease})
if handoff:
return {
"safe_next_action": SAFE_CONSUME_CROSS_ROLE,
"reasons": [
f"controller allocation pending handoff (freshness={status}); "
"required-role worker may consume without abandon/reassign; "
f"required_role={handoff['required_role']}"
],
"block": False,
"same_owner": False,
"owner_session_id": owner,
"required_role": handoff["required_role"],
"cross_role_handoff": True,
"handoff_status": "pending",
"also_allowed": [SAFE_ABANDON_ALLOWED],
}
return {
"safe_next_action": SAFE_ABANDON_ALLOWED,
"reasons": [
@@ -290,24 +272,6 @@ def decide_safe_next_action(
}
if not same_owner and status == "active":
# #843: pending cross-role handoff is consumable by required role
handoff = is_pending_cross_role_handoff({"lease": lease})
if handoff:
return {
"safe_next_action": SAFE_CONSUME_CROSS_ROLE,
"reasons": [
"controller cross-role allocation pending handoff; "
f"required_role={handoff['required_role']}; "
"consume via gitea_adopt_workflow_lease without "
"abandonment or sharing the controller session"
],
"block": False,
"same_owner": False,
"owner_session_id": owner,
"required_role": handoff["required_role"],
"cross_role_handoff": True,
"handoff_status": "pending",
}
return {
"safe_next_action": SAFE_WAIT_FOREIGN,
"reasons": [
@@ -476,84 +440,6 @@ def list_active_leases(
}
def parse_lease_provenance(lease_or_state: Mapping[str, Any] | None) -> dict[str, Any]:
"""Return durable lease provenance dict (empty when absent/unparseable)."""
if not lease_or_state:
return {}
if "provenance" in lease_or_state and isinstance(lease_or_state.get("provenance"), dict):
return dict(lease_or_state["provenance"])
raw = None
if "provenance_json" in lease_or_state:
raw = lease_or_state.get("provenance_json")
elif "lease" in lease_or_state and isinstance(lease_or_state.get("lease"), Mapping):
raw = lease_or_state["lease"].get("provenance_json")
if not raw:
return {}
if isinstance(raw, dict):
return dict(raw)
try:
loaded = json.loads(raw)
except (TypeError, json.JSONDecodeError):
return {}
return dict(loaded) if isinstance(loaded, dict) else {}
def is_pending_cross_role_handoff(
state: Mapping[str, Any] | None,
) -> dict[str, Any] | None:
"""Return handoff evidence when a controller allocation awaits consume (#843).
A pending handoff is identified by durable provenance written at
cross-role apply time — not by title heuristics or session-id guessing.
"""
if not state:
return None
lease = state.get("lease") if isinstance(state.get("lease"), Mapping) else state
if not isinstance(lease, Mapping):
return None
status = str(lease.get("status") or "").strip().lower()
if status in (LEASE_STATUS_ABANDONED, LEASE_STATUS_RELEASED, LEASE_STATUS_EXPIRED):
return None
prov = parse_lease_provenance(state)
if not prov and isinstance(lease, Mapping):
prov = parse_lease_provenance(lease)
if not prov.get("cross_role_handoff"):
return None
handoff_status = str(prov.get("handoff_status") or "pending").strip().lower()
if handoff_status != "pending":
return None
adopted_by = (
lease.get("adopted_by_session_id")
or prov.get("adopted_by_session_id")
or ""
)
if str(adopted_by).strip():
return None
required_role = str(
prov.get("required_role") or lease.get("role") or ""
).strip().lower()
if not required_role:
return None
return {
"cross_role_handoff": True,
"handoff_status": "pending",
"required_role": required_role,
"allocating_session_id": str(
prov.get("allocating_session_id") or lease.get("session_id") or ""
),
"allocating_role": str(prov.get("allocating_role") or "controller"),
"lease_id": str(lease.get("lease_id") or ""),
"assignment_id": (
str(state["assignment"]["assignment_id"])
if isinstance(state.get("assignment"), Mapping)
and state["assignment"].get("assignment_id")
else None
),
"provenance": prov,
}
def adopt_lease(
db: cpd.ControlPlaneDB,
*,
@@ -564,17 +450,8 @@ def adopt_lease(
expected_head_sha: str | None = None,
owner_pid: int | None = None,
operator_authorized: bool = False,
adopter_profile_name: str | None = None,
adopter_namespace: str | None = None,
) -> dict[str, Any]:
"""Sanctioned adopt path with provenance; never silent foreign steal.
#843 F1: for a pending cross-role handoff, ``role`` must be the
authoritative profile-derived role supplied by the MCP boundary — never
caller-asserted authority. When the handoff provenance declares
``required_profile`` / ``required_namespace`` and the caller context is
provided, both are validated exactly; a mismatch fails closed.
"""
"""Sanctioned adopt path with provenance; never silent foreign steal."""
state = db.get_lease_workflow_state(lease_id)
if not state:
raise LeaseLifecycleError(
@@ -586,8 +463,11 @@ def adopt_lease(
owner = str(lease.get("session_id") or "")
same_owner = owner == str(adopter_session_id)
handoff = is_pending_cross_role_handoff(state)
adopter_role = (role or "").strip().lower()
if freshness["freshness"] == "active" and not same_owner:
raise LeaseLifecycleError(
f"refusing to steal active foreign lease {lease_id} owned by "
f"{owner} (fail closed)"
)
if freshness["freshness"] in ("abandoned", "released"):
raise LeaseLifecycleError(
@@ -595,64 +475,13 @@ def adopt_lease(
"(fail closed)"
)
if handoff and not same_owner:
# Terminal statuses already rejected above. Freshness may be
# active OR stale_dead_process (controller exited) — both are
# consumable without abandonment when handoff is still pending.
if freshness["freshness"] not in (
"active",
"stale_dead_process",
"stale_missing_worktree",
):
raise LeaseLifecycleError(
f"lease {lease_id} freshness={freshness['freshness']}; "
"terminal or non-active allocation cannot be handoff-consumed "
"(fail closed)"
)
required = handoff["required_role"]
if adopter_role != required:
raise LeaseLifecycleError(
f"wrong role for cross-role handoff consume of {lease_id}: "
f"required={required} adopter={adopter_role or 'none'} "
"(fail closed)"
)
# #843 F1: provenance profile/namespace restrictions are validated
# against the authoritative caller context when declared. Caller
# input can never widen authority; a mismatch fails closed.
handoff_prov = handoff.get("provenance") or {}
required_profile = str(
handoff_prov.get("required_profile") or ""
).strip()
if required_profile and adopter_profile_name is not None:
if str(adopter_profile_name).strip() != required_profile:
raise LeaseLifecycleError(
f"wrong profile for cross-role handoff consume of "
f"{lease_id}: required_profile={required_profile} "
f"adopter_profile={adopter_profile_name} (fail closed)"
)
required_namespace = str(
handoff_prov.get("required_namespace") or ""
).strip()
if required_namespace and adopter_namespace is not None:
if str(adopter_namespace).strip() != required_namespace:
raise LeaseLifecycleError(
f"wrong namespace for cross-role handoff consume of "
f"{lease_id}: required_namespace={required_namespace} "
f"adopter_namespace={adopter_namespace} (fail closed)"
)
reason = "cross-role-handoff-consume"
elif freshness["freshness"] == "active" and not same_owner:
raise LeaseLifecycleError(
f"refusing to steal active foreign lease {lease_id} owned by "
f"{owner} (fail closed)"
)
elif not same_owner and freshness["freshness"] in (
# Expired or stale: require abandon-style safety before ownership transfer
# when not same owner; same owner may reclaim.
if not same_owner and freshness["freshness"] in (
"expired",
"stale_dead_process",
"stale_missing_worktree",
):
# Expired or stale (non-handoff): require abandon-style safety before
# ownership transfer when not same owner; same owner may reclaim.
if not operator_authorized and freshness["freshness"] == "expired":
# Deterministic reclaim of expired foreign lease is allowed
# without operator flag (sanctioned expire reclaim).
@@ -663,9 +492,6 @@ def adopt_lease(
f"lease {lease_id} freshness={freshness['freshness']}; "
"use abandon with proof before foreign adopt (fail closed)"
)
reason = "sanctioned-reclaim-adopt"
else:
reason = "owner-resume-adopt" if same_owner else "sanctioned-reclaim-adopt"
provenance = build_adopt_provenance(
adopted_from_session_id=owner,
@@ -678,14 +504,10 @@ def adopt_lease(
worktree_path=worktree_path,
expected_head_sha=expected_head_sha or lease.get("expected_head_sha"),
prior_lease_id=lease_id,
reason=reason,
reason=(
"owner-resume-adopt" if same_owner else "sanctioned-reclaim-adopt"
),
)
if handoff and not same_owner:
provenance["cross_role_handoff"] = True
provenance["handoff_status"] = "adopted"
provenance["required_role"] = handoff["required_role"]
provenance["allocating_session_id"] = handoff["allocating_session_id"]
provenance["allocating_role"] = handoff["allocating_role"]
result = db.adopt_lease(
lease_id=lease_id,
@@ -696,7 +518,7 @@ def adopt_lease(
owner_pid=owner_pid if owner_pid is not None else os.getpid(),
provenance=provenance,
)
out = {
return {
"success": True,
"outcome": result.get("outcome"),
"same_owner": same_owner,
@@ -709,24 +531,6 @@ def adopt_lease(
"comment_lease_only": False,
"reasons": result.get("reasons") or [],
}
if handoff and not same_owner:
out["cross_role_handoff"] = True
out["handoff_status"] = "adopted"
out["required_role"] = handoff["required_role"]
out["adopted_by_session_id"] = adopter_session_id
out["adopted_from_session_id"] = owner
lease_row = result.get("lease") or {}
if isinstance(lease_row, Mapping):
out["read_after_write"] = {
"lease_id": lease_row.get("lease_id"),
"session_id": lease_row.get("session_id"),
"role": lease_row.get("role"),
"status": lease_row.get("status"),
"adopted_by_session_id": lease_row.get("adopted_by_session_id"),
"adopted_from_session_id": lease_row.get("adopted_from_session_id"),
"phase": lease_row.get("phase"),
}
return out
def release_lease(
+212
View File
@@ -0,0 +1,212 @@
"""Central lease policy configuration (#790 Slice A, AC-N7).
The single authoritative source for every lease duration in the project. Before
this module the numbers were scattered: a four-hour author TTL was declared
twice (``issue_lock_store`` and ``gitea_mcp_server``), the reviewer/merger
sliding window lived in ``reviewer_pr_lease``, the conflict-fix window in
``pr_work_lease``, and the control-plane default in ``control_plane_db``.
Nothing tied them together, so tuning one class silently diverged from the
others and no reader could answer "how long does a lease live?" without
grepping four files.
AC-N7 requires that this configuration exist *before* the first heartbeat and
TTL behavior that reads from it, so it ships in Slice A rather than trailing the
code it governs.
Deliberate boundaries:
* **Declaration is not rewiring.** Every task class is declared here, but only
those with ``heartbeat_lifecycle_active`` were migrated onto the shared
heartbeat lifecycle in Slice A — currently ``author_issue_work`` alone.
Reviewer, merger, and conflict-fix leases keep their own existing behavior
until Slice C moves them; their numbers are recorded here so the two cannot
drift apart unnoticed, and ``tests/test_issue_790_lease_policy.py`` asserts
the recorded values still equal the constants those modules use.
* **No policy decision lives here.** This module answers "how long", never "may
this session proceed". Freshness, reclaim, and renewal dispositions stay in
``issue_lock_store``.
"""
from __future__ import annotations
import os
from dataclasses import dataclass
from typing import Any
# Task classes. Only the first is migrated onto the shared lifecycle in Slice A.
TASK_CLASS_AUTHOR_ISSUE_WORK = "author_issue_work"
TASK_CLASS_REVIEWER_PR = "reviewer_pr"
TASK_CLASS_MERGER_PR = "merger_pr"
TASK_CLASS_CONFLICT_FIX = "conflict_fix"
# Durable marker for a lease minted under the shared heartbeat lifecycle.
#
# #790 AC-N8: this explicit marker — never a timestamp comparison — is what
# distinguishes a heartbeat-lifecycle lease from a legacy one. A lock written
# before this lifecycle existed carries no marker and reads as
# ``LIFECYCLE_LEGACY``.
LIFECYCLE_HEARTBEAT_V1 = "heartbeat-v1"
LIFECYCLE_LEGACY = "legacy"
_ENV_PREFIX = "GITEA_LEASE_POLICY"
@dataclass(frozen=True)
class LeasePolicy:
"""Durations governing one task class.
All intervals are minutes except ``absolute_cap_hours``. ``None`` for the
cap means the class has no maximum continuous duration.
"""
task_class: str
initial_ttl_minutes: float
heartbeat_cadence_minutes: float
stale_warning_minutes: float
missed_heartbeat_grace_minutes: float
absolute_cap_hours: float | None
recovery_grace_minutes: float
terminal_race_drain_minutes: float
terminal_retirement_eligible: bool
heartbeat_lifecycle_active: bool
# Defaults. ``author_issue_work`` adopts the reviewer window proven by #747
# rather than inventing new numbers: a lease expires 10 minutes after its last
# valid heartbeat, warns at half that, and an actively heartbeating session is
# never evicted. The prior value was a fixed four hours (240 minutes) that no
# heartbeat could shorten — the defect this issue exists to correct.
_DEFAULTS: dict[str, LeasePolicy] = {
TASK_CLASS_AUTHOR_ISSUE_WORK: LeasePolicy(
task_class=TASK_CLASS_AUTHOR_ISSUE_WORK,
initial_ttl_minutes=10.0,
heartbeat_cadence_minutes=2.0,
stale_warning_minutes=5.0,
missed_heartbeat_grace_minutes=10.0,
absolute_cap_hours=8.0,
recovery_grace_minutes=10.0,
terminal_race_drain_minutes=2.0,
terminal_retirement_eligible=True,
heartbeat_lifecycle_active=True,
),
# Declared, not rewired. These mirror reviewer_pr_lease.LEASE_TTL_MINUTES
# and STALE_WARNING_MINUTES; Slice C migrates the call sites.
TASK_CLASS_REVIEWER_PR: LeasePolicy(
task_class=TASK_CLASS_REVIEWER_PR,
initial_ttl_minutes=10.0,
heartbeat_cadence_minutes=2.0,
stale_warning_minutes=5.0,
missed_heartbeat_grace_minutes=10.0,
absolute_cap_hours=None,
recovery_grace_minutes=10.0,
terminal_race_drain_minutes=2.0,
terminal_retirement_eligible=False,
heartbeat_lifecycle_active=False,
),
TASK_CLASS_MERGER_PR: LeasePolicy(
task_class=TASK_CLASS_MERGER_PR,
initial_ttl_minutes=10.0,
heartbeat_cadence_minutes=2.0,
stale_warning_minutes=5.0,
missed_heartbeat_grace_minutes=10.0,
absolute_cap_hours=None,
recovery_grace_minutes=10.0,
terminal_race_drain_minutes=2.0,
terminal_retirement_eligible=False,
heartbeat_lifecycle_active=False,
),
# Mirrors pr_work_lease.DEFAULT_CONFLICT_FIX_TTL_MINUTES. Deliberately left
# at its current window; shortening it is Slice C's call, not this slice's.
TASK_CLASS_CONFLICT_FIX: LeasePolicy(
task_class=TASK_CLASS_CONFLICT_FIX,
initial_ttl_minutes=120.0,
heartbeat_cadence_minutes=2.0,
stale_warning_minutes=5.0,
missed_heartbeat_grace_minutes=10.0,
absolute_cap_hours=None,
recovery_grace_minutes=10.0,
terminal_race_drain_minutes=2.0,
terminal_retirement_eligible=False,
heartbeat_lifecycle_active=False,
),
}
_NUMERIC_FIELDS = (
"initial_ttl_minutes",
"heartbeat_cadence_minutes",
"stale_warning_minutes",
"missed_heartbeat_grace_minutes",
"absolute_cap_hours",
"recovery_grace_minutes",
"terminal_race_drain_minutes",
)
def env_var_name(task_class: str, field: str) -> str:
"""Environment variable that overrides one field of one task class."""
return f"{_ENV_PREFIX}_{task_class.upper()}_{field.upper()}"
def _override(task_class: str, field: str, default: float | None) -> float | None:
"""Read one override, falling back to *default* on anything unusable.
A malformed or non-positive override is ignored rather than raised: a typo
in an environment variable must not be able to mint a zero-length lease that
makes every claim instantly reclaimable, nor crash the server at import.
"""
raw = (os.environ.get(env_var_name(task_class, field)) or "").strip()
if not raw:
return default
try:
value = float(raw)
except (TypeError, ValueError):
return default
if value <= 0:
return default
return value
def policy_for(task_class: str) -> LeasePolicy:
"""Return the effective policy for *task_class*.
Unknown task classes fall back to the author policy, which is the most
conservative migrated class, rather than raising — a new caller must never
be able to crash a lock write by naming a class this table has not learned.
"""
key = str(task_class or "").strip() or TASK_CLASS_AUTHOR_ISSUE_WORK
base = _DEFAULTS.get(key) or _DEFAULTS[TASK_CLASS_AUTHOR_ISSUE_WORK]
resolved = {
field: _override(base.task_class, field, getattr(base, field))
for field in _NUMERIC_FIELDS
}
if all(resolved[field] == getattr(base, field) for field in _NUMERIC_FIELDS):
return base
return LeasePolicy(
task_class=base.task_class,
terminal_retirement_eligible=base.terminal_retirement_eligible,
heartbeat_lifecycle_active=base.heartbeat_lifecycle_active,
**resolved,
)
def known_task_classes() -> tuple[str, ...]:
"""Every declared task class, migrated or not."""
return tuple(_DEFAULTS)
def describe(task_class: str) -> dict[str, Any]:
"""Serializable view of a policy, for audit records and tool payloads."""
policy = policy_for(task_class)
return {
"task_class": policy.task_class,
"initial_ttl_minutes": policy.initial_ttl_minutes,
"heartbeat_cadence_minutes": policy.heartbeat_cadence_minutes,
"stale_warning_minutes": policy.stale_warning_minutes,
"missed_heartbeat_grace_minutes": policy.missed_heartbeat_grace_minutes,
"absolute_cap_hours": policy.absolute_cap_hours,
"recovery_grace_minutes": policy.recovery_grace_minutes,
"terminal_race_drain_minutes": policy.terminal_race_drain_minutes,
"terminal_retirement_eligible": policy.terminal_retirement_eligible,
"heartbeat_lifecycle_active": policy.heartbeat_lifecycle_active,
"lifecycle_version": LIFECYCLE_HEARTBEAT_V1,
}
+9
View File
@@ -32,6 +32,15 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = {
"permission": "gitea.issue.comment",
"role": "author",
},
# #790 Slice A: prove an owned author lease is still active. Strictly
# narrower than lock_issue — it can only slide a lease this exact session
# already owns, never acquire, take over, or revive one — so it gates on the
# same authority rather than introducing an operation name that every
# already-configured author profile would be missing.
"heartbeat_issue_lock": {
"permission": "gitea.issue.comment",
"role": "author",
},
"set_issue_labels": {
"permission": "gitea.issue.comment",
"role": "author",
@@ -1,243 +0,0 @@
"""Allocator epic / child-only container pre-rank exclusion (#844).
Covers:
* Issue #631-shaped child-only epic is excluded before ranking.
* Implementable child issues remain eligible and can be selected.
* Ordinary issues that merely mention "epic" in title/body are not excluded.
* Excluded containers never receive assignments or workflow leases.
* Structured skip reason ``epic_or_child_only_container`` is reported.
"""
from __future__ import annotations
import os
import tempfile
import unittest
from allocator_service import (
OUTCOME_ASSIGNED,
OUTCOME_PREVIEW,
SKIP_EPIC_OR_CHILD_ONLY_CONTAINER,
WorkCandidate,
allocate_next_work,
classify_epic_or_child_only_container,
)
from control_plane_db import ControlPlaneDB
REMOTE = "prgs"
ORG = "Scaled-Tech-Consulting"
REPO = "Gitea-Tools"
# Minimal body mirroring issue #631 authoritative scope language.
_EPIC_631_BODY = """
## Scope (umbrella)
This epic owns the **product roadmap and linkage** for the Web Console.
Implementation is delivered via child issues only.
## Explicit non-goals
* Do not implement product features in this epic issue itself.
* No product feature implementation is claimed complete solely on this epic.
"""
_CHILD_BODY = """
## Problem
Operators need a workflow-event timeline model for Phase 1.
## Acceptance criteria
- [ ] Timeline model API exists
"""
def _issue(
number: int,
*,
title: str = "",
body: str = "",
labels: tuple[str, ...] = ("status:ready", "type:feature"),
priority: int = 20,
) -> WorkCandidate:
return WorkCandidate(
kind="issue",
number=number,
state="open",
labels=labels,
title=title or f"issue {number}",
body=body,
priority=priority,
)
class ClassifyEpicContainerTest(unittest.TestCase):
def test_631_shaped_body_and_title_is_container(self) -> None:
c = _issue(
631,
title="Epic: MCP Control Plane Web Console",
body=_EPIC_631_BODY,
)
is_c, detail = classify_epic_or_child_only_container(c)
self.assertTrue(is_c)
self.assertIsNotNone(detail)
self.assertIn("body_marker", detail or "")
def test_body_markers_without_epic_title(self) -> None:
c = _issue(
900,
title="Control plane roadmap tracker",
body="Implementation is delivered via child issues only.",
)
is_c, _ = classify_epic_or_child_only_container(c)
self.assertTrue(is_c)
def test_epic_label_alone_is_container(self) -> None:
c = _issue(
901,
title="Roadmap linkage",
body="Track children.",
labels=("status:ready", "type:epic"),
)
is_c, detail = classify_epic_or_child_only_container(c)
self.assertTrue(is_c)
self.assertIn("type:epic", detail or "")
def test_title_epic_prefix_alone_not_container(self) -> None:
"""Title-only 'Epic:' without body scope evidence stays eligible (#844)."""
c = _issue(
902,
title="Epic: something mentioned only in title",
body="Implement a concrete fix for the allocator skip list.",
)
is_c, detail = classify_epic_or_child_only_container(c)
self.assertFalse(is_c)
self.assertIsNone(detail)
def test_incidental_epic_word_not_container(self) -> None:
c = _issue(
903,
title="Document epic handoff conventions",
body=(
"Update the docs so implementable issues that mention an epic "
"remain independently executable."
),
)
is_c, _ = classify_epic_or_child_only_container(c)
self.assertFalse(is_c)
def test_prs_never_classified(self) -> None:
pr = WorkCandidate(
kind="pr",
number=10,
state="open",
title="Epic: fake",
body="Implementation is delivered via child issues only.",
head_sha="a" * 40,
priority=5,
)
is_c, _ = classify_epic_or_child_only_container(pr)
self.assertFalse(is_c)
class AllocateEpicContainerExclusionTest(unittest.TestCase):
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.addCleanup(self._tmp.cleanup)
self.db = ControlPlaneDB(os.path.join(self._tmp.name, "cp.sqlite3"))
def _alloc(self, candidates, **kwargs):
defaults = dict(
session_id="sess-844",
role="author",
remote=REMOTE,
org=ORG,
repo=REPO,
profile_name="prgs-author",
username="jcwalker3",
claims={},
apply=False,
)
defaults.update(kwargs)
return allocate_next_work(self.db, candidates=candidates, **defaults)
def test_631_shaped_epic_excluded_child_selected(self) -> None:
epic = _issue(
631,
title="Epic: MCP Control Plane Web Console",
body=_EPIC_631_BODY,
)
child = _issue(
637,
title="Web Console: Workflow-event timeline model (Phase 1)",
body=_CHILD_BODY,
)
res = self._alloc([epic, child], apply=False)
self.assertTrue(res["success"], res)
self.assertEqual(res["outcome"], OUTCOME_PREVIEW)
self.assertEqual(res["selected"]["number"], 637)
skipped = {s["number"]: s for s in res["skipped"]}
self.assertIn(631, skipped)
self.assertEqual(
skipped[631]["reason_code"], SKIP_EPIC_OR_CHILD_ONLY_CONTAINER
)
self.assertIn(SKIP_EPIC_OR_CHILD_ONLY_CONTAINER, skipped[631]["reason"])
def test_container_cannot_receive_assignment_or_lease(self) -> None:
epic = _issue(
631,
title="Epic: MCP Control Plane Web Console",
body=_EPIC_631_BODY,
)
res = self._alloc([epic], apply=True)
self.assertTrue(res["success"], res)
# Only container present → no safe work; never assigned_work.
self.assertNotEqual(res["outcome"], OUTCOME_ASSIGNED)
self.assertIsNone(res.get("assignment"))
self.assertIsNone(res.get("selected"))
skipped = {s["number"]: s for s in res["skipped"]}
self.assertEqual(
skipped[631]["reason_code"], SKIP_EPIC_OR_CHILD_ONLY_CONTAINER
)
# No lease row for the epic.
leases = self.db.list_active_leases(
remote=REMOTE, org=ORG, repo=REPO
) if hasattr(self.db, "list_active_leases") else []
# Prefer generic inventory if available.
if not leases and hasattr(self.db, "list_leases"):
leases = self.db.list_leases(remote=REMOTE, org=ORG, repo=REPO)
for lease in leases or []:
work_number = lease.get("work_number") if isinstance(lease, dict) else None
self.assertNotEqual(work_number, 631)
def test_incidental_epic_title_remains_eligible(self) -> None:
ordinary = _issue(
700,
title="Document epic handoff conventions",
body="Write runbook text about epic vs child issues.",
)
res = self._alloc([ordinary], apply=False)
self.assertTrue(res["success"], res)
self.assertEqual(res["selected"]["number"], 700)
self.assertEqual(res["skipped"], [])
def test_apply_selects_child_not_epic(self) -> None:
epic = _issue(
631,
title="Epic: MCP Control Plane Web Console",
body=_EPIC_631_BODY,
)
child = _issue(
637,
title="Web Console: Workflow-event timeline model (Phase 1)",
body=_CHILD_BODY,
)
res = self._alloc([epic, child], apply=True)
self.assertTrue(res["success"], res)
self.assertEqual(res["outcome"], OUTCOME_ASSIGNED)
self.assertEqual(res["selected"]["number"], 637)
self.assertEqual(res["assignment"]["work_number"], 637)
if __name__ == "__main__":
unittest.main()
+444
View File
@@ -0,0 +1,444 @@
"""Task heartbeat through the native MCP author path (#790 Slice A, AC-N6).
Assessor-level coverage is not sufficient here, and this project has already
paid for learning that: in review #499 on PR #791 the #760 renewal waiver was
computed correctly and then *discarded* at two later gates, so every real
renewal still failed while the unit suite stayed green. AC-N6 exists because of
that, and requires driving the real tools against a real git repository and a
real durable lock file, composing the gates in production order.
These tests therefore call ``gitea_lock_issue`` and
``gitea_heartbeat_issue_lock`` themselves and assert on what lands on disk,
never on an assessor's return value alone.
"""
from __future__ import annotations
import os
import subprocess
import sys
import tempfile
import unittest
from datetime import datetime, timedelta, timezone
from unittest.mock import patch
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
sys.path.insert(0, os.path.dirname(os.path.abspath(__file__)))
from mutation_profile_fixture import shared_mutation_env # noqa: E402
import issue_lock_provenance # noqa: E402
import issue_lock_store # noqa: E402
import lease_policy # noqa: E402
import mcp_server # noqa: E402
ISSUE = 9791
BRANCH = f"fix/issue-{ISSUE}-heartbeat-mcp"
IDENTITY = "example-user"
PROFILE = "test-author-prgs"
ORG = "Scaled-Tech-Consulting"
REPO = "Gitea-Tools"
def _ts(moment: datetime) -> str:
return (
moment.astimezone(timezone.utc)
.replace(microsecond=0)
.isoformat()
.replace("+00:00", "Z")
)
class _HeartbeatMcpBase(unittest.TestCase):
"""Real git repo plus a real durable lock, driven through the real tools."""
def setUp(self):
self.lock_dir = tempfile.TemporaryDirectory()
self.addCleanup(self.lock_dir.cleanup)
self.repo = tempfile.mkdtemp(prefix="issue790-mcp-")
self.addCleanup(lambda: subprocess.run(["rm", "-rf", self.repo], check=False))
self._init_worktree()
self.remotes = patch.dict(
mcp_server.REMOTES,
{"prgs": {"host": "gitea.prgs.cc", "org": ORG, "repo": REPO}},
)
self.remotes.start()
self.addCleanup(patch.stopall)
mcp_server._IDENTITY_CACHE.clear()
def _git(self, *args):
return subprocess.run(
["git", "-C", self.repo, *args], capture_output=True, text=True, check=True
)
def _init_worktree(self):
self._git("init", "-q", "-b", "master")
self._git("config", "user.email", "[email protected]")
self._git("config", "user.name", "Test")
with open(os.path.join(self.repo, "seed.txt"), "w") as fh:
fh.write("seed\n")
self._git("add", "seed.txt")
self._git("commit", "-q", "-m", "seed")
self.base_sha = self._git("rev-parse", "HEAD").stdout.strip()
# A fresh claim starts base-equivalent, which is the ordinary first-lock
# shape and exercises assess_issue_lock_worktree on its normal path.
self._git("checkout", "-q", "-b", BRANCH)
self.head_sha = self.base_sha
self.worktree = os.path.realpath(self.repo)
def _lock_path(self):
return issue_lock_store.lock_file_path(
remote="prgs",
org=ORG,
repo=REPO,
issue_number=ISSUE,
lock_dir=self.lock_dir.name,
)
def _tool_env(self):
env = shared_mutation_env(
PROFILE, include_example_repo=True, GITEA_ISSUE_LOCK_DIR=self.lock_dir.name
)
env["GITEA_ISSUE_LOCK_DIR"] = self.lock_dir.name
return env
def _git_state(self, *, porcelain="", base_equivalent=True):
return {
"current_branch": BRANCH,
"porcelain_status": porcelain,
"base_equivalent": base_equivalent,
"head_sha": self.head_sha,
"inspected_git_root": self.worktree,
"base_branch": "master",
}
def run_lock_issue(
self,
*,
branch_entries=None,
open_prs=None,
git_state=None,
identity=IDENTITY,
profile=PROFILE,
):
branch_entries = branch_entries if branch_entries is not None else []
open_prs = open_prs if open_prs is not None else []
git_state = git_state or self._git_state()
env = self._tool_env()
with patch(
"mcp_server.api_get_all", return_value=list(branch_entries)
), patch(
"mcp_server._list_open_pulls", return_value=list(open_prs)
), patch(
"mcp_server.get_auth_header", return_value="token x"
), patch(
"mcp_server._work_lease_claimant",
return_value={"username": identity, "profile": profile},
), patch(
"mcp_server.issue_lock_worktree.read_worktree_git_state",
return_value=git_state,
), patch(
"mcp_server.issue_duplicate_context_fetcher",
side_effect=lambda h, o, r, auth, issue_number: (
list(open_prs),
[b.get("name") for b in branch_entries if isinstance(b, dict)],
{"status": "not_claimed"},
),
), patch.dict(os.environ, env, clear=True):
os.environ["GITEA_ISSUE_LOCK_DIR"] = self.lock_dir.name
return mcp_server.gitea_lock_issue(
issue_number=ISSUE,
branch_name=BRANCH,
remote="prgs",
worktree_path=self.worktree,
)
def run_heartbeat(
self, *, task_session_id, identity=IDENTITY, profile=PROFILE, **kwargs
):
env = self._tool_env()
with patch(
"mcp_server._work_lease_claimant",
return_value={"username": identity, "profile": profile},
), patch("mcp_server.get_auth_header", return_value="token x"), patch.dict(
os.environ, env, clear=True
):
os.environ["GITEA_ISSUE_LOCK_DIR"] = self.lock_dir.name
return mcp_server.gitea_heartbeat_issue_lock(
issue_number=ISSUE,
branch_name=kwargs.pop("branch_name", BRANCH),
task_session_id=task_session_id,
remote="prgs",
worktree_path=kwargs.pop("worktree_path", self.worktree),
**kwargs,
)
def write_legacy_lock(self, *, hours_old: float = 3.0, ttl_hours: float = 4.0):
"""A durable lock in the shape the store wrote before this slice."""
now = datetime.now(timezone.utc)
claimant = {"username": IDENTITY, "profile": PROFILE}
created = now - timedelta(hours=hours_old)
record = {
"issue_number": ISSUE,
"branch_name": BRANCH,
"remote": "prgs",
"org": ORG,
"repo": REPO,
"worktree_path": self.worktree,
"session_pid": os.getpid(),
"pid": os.getpid(),
"lock_generation": 1,
"work_lease": {
"operation_type": issue_lock_store.AUTHOR_ISSUE_WORK_LEASE,
"issue_number": ISSUE,
"pr_number": None,
"branch": BRANCH,
"worktree_path": self.worktree,
"claimant": claimant,
"created_at": _ts(created),
# The legacy signature: never advanced past creation.
"last_heartbeat_at": _ts(created),
"expires_at": _ts(created + timedelta(hours=ttl_hours)),
},
"lock_provenance": issue_lock_provenance.build_sanctioned_lock_provenance(
tool="gitea_lock_issue", claimant=claimant
),
}
path = self._lock_path()
record["lock_file_path"] = path
issue_lock_store.save_lock_file(path, record)
return record
class TestLockIssueMintsTheLifecycle(_HeartbeatMcpBase):
"""Durable lock creation and read-back through the real tool."""
def test_native_lock_writes_the_marker_and_a_task_session_id(self):
result = self.run_lock_issue()
self.assertTrue(result["success"], result)
written = issue_lock_store.read_lock_file(result["lock_file_path"])
lease = written["work_lease"]
self.assertEqual(
lease["lifecycle_version"], lease_policy.LIFECYCLE_HEARTBEAT_V1
)
self.assertTrue(lease["task_session_id"])
self.assertFalse(issue_lock_store.is_legacy_lease(written))
# AC-N1: the ownership key is not the daemon pid, which is recorded
# separately as evidence.
self.assertNotIn(str(written["session_pid"]), lease["task_session_id"])
self.assertEqual(written["session_pid"], os.getpid())
def test_native_lease_uses_the_policy_window_not_four_hours(self):
result = self.run_lock_issue()
lease = result["work_lease"]
created = datetime.fromisoformat(lease["created_at"].replace("Z", "+00:00"))
expires = datetime.fromisoformat(lease["expires_at"].replace("Z", "+00:00"))
policy = lease_policy.policy_for(lease_policy.TASK_CLASS_AUTHOR_ISSUE_WORK)
self.assertEqual(
(expires - created).total_seconds() / 60.0, policy.initial_ttl_minutes
)
def test_freshness_of_a_new_native_lock_is_live(self):
result = self.run_lock_issue()
self.assertEqual(
result["lock_freshness"]["status"], issue_lock_store.STATUS_LIVE
)
self.assertTrue(result["lock_freshness"]["live"])
class TestHeartbeatThroughTheTool(_HeartbeatMcpBase):
def _lock_and_session(self):
result = self.run_lock_issue()
self.assertTrue(result["success"], result)
return result, result["work_lease"]["task_session_id"]
def test_heartbeat_slides_the_lease_and_advances_the_generation(self):
locked, session = self._lock_and_session()
before = issue_lock_store.read_lock_file(locked["lock_file_path"])
beat = self.run_heartbeat(task_session_id=session)
self.assertTrue(beat["success"], beat)
self.assertEqual(beat["operation"], "heartbeat")
after = issue_lock_store.read_lock_file(locked["lock_file_path"])
self.assertGreater(
issue_lock_store.lock_generation(after),
issue_lock_store.lock_generation(before),
)
self.assertGreaterEqual(
after["work_lease"]["expires_at"], before["work_lease"]["expires_at"]
)
self.assertEqual(after["work_lease"]["heartbeat_count"], 2)
def test_heartbeat_evidence_survives_the_downstream_mutation_gate(self):
"""The #499 F2 lesson, applied.
A sanction that is computed and then discarded downstream is worthless.
After a heartbeat the lock must still satisfy the gate every author
mutation runs through.
"""
locked, session = self._lock_and_session()
self.run_heartbeat(task_session_id=session)
written = issue_lock_store.read_lock_file(locked["lock_file_path"])
verdict = issue_lock_store.verify_lock_for_mutation(
written,
issue_number=ISSUE,
branch_name=BRANCH,
worktree_path=self.worktree,
)
self.assertTrue(verdict["proven"], verdict)
self.assertFalse(verdict["block"])
def _duplicate_gate(self, *, open_prs, branches):
env = self._tool_env()
with patch("mcp_server.get_auth_header", return_value="token x"), patch(
"mcp_server.issue_duplicate_context_fetcher",
side_effect=lambda h, o, r, auth, issue_number: (
list(open_prs),
list(branches),
{"status": "not_claimed"},
),
), patch.dict(os.environ, env, clear=True):
os.environ["GITEA_ISSUE_LOCK_DIR"] = self.lock_dir.name
return mcp_server.gitea_assess_work_issue_duplicate(
issue_number=ISSUE, branch_name=BRANCH, remote="prgs"
)
def test_heartbeat_does_not_change_the_duplicate_gate_verdict(self):
"""The gate must be invariant under heartbeating.
The point is not that the gate passes — with a linked open PR at the
lock phase it correctly blocks (#400), heartbeat or not. The property
that matters is that sliding a lease neither loosens the gate nor
corrupts the lock state it reads: the verdict before and after a
heartbeat must be identical, for both the clear and the blocking shape.
"""
_, session = self._lock_and_session()
linked = [{"number": 4242, "head": {"ref": BRANCH, "sha": self.head_sha}}]
clear_before = self._duplicate_gate(open_prs=[], branches=[])
blocked_before = self._duplicate_gate(open_prs=linked, branches=[BRANCH])
self.assertTrue(self.run_heartbeat(task_session_id=session)["success"])
clear_after = self._duplicate_gate(open_prs=[], branches=[])
blocked_after = self._duplicate_gate(open_prs=linked, branches=[BRANCH])
self.assertEqual(clear_before["outcome"], clear_after["outcome"])
self.assertFalse(clear_after["block"])
self.assertEqual(blocked_before["outcome"], blocked_after["outcome"])
self.assertTrue(blocked_after["block"])
self.assertEqual(blocked_after["linked_open_pr"], 4242)
def test_foreign_session_id_is_refused_through_the_tool(self):
self._lock_and_session()
beat = self.run_heartbeat(task_session_id="author_issue_work-ffffffffffffffff")
self.assertFalse(beat["success"])
self.assertIn("task_session_id does not match", " ".join(beat["reasons"]))
def test_stale_generation_is_refused_through_the_tool(self):
locked, session = self._lock_and_session()
current = issue_lock_store.lock_generation(
issue_lock_store.read_lock_file(locked["lock_file_path"])
)
beat = self.run_heartbeat(
task_session_id=session, expected_generation=current + 5
)
self.assertFalse(beat["success"])
self.assertIn("generation changed", beat["reasons"][0])
def test_foreign_claimant_is_refused_through_the_tool(self):
_, session = self._lock_and_session()
beat = self.run_heartbeat(task_session_id=session, identity="someone-else")
self.assertFalse(beat["success"])
def test_heartbeat_cannot_acquire_a_missing_lock(self):
beat = self.run_heartbeat(task_session_id="author_issue_work-000000000000")
self.assertFalse(beat["success"])
self.assertIn("no durable lock", beat["reasons"][0])
def test_alive_pid_alone_does_not_keep_a_lease_live_through_the_tool(self):
"""PID-only refusal, end to end.
The recorded pid is this live process. The lock is aged past its grace
with no heartbeat, so the tool must refuse to slide it and the durable
record must classify as a missed heartbeat rather than as live.
"""
locked, session = self._lock_and_session()
record = issue_lock_store.read_lock_file(locked["lock_file_path"])
record["work_lease"]["last_heartbeat_at"] = _ts(
datetime.now(timezone.utc) - timedelta(minutes=30)
)
record["work_lease"]["expires_at"] = _ts(
datetime.now(timezone.utc) + timedelta(hours=2)
)
issue_lock_store.save_lock_file(locked["lock_file_path"], record)
self.assertTrue(issue_lock_store.is_process_alive(record["session_pid"]))
fresh = issue_lock_store.assess_lock_freshness(record)
self.assertEqual(
fresh["status"], issue_lock_store.STATUS_STALE_MISSED_HEARTBEAT
)
self.assertTrue(fresh["pid_alive"])
beat = self.run_heartbeat(task_session_id=session)
self.assertFalse(beat["success"])
self.assertIn("reclaimed", " ".join(beat["reasons"]))
class TestLegacyLocksThroughTheTool(_HeartbeatMcpBase):
"""AC-N8 end to end: protected on deployment, and rebindable."""
def test_legacy_lock_stays_protected_after_deployment(self):
record = self.write_legacy_lock(hours_old=3.0, ttl_hours=4.0)
fresh = issue_lock_store.assess_lock_freshness(record)
self.assertEqual(fresh["status"], issue_lock_store.STATUS_LIVE)
self.assertTrue(fresh["legacy_lease"])
self.assertTrue(fresh["legacy_expiry_preserved"])
# It had never heartbeated, so under the new grace alone it would be
# long gone; the preserved absolute expiry is what protects it.
self.assertEqual(
record["work_lease"]["created_at"],
record["work_lease"]["last_heartbeat_at"],
)
def test_tool_rebinds_a_legacy_lock_and_mints_a_first_heartbeat(self):
self.write_legacy_lock(hours_old=3.0, ttl_hours=4.0)
result = self.run_heartbeat(task_session_id=None)
self.assertTrue(result["success"], result)
self.assertEqual(result["operation"], "legacy_rebind")
self.assertTrue(result["task_session_id"])
written = issue_lock_store.read_lock_file(self._lock_path())
lease = written["work_lease"]
self.assertEqual(
lease["lifecycle_version"], lease_policy.LIFECYCLE_HEARTBEAT_V1
)
self.assertEqual(lease["heartbeat_count"], 1)
self.assertNotEqual(
lease["created_at"],
written["legacy_rebind"]["legacy_origin"]["created_at"],
)
self.assertFalse(issue_lock_store.is_legacy_lease(written))
def test_rebound_lock_then_heartbeats_through_the_tool(self):
self.write_legacy_lock(hours_old=3.0, ttl_hours=4.0)
rebound = self.run_heartbeat(task_session_id=None)
beat = self.run_heartbeat(task_session_id=rebound["task_session_id"])
self.assertTrue(beat["success"], beat)
self.assertEqual(beat["operation"], "heartbeat")
self.assertEqual(beat["heartbeat_count"], 2)
def test_rebind_refuses_a_foreign_owner_through_the_tool(self):
self.write_legacy_lock(hours_old=3.0, ttl_hours=4.0)
result = self.run_heartbeat(task_session_id=None, identity="someone-else")
self.assertFalse(result["success"])
self.assertEqual(result["operation"], "legacy_rebind")
if __name__ == "__main__":
unittest.main()
+594
View File
@@ -0,0 +1,594 @@
"""Central lease policy and load-bearing heartbeat freshness (#790 Slice A).
Before this slice, ``issue_lock_store.assess_lock_freshness`` parsed
``last_heartbeat_at`` and then never consulted it: liveness was decided by an
absolute four-hour ``expires_at`` and by PID liveness. Because the recorded PID
is the long-lived MCP daemon rather than the authoring task, an abandoned claim
stayed "live" for the full four hours, and a claim whose work had already landed
blocked reconciliation for just as long (Issue #787 / PR #789, and again Issue
#760 / PR #791).
These tests pin the corrected semantics, including the two asymmetries that are
easy to lose in a refactor:
* an **alive** PID must never make anything live (AC-N2), while
* a **dead** PID must still mark a lease stale, because #753 dead-session
recovery keys on exactly that classification.
Durable-state helpers here write real lock files through the real flock path;
they are not mocks of the store.
"""
from __future__ import annotations
import os
import sys
import tempfile
import unittest
from datetime import datetime, timedelta, timezone
from unittest.mock import patch
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
import issue_lock_store # noqa: E402
import lease_policy # noqa: E402
import pr_work_lease # noqa: E402
import reviewer_pr_lease # noqa: E402
ISSUE = 9790
BRANCH = f"fix/issue-{ISSUE}-heartbeat"
IDENTITY = "example-user"
PROFILE = "test-author-prgs"
ORG = "Example-Org"
REPO = "Example-Repo"
REMOTE = "prgs"
DEAD_PID = 2**22 # far above any live pid on a test host
def _ts(moment: datetime) -> str:
return (
moment.astimezone(timezone.utc)
.replace(microsecond=0)
.isoformat()
.replace("+00:00", "Z")
)
class _LockFixture(unittest.TestCase):
def setUp(self):
self.lock_dir = tempfile.TemporaryDirectory()
self.addCleanup(self.lock_dir.cleanup)
self.now = datetime.now(timezone.utc)
self.worktree = os.path.realpath(tempfile.mkdtemp(prefix="issue790-"))
self.addCleanup(patch.stopall)
def _path(self):
return issue_lock_store.lock_file_path(
remote=REMOTE,
org=ORG,
repo=REPO,
issue_number=ISSUE,
lock_dir=self.lock_dir.name,
)
def write_lock(
self,
*,
lifecycle: str | None = lease_policy.LIFECYCLE_HEARTBEAT_V1,
created_delta: timedelta = timedelta(minutes=1),
heartbeat_delta: timedelta = timedelta(minutes=1),
expires_delta: timedelta = timedelta(minutes=9),
pid: int | None = None,
task_session_id: str | None = "author_issue_work-aaaabbbbccccdddd",
generation: int = 1,
identity: str = IDENTITY,
profile: str = PROFILE,
branch: str = BRANCH,
worktree: str | None = None,
) -> dict:
"""Write a real durable lock and return the record.
Deltas are relative to ``self.now``; ``expires_delta`` is added, the
others subtracted, so "in the past" reads naturally at each call site.
"""
lease: dict = {
"operation_type": issue_lock_store.AUTHOR_ISSUE_WORK_LEASE,
"issue_number": ISSUE,
"pr_number": None,
"branch": branch,
"worktree_path": worktree or self.worktree,
"claimant": {"username": identity, "profile": profile},
"created_at": _ts(self.now - created_delta),
"last_heartbeat_at": _ts(self.now - heartbeat_delta),
"expires_at": _ts(self.now + expires_delta),
}
if lifecycle is not None:
lease["lifecycle_version"] = lifecycle
if task_session_id is not None:
lease["task_session_id"] = task_session_id
pid_value = os.getpid() if pid is None else pid
record = {
"issue_number": ISSUE,
"branch_name": branch,
"remote": REMOTE,
"org": ORG,
"repo": REPO,
"worktree_path": worktree or self.worktree,
"session_pid": pid_value,
"pid": pid_value,
"lock_generation": generation,
"work_lease": lease,
}
path = self._path()
record["lock_file_path"] = path
issue_lock_store.save_lock_file(path, record)
return record
class TestPolicyIsTheSingleSource(unittest.TestCase):
"""AC-N7: one authoritative configuration source for every duration."""
def test_author_policy_carries_the_agreed_values(self):
policy = lease_policy.policy_for(lease_policy.TASK_CLASS_AUTHOR_ISSUE_WORK)
self.assertEqual(policy.initial_ttl_minutes, 10.0)
self.assertEqual(policy.heartbeat_cadence_minutes, 2.0)
self.assertEqual(policy.stale_warning_minutes, 5.0)
self.assertEqual(policy.missed_heartbeat_grace_minutes, 10.0)
self.assertEqual(policy.absolute_cap_hours, 8.0)
self.assertEqual(policy.recovery_grace_minutes, 10.0)
self.assertEqual(policy.terminal_race_drain_minutes, 2.0)
self.assertTrue(policy.terminal_retirement_eligible)
self.assertTrue(policy.heartbeat_lifecycle_active)
def test_the_four_hour_author_ttl_literal_is_gone(self):
"""The duplicated literal AC-N7 exists to remove."""
self.assertFalse(hasattr(issue_lock_store, "WORK_LEASE_TTL_HOURS"))
import gitea_mcp_server
self.assertFalse(hasattr(gitea_mcp_server, "WORK_LEASE_TTL_HOURS"))
def test_declared_reviewer_values_match_the_module_still_using_them(self):
"""Slice A declares reviewer/merger numbers without rewiring them.
Recording a value in two places is only safe if drift is detectable, so
this asserts the declaration still equals the constants #747 owns. When
Slice C migrates those call sites, this test becomes the proof the
migration changed nothing.
"""
policy = lease_policy.policy_for(lease_policy.TASK_CLASS_REVIEWER_PR)
self.assertEqual(
policy.initial_ttl_minutes, float(reviewer_pr_lease.LEASE_TTL_MINUTES)
)
self.assertEqual(
policy.stale_warning_minutes,
float(reviewer_pr_lease.STALE_WARNING_MINUTES),
)
self.assertFalse(policy.heartbeat_lifecycle_active)
def test_declared_conflict_fix_value_matches_its_module(self):
policy = lease_policy.policy_for(lease_policy.TASK_CLASS_CONFLICT_FIX)
self.assertEqual(
policy.initial_ttl_minutes,
float(pr_work_lease.DEFAULT_CONFLICT_FIX_TTL_MINUTES),
)
self.assertFalse(policy.heartbeat_lifecycle_active)
def test_environment_override_applies(self):
var = lease_policy.env_var_name(
lease_policy.TASK_CLASS_AUTHOR_ISSUE_WORK, "initial_ttl_minutes"
)
with patch.dict(os.environ, {var: "7"}):
self.assertEqual(
lease_policy.policy_for(
lease_policy.TASK_CLASS_AUTHOR_ISSUE_WORK
).initial_ttl_minutes,
7.0,
)
def test_unusable_override_falls_back_instead_of_minting_a_zero_lease(self):
"""A typo must not make every claim instantly reclaimable."""
var = lease_policy.env_var_name(
lease_policy.TASK_CLASS_AUTHOR_ISSUE_WORK, "initial_ttl_minutes"
)
for bad in ("0", "-5", "not-a-number", " "):
with self.subTest(value=bad), patch.dict(os.environ, {var: bad}):
self.assertEqual(
lease_policy.policy_for(
lease_policy.TASK_CLASS_AUTHOR_ISSUE_WORK
).initial_ttl_minutes,
10.0,
)
def test_unknown_task_class_does_not_raise(self):
policy = lease_policy.policy_for("something-new")
self.assertEqual(policy.task_class, lease_policy.TASK_CLASS_AUTHOR_ISSUE_WORK)
class TestLifecycleDiscrimination(_LockFixture):
"""AC-N8: the marker, never a timestamp, decides legacy vs heartbeat."""
def test_missing_marker_reads_as_legacy(self):
record = self.write_lock(lifecycle=None)
self.assertTrue(issue_lock_store.is_legacy_lease(record))
self.assertEqual(
issue_lock_store.lease_lifecycle_version(record),
lease_policy.LIFECYCLE_LEGACY,
)
def test_marker_present_reads_as_heartbeat_lifecycle(self):
record = self.write_lock()
self.assertFalse(issue_lock_store.is_legacy_lease(record))
def test_equal_created_and_heartbeat_never_implies_a_fresh_heartbeat(self):
"""The exact inversion AC-N8 forbids.
A legacy lock has ``last_heartbeat_at == created_at`` forever because
nothing ever advanced it. Reading that equality as "recently
heartbeated" would classify every never-heartbeated lock as fresh.
"""
legacy = self.write_lock(
lifecycle=None,
created_delta=timedelta(hours=3),
heartbeat_delta=timedelta(hours=3),
)
lease = legacy["work_lease"]
self.assertEqual(lease["created_at"], lease["last_heartbeat_at"])
self.assertTrue(issue_lock_store.is_legacy_lease(legacy))
# A brand-new heartbeat lease has them equal too, so the equality
# carries no information in either direction.
fresh = self.write_lock(
created_delta=timedelta(seconds=0), heartbeat_delta=timedelta(seconds=0)
)
self.assertEqual(
fresh["work_lease"]["created_at"],
fresh["work_lease"]["last_heartbeat_at"],
)
self.assertFalse(issue_lock_store.is_legacy_lease(fresh))
def test_minted_session_id_contains_no_pid(self):
"""AC-N1: the ownership key must not be derived from the daemon pid."""
minted = issue_lock_store.mint_task_session_id()
self.assertNotIn(str(os.getpid()), minted)
self.assertNotEqual(minted, issue_lock_store.mint_task_session_id())
class TestFreshnessIsHeartbeatDriven(_LockFixture):
"""AC-N2 and the new bands."""
def test_fresh_heartbeat_is_live(self):
record = self.write_lock(heartbeat_delta=timedelta(minutes=1))
fresh = issue_lock_store.assess_lock_freshness(record, now=self.now)
self.assertEqual(fresh["status"], issue_lock_store.STATUS_LIVE)
self.assertTrue(fresh["live"])
self.assertFalse(fresh["heartbeat_warning"])
def test_heartbeat_past_warning_is_still_live_but_flagged(self):
record = self.write_lock(heartbeat_delta=timedelta(minutes=6))
fresh = issue_lock_store.assess_lock_freshness(record, now=self.now)
self.assertEqual(fresh["status"], issue_lock_store.STATUS_LIVE)
self.assertTrue(fresh["heartbeat_warning"])
def test_missed_heartbeat_past_grace_is_classified_explicitly(self):
record = self.write_lock(
heartbeat_delta=timedelta(minutes=11),
expires_delta=timedelta(minutes=30),
)
fresh = issue_lock_store.assess_lock_freshness(record, now=self.now)
self.assertEqual(
fresh["status"], issue_lock_store.STATUS_STALE_MISSED_HEARTBEAT
)
self.assertFalse(fresh["live"])
self.assertTrue(fresh["stale"])
def test_alive_pid_never_establishes_freshness(self):
"""The defect in one assertion.
The recorded PID is this very process, so it is unambiguously alive —
and the lease is still not live, because the task stopped heartbeating.
"""
record = self.write_lock(
pid=os.getpid(),
heartbeat_delta=timedelta(hours=4),
expires_delta=timedelta(hours=4),
)
fresh = issue_lock_store.assess_lock_freshness(record, now=self.now)
self.assertTrue(fresh["pid_alive"])
self.assertFalse(fresh["live"])
self.assertEqual(
fresh["status"], issue_lock_store.STATUS_STALE_MISSED_HEARTBEAT
)
def test_dead_pid_still_marks_stale_for_issue_753(self):
"""The opposite asymmetry: dead-PID corroboration is preserved."""
record = self.write_lock(pid=DEAD_PID, heartbeat_delta=timedelta(minutes=1))
fresh = issue_lock_store.assess_lock_freshness(record, now=self.now)
self.assertEqual(fresh["status"], issue_lock_store.STATUS_STALE)
self.assertFalse(fresh["live"])
self.assertIn("not alive", fresh["reason"])
def test_absolute_cap_requires_readoption(self):
record = self.write_lock(
created_delta=timedelta(hours=9), heartbeat_delta=timedelta(minutes=1)
)
fresh = issue_lock_store.assess_lock_freshness(record, now=self.now)
self.assertEqual(fresh["status"], issue_lock_store.STATUS_STALE_ABSOLUTE_CAP)
self.assertIn("re-adoption", fresh["reason"])
def test_heartbeat_lifecycle_without_a_heartbeat_fails_closed(self):
record = self.write_lock()
del record["work_lease"]["last_heartbeat_at"]
issue_lock_store.save_lock_file(self._path(), record)
fresh = issue_lock_store.assess_lock_freshness(record, now=self.now)
self.assertEqual(
fresh["status"], issue_lock_store.STATUS_STALE_MISSED_HEARTBEAT
)
self.assertIn("fail closed", fresh["reason"])
def test_absent_lock(self):
fresh = issue_lock_store.assess_lock_freshness(None)
self.assertEqual(fresh["status"], issue_lock_store.STATUS_ABSENT)
self.assertFalse(fresh["stale"])
class TestLegacyLocksStayProtected(_LockFixture):
"""AC-N8: deployment must not retroactively shorten an existing claim."""
def test_legacy_lock_with_a_stale_heartbeat_remains_live(self):
"""The deployment-safety case.
A four-hour legacy lease minted three hours ago has not heartbeated
once. Under the new grace it would be long gone; under its preserved
absolute expiry it is still live, and must stay that way.
"""
record = self.write_lock(
lifecycle=None,
created_delta=timedelta(hours=3),
heartbeat_delta=timedelta(hours=3),
expires_delta=timedelta(hours=1),
)
fresh = issue_lock_store.assess_lock_freshness(record, now=self.now)
self.assertEqual(fresh["status"], issue_lock_store.STATUS_LIVE)
self.assertTrue(fresh["live"])
self.assertTrue(fresh["legacy_lease"])
self.assertTrue(fresh["legacy_expiry_preserved"])
def test_legacy_lock_past_its_absolute_expiry_is_expired_as_before(self):
record = self.write_lock(
lifecycle=None,
created_delta=timedelta(hours=5),
heartbeat_delta=timedelta(hours=5),
expires_delta=timedelta(hours=-1),
)
fresh = issue_lock_store.assess_lock_freshness(record, now=self.now)
self.assertEqual(fresh["status"], issue_lock_store.STATUS_EXPIRED)
def test_legacy_lock_is_never_reclaimed_by_the_heartbeat_band(self):
record = self.write_lock(
lifecycle=None,
created_delta=timedelta(hours=3),
heartbeat_delta=timedelta(hours=3),
expires_delta=timedelta(hours=1),
)
reclaim = issue_lock_store.assess_expired_lock_reclaim(record, now=self.now)
self.assertFalse(reclaim["reclaim_allowed"])
class TestReclaimAfterMissedHeartbeat(_LockFixture):
def test_missed_heartbeat_makes_ownership_reclaimable(self):
record = self.write_lock(
pid=os.getpid(),
heartbeat_delta=timedelta(minutes=15),
expires_delta=timedelta(hours=3),
)
reclaim = issue_lock_store.assess_expired_lock_reclaim(record, now=self.now)
self.assertTrue(reclaim["reclaim_allowed"])
self.assertIn("stale_missed_heartbeat", reclaim["reasons"][0])
def test_live_lease_is_never_reclaimable(self):
record = self.write_lock(heartbeat_delta=timedelta(minutes=1))
reclaim = issue_lock_store.assess_expired_lock_reclaim(record, now=self.now)
self.assertFalse(reclaim["reclaim_allowed"])
def test_dead_pid_reclaim_path_is_unchanged(self):
"""#753 must keep working through its original conditions."""
record = self.write_lock(pid=DEAD_PID, heartbeat_delta=timedelta(minutes=1))
reclaim = issue_lock_store.assess_expired_lock_reclaim(record, now=self.now)
self.assertTrue(reclaim["reclaim_allowed"])
self.assertTrue(reclaim["owner_pid_dead"])
class TestHeartbeatWriter(_LockFixture):
"""A4: flock + CAS + exact verification, and no revival path."""
def _heartbeat(self, **kwargs):
params = {
"remote": REMOTE,
"org": ORG,
"repo": REPO,
"issue_number": ISSUE,
"branch_name": BRANCH,
"worktree_path": self.worktree,
"identity": IDENTITY,
"profile": PROFILE,
"task_session_id": "author_issue_work-aaaabbbbccccdddd",
"lock_dir": self.lock_dir.name,
"now": self.now,
}
params.update(kwargs)
return issue_lock_store.heartbeat_session_lock(**params)
def test_heartbeat_slides_expiry_and_advances_generation(self):
self.write_lock(heartbeat_delta=timedelta(minutes=4), generation=5)
result = self._heartbeat()
self.assertTrue(result["success"], result)
self.assertEqual(result["prior_generation"], 5)
self.assertEqual(result["lock_generation"], 6)
self.assertEqual(result["heartbeat_count"], 1)
self.assertEqual(result["last_heartbeat_at"], _ts(self.now))
self.assertEqual(result["expires_at"], _ts(self.now + timedelta(minutes=10)))
self.assertTrue(result["freshness"]["live"])
def test_heartbeat_is_durable_and_repeatable(self):
self.write_lock(heartbeat_delta=timedelta(minutes=4))
self._heartbeat()
second = self._heartbeat(now=self.now + timedelta(minutes=1))
self.assertTrue(second["success"], second)
self.assertEqual(second["heartbeat_count"], 2)
written = issue_lock_store.read_lock_file(self._path())
self.assertEqual(written["work_lease"]["heartbeat_count"], 2)
def test_stale_generation_is_refused(self):
self.write_lock(generation=5)
result = self._heartbeat(expected_generation=4)
self.assertFalse(result["success"])
self.assertIn("generation changed", result["reasons"][0])
def test_foreign_session_is_refused(self):
self.write_lock()
result = self._heartbeat(task_session_id="author_issue_work-ffffffffffffffff")
self.assertFalse(result["success"])
self.assertIn("task_session_id does not match", " ".join(result["reasons"]))
def test_missing_session_id_is_refused(self):
self.write_lock()
result = self._heartbeat(task_session_id="")
self.assertFalse(result["success"])
def test_foreign_claimant_is_refused(self):
self.write_lock()
for field, value in (
("identity", "someone-else"),
("profile", "other-profile"),
):
with self.subTest(field=field):
result = self._heartbeat(**{field: value})
self.assertFalse(result["success"])
def test_branch_and_worktree_mismatch_are_refused(self):
self.write_lock()
wrong_branch = self._heartbeat(branch_name=f"fix/issue-{ISSUE}-other")
self.assertFalse(wrong_branch["success"])
wrong_worktree = self._heartbeat(worktree_path="/tmp/not-the-worktree")
self.assertFalse(wrong_worktree["success"])
def test_lapsed_lease_cannot_be_heartbeated_back_to_life(self):
"""No revival path (A4).
A session that stopped proving liveness must reclaim under a fresh
generation, not restore ownership retroactively.
"""
self.write_lock(
heartbeat_delta=timedelta(minutes=30), expires_delta=timedelta(hours=1)
)
result = self._heartbeat()
self.assertFalse(result["success"])
self.assertIn("reclaimed", " ".join(result["reasons"]))
def test_absent_lock_cannot_be_created_by_heartbeat(self):
result = self._heartbeat()
self.assertFalse(result["success"])
self.assertIn("no durable lock", result["reasons"][0])
def test_legacy_lock_is_refused_until_rebound(self):
self.write_lock(lifecycle=None)
result = self._heartbeat()
self.assertFalse(result["success"])
self.assertTrue(result["legacy_lease"])
self.assertIn("rebound", " ".join(result["reasons"]))
class TestLegacyRebind(_LockFixture):
"""AC-N8 exit route: canonical exact-owner rebinding."""
def _rebind(self, **kwargs):
params = {
"remote": REMOTE,
"org": ORG,
"repo": REPO,
"issue_number": ISSUE,
"branch_name": BRANCH,
"worktree_path": self.worktree,
"identity": IDENTITY,
"profile": PROFILE,
"lock_dir": self.lock_dir.name,
"now": self.now,
}
params.update(kwargs)
return issue_lock_store.rebind_legacy_lock(**params)
def test_rebind_mints_a_session_and_a_genuine_first_heartbeat(self):
self.write_lock(
lifecycle=None,
created_delta=timedelta(hours=3),
heartbeat_delta=timedelta(hours=3),
expires_delta=timedelta(hours=1),
generation=2,
)
result = self._rebind()
self.assertTrue(result["success"], result)
self.assertTrue(result["task_session_id"])
self.assertEqual(result["lock_generation"], 3)
written = issue_lock_store.read_lock_file(self._path())
lease = written["work_lease"]
self.assertEqual(
lease["lifecycle_version"], lease_policy.LIFECYCLE_HEARTBEAT_V1
)
self.assertEqual(lease["last_heartbeat_at"], _ts(self.now))
self.assertEqual(lease["expires_at"], _ts(self.now + timedelta(minutes=10)))
self.assertFalse(issue_lock_store.is_legacy_lease(written))
# The original claim is preserved for audit rather than overwritten.
origin = written["legacy_rebind"]["legacy_origin"]
self.assertTrue(origin["created_at"])
self.assertEqual(origin["lifecycle"], lease_policy.LIFECYCLE_LEGACY)
def test_rebound_lock_can_then_heartbeat(self):
self.write_lock(
lifecycle=None,
created_delta=timedelta(hours=3),
heartbeat_delta=timedelta(hours=3),
expires_delta=timedelta(hours=1),
)
rebound = self._rebind()
beat = issue_lock_store.heartbeat_session_lock(
remote=REMOTE,
org=ORG,
repo=REPO,
issue_number=ISSUE,
branch_name=BRANCH,
worktree_path=self.worktree,
identity=IDENTITY,
profile=PROFILE,
task_session_id=rebound["task_session_id"],
lock_dir=self.lock_dir.name,
now=self.now + timedelta(minutes=1),
)
self.assertTrue(beat["success"], beat)
def test_rebind_refuses_a_foreign_owner(self):
self.write_lock(lifecycle=None, expires_delta=timedelta(hours=1))
result = self._rebind(identity="someone-else")
self.assertFalse(result["success"])
def test_rebind_refuses_a_lock_already_on_the_lifecycle(self):
self.write_lock()
result = self._rebind()
self.assertFalse(result["success"])
self.assertFalse(result["legacy_lease"])
def test_rebind_is_not_a_recovery_path_for_a_lapsed_legacy_lease(self):
"""An expired legacy lease belongs to #760 renewal or #601 reclaim."""
self.write_lock(
lifecycle=None,
created_delta=timedelta(hours=5),
heartbeat_delta=timedelta(hours=5),
expires_delta=timedelta(hours=-1),
)
result = self._rebind()
self.assertFalse(result["success"])
self.assertIn("not a recovery path", " ".join(result["reasons"]))
if __name__ == "__main__":
unittest.main()
-682
View File
@@ -1,682 +0,0 @@
"""Cross-role allocation handoff consumable by independent workers (#843).
Regression coverage for the controller→required-role consume path:
* controller allocates author work; independent author adopts successfully
* author adoption succeeds after allocating controller process exits
* author adoption without sharing controller session identity
* wrong-role adoption rejected
* concurrent/second adoption rejected without state corruption
* terminal allocation adoption rejected
* successful adoption produces authoritative ownership evidence
* genuine abandoned-lease recovery remains valid
* process_work_queue / allocate results include consume identifiers
* same-role allocation behavior remains compatible
"""
from __future__ import annotations
import os
import tempfile
import unittest
from datetime import timedelta
from unittest.mock import patch
from allocator_service import (
ALLOCATION_MODE_CROSS_ROLE,
ALLOCATION_MODE_ROLE_SCOPED,
OUTCOME_ASSIGNED,
ROLE_AUTHOR,
ROLE_CONTROLLER,
ROLE_REVIEWER,
WorkCandidate,
allocate_next_work,
)
from control_plane_db import ControlPlaneDB, ForeignLeaseError, _ts, _utc_now
import lease_lifecycle as ll
class CrossRoleHandoffTest(unittest.TestCase):
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.db_path = os.path.join(self._tmp.name, "cp.sqlite3")
self.db = ControlPlaneDB(self.db_path)
self.db.upsert_session(
session_id="ctrl-session",
role="controller",
profile="prgs-controller",
pid=99999999, # dead-looking pid
)
self.db.upsert_session(
session_id="author-worker",
role="author",
profile="prgs-author",
pid=os.getpid(),
)
self.db.upsert_session(
session_id="author-worker-2",
role="author",
profile="prgs-author",
pid=os.getpid(),
)
self.db.upsert_session(
session_id="reviewer-worker",
role="reviewer",
profile="prgs-reviewer",
pid=os.getpid(),
)
self.wt = self._tmp.name
def tearDown(self) -> None:
self._tmp.cleanup()
def _ready_issue(self, number: int = 843, title: str = "handoff target") -> WorkCandidate:
return WorkCandidate(
kind="issue",
number=number,
labels=("status:ready", "type:bug"),
title=title,
priority=20,
)
def _controller_allocate(self, number: int = 843, **kwargs):
defaults = dict(
db=self.db,
session_id="ctrl-session",
role=ROLE_CONTROLLER,
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
candidates=[self._ready_issue(number)],
apply=True,
profile_name="prgs-controller",
username="controller-user",
allocation_mode=ALLOCATION_MODE_CROSS_ROLE,
)
defaults.update(kwargs)
return allocate_next_work(**defaults)
def test_controller_allocates_author_independent_author_adopts(self) -> None:
res = self._controller_allocate()
self.assertEqual(res["outcome"], OUTCOME_ASSIGNED)
self.assertEqual(res["required_role"], ROLE_AUTHOR)
self.assertIn("consume_allocation", res)
consume = res["consume_allocation"]
self.assertEqual(consume["tool"], "gitea_adopt_workflow_lease")
self.assertEqual(consume["required_role"], ROLE_AUTHOR)
self.assertFalse(consume["controller_session_required"])
lid = res["assignment"]["lease_id"]
self.assertEqual(consume["lease_id"], lid)
self.assertIn(lid, res["next_valid_command"])
adopted = ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="author-worker",
role=ROLE_AUTHOR,
worktree_path=self.wt,
)
self.assertTrue(adopted["success"])
self.assertEqual(adopted["outcome"], "adopted_cross_role_handoff")
self.assertEqual(adopted["adopted_by_session_id"], "author-worker")
self.assertEqual(adopted["adopted_from_session_id"], "ctrl-session")
raw = adopted["read_after_write"]
self.assertEqual(raw["session_id"], "author-worker")
self.assertEqual(raw["adopted_by_session_id"], "author-worker")
self.assertEqual(raw["status"], "active")
self.assertEqual(raw["phase"], "adopted")
# Authoritative re-read
state = self.db.get_lease_workflow_state(lid)
self.assertEqual(state["lease"]["session_id"], "author-worker")
self.assertEqual(state["lease"]["adopted_by_session_id"], "author-worker")
self.assertEqual(state["assignment"]["session_id"], "author-worker")
self.assertEqual(state["provenance"]["handoff_status"], "adopted")
def test_author_adoption_after_controller_process_exits(self) -> None:
res = self._controller_allocate(number=900)
lid = res["assignment"]["lease_id"]
# Force owner_pid dead + freshness stale_dead_process
import sqlite3
conn = sqlite3.connect(self.db_path)
try:
conn.execute(
"UPDATE leases SET owner_pid = 99999999 WHERE lease_id = ?",
(lid,),
)
conn.commit()
finally:
conn.close()
state = self.db.get_lease_workflow_state(lid)
fr = ll.classify_lease_freshness(
state["lease"], pid_checker=lambda _p: False
)
self.assertEqual(fr["freshness"], "stale_dead_process")
adopted = ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="author-worker",
role=ROLE_AUTHOR,
worktree_path=self.wt,
)
self.assertEqual(adopted["outcome"], "adopted_cross_role_handoff")
self.assertEqual(adopted["adopted_by_session_id"], "author-worker")
# No abandon required
state2 = self.db.get_lease_workflow_state(lid)
self.assertEqual(state2["lease"]["status"], "active")
self.assertNotEqual(state2["lease"]["status"], "abandoned")
def test_adoption_without_sharing_controller_session_identity(self) -> None:
res = self._controller_allocate(number=901)
lid = res["assignment"]["lease_id"]
adopted = ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="author-worker",
role=ROLE_AUTHOR,
worktree_path=self.wt,
)
self.assertNotEqual(adopted["adopted_by_session_id"], "ctrl-session")
self.assertFalse(adopted["same_owner"])
self.assertEqual(adopted["adopted_from_session_id"], "ctrl-session")
def test_wrong_role_adoption_rejected(self) -> None:
res = self._controller_allocate(number=902)
lid = res["assignment"]["lease_id"]
with self.assertRaises(ll.LeaseLifecycleError) as ctx:
ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="reviewer-worker",
role=ROLE_REVIEWER,
)
self.assertIn("wrong role", str(ctx.exception).lower())
# State unchanged
state = self.db.get_lease_workflow_state(lid)
self.assertEqual(state["lease"]["session_id"], "ctrl-session")
self.assertIsNone(state["lease"].get("adopted_by_session_id") or None)
self.assertEqual(state["provenance"]["handoff_status"], "pending")
def test_second_adoption_rejected_without_corruption(self) -> None:
res = self._controller_allocate(number=903)
lid = res["assignment"]["lease_id"]
first = ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="author-worker",
role=ROLE_AUTHOR,
worktree_path=self.wt,
)
self.assertEqual(first["outcome"], "adopted_cross_role_handoff")
with self.assertRaises(ll.LeaseLifecycleError):
ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="author-worker-2",
role=ROLE_AUTHOR,
worktree_path=self.wt,
)
state = self.db.get_lease_workflow_state(lid)
self.assertEqual(state["lease"]["session_id"], "author-worker")
self.assertEqual(state["lease"]["adopted_by_session_id"], "author-worker")
self.assertEqual(state["assignment"]["session_id"], "author-worker")
self.assertEqual(state["lease"]["status"], "active")
def test_terminal_allocation_adoption_rejected(self) -> None:
res = self._controller_allocate(number=904)
lid = res["assignment"]["lease_id"]
# Abandon as terminal
proof = ll.AbandonProof(
dead_process=True,
missing_worktree=True,
no_open_pr=True,
no_live_mutation_risk=True,
owner_pid=99999999,
worktree_path="/nonexistent/for-843",
)
# Attach dead pid / missing wt for abandon eligibility
import sqlite3
conn = sqlite3.connect(self.db_path)
try:
conn.execute(
"UPDATE leases SET owner_pid = 99999999, worktree_path = ? WHERE lease_id = ?",
("/nonexistent/for-843", lid),
)
conn.commit()
finally:
conn.close()
abandoned = ll.abandon_lease(
self.db,
lease_id=lid,
requester_session_id="author-worker",
proof=proof,
)
self.assertEqual(abandoned["outcome"], "abandoned")
with self.assertRaises(ll.LeaseLifecycleError) as ctx:
ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="author-worker",
role=ROLE_AUTHOR,
)
self.assertIn("abandoned", str(ctx.exception).lower())
def test_successful_adoption_read_after_write_ownership(self) -> None:
res = self._controller_allocate(number=905)
lid = res["assignment"]["lease_id"]
adopted = ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="author-worker",
role=ROLE_AUTHOR,
worktree_path=self.wt,
)
raw = adopted["read_after_write"]
self.assertEqual(raw["lease_id"], lid)
self.assertEqual(raw["session_id"], "author-worker")
self.assertEqual(raw["adopted_by_session_id"], "author-worker")
self.assertEqual(raw["adopted_from_session_id"], "ctrl-session")
# Re-fetch proves durable write
state = self.db.get_lease_workflow_state(lid)
self.assertEqual(state["lease"]["session_id"], raw["session_id"])
self.assertEqual(
state["lease"]["adopted_by_session_id"], raw["adopted_by_session_id"]
)
def test_genuine_abandoned_recovery_still_valid(self) -> None:
"""Same-role author lease abandoned remains reclaimable via abandon path."""
same = allocate_next_work(
self.db,
session_id="author-worker",
role=ROLE_AUTHOR,
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
candidates=[self._ready_issue(906, "same-role")],
apply=True,
profile_name="prgs-author",
username="author-user",
allocation_mode=ALLOCATION_MODE_ROLE_SCOPED,
)
self.assertEqual(same["outcome"], OUTCOME_ASSIGNED)
lid = same["assignment"]["lease_id"]
import sqlite3
conn = sqlite3.connect(self.db_path)
try:
conn.execute(
"UPDATE leases SET owner_pid = 99999999, worktree_path = ? WHERE lease_id = ?",
("/nonexistent/same-role", lid),
)
conn.commit()
finally:
conn.close()
proof = ll.AbandonProof(
dead_process=True,
missing_worktree=True,
no_open_pr=True,
no_live_mutation_risk=True,
owner_pid=99999999,
worktree_path="/nonexistent/same-role",
)
abandoned = ll.abandon_lease(
self.db,
lease_id=lid,
requester_session_id="author-worker-2",
proof=proof,
)
self.assertEqual(abandoned["outcome"], "abandoned")
# Foreign author cannot handoff-consume an abandoned non-handoff lease
with self.assertRaises(ll.LeaseLifecycleError):
ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="author-worker-2",
role=ROLE_AUTHOR,
)
# Reclaim path still works for expired/abandoned after force-expire
reclaimed = ll.reclaim_expired_lease(
self.db,
lease_id=lid,
session_id="author-worker-2",
role=ROLE_AUTHOR,
worktree_path=self.wt,
)
self.assertEqual(reclaimed["outcome"], "reclaimed")
self.assertEqual(reclaimed["assignment"]["session_id"], "author-worker-2")
def test_allocate_payload_includes_consume_identifiers(self) -> None:
res = self._controller_allocate(number=907)
self.assertIn("consume_allocation", res)
c = res["consume_allocation"]
for key in (
"tool",
"lease_id",
"assignment_id",
"required_role",
"required_profile",
"required_namespace",
"instructions",
"handoff_status",
):
self.assertIn(key, c)
self.assertEqual(c["required_namespace"], "gitea-author")
self.assertEqual(c["required_profile"], "prgs-author")
self.assertIn("gitea_adopt_workflow_lease", c["instructions"])
self.assertTrue(res["lease_proof"]["cross_role_handoff"])
self.assertEqual(res["lease_proof"]["handoff_status"], "pending")
def test_same_role_allocation_remains_compatible(self) -> None:
res = allocate_next_work(
self.db,
session_id="author-worker",
role=ROLE_AUTHOR,
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
candidates=[self._ready_issue(908)],
apply=True,
profile_name="prgs-author",
username="author-user",
)
self.assertEqual(res["outcome"], OUTCOME_ASSIGNED)
self.assertNotIn("consume_allocation", res)
lid = res["assignment"]["lease_id"]
state = self.db.get_lease_workflow_state(lid)
# No cross-role handoff provenance
prov = state.get("provenance") or {}
self.assertFalse(prov.get("cross_role_handoff"))
# Owner resume still works
resume = ll.adopt_lease(
self.db,
lease_id=lid,
adopter_session_id="author-worker",
role=ROLE_AUTHOR,
worktree_path=self.wt,
)
self.assertTrue(resume["same_owner"])
self.assertEqual(resume["outcome"], "adopted_owner_resume")
def test_inspect_points_required_role_at_consume(self) -> None:
res = self._controller_allocate(number=909)
lid = res["assignment"]["lease_id"]
decision = ll.inspect_lease(
self.db, lid, caller_session_id="author-worker"
)
self.assertEqual(
decision["safe_next_action"], ll.SAFE_CONSUME_CROSS_ROLE
)
self.assertFalse(decision["block"])
self.assertEqual(decision["required_role"], ROLE_AUTHOR)
def test_db_cas_rejects_concurrent_second_consume(self) -> None:
res = self._controller_allocate(number=910)
lid = res["assignment"]["lease_id"]
# First consume via DB layer directly
first = self.db.adopt_lease(
lease_id=lid,
adopter_session_id="author-worker",
role=ROLE_AUTHOR,
worktree_path=self.wt,
provenance={
"cross_role_handoff": True,
"handoff_status": "adopted",
"required_role": "author",
},
)
self.assertEqual(first["outcome"], "adopted_cross_role_handoff")
# Second CAS must fail
with self.assertRaises(ForeignLeaseError):
self.db.adopt_lease(
lease_id=lid,
adopter_session_id="author-worker-2",
role=ROLE_AUTHOR,
worktree_path=self.wt,
provenance={
"cross_role_handoff": True,
"handoff_status": "pending",
"required_role": "author",
},
)
state = self.db.get_lease_workflow_state(lid)
self.assertEqual(state["lease"]["session_id"], "author-worker")
class MCPBoundaryAdoptRoleBindingTest(unittest.TestCase):
"""#843 F1: MCP-boundary role binding for ``gitea_adopt_workflow_lease``.
The library-level wrong-role test calls ``lease_lifecycle.adopt_lease``
directly. These tests prove the MCP entry point derives the adopter role
authoritatively from the active authenticated profile and rejects any
caller-supplied role that disagrees, so a reviewer/merger profile cannot
consume an author handoff by passing ``role="author"``.
"""
AUTHOR_PROFILE = {
"profile_name": "prgs-author",
"role": "author",
"allowed_operations": [
"gitea.read",
"gitea.pr.create",
"gitea.branch.push",
],
"forbidden_operations": [],
}
REVIEWER_PROFILE = {
"profile_name": "prgs-reviewer",
"role": "reviewer",
"allowed_operations": [
"gitea.read",
"gitea.pr.review",
"gitea.pr.approve",
"gitea.pr.request_changes",
],
"forbidden_operations": ["gitea.pr.create", "gitea.branch.push"],
}
MERGER_PROFILE = {
"profile_name": "prgs-merger",
"role": "merger",
"allowed_operations": ["gitea.read", "gitea.pr.merge"],
"forbidden_operations": ["gitea.pr.create", "gitea.branch.push"],
}
FOREIGN_AUTHOR_PROFILE = {
"profile_name": "dadeschools-author",
"role": "author",
"allowed_operations": [
"gitea.read",
"gitea.pr.create",
"gitea.branch.push",
],
"forbidden_operations": [],
}
def setUp(self) -> None:
self._tmp = tempfile.TemporaryDirectory()
self.db_path = os.path.join(self._tmp.name, "cp.sqlite3")
self.db = ControlPlaneDB(self.db_path)
self.db.upsert_session(
session_id="ctrl-session",
role="controller",
profile="prgs-controller",
pid=99999999,
)
self.db.upsert_session(
session_id="author-worker",
role="author",
profile="prgs-author",
pid=os.getpid(),
)
self.wt = self._tmp.name
def tearDown(self) -> None:
self._tmp.cleanup()
def _ready_issue(self, number: int) -> WorkCandidate:
return WorkCandidate(
kind="issue",
number=number,
labels=("status:ready", "type:bug"),
title="handoff target",
priority=20,
)
def _handoff_lease(self, number: int = 843) -> str:
res = allocate_next_work(
db=self.db,
session_id="ctrl-session",
role=ROLE_CONTROLLER,
remote="prgs",
org="Scaled-Tech-Consulting",
repo="Gitea-Tools",
candidates=[self._ready_issue(number)],
apply=True,
profile_name="prgs-controller",
username="controller-user",
allocation_mode=ALLOCATION_MODE_CROSS_ROLE,
)
self.assertEqual(res["outcome"], OUTCOME_ASSIGNED)
self.assertEqual(res["required_role"], ROLE_AUTHOR)
return res["assignment"]["lease_id"]
def _call_adopt_tool(self, profile: dict, **kwargs):
import gitea_mcp_server as mcp_server
with (
patch.object(mcp_server, "get_profile", return_value=profile),
patch.object(
mcp_server,
"_control_plane_db_or_error",
return_value=(self.db, []),
),
):
return mcp_server.gitea_adopt_workflow_lease(
remote="prgs", **kwargs
)
def _assert_handoff_untouched(self, lease_id: str) -> None:
state = self.db.get_lease_workflow_state(lease_id)
self.assertEqual(state["lease"]["session_id"], "ctrl-session")
self.assertIsNone(state["lease"].get("adopted_by_session_id") or None)
self.assertEqual(state["lease"]["status"], "active")
self.assertEqual(state["provenance"]["handoff_status"], "pending")
def test_reviewer_profile_cannot_consume_author_handoff_via_role_author(
self,
) -> None:
lid = self._handoff_lease(920)
result = self._call_adopt_tool(
self.REVIEWER_PROFILE,
lease_id=lid,
session_id="reviewer-worker",
role="author",
worktree_path=self.wt,
)
self.assertFalse(result["success"])
self.assertEqual(result["outcome"], "blocked")
self.assertEqual(result["profile_role_kind"], "reviewer")
self.assertEqual(result["supplied_role"], "author")
self.assertIn("does not match", result["reasons"][0])
self._assert_handoff_untouched(lid)
def test_merger_profile_cannot_consume_author_handoff_via_role_author(
self,
) -> None:
lid = self._handoff_lease(921)
result = self._call_adopt_tool(
self.MERGER_PROFILE,
lease_id=lid,
session_id="merger-worker",
role="author",
worktree_path=self.wt,
)
self.assertFalse(result["success"])
self.assertEqual(result["outcome"], "blocked")
self.assertEqual(result["profile_role_kind"], "merger")
self._assert_handoff_untouched(lid)
def test_reviewer_profile_rejected_without_role_argument(self) -> None:
"""Even without a spoofed role, the profile-derived role binds."""
lid = self._handoff_lease(922)
result = self._call_adopt_tool(
self.REVIEWER_PROFILE,
lease_id=lid,
session_id="reviewer-worker",
worktree_path=self.wt,
)
self.assertFalse(result["success"])
self.assertEqual(result["outcome"], "blocked")
self.assertIn("wrong role", result["reasons"][0].lower())
self._assert_handoff_untouched(lid)
def test_author_profile_mismatching_supplied_role_rejected(self) -> None:
lid = self._handoff_lease(923)
result = self._call_adopt_tool(
self.AUTHOR_PROFILE,
lease_id=lid,
session_id="author-worker",
role="reviewer",
worktree_path=self.wt,
)
self.assertFalse(result["success"])
self.assertEqual(result["outcome"], "blocked")
self.assertEqual(result["profile_role_kind"], "author")
self.assertEqual(result["supplied_role"], "reviewer")
self._assert_handoff_untouched(lid)
def test_foreign_profile_name_rejected_for_author_handoff(self) -> None:
"""Provenance required_profile binds even when the role matches."""
lid = self._handoff_lease(924)
result = self._call_adopt_tool(
self.FOREIGN_AUTHOR_PROFILE,
lease_id=lid,
session_id="foreign-author-worker",
worktree_path=self.wt,
)
self.assertFalse(result["success"])
self.assertEqual(result["outcome"], "blocked")
self.assertIn("wrong profile", result["reasons"][0].lower())
self._assert_handoff_untouched(lid)
def test_author_profile_consumes_author_handoff(self) -> None:
lid = self._handoff_lease(925)
result = self._call_adopt_tool(
self.AUTHOR_PROFILE,
lease_id=lid,
session_id="author-worker",
role="author",
worktree_path=self.wt,
)
self.assertTrue(result["success"])
self.assertEqual(result["outcome"], "adopted_cross_role_handoff")
self.assertEqual(result["adopted_by_session_id"], "author-worker")
self.assertEqual(result["adopted_from_session_id"], "ctrl-session")
state = self.db.get_lease_workflow_state(lid)
self.assertEqual(state["lease"]["session_id"], "author-worker")
self.assertEqual(
state["lease"]["adopted_by_session_id"], "author-worker"
)
self.assertEqual(state["assignment"]["session_id"], "author-worker")
self.assertEqual(state["provenance"]["handoff_status"], "adopted")
def test_author_profile_consumes_author_handoff_without_role_argument(
self,
) -> None:
lid = self._handoff_lease(926)
result = self._call_adopt_tool(
self.AUTHOR_PROFILE,
lease_id=lid,
session_id="author-worker",
worktree_path=self.wt,
)
self.assertTrue(result["success"])
self.assertEqual(result["outcome"], "adopted_cross_role_handoff")
state = self.db.get_lease_workflow_state(lid)
self.assertEqual(state["lease"]["session_id"], "author-worker")
self.assertEqual(state["lease"]["role"], "author")
if __name__ == "__main__":
unittest.main()
-618
View File
@@ -1,618 +0,0 @@
"""Tests for the workflow-event timeline model and read API (#637).
Covers the acceptance criteria: versioned schema, adaptation of control-plane
events and Gitea handoff comments, filter by issue/PR/session, redaction of
secret-like payloads, and stable pagination.
"""
import json
import os
import sqlite3
import sys
import unittest
from pathlib import Path
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
from starlette.testclient import TestClient
import control_plane_db
from canonical_thread_handoff import format_cth_body
from webui import timeline
from webui.app import create_app
def _seed_db(path: str) -> None:
"""Create a control-plane DB and seed scoped work_items + events."""
# Constructing ControlPlaneDB runs the schema migration once.
control_plane_db.ControlPlaneDB(db_path=path)
conn = sqlite3.connect(path)
try:
conn.execute(
"INSERT INTO work_items(remote, org, repo, kind, number, state, updated_at) "
"VALUES (?,?,?,?,?,?,?)",
("prgs", "Scaled-Tech-Consulting", "Gitea-Tools", "issue", 637, "open", "2026-07-23T00:00:00Z"),
)
issue_wid = conn.execute("SELECT last_insert_rowid()").fetchone()[0]
conn.execute(
"INSERT INTO work_items(remote, org, repo, kind, number, state, updated_at) "
"VALUES (?,?,?,?,?,?,?)",
("prgs", "Scaled-Tech-Consulting", "Gitea-Tools", "pr", 813, "open", "2026-07-23T00:00:00Z"),
)
pr_wid = conn.execute("SELECT last_insert_rowid()").fetchone()[0]
# A work item for a different repo — must never appear in prgs/Gitea-Tools scope.
conn.execute(
"INSERT INTO work_items(remote, org, repo, kind, number, state, updated_at) "
"VALUES (?,?,?,?,?,?,?)",
("dadeschools", "Other", "Elsewhere", "issue", 1, "open", "2026-07-23T00:00:00Z"),
)
other_wid = conn.execute("SELECT last_insert_rowid()").fetchone()[0]
rows = [
(issue_wid, "allocation", "assigned author work", "2026-07-23T01:00:00Z"),
(issue_wid, "lease.renew", "token=ghs_ABCDEF1234567890abcdef lease renewed", "2026-07-23T02:00:00Z"),
(pr_wid, "pr.opened", "PR opened for review", "2026-07-23T03:00:00Z"),
(other_wid, "allocation", "off-scope event", "2026-07-23T04:00:00Z"),
]
conn.executemany(
"INSERT INTO events(work_item_id, event_type, message, created_at) VALUES (?,?,?,?)",
rows,
)
conn.commit()
finally:
conn.close()
class TestSchema(unittest.TestCase):
def test_schema_is_versioned(self):
self.assertIsInstance(timeline.TIMELINE_SCHEMA_VERSION, int)
self.assertGreaterEqual(timeline.TIMELINE_SCHEMA_VERSION, 1)
def test_event_to_dict_shape(self):
ev = timeline.WorkflowEvent(
source=timeline.SOURCE_CONTROL_PLANE,
event_type="allocation",
event_key="cp:1",
timestamp="2026-07-23T01:00:00Z",
issue_number=637,
)
d = ev.to_dict()
for key in (
"source", "event_type", "event_key", "timestamp", "actor", "role",
"issue_number", "pr_number", "session_id", "tool_name", "decision",
"message", "correlation_id", "evidence_refs", "sensitive",
):
self.assertIn(key, d)
self.assertEqual(d["evidence_refs"], [])
class TestCpAdapter(unittest.TestCase):
def test_issue_and_pr_mapping(self):
rows = [
{"event_id": 1, "event_type": "allocation", "message": "x", "created_at": "2026-07-23T01:00:00Z", "kind": "issue", "number": 637},
{"event_id": 2, "event_type": "pr.opened", "message": "y", "created_at": "2026-07-23T02:00:00Z", "kind": "pr", "number": 813},
]
events = timeline.adapt_cp_events(rows)
self.assertEqual(len(events), 2)
self.assertEqual(events[0].issue_number, 637)
self.assertIsNone(events[0].pr_number)
self.assertEqual(events[0].correlation_id, "issue#637")
self.assertIsNone(events[1].issue_number)
self.assertEqual(events[1].pr_number, 813)
def test_malformed_rows_skipped(self):
rows = [
{"event_id": None, "event_type": "x", "kind": "issue", "number": 1},
{"event_id": 5, "event_type": "", "kind": "issue", "number": 1},
{"event_id": 6, "event_type": "ok", "message": "m", "created_at": None, "kind": "issue", "number": 1},
]
events = timeline.adapt_cp_events(rows)
self.assertEqual(len(events), 1)
self.assertIsNone(events[0].timestamp)
def test_sensitive_event_flagged(self):
rows = [{"event_id": 1, "event_type": "lease.renew", "message": "m", "created_at": "2026-07-23T01:00:00Z", "kind": "issue", "number": 1}]
events = timeline.adapt_cp_events(rows)
self.assertTrue(events[0].sensitive)
class TestCthAdapter(unittest.TestCase):
def test_cth_comment_becomes_event(self):
body = format_cth_body(
cth_type="Author Handoff",
status="ready",
next_owner="reviewer",
decision="implement timeline",
proof="commit abc1234 closes #637",
next_action="review PR",
ready_to_paste_prompt="Review PR #900 as reviewer",
)
comments = [{"id": 42, "body": body, "created_at": "2026-07-23T05:00:00Z", "user": {"login": "jcwalker3"}}]
events = timeline.adapt_cth_comments(comments, kind="issue", number=637)
self.assertEqual(len(events), 1)
ev = events[0]
self.assertEqual(ev.source, timeline.SOURCE_GITEA_HANDOFF)
self.assertEqual(ev.event_type, "handoff:Author Handoff")
self.assertEqual(ev.actor, "jcwalker3")
self.assertEqual(ev.issue_number, 637)
self.assertEqual(ev.event_key, "cth:issue:637:42")
self.assertIn("#637", ev.evidence_refs)
self.assertIn("abc1234", ev.evidence_refs)
def test_non_cth_comment_ignored(self):
comments = [{"id": 1, "body": "just a normal comment", "created_at": "2026-07-23T05:00:00Z", "user": {"login": "x"}}]
self.assertEqual(timeline.adapt_cth_comments(comments, kind="issue", number=1), [])
class TestRedaction(unittest.TestCase):
def test_cp_message_redacted(self):
rows = [{"event_id": 1, "event_type": "lease", "message": "token=ghs_ABCDEF1234567890abcdef here", "created_at": "2026-07-23T01:00:00Z", "kind": "issue", "number": 1}]
events = timeline.adapt_cp_events(rows)
self.assertNotIn("ghs_ABCDEF1234567890abcdef", events[0].message or "")
def test_handoff_decision_redacted(self):
body = format_cth_body(
cth_type="Blocker",
status="blocked",
next_owner="author",
decision="password=SuperSecret123! must rotate",
proof="none",
next_action="rotate",
ready_to_paste_prompt="Rotate the credential and retry",
)
comments = [{"id": 7, "body": body, "created_at": "2026-07-23T05:00:00Z", "user": {"login": "x"}}]
events = timeline.adapt_cth_comments(comments, kind="issue", number=1)
self.assertNotIn("SuperSecret123!", events[0].decision or "")
class TestFilterSortPaginate(unittest.TestCase):
def _events(self):
return [
timeline.WorkflowEvent(source="control_plane", event_type="a", event_key="cp:3", timestamp="2026-07-23T03:00:00Z", pr_number=813),
timeline.WorkflowEvent(source="control_plane", event_type="b", event_key="cp:1", timestamp="2026-07-23T01:00:00Z", issue_number=637),
timeline.WorkflowEvent(source="control_plane", event_type="c", event_key="cp:2", timestamp="2026-07-23T02:00:00Z", issue_number=637, session_id="sess-1"),
]
def test_filter_by_issue(self):
out = timeline.filter_events(self._events(), issue=637)
self.assertEqual({e.event_key for e in out}, {"cp:1", "cp:2"})
def test_filter_by_pr(self):
out = timeline.filter_events(self._events(), pr=813)
self.assertEqual([e.event_key for e in out], ["cp:3"])
def test_filter_by_session(self):
out = timeline.filter_events(self._events(), session="sess-1")
self.assertEqual([e.event_key for e in out], ["cp:2"])
def test_stable_sort_ascending(self):
out = timeline.sort_events(self._events())
self.assertEqual([e.event_key for e in out], ["cp:1", "cp:2", "cp:3"])
def test_missing_timestamp_sorts_last(self):
evs = self._events() + [
timeline.WorkflowEvent(source="control_plane", event_type="z", event_key="cp:9", timestamp=None)
]
out = timeline.sort_events(evs)
self.assertEqual(out[-1].event_key, "cp:9")
def test_pagination_windows_and_next_offset(self):
evs = timeline.sort_events(self._events())
page1 = timeline.paginate(evs, limit=2, offset=0)
self.assertEqual(len(page1.events), 2)
self.assertEqual(page1.total, 3)
self.assertEqual(page1.next_offset, 2)
page2 = timeline.paginate(evs, limit=2, offset=2)
self.assertEqual(len(page2.events), 1)
self.assertIsNone(page2.next_offset)
def test_pagination_bounds_coerced(self):
evs = self._events()
page = timeline.paginate(evs, limit=-5, offset=-3)
self.assertGreaterEqual(page.limit, 1)
self.assertEqual(page.offset, 0)
class TestCpReader(unittest.TestCase):
def test_reads_scoped_events_only(self):
import tempfile
with tempfile.TemporaryDirectory() as tmp:
db = os.path.join(tmp, "cp.sqlite3")
_seed_db(db)
events, status = timeline.read_cp_events(
remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools", db_path=db
)
self.assertTrue(status.ok)
# 3 scoped events; the dadeschools/Other event is excluded.
self.assertEqual(len(events), 3)
self.assertTrue(all(e.source == "control_plane" for e in events))
# Redaction applied to the token-bearing message.
joined = " ".join(e.message or "" for e in events)
self.assertNotIn("ghs_ABCDEF1234567890abcdef", joined)
def test_missing_db_degrades(self):
events, status = timeline.read_cp_events(
remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools",
db_path="/nonexistent/path/to/cp.sqlite3",
)
self.assertEqual(events, [])
self.assertFalse(status.ok)
self.assertIsNotNone(status.reason)
class TestLoadTimeline(unittest.TestCase):
def test_handoff_not_run_without_thread_filter(self):
import tempfile
with tempfile.TemporaryDirectory() as tmp:
db = os.path.join(tmp, "cp.sqlite3")
_seed_db(db)
snap = timeline.load_timeline(
remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools", db_path=db
)
d = snap.to_dict()
handoff = [s for s in d["sources"] if s["name"] == "gitea_handoff"][0]
self.assertFalse(handoff["ok"])
self.assertIn("thread-scoped", handoff["reason"])
self.assertEqual(d["schema_version"], timeline.TIMELINE_SCHEMA_VERSION)
def test_handoff_included_via_injected_source(self):
import tempfile
body = format_cth_body(
cth_type="Author Handoff", status="ready", next_owner="reviewer",
decision="d", proof="#637", next_action="review", ready_to_paste_prompt="Review PR #1 now",
)
def source(kind, number):
return [{"id": 1, "body": body, "created_at": "2026-07-23T09:00:00Z", "user": {"login": "jcwalker3"}}]
with tempfile.TemporaryDirectory() as tmp:
db = os.path.join(tmp, "cp.sqlite3")
_seed_db(db)
snap = timeline.load_timeline(
remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools",
issue=637, db_path=db, comment_source=source,
)
d = snap.to_dict()
handoff = [s for s in d["sources"] if s["name"] == "gitea_handoff"][0]
self.assertTrue(handoff["ok"])
self.assertEqual(handoff["count"], 1)
# Both a CP event and the handoff event for issue 637 appear, sorted.
kinds = {e["source"] for e in d["events"]}
self.assertEqual(kinds, {"control_plane", "gitea_handoff"})
def test_failing_comment_source_degrades_only_handoff(self):
import tempfile
def boom(kind, number):
raise RuntimeError("network down")
with tempfile.TemporaryDirectory() as tmp:
db = os.path.join(tmp, "cp.sqlite3")
_seed_db(db)
snap = timeline.load_timeline(
remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools",
issue=637, db_path=db, comment_source=boom,
)
d = snap.to_dict()
cp = [s for s in d["sources"] if s["name"] == "control_plane"][0]
handoff = [s for s in d["sources"] if s["name"] == "gitea_handoff"][0]
self.assertTrue(cp["ok"])
self.assertFalse(handoff["ok"])
self.assertIn("network down", handoff["reason"])
# A fabricated 40-character lowercase hex value with the shape of a Gitea
# personal access token. Never a real credential — its only job is to prove it
# cannot reach any part of a serialized timeline payload.
SYNTHETIC_SECRET_40_HEX = "a3f9c17be44d2058e6b17c9d0f5321ab77c4e9d1"
def _cth_comment(comment_id, *, created_at, session=None, decision="d", proof="none", **kw):
"""Build one real CTH comment record, optionally declaring a session."""
extra = {"Session": session} if session is not None else None
body = format_cth_body(
cth_type=kw.pop("cth_type", "Author Handoff"),
status=kw.pop("status", "ready"),
next_owner=kw.pop("next_owner", "reviewer"),
decision=decision,
proof=proof,
next_action=kw.pop("next_action", "review"),
ready_to_paste_prompt=kw.pop("ready_to_paste_prompt", "Review PR #1 now"),
extra_fields=extra,
)
return {
"id": comment_id,
"body": body,
"created_at": created_at,
"user": {"login": kw.pop("login", "jcwalker3")},
}
class TestSessionFilterThroughAdapter(unittest.TestCase):
"""F1: the session dimension must be real, or explicitly refused.
These drive the filter through the CTH adapter and the composed
``load_timeline``/API path, never through a hand-built ``WorkflowEvent``.
"""
def _source(self, comments):
return lambda kind, number: list(comments)
def test_adapter_populates_declared_session(self):
events = timeline.adapt_cth_comments(
[_cth_comment(1, created_at="2026-07-23T05:00:00Z", session="sess-alpha")],
kind="issue",
number=637,
)
self.assertEqual(len(events), 1)
self.assertEqual(events[0].session_id, "sess-alpha")
def test_session_filter_matches_through_adapter(self):
import tempfile
comments = [
_cth_comment(1, created_at="2026-07-23T09:00:00Z", session="sess-alpha"),
_cth_comment(2, created_at="2026-07-23T08:00:00Z", session="sess-alpha"),
_cth_comment(3, created_at="2026-07-23T10:00:00Z", session="sess-beta"),
]
with tempfile.TemporaryDirectory() as tmp:
db = os.path.join(tmp, "cp.sqlite3")
_seed_db(db)
snap = timeline.load_timeline(
remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools",
issue=637, session="sess-alpha", db_path=db,
comment_source=self._source(comments),
)
d = snap.to_dict()
self.assertTrue(d["ok"])
self.assertIsNone(d["error"])
keys = [e["event_key"] for e in d["events"]]
# Only the two sess-alpha events, still in ascending timestamp order.
self.assertEqual(keys, ["cth:issue:637:2", "cth:issue:637:1"])
self.assertTrue(all(e["session_id"] == "sess-alpha" for e in d["events"]))
self.assertEqual(d["pagination"]["total"], 2)
def test_session_filter_paginates_and_keeps_order(self):
import tempfile
comments = [
_cth_comment(i, created_at=f"2026-07-23T0{i}:00:00Z", session="sess-alpha")
for i in range(1, 4)
]
with tempfile.TemporaryDirectory() as tmp:
db = os.path.join(tmp, "cp.sqlite3")
_seed_db(db)
kwargs = dict(
remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools",
issue=637, session="sess-alpha", db_path=db,
comment_source=self._source(comments),
)
page1 = timeline.load_timeline(limit=2, offset=0, **kwargs).to_dict()
page2 = timeline.load_timeline(limit=2, offset=2, **kwargs).to_dict()
self.assertEqual(
[e["event_key"] for e in page1["events"]],
["cth:issue:637:1", "cth:issue:637:2"],
)
self.assertEqual(page1["pagination"]["next_offset"], 2)
self.assertEqual([e["event_key"] for e in page2["events"]], ["cth:issue:637:3"])
self.assertIsNone(page2["pagination"]["next_offset"])
def test_unknown_session_is_honestly_empty_when_supported(self):
"""A source that *can* answer the dimension may legitimately match nothing."""
import tempfile
comments = [_cth_comment(1, created_at="2026-07-23T09:00:00Z", session="sess-alpha")]
with tempfile.TemporaryDirectory() as tmp:
db = os.path.join(tmp, "cp.sqlite3")
_seed_db(db)
d = timeline.load_timeline(
remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools",
issue=637, session="sess-nope", db_path=db,
comment_source=self._source(comments),
).to_dict()
self.assertTrue(d["ok"])
self.assertEqual(d["events"], [])
def test_session_filter_refused_when_no_source_can_answer(self):
"""The F1 defect: an empty-and-healthy page for an unanswerable filter."""
import tempfile
with tempfile.TemporaryDirectory() as tmp:
db = os.path.join(tmp, "cp.sqlite3")
_seed_db(db)
d = timeline.load_timeline(
remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools",
issue=637, session="sess-alpha", db_path=db, comment_source=None,
).to_dict()
self.assertFalse(d["ok"])
self.assertEqual(d["error"]["code"], "filter_not_supported")
self.assertEqual(d["error"]["unsupported_filters"], ["session"])
self.assertEqual(d["events"], [])
self.assertEqual(d["pagination"]["total"], 0)
# The refusal states which source could not answer, and why.
explained = {r["source"] for r in d["error"]["sources"]}
self.assertEqual(explained, {"control_plane", "gitea_handoff"})
def test_control_plane_declares_session_unsupported(self):
import tempfile
with tempfile.TemporaryDirectory() as tmp:
db = os.path.join(tmp, "cp.sqlite3")
_seed_db(db)
d = timeline.load_timeline(
remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools",
issue=637, session="sess-alpha", db_path=db,
comment_source=self._source([]),
).to_dict()
cp = [s for s in d["sources"] if s["name"] == "control_plane"][0]
handoff = [s for s in d["sources"] if s["name"] == "gitea_handoff"][0]
self.assertNotIn("session", cp["supported_filters"])
self.assertEqual(cp["unsupported_filters"], ["session"])
self.assertIn("session", handoff["supported_filters"])
self.assertEqual(handoff["unsupported_filters"], [])
def test_secret_shaped_session_value_is_dropped(self):
events = timeline.adapt_cth_comments(
[_cth_comment(1, created_at="2026-07-23T05:00:00Z", session=SYNTHETIC_SECRET_40_HEX)],
kind="issue",
number=637,
)
self.assertIsNone(events[0].session_id)
self.assertNotIn(SYNTHETIC_SECRET_40_HEX, json.dumps(events[0].to_dict()))
class TestEvidenceRefRedaction(unittest.TestCase):
"""F2: evidence_refs must not be a hole in the redaction boundary."""
def _refs_for(self, *, proof="none", decision="d"):
events = timeline.adapt_cth_comments(
[_cth_comment(1, created_at="2026-07-23T05:00:00Z", proof=proof, decision=decision)],
kind="issue",
number=637,
)
self.assertEqual(len(events), 1)
return events[0]
def test_assigned_secret_never_reaches_evidence_refs(self):
ev = self._refs_for(proof=f"authenticated with token={SYNTHETIC_SECRET_40_HEX}")
self.assertNotIn(SYNTHETIC_SECRET_40_HEX, ev.evidence_refs)
self.assertNotIn(SYNTHETIC_SECRET_40_HEX, json.dumps(ev.to_dict()))
def test_bare_secret_shaped_value_never_reaches_evidence_refs(self):
# An undeclared hex run in proof text is not evidence of anything, and
# proof itself is never serialized — so the value has no way out.
ev = self._refs_for(proof=f"proof {SYNTHETIC_SECRET_40_HEX}", decision="rotate the credential")
self.assertNotIn(SYNTHETIC_SECRET_40_HEX, ev.evidence_refs)
self.assertNotIn(SYNTHETIC_SECRET_40_HEX, json.dumps(ev.to_dict()))
def test_secret_absent_from_complete_serialized_payload(self):
import tempfile
comments = [
_cth_comment(
1,
created_at="2026-07-23T09:00:00Z",
proof=f"lease token={SYNTHETIC_SECRET_40_HEX} and bare {SYNTHETIC_SECRET_40_HEX}",
decision=f"rotate api_key={SYNTHETIC_SECRET_40_HEX}",
)
]
with tempfile.TemporaryDirectory() as tmp:
db = os.path.join(tmp, "cp.sqlite3")
_seed_db(db)
snap = timeline.load_timeline(
remote="prgs", org="Scaled-Tech-Consulting", repo="Gitea-Tools",
issue=637, db_path=db, comment_source=lambda k, n: comments,
)
payload = json.dumps(snap.to_dict())
self.assertNotIn(SYNTHETIC_SECRET_40_HEX, payload)
# And the surface is clean by the redaction policy's own detectors.
from webui import console_redaction
self.assertEqual(console_redaction.scan_for_secrets(snap.to_dict()), [])
def test_legitimate_references_still_usable(self):
head = "4f3a464a1c455b6ecaaae9b6eef496c8f8ed451a"
ev = self._refs_for(
proof=f"closes #637, PR #849, commit abc1234, at head {head}",
decision="none",
)
for token in ("#637", "#849", "abc1234", head):
self.assertIn(token, ev.evidence_refs)
def test_undeclared_hex_words_are_not_references(self):
ev = self._refs_for(proof="the record was defaced and the facade decayed")
self.assertEqual(ev.evidence_refs, ())
def test_refs_revalidated_independently_before_serialization(self):
"""Extraction is not trusted: the validator drops anything unproven."""
safe, dropped = timeline._validated_evidence_refs(
["#637", "abc1234", "not-a-ref", f"token={SYNTHETIC_SECRET_40_HEX}", ""]
)
self.assertEqual(safe, ("#637", "abc1234"))
self.assertTrue(dropped)
self.assertNotIn(SYNTHETIC_SECRET_40_HEX, json.dumps(list(safe)))
def test_clean_event_is_not_marked_sensitive(self):
ev = self._refs_for(proof="closes #637")
self.assertEqual(ev.evidence_refs, ("#637",))
self.assertFalse(ev.sensitive)
class TestTimelineApi(unittest.TestCase):
def setUp(self):
self._prev_db = os.environ.get(control_plane_db.DB_PATH_ENV)
self._prev_offline = os.environ.get("WEBUI_TEST_OFFLINE")
import tempfile
self._tmpdir = tempfile.TemporaryDirectory()
self._db = os.path.join(self._tmpdir.name, "cp.sqlite3")
_seed_db(self._db)
os.environ[control_plane_db.DB_PATH_ENV] = self._db
os.environ["WEBUI_TEST_OFFLINE"] = "1"
self.client = TestClient(create_app())
def tearDown(self):
if self._prev_db is None:
os.environ.pop(control_plane_db.DB_PATH_ENV, None)
else:
os.environ[control_plane_db.DB_PATH_ENV] = self._prev_db
if self._prev_offline is None:
os.environ.pop("WEBUI_TEST_OFFLINE", None)
else:
os.environ["WEBUI_TEST_OFFLINE"] = self._prev_offline
self._tmpdir.cleanup()
def test_api_returns_timeline(self):
resp = self.client.get("/api/v1/timeline")
self.assertEqual(resp.status_code, 200)
body = resp.json()
self.assertEqual(body["schema_version"], timeline.TIMELINE_SCHEMA_VERSION)
self.assertIn("events", body)
self.assertIn("pagination", body)
self.assertGreaterEqual(body["pagination"]["total"], 1)
def test_api_filter_by_issue(self):
resp = self.client.get("/api/v1/timeline?issue=637")
self.assertEqual(resp.status_code, 200)
events = resp.json()["events"]
self.assertTrue(events)
self.assertTrue(all(e["issue_number"] == 637 for e in events))
def test_api_pagination(self):
resp = self.client.get("/api/v1/timeline?limit=1&offset=0")
self.assertEqual(resp.status_code, 200)
pg = resp.json()["pagination"]
self.assertEqual(pg["limit"], 1)
self.assertEqual(len(resp.json()["events"]), 1)
if pg["total"] > 1:
self.assertTrue(pg["has_more"])
def test_api_is_read_only(self):
resp = self.client.post("/api/v1/timeline")
self.assertIn(resp.status_code, (404, 405))
def test_api_no_secret_leak(self):
resp = self.client.get("/api/v1/timeline?issue=637")
self.assertNotIn("ghs_ABCDEF1234567890abcdef", resp.text)
def test_api_refuses_unanswerable_session_filter(self):
"""No handoff source is configured offline, so nothing can carry a session."""
resp = self.client.get("/api/v1/timeline?issue=637&session=sess-alpha")
self.assertEqual(resp.status_code, 422)
body = resp.json()
self.assertFalse(body["ok"])
self.assertEqual(body["error"]["code"], "filter_not_supported")
self.assertEqual(body["error"]["unsupported_filters"], ["session"])
self.assertEqual(body["events"], [])
self.assertEqual(body["pagination"]["total"], 0)
def test_api_unfiltered_read_stays_ok(self):
body = self.client.get("/api/v1/timeline?issue=637").json()
self.assertTrue(body["ok"])
self.assertIsNone(body["error"])
if __name__ == "__main__":
unittest.main()
-100
View File
@@ -45,7 +45,6 @@ from webui.worktree_scanner import load_hygiene_snapshot, snapshot_to_dict as wo
from webui.worktree_views import render_worktrees_page
from webui.runtime_health import load_runtime_snapshot, snapshot_to_dict as runtime_snapshot_to_dict
from webui.runtime_views import render_runtime_page
from webui.timeline import load_timeline, snapshot_to_dict as timeline_snapshot_to_dict
from webui.system_health import (
API_PATH as SYSTEM_HEALTH_API_PATH,
load_system_health,
@@ -411,104 +410,6 @@ async def api_console_security_model(_request: Request) -> JSONResponse:
})
def _query_int(request: Request, key: str) -> int | None:
"""Parse an optional integer query parameter; None when absent/invalid."""
raw = request.query_params.get(key)
if raw is None or not str(raw).strip():
return None
try:
return int(str(raw).strip())
except (TypeError, ValueError):
return None
def _derive_remote(host: str) -> str:
"""Map a Gitea host to its known short remote name (control-plane scope key)."""
text = (host or "").lower()
if "prgs" in text:
return "prgs"
if "dadeschools" in text:
return "dadeschools"
return text.split(".")[0] if text else ""
def _timeline_comment_source(host: str, org: str, repo: str):
"""Build a fail-soft CTH-comment fetcher for one repo, or None when offline.
Returns a callable ``(kind, number) -> list[comment]``. Credentials or
network failures raise inside the callable so ``load_timeline`` degrades the
handoff source rather than the whole timeline. Offline test mode yields no
live source so the handoff section reports ``not run``.
"""
import os
from gitea_auth import api_fetch_page, get_auth_header, repo_api_url
offline = (os.environ.get("WEBUI_TEST_OFFLINE") or "").strip().lower() in {"1", "true", "yes"}
if offline:
return None
auth = get_auth_header(host)
if not auth:
return None
def _fetch(kind: str, number: int) -> list:
segment = "pulls" if kind == "pr" else "issues"
url = f"{repo_api_url(host, org, repo)}/{segment}/{int(number)}/comments"
comments: list = []
page = 1
while page <= 20:
raw, meta = api_fetch_page(url, auth, page=page, limit=50)
comments.extend(raw)
if bool(meta["is_final_page"]):
break
page += 1
return comments
return _fetch
async def api_v1_timeline(request: Request) -> JSONResponse:
"""Read-only workflow-event timeline (#637). Filter by issue/PR/session."""
from webui.queue_loader import _host_from_url # host normalisation helper
registry, error = _load_project_registry()
if error is not None:
return JSONResponse(error.to_dict(), status_code=500)
project = registry.projects[0] if registry.projects else None
org = request.query_params.get("org") or (project.gitea_owner if project else "")
repo = request.query_params.get("repo") or (project.repo_name if project else "")
host = _host_from_url(project.remote_host) if project else ""
remote = request.query_params.get("remote") or _derive_remote(host)
if not (remote and org and repo):
return JSONResponse(
{
"error": "timeline_scope_unresolved",
"detail": "no project in registry and no remote/org/repo query params provided",
},
status_code=400,
)
comment_source = _timeline_comment_source(host, org, repo) if (host and org and repo) else None
snapshot = load_timeline(
remote=remote,
org=org,
repo=repo,
issue=_query_int(request, "issue"),
pr=_query_int(request, "pr"),
session=(request.query_params.get("session") or None),
limit=_query_int(request, "limit"),
offset=_query_int(request, "offset"),
comment_source=comment_source,
)
# A filter no surviving source can carry is refused, not answered empty:
# a 200 with zero events would tell the operator no such activity exists.
status_code = 200 if snapshot.ok else 422
return JSONResponse(timeline_snapshot_to_dict(snapshot), status_code=status_code)
async def method_not_allowed(request: Request, _exc: Exception) -> Response:
path = request.url.path
if path in _AUDIT_MUTATION_PATHS and request.method == "POST":
@@ -548,7 +449,6 @@ def create_app(*, bind_host: str | None = None) -> Starlette:
Route("/api/prompts", api_prompts, methods=["GET"]),
Route("/runtime", runtime, methods=["GET"]),
Route("/api/runtime", api_runtime, methods=["GET"]),
Route("/api/v1/timeline", api_v1_timeline, methods=["GET"]),
Route("/audit", audit, methods=["GET", "POST"]),
Route("/api/audit", api_audit, methods=["GET", "POST"]),
Route("/worktrees", worktrees, methods=["GET"]),
-784
View File
@@ -1,784 +0,0 @@
"""Workflow-event and conversation timeline model (#637, Phase 1).
Operators cannot browse a unified timeline of workflow events, decisions,
tool calls, and handoffs: the evidence is scattered across control-plane
events, Gitea canonical handoff comments, and local logs. This module defines
one durable, versioned event schema and per-source adapters that normalise
those scattered records into a single ``WorkflowEvent`` stream, plus a
read-only query layer (filter by issue / PR / session, stable ordering,
pagination) that the ``/api/v1/timeline`` route serves.
Design rules honoured here:
- **Read-only.** Sources are read; nothing is mutated. The control-plane
database is opened through a ``mode=ro`` URI so a missing or unwritable DB
degrades to a reason instead of creating directories or running migrations.
- **Fail-soft per source.** An unavailable source degrades to a status with a
reason rather than raising, and a source that could not run is never
rendered as an empty-and-healthy timeline.
- **Answerable filters only.** Each source declares which filter dimensions it
can actually answer. A filter dimension no source that ran can carry is
refused with an explicit reason rather than silently matching nothing: an
empty page from an unanswerable filter reads to an operator as "no such
activity", which is a different — and false — statement.
- **Redaction at the boundary, fail closed.** Every free-text field (event
messages, redacted tool arguments, decision/proof text) is run through the
console redaction policy before it leaves this module, and *before* any
structured value is derived from it evidence references are extracted from
redacted text, then independently revalidated before serialization. An
unredactable value becomes the placeholder, and a value that cannot be proven
safe is dropped an unredacted payload is never emitted, and a generation
error never drops raw data to a caller or a log.
- **Stable ordering.** Events sort by ``(timestamp, source_rank, event_key)``
with a deterministic tiebreak, so pagination is stable across calls and
events with equal or missing timestamps keep a fixed order.
Non-goals (from the issue): no full chat replay, no mutation of historical
events, no unredacted tool-argument storage.
"""
from __future__ import annotations
import re
import sqlite3
from dataclasses import dataclass, replace
from datetime import datetime, timezone
from typing import Any, Callable, Iterable
import control_plane_db
from webui import console_redaction
# The schema is versioned so consumers can branch on shape. Bump on any
# breaking change to WorkflowEvent's serialized form.
TIMELINE_SCHEMA_VERSION = 1
# Known event sources and their deterministic ordering rank. When two events
# carry the same timestamp, the source rank breaks the tie before the
# per-source event key, so a control-plane event and a handoff comment minted
# in the same second always sort in a fixed order.
SOURCE_CONTROL_PLANE = "control_plane"
SOURCE_GITEA_HANDOFF = "gitea_handoff"
_SOURCE_RANK = {
SOURCE_CONTROL_PLANE: 0,
SOURCE_GITEA_HANDOFF: 1,
}
# The filter dimensions the query layer accepts.
FILTER_ISSUE = "issue"
FILTER_PR = "pr"
FILTER_SESSION = "session"
# Which dimensions each source can actually answer. This is a property of the
# underlying records, not of the query code: the control-plane ``events`` table
# is (event_id, work_item_id, event_type, message, created_at) and carries no
# session identity at all, so no control-plane event can ever match a session
# filter. A CTH handoff comment can declare its session as a field, so the
# handoff source answers all three. Filtering on a dimension the surviving
# sources cannot carry is refused in ``load_timeline`` rather than answered
# with an empty page.
_SOURCE_FILTER_SUPPORT: dict[str, tuple[str, ...]] = {
SOURCE_CONTROL_PLANE: (FILTER_ISSUE, FILTER_PR),
SOURCE_GITEA_HANDOFF: (FILTER_ISSUE, FILTER_PR, FILTER_SESSION),
}
# Why a source cannot answer a dimension, for the refusal reason an operator reads.
_SOURCE_FILTER_LIMITS: dict[tuple[str, str], str] = {
(SOURCE_CONTROL_PLANE, FILTER_SESSION): (
"control-plane events carry no session identity "
"(the events table has no session column)"
),
}
# A timestamp far in the future so events with no parseable timestamp sort
# last (after everything real) instead of first, without raising.
_MISSING_TS_SORT = "9999-12-31T23:59:59Z"
def _parse_ts(value: str | None) -> str | None:
"""Normalise a timestamp to ``...Z`` UTC, or None when unparseable."""
if not value:
return None
text = str(value).strip()
if not text:
return None
candidate = text[:-1] + "+00:00" if text.endswith("Z") else text
try:
parsed = datetime.fromisoformat(candidate)
except ValueError:
return None
if parsed.tzinfo is None:
parsed = parsed.replace(tzinfo=timezone.utc)
return parsed.astimezone(timezone.utc).replace(microsecond=0).isoformat().replace("+00:00", "Z")
def _redact(value: Any) -> Any:
"""Redact a single free-text field, failing closed to the placeholder."""
if value is None:
return None
return console_redaction.redact_text(str(value))
@dataclass(frozen=True)
class WorkflowEvent:
"""One normalised timeline event.
Every field is optional except ``source``/``event_type``/``event_key``
because sources carry different subsets. The class is frozen so an adapted
event is an immutable record; a consumer that needs a variant builds a new
one rather than mutating history.
"""
source: str
event_type: str
event_key: str
timestamp: str | None = None
actor: str | None = None
role: str | None = None
issue_number: int | None = None
pr_number: int | None = None
session_id: str | None = None
tool_name: str | None = None
decision: str | None = None
message: str | None = None
correlation_id: str | None = None
evidence_refs: tuple[str, ...] = ()
sensitive: bool = False
def sort_key(self) -> tuple[str, int, str]:
return (
self.timestamp or _MISSING_TS_SORT,
_SOURCE_RANK.get(self.source, 99),
self.event_key,
)
def to_dict(self) -> dict[str, Any]:
return {
"source": self.source,
"event_type": self.event_type,
"event_key": self.event_key,
"timestamp": self.timestamp,
"actor": self.actor,
"role": self.role,
"issue_number": self.issue_number,
"pr_number": self.pr_number,
"session_id": self.session_id,
"tool_name": self.tool_name,
"decision": self.decision,
"message": self.message,
"correlation_id": self.correlation_id,
"evidence_refs": list(self.evidence_refs),
"sensitive": self.sensitive,
}
# --------------------------------------------------------------------------- #
# Adapters — pure functions from a source's raw records to WorkflowEvents. #
# Each is total: a malformed record is skipped, never raised on. #
# --------------------------------------------------------------------------- #
# Event types whose payload is treated as sensitive and always redaction-hard
# (they can carry lease/session provenance or tool arguments).
_SENSITIVE_EVENT_HINTS = ("lease", "capability", "token", "auth", "secret")
# Reference tokens (issue/PR/comment ids) and SHAs parsed out of proof text.
_EVIDENCE_REF_RE = re.compile(r"(?:#|PR\s*#?|issue\s*#?|comment\s*#?)(\d+)", re.IGNORECASE)
# A commit reference is only recognised when the text *declares* it as one.
# A bare lowercase hex run is not evidence of anything: at 40 characters it is
# exactly the shape of a Gitea personal access token, and at 7 it also matches
# ordinary words such as "defaced". Requiring an anchoring keyword keeps real
# references ("commit abc1234", "at head a209756...", "base caaae9b6") usable
# while refusing to lift an undeclared secret-shaped run out of free text.
_SHA_RE = re.compile(
r"(?i:\b(?:commit|sha|head|base|parent|revision|rev|merge[- ]base)\b[\s:=@#]*)"
r"([0-9a-f]{7,40})\b"
)
# Shapes a serialized evidence reference is allowed to take. Anything else is
# dropped rather than emitted.
_REF_ISSUE_SHAPE = re.compile(r"^#[0-9]{1,9}$")
_REF_SHA_SHAPE = re.compile(r"^[0-9a-f]{7,40}$")
# A long undelimited hex run with no declaring context is treated as credential
# material wherever it appears, never as an identifier.
_BARE_SECRET_SHAPE = re.compile(r"^[0-9a-f]{32,}$")
def _kind_to_numbers(kind: str | None, number: int | None) -> tuple[int | None, int | None]:
"""Map a control-plane work-item (kind, number) to (issue_no, pr_no)."""
if number is None:
return (None, None)
if kind == "pr":
return (None, int(number))
if kind == "issue":
return (int(number), None)
return (None, None)
def _correlation_for(kind: str | None, number: int | None) -> str | None:
if number is None or kind not in ("issue", "pr"):
return None
return f"{kind}#{number}"
def _extract_evidence_refs(*texts: str | None) -> tuple[str, ...]:
"""Extract issue/PR and declared-commit references from **redacted** text.
Callers must pass text that has already been through :func:`_redact`; this
function derives a structured field from its input, so extracting ahead of
redaction would republish whatever redaction was about to remove. Every
reference is revalidated by :func:`_validated_evidence_refs` before it is
serialized.
"""
refs: list[str] = []
for text in texts:
if not text:
continue
for match in _EVIDENCE_REF_RE.finditer(text):
token = f"#{match.group(1)}"
if token not in refs:
refs.append(token)
for match in _SHA_RE.finditer(text):
token = match.group(1)
if token not in refs:
refs.append(token)
return tuple(refs)
def _validated_evidence_refs(refs: Iterable[str]) -> tuple[tuple[str, ...], bool]:
"""Independently revalidate references immediately before serialization.
Extraction is not trusted on its own. A reference survives only when it has
a known reference shape and is unchanged by a second redaction pass a
value the redaction policy would alter is credential material that must not
be emitted as a structured field. A full 40-character SHA stays usable
because extraction only accepts a hex run the source text explicitly
declared as a commit. Returns ``(safe_refs, dropped_any)``; ``dropped_any``
marks the event sensitive so the drop is visible rather than silent.
"""
safe: list[str] = []
dropped = False
for ref in refs or ():
try:
token = str(ref).strip()
if not token:
continue
recognised = bool(_REF_ISSUE_SHAPE.match(token) or _REF_SHA_SHAPE.match(token))
if not recognised:
dropped = True
continue
if _redact(token) != token:
dropped = True
continue
if token not in safe:
safe.append(token)
except Exception:
# Fail closed: a reference that cannot be proven safe is dropped.
dropped = True
continue
return (tuple(safe), dropped)
def _safe_session_id(value: Any) -> str | None:
"""Return a session identifier only when it is safe to emit.
The value is authoritative source data a session the record names for
itself but it is still free text. It is dropped when redaction alters it
or when it is a bare secret-shaped hex run, so a credential parked in a
session field can never reach the payload or be echoed back by a filter.
"""
if value is None:
return None
text = str(value).strip()
if not text:
return None
if _BARE_SECRET_SHAPE.match(text):
return None
return text if _redact(text) == text else None
def adapt_cp_events(rows: Iterable[dict[str, Any]]) -> list[WorkflowEvent]:
"""Adapt control-plane ``events`` rows (joined to work_items) into events.
Each row is expected to carry ``event_id``, ``event_type``, ``message``,
``created_at`` and the joined work-item ``kind``/``number``. Rows missing
an id or type are skipped so a partially written table never raises.
"""
events: list[WorkflowEvent] = []
for row in rows or []:
try:
event_id = row.get("event_id")
event_type = (row.get("event_type") or "").strip()
if event_id is None or not event_type:
continue
kind = row.get("kind")
number = row.get("number")
issue_no, pr_no = _kind_to_numbers(kind, number)
sensitive = any(hint in event_type.lower() for hint in _SENSITIVE_EVENT_HINTS)
events.append(
WorkflowEvent(
source=SOURCE_CONTROL_PLANE,
event_type=event_type,
event_key=f"cp:{event_id}",
timestamp=_parse_ts(row.get("created_at")),
issue_number=issue_no,
pr_number=pr_no,
# No session_id: the control-plane events table is
# (event_id, work_item_id, event_type, message, created_at)
# and records no session. Inventing one from the work item
# or the message text would be a guess, so this source
# declares the session dimension unsupported instead
# (_SOURCE_FILTER_SUPPORT) and the query layer refuses a
# session filter it cannot honestly answer.
message=_redact(row.get("message")),
correlation_id=_correlation_for(kind, number),
sensitive=sensitive,
)
)
except Exception:
# A single malformed row must not sink the whole adaptation.
continue
return events
def adapt_cth_comments(
comments: Iterable[dict[str, Any]],
*,
kind: str,
number: int,
) -> list[WorkflowEvent]:
"""Adapt Gitea Canonical Thread Handoff (CTH) comments into events.
Only comments that parse as a CTH (``canonical_thread_handoff.parse_cth_comment``)
become events; ordinary comments are ignored. ``kind``/``number`` scope the
events to the issue or PR the comments belong to.
"""
# Imported lazily so this module has no import-time dependency on the
# handoff parser when only the control-plane adapter is used.
from canonical_thread_handoff import parse_cth_comment
issue_no, pr_no = _kind_to_numbers(kind, number)
correlation = _correlation_for(kind, number)
events: list[WorkflowEvent] = []
for comment in comments or []:
try:
body = comment.get("body") or ""
parsed = parse_cth_comment(body)
if not parsed:
continue
fields = parsed.get("fields") or {}
cth_type = parsed.get("cth_type") or "handoff"
comment_id = comment.get("id")
# Redaction runs first, and every derived value is taken from the
# redacted text — deriving evidence refs from the raw proof would
# re-emit exactly what redaction was about to remove.
decision = _redact(fields.get("decision"))
proof = _redact(fields.get("proof"))
next_action = _redact(fields.get("next action"))
refs, refs_dropped = _validated_evidence_refs(
_extract_evidence_refs(proof, decision)
)
events.append(
WorkflowEvent(
source=SOURCE_GITEA_HANDOFF,
event_type=f"handoff:{cth_type}",
event_key=f"cth:{kind}:{number}:{comment_id}",
timestamp=_parse_ts(comment.get("created_at")),
actor=_redact((comment.get("user") or {}).get("login")),
role=_redact(fields.get("next owner")),
issue_number=issue_no,
pr_number=pr_no,
# A CTH names its own session when the producer records one;
# it is read from that declared field, never inferred from
# unrelated text.
session_id=_safe_session_id(fields.get("session")),
decision=decision,
message=next_action or _redact(fields.get("status")),
correlation_id=correlation,
evidence_refs=refs,
sensitive=refs_dropped,
)
)
except Exception:
continue
return events
# --------------------------------------------------------------------------- #
# Read-only control-plane event source. #
# --------------------------------------------------------------------------- #
_CP_EVENTS_QUERY = """
SELECT e.event_id AS event_id,
e.event_type AS event_type,
e.message AS message,
e.created_at AS created_at,
w.kind AS kind,
w.number AS number
FROM events e
JOIN work_items w ON e.work_item_id = w.work_item_id
WHERE w.remote = ? AND w.org = ? AND w.repo = ?
"""
@dataclass(frozen=True)
class SourceStatus:
"""Fail-soft status for one timeline source.
``supported_filters`` states which filter dimensions this source's records
can carry; ``unsupported_filters`` names the requested dimensions it cannot,
so an operator can see *why* a source contributed nothing rather than being
left to read an empty list as an absence of activity.
"""
name: str
ok: bool
reason: str | None = None
count: int = 0
supported_filters: tuple[str, ...] = ()
unsupported_filters: tuple[str, ...] = ()
def to_dict(self) -> dict[str, Any]:
return {
"name": self.name,
"ok": self.ok,
"reason": self.reason,
"count": self.count,
"supported_filters": list(self.supported_filters),
"unsupported_filters": list(self.unsupported_filters),
}
def _cp_status(*, ok: bool, reason: str | None = None, count: int = 0) -> SourceStatus:
return SourceStatus(
SOURCE_CONTROL_PLANE,
ok=ok,
reason=reason,
count=count,
supported_filters=_SOURCE_FILTER_SUPPORT[SOURCE_CONTROL_PLANE],
)
def _handoff_status(*, ok: bool, reason: str | None = None, count: int = 0) -> SourceStatus:
return SourceStatus(
SOURCE_GITEA_HANDOFF,
ok=ok,
reason=reason,
count=count,
supported_filters=_SOURCE_FILTER_SUPPORT[SOURCE_GITEA_HANDOFF],
)
def read_cp_events(
*,
remote: str,
org: str,
repo: str,
db_path: str | None = None,
) -> tuple[list[WorkflowEvent], SourceStatus]:
"""Read scoped control-plane events read-only. Never creates the DB.
Opens the SQLite file through a ``mode=ro`` URI: a health/timeline read
must never create directories or run the schema migration that
``ControlPlaneDB()`` performs on construction. A missing or unreadable DB
degrades to a status with a reason.
"""
path = (db_path or control_plane_db.default_db_path()).strip()
conn: sqlite3.Connection | None = None
try:
conn = sqlite3.connect(f"file:{path}?mode=ro", uri=True)
conn.row_factory = sqlite3.Row
cursor = conn.execute(_CP_EVENTS_QUERY, (remote, org, repo))
rows = [dict(r) for r in cursor.fetchall()]
except sqlite3.OperationalError as exc:
return ([], _cp_status(ok=False, reason=f"control-plane DB unavailable: {exc}"))
except sqlite3.Error as exc:
return ([], _cp_status(ok=False, reason=f"control-plane read failed: {exc}"))
finally:
if conn is not None:
conn.close()
events = adapt_cp_events(rows)
return (events, _cp_status(ok=True, count=len(events)))
# --------------------------------------------------------------------------- #
# Filter, sort, paginate. #
# --------------------------------------------------------------------------- #
def filter_events(
events: Iterable[WorkflowEvent],
*,
issue: int | None = None,
pr: int | None = None,
session: str | None = None,
) -> list[WorkflowEvent]:
"""Filter events by issue number, PR number, and/or session id.
Filters are conjunctive. A filter that names a dimension an event does not
carry excludes that event (an issue filter excludes PR-only events).
"""
out: list[WorkflowEvent] = []
for ev in events:
if issue is not None and ev.issue_number != issue:
continue
if pr is not None and ev.pr_number != pr:
continue
if session is not None and ev.session_id != session:
continue
out.append(ev)
return out
def sort_events(events: Iterable[WorkflowEvent]) -> list[WorkflowEvent]:
"""Return events in stable timeline order (ascending)."""
return sorted(events, key=lambda ev: ev.sort_key())
@dataclass(frozen=True)
class TimelinePage:
"""One page of the sorted, filtered timeline."""
events: tuple[WorkflowEvent, ...]
total: int
limit: int
offset: int
@property
def next_offset(self) -> int | None:
nxt = self.offset + len(self.events)
return nxt if nxt < self.total else None
def to_dict(self) -> dict[str, Any]:
return {
"events": [ev.to_dict() for ev in self.events],
"pagination": {
"total": self.total,
"limit": self.limit,
"offset": self.offset,
"returned": len(self.events),
"next_offset": self.next_offset,
"has_more": self.next_offset is not None,
},
}
_MAX_LIMIT = 500
_DEFAULT_LIMIT = 50
def _coerce_bounds(limit: int | None, offset: int | None) -> tuple[int, int]:
try:
lim = int(limit) if limit is not None else _DEFAULT_LIMIT
except (TypeError, ValueError):
lim = _DEFAULT_LIMIT
try:
off = int(offset) if offset is not None else 0
except (TypeError, ValueError):
off = 0
lim = max(1, min(lim, _MAX_LIMIT))
off = max(0, off)
return (lim, off)
def paginate(events: list[WorkflowEvent], *, limit: int | None, offset: int | None) -> TimelinePage:
lim, off = _coerce_bounds(limit, offset)
window = events[off : off + lim]
return TimelinePage(events=tuple(window), total=len(events), limit=lim, offset=off)
# --------------------------------------------------------------------------- #
# Composition — load_timeline aggregates all sources, fail-soft. #
# --------------------------------------------------------------------------- #
# A comment source is a callable that, given (kind, number), returns the raw
# Gitea comment list for that issue/PR. The route supplies a live fail-soft
# fetcher; tests supply a fixture. When None, the handoff source is reported as
# not-run (never silently empty-and-healthy).
CommentSource = Callable[[str, int], list[dict[str, Any]]]
@dataclass(frozen=True)
class TimelineSnapshot:
"""One answered timeline query.
``ok`` is False when the query could not be answered as asked currently
when a requested filter dimension no surviving source can carry was
supplied. The page is then empty *and* the snapshot says so, because an
``ok`` empty page is a claim that no such activity exists.
"""
schema_version: int
remote: str
org: str
repo: str
filters: dict[str, Any]
page: TimelinePage
sources: tuple[SourceStatus, ...]
ok: bool = True
error: dict[str, Any] | None = None
def to_dict(self) -> dict[str, Any]:
return {
"ok": self.ok,
"error": self.error,
"schema_version": self.schema_version,
"scope": {"remote": self.remote, "org": self.org, "repo": self.repo},
"filters": self.filters,
"sources": [s.to_dict() for s in self.sources],
**self.page.to_dict(),
}
def _unanswerable_reasons(
statuses: Iterable[SourceStatus], unanswerable: Iterable[str]
) -> list[dict[str, str]]:
"""Explain, per source, why each unanswerable dimension went unanswered."""
out: list[dict[str, str]] = []
for status in statuses:
for dim in unanswerable:
if dim not in status.supported_filters:
reason = _SOURCE_FILTER_LIMITS.get(
(status.name, dim), f"this source's records carry no {dim} identity"
)
elif not status.ok:
reason = (
f"this source can carry {dim} but did not run: "
f"{status.reason or 'unavailable'}"
)
else:
continue
out.append({"source": status.name, "filter": dim, "reason": reason})
return out
def load_timeline(
*,
remote: str,
org: str,
repo: str,
issue: int | None = None,
pr: int | None = None,
session: str | None = None,
limit: int | None = None,
offset: int | None = None,
db_path: str | None = None,
comment_source: CommentSource | None = None,
) -> TimelineSnapshot:
"""Aggregate every timeline source into one filtered, paginated snapshot.
Sources are read independently and fail soft: an unavailable source
contributes a ``SourceStatus`` with ``ok=False`` and a reason, and never
collapses the whole timeline. The handoff source only runs when a specific
issue or PR is requested (a handoff comment belongs to one thread) and a
``comment_source`` is available; otherwise it is reported as ``not run``
rather than as an empty-and-healthy source.
A filter dimension that no surviving source can carry a ``session``
filter when the only source that ran is the control plane, whose events
record no session is refused with ``ok=False`` and a structured error
instead of being answered with an empty page.
"""
all_events: list[WorkflowEvent] = []
statuses: list[SourceStatus] = []
cp_events, cp_status = read_cp_events(remote=remote, org=org, repo=repo, db_path=db_path)
all_events.extend(cp_events)
statuses.append(cp_status)
# Gitea handoff comments are thread-scoped: only fetch when the caller
# narrowed to one issue or PR, and only when a source was provided.
handoff_target: tuple[str, int] | None = None
if pr is not None:
handoff_target = ("pr", pr)
elif issue is not None:
handoff_target = ("issue", issue)
if handoff_target is None:
statuses.append(
_handoff_status(
ok=False,
reason="not run: handoff comments are thread-scoped; filter by issue or pr to include them",
)
)
elif comment_source is None:
statuses.append(
_handoff_status(
ok=False,
reason="not run: no comment source configured for this timeline read",
)
)
else:
kind, number = handoff_target
try:
comments = comment_source(kind, number) or []
handoff_events = adapt_cth_comments(comments, kind=kind, number=number)
all_events.extend(handoff_events)
statuses.append(_handoff_status(ok=True, count=len(handoff_events)))
except Exception as exc: # fail soft: a fetch/parse error degrades this source only
statuses.append(_handoff_status(ok=False, reason=f"handoff source failed: {exc}"))
requested = tuple(
name
for name, value in ((FILTER_ISSUE, issue), (FILTER_PR, pr), (FILTER_SESSION, session))
if value is not None
)
statuses = [
replace(
status,
unsupported_filters=tuple(
dim for dim in requested if dim not in status.supported_filters
),
)
for status in statuses
]
filters = {"issue": issue, "pr": pr, "session": session}
# A dimension is answerable only if a source that actually ran can carry it.
# If none can, refuse: an empty page would assert "no such activity", which
# is a claim this timeline is not in a position to make.
answerable: set[str] = set()
for status in statuses:
if status.ok:
answerable.update(status.supported_filters)
unanswerable = tuple(dim for dim in requested if dim not in answerable)
if unanswerable:
return TimelineSnapshot(
schema_version=TIMELINE_SCHEMA_VERSION,
remote=remote,
org=org,
repo=repo,
filters=filters,
page=paginate([], limit=limit, offset=offset),
sources=tuple(statuses),
ok=False,
error={
"code": "filter_not_supported",
"unsupported_filters": list(unanswerable),
"detail": (
"no timeline source that ran can answer "
+ ", ".join(f"'{dim}'" for dim in unanswerable)
+ "; the result is refused rather than returned empty"
),
"sources": _unanswerable_reasons(statuses, unanswerable),
},
)
filtered = filter_events(all_events, issue=issue, pr=pr, session=session)
ordered = sort_events(filtered)
page = paginate(ordered, limit=limit, offset=offset)
return TimelineSnapshot(
schema_version=TIMELINE_SCHEMA_VERSION,
remote=remote,
org=org,
repo=repo,
filters=filters,
page=page,
sources=tuple(statuses),
)
def snapshot_to_dict(snapshot: TimelineSnapshot) -> dict[str, Any]:
return snapshot.to_dict()