diff --git a/branch_cleanup_guard.py b/branch_cleanup_guard.py index e91ea02..584d34a 100644 --- a/branch_cleanup_guard.py +++ b/branch_cleanup_guard.py @@ -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: diff --git a/control_plane_db.py b/control_plane_db.py index 4616cb3..d7ad05b 100644 --- a/control_plane_db.py +++ b/control_plane_db.py @@ -31,7 +31,7 @@ from typing import Any, Iterator, Sequence import dependency_graph -SCHEMA_VERSION = 4 +SCHEMA_VERSION = 5 # Assignable work kinds only — raw monitoring incidents are never work items. WORK_KINDS = frozenset({"issue", "pr"}) @@ -186,6 +186,34 @@ CREATE INDEX IF NOT EXISTS idx_dependency_edges_target ON dependency_edges(remote, org, repo, target_kind, target_number); CREATE INDEX IF NOT EXISTS idx_assignments_session ON assignments(session_id, status); CREATE INDEX IF NOT EXISTS idx_incident_gitea ON incident_links(gitea_org, gitea_repo, gitea_issue_number); + +-- Model usage, token cost, latency, and performance events (#651) +CREATE TABLE IF NOT EXISTS usage_events ( + usage_id INTEGER PRIMARY KEY AUTOINCREMENT, + session_id TEXT, + remote TEXT NOT NULL DEFAULT 'dadeschools', + org TEXT NOT NULL DEFAULT '', + repo TEXT NOT NULL DEFAULT '', + project_id TEXT, + role TEXT NOT NULL DEFAULT 'unknown', + model TEXT NOT NULL DEFAULT 'unknown', + issue_number INTEGER, + pr_number INTEGER, + stage TEXT NOT NULL DEFAULT 'unknown', + input_tokens INTEGER, + output_tokens INTEGER, + total_tokens INTEGER, + estimated_cost_usd REAL, + latency_ms INTEGER, + duration_ms INTEGER, + status TEXT NOT NULL DEFAULT 'success', + metadata TEXT, + created_at TEXT NOT NULL +); + +CREATE INDEX IF NOT EXISTS idx_usage_events_scope ON usage_events(remote, org, repo); +CREATE INDEX IF NOT EXISTS idx_usage_events_role_model ON usage_events(role, model); +CREATE INDEX IF NOT EXISTS idx_usage_events_stage ON usage_events(stage); """ @@ -339,6 +367,7 @@ class ControlPlaneDB: self._migrate_incident_links_null_scope(conn) self._migrate_lease_lifecycle_columns(conn) self._migrate_session_ownership_columns(conn) + self._migrate_usage_events_table(conn) conn.execute( "INSERT OR REPLACE INTO schema_meta(key, value) VALUES (?, ?)", ("schema_version", str(SCHEMA_VERSION)), @@ -518,6 +547,207 @@ class ControlPlaneDB: f"UPDATE incident_links SET {col} = '' WHERE {col} IS NULL" ) + def _migrate_usage_events_table(self, conn: sqlite3.Connection) -> None: + """Create usage_events table and indexes if they do not exist (#651).""" + conn.execute(""" + CREATE TABLE IF NOT EXISTS usage_events ( + usage_id INTEGER PRIMARY KEY AUTOINCREMENT, + session_id TEXT, + remote TEXT NOT NULL DEFAULT 'dadeschools', + org TEXT NOT NULL DEFAULT '', + repo TEXT NOT NULL DEFAULT '', + project_id TEXT, + role TEXT NOT NULL DEFAULT 'unknown', + model TEXT NOT NULL DEFAULT 'unknown', + issue_number INTEGER, + pr_number INTEGER, + stage TEXT NOT NULL DEFAULT 'unknown', + input_tokens INTEGER, + output_tokens INTEGER, + total_tokens INTEGER, + estimated_cost_usd REAL, + latency_ms INTEGER, + duration_ms INTEGER, + status TEXT NOT NULL DEFAULT 'success', + metadata TEXT, + created_at TEXT NOT NULL + ); + """) + conn.execute("CREATE INDEX IF NOT EXISTS idx_usage_events_scope ON usage_events(remote, org, repo);") + conn.execute("CREATE INDEX IF NOT EXISTS idx_usage_events_role_model ON usage_events(role, model);") + conn.execute("CREATE INDEX IF NOT EXISTS idx_usage_events_stage ON usage_events(stage);") + + # #651 retention: cap growth so unauthenticated or high-volume ingest + # cannot DoS the control-plane DB (PR #876 F3). Applied after every write. + USAGE_EVENTS_MAX_ROWS = 10_000 + USAGE_EVENTS_RETENTION_DAYS = 90 + + def record_usage_event( + self, + *, + session_id: str | None = None, + remote: str = "dadeschools", + org: str = "", + repo: str = "", + project_id: str | None = None, + role: str = "unknown", + model: str = "unknown", + issue_number: int | None = None, + pr_number: int | None = None, + stage: str = "unknown", + input_tokens: int | None = None, + output_tokens: int | None = None, + total_tokens: int | None = None, + estimated_cost_usd: float | None = None, + latency_ms: int | None = None, + duration_ms: int | None = None, + status: str = "success", + metadata: str | dict[str, Any] | None = None, + created_at: str | None = None, + ) -> int: + """Record a model usage, token cost, latency, or stage performance event (#651).""" + ts = created_at or _ts() + meta_str: str | None = None + if metadata is not None: + from webui import console_redaction + redacted_meta = console_redaction.redact_payload(metadata) + if isinstance(redacted_meta, str): + meta_str = redacted_meta + else: + try: + meta_str = json.dumps(redacted_meta, default=str) + except Exception: + meta_str = str(redacted_meta) + + if total_tokens is None and (input_tokens is not None or output_tokens is not None): + total_tokens = (input_tokens or 0) + (output_tokens or 0) + + with self._tx(immediate=True) as conn: + cursor = conn.execute( + """ + INSERT INTO usage_events ( + session_id, remote, org, repo, project_id, role, model, + issue_number, pr_number, stage, input_tokens, output_tokens, + total_tokens, estimated_cost_usd, latency_ms, duration_ms, + status, metadata, created_at + ) VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?) + """, + ( + session_id, + remote, + org, + repo, + project_id, + role, + model, + issue_number, + pr_number, + stage, + input_tokens, + output_tokens, + total_tokens, + estimated_cost_usd, + latency_ms, + duration_ms, + status, + meta_str, + ts, + ), + ) + usage_id = cursor.lastrowid + self._enforce_usage_events_retention(conn) + return usage_id + + def _enforce_usage_events_retention(self, conn: sqlite3.Connection) -> None: + """Drop aged and excess usage_events rows (PR #876 F3).""" + # Age-based: ISO-8601 UTC timestamps compare lexicographically. + cutoff = ( + datetime.now(timezone.utc) + - timedelta(days=int(self.USAGE_EVENTS_RETENTION_DAYS)) + ).strftime("%Y-%m-%dT%H:%M:%SZ") + conn.execute( + "DELETE FROM usage_events WHERE created_at < ?", + (cutoff,), + ) + # Count-based: keep the newest USAGE_EVENTS_MAX_ROWS by usage_id. + max_rows = int(self.USAGE_EVENTS_MAX_ROWS) + if max_rows > 0: + conn.execute( + """ + DELETE FROM usage_events + WHERE usage_id NOT IN ( + SELECT usage_id FROM usage_events + ORDER BY usage_id DESC + LIMIT ? + ) + """, + (max_rows,), + ) + + def query_usage_events( + self, + *, + remote: str | None = None, + org: str | None = None, + repo: str | None = None, + project_id: str | None = None, + role: str | None = None, + model: str | None = None, + issue_number: int | None = None, + pr_number: int | None = None, + stage: str | None = None, + session_id: str | None = None, + limit: int = 500, + offset: int = 0, + ) -> list[dict[str, Any]]: + """Query stored usage events matching filters (#651).""" + conditions = [] + params = [] + if remote: + conditions.append("remote = ?") + params.append(remote) + if org: + conditions.append("org = ?") + params.append(org) + if repo: + conditions.append("repo = ?") + params.append(repo) + if project_id: + conditions.append("project_id = ?") + params.append(project_id) + if role: + conditions.append("role = ?") + params.append(role) + if model: + conditions.append("model = ?") + params.append(model) + if issue_number is not None: + conditions.append("issue_number = ?") + params.append(issue_number) + if pr_number is not None: + conditions.append("pr_number = ?") + params.append(pr_number) + if stage: + conditions.append("stage = ?") + params.append(stage) + if session_id: + conditions.append("session_id = ?") + params.append(session_id) + + where_clause = f"WHERE {' AND '.join(conditions)}" if conditions else "" + sql = f""" + SELECT * FROM usage_events + {where_clause} + ORDER BY usage_id ASC + LIMIT ? OFFSET ? + """ + params.extend([limit, offset]) + + with self._tx(immediate=False) as conn: + cursor = conn.execute(sql, params) + rows = cursor.fetchall() + return [dict(row) for row in rows] + # ── sessions ────────────────────────────────────────────────────────── def upsert_session( diff --git a/docs/mcp-tool-inventory.md b/docs/mcp-tool-inventory.md index 47b8021..66094ec 100644 --- a/docs/mcp-tool-inventory.md +++ b/docs/mcp-tool-inventory.md @@ -101,6 +101,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` diff --git a/docs/observability/analytics-instrumentation.md b/docs/observability/analytics-instrumentation.md new file mode 100644 index 0000000..55d7057 --- /dev/null +++ b/docs/observability/analytics-instrumentation.md @@ -0,0 +1,143 @@ +# Model Usage, Token Cost, Latency, and Workflow Analytics (Phase 4) + +- **Tracking Issue:** [#651](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/651) +- **Parent Epic:** [#631](https://gitea.prgs.cc/Scaled-Tech-Consulting/Gitea-Tools/issues/651) +- **Console Surface:** `/analytics`, `/api/v1/analytics`, `/api/v1/analytics/usage` + +## 1. Overview + +The Web Console Analytics module provides durable, aggregate visibility into **model usage, token cost, latency percentiles, and workflow-stage performance** across projects, worker roles, AI models, issues, and PRs. + +### Non-Goals +- No mandatory client-side telemetry that leaks prompts or secret keys. +- No third-party payment provider or billing integration. +- No automatic model routing changes without controller policy (#647). + +--- + +## 2. Event Schema (`usage_events`) + +Usage metrics are stored in the control-plane database under table `usage_events`. + +| Column | Type | Description | +|---|---|---| +| `usage_id` | `INTEGER` | Primary key (autoincrement) | +| `session_id` | `TEXT` | Optional active session identifier | +| `remote` | `TEXT` | Known Gitea instance (`dadeschools` or `prgs`) | +| `org` | `TEXT` | Repository owner / organization | +| `repo` | `TEXT` | Repository name | +| `project_id` | `TEXT` | Optional project identifier | +| `role` | `TEXT` | Active worker role (`author`, `reviewer`, `merger`, `reconciler`, `controller`) | +| `model` | `TEXT` | LLM model identifier (e.g. `gemini-3.6-flash`, `claude-3-5-sonnet`) | +| `issue_number` | `INTEGER` | Correlated Gitea issue number (optional) | +| `pr_number` | `INTEGER` | Correlated Gitea PR number (optional) | +| `stage` | `TEXT` | Workflow stage (`preflight`, `implementation`, `review`, `merge`, `reconciliation`) | +| `input_tokens` | `INTEGER` | Input token count (optional / nullable) | +| `output_tokens` | `INTEGER` | Output token count (optional / nullable) | +| `total_tokens` | `INTEGER` | Total token count (optional / nullable) | +| `estimated_cost_usd` | `REAL` | Estimated USD cost (optional / nullable) | +| `latency_ms` | `INTEGER` | Request latency in milliseconds (optional / nullable) | +| `duration_ms` | `INTEGER` | Stage execution duration in milliseconds (optional / nullable) | +| `status` | `TEXT` | Outcome status (`success`, `failure`, `timeout`) | +| `metadata` | `TEXT` | Redacted metadata or summary string | +| `created_at` | `TEXT` | ISO 8601 UTC timestamp | + +--- + +## 3. Handling of Missing Data ("Unknown" vs. Zero Fabrication) + +To ensure operational metrics accurately reflect evidence: +- **Untracked or missing metrics are displayed as `Unknown`**, never zero-fabricated. +- If an event omits `estimated_cost_usd`, `latency_ms`, or token counts, the aggregator marks those fields as missing (`None`) rather than defaulting to `0` or `$0.00`. +- Summary tables and KPI cards explicitly indicate when data is unmeasured or partially reported. + +--- + +## 4. Redaction & Security Rules + +Per `#633` security policy: +- Free-text fields (`metadata`, `prompt_summary`, `session_id`) are run through `console_redaction.redact_text` before persistence and output serialization. +- Secret tokens, keychain commands, authorization headers, passwords, and JWTs are stripped automatically. + +--- + +## 5. Opt-in Instrumentation Guide + +Applications, MCP servers, and background sessions can report usage metrics through either Python API or HTTP ingestion. + +### Python Ingestion + +```python +from webui.analytics_loader import record_usage + +record_usage( + remote="dadeschools", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + role="author", + model="gemini-3.6-flash", + issue_number=651, + stage="implementation", + input_tokens=1420, + output_tokens=380, + total_tokens=1800, + estimated_cost_usd=0.00045, + latency_ms=320, + duration_ms=4500, + status="success", + metadata={"note": "Implementation of analytics module"}, +) +``` + +### HTTP Ingestion API (authorized write) + +`POST /api/v1/analytics/usage` is a **gated write**. It runs through +`console_authz` action `record_analytics_usage` (operator+, Phase 2 execution). +Unauthenticated or phase-inactive requests receive **403** and do not write. +Prefer in-process `record_usage` for MCP / session instrumentation. + +```http +POST /api/v1/analytics/usage HTTP/1.1 +Content-Type: application/json +# Requires authenticated principal with record_analytics_usage execution enabled + +{ + "remote": "dadeschools", + "org": "Scaled-Tech-Consulting", + "repo": "Gitea-Tools", + "role": "author", + "model": "gemini-3.6-flash", + "issue_number": 651, + "stage": "implementation", + "input_tokens": 1420, + "output_tokens": 380, + "total_tokens": 1800, + "estimated_cost_usd": 0.00045, + "latency_ms": 320, + "duration_ms": 4500, + "status": "success", + "metadata": "Analytics schema landed" +} +``` + +### Retention + +`usage_events` is retained with hard caps applied on every write: + +| Limit | Default | +|---|---| +| Max rows | 10,000 (`ControlPlaneDB.USAGE_EVENTS_MAX_ROWS`) | +| Max age | 90 days (`ControlPlaneDB.USAGE_EVENTS_RETENTION_DAYS`) | + +Older rows (by `created_at`) and excess oldest rows (by `usage_id`) are deleted +after each insert so unbounded growth / DoS-by-volume cannot fill the DB. + +--- + +## 6. Querying Analytics API + +```http +GET /api/v1/analytics?role=author&stage=implementation HTTP/1.1 +``` + +Returns `AnalyticsSnapshot` JSON containing aggregations (`by_model`, `by_stage`, `by_role`, `by_work_item`, `by_project`) and latency percentiles (`p50`, `p90`, `p95`, `p99`). diff --git a/docs/webui-authz-audit.md b/docs/webui-authz-audit.md index 9e010e9..2dcffac 100644 --- a/docs/webui-authz-audit.md +++ b/docs/webui-authz-audit.md @@ -91,6 +91,7 @@ already define, and a regression test asserts each mapping matches. | `close_pr` | controller | privileged | `gitea.pr.close` | Yes | No | No | 3 | | `merge_pr` | controller | privileged | `gitea.pr.merge` | Yes | **Yes** | **Yes** | 3 | | `delete_branch` | admin | destructive | `gitea.branch.delete` | Yes | **Yes** | **Yes** | 3 | +| `record_analytics_usage` | operator | gated_write | `runtime.record_analytics_usage` | Yes | No | No | 2 | | `system.reload_namespace` | controller | privileged | `runtime.reload_namespace` | Yes | No | No | 2 | | `system.restart_namespace` | admin | destructive | `runtime.restart_namespace` | Yes | **Yes** | **Yes** | 2 | diff --git a/gitea_mcp_server.py b/gitea_mcp_server.py index c8a9461..5fb0988 100644 --- a/gitea_mcp_server.py +++ b/gitea_mcp_server.py @@ -2065,6 +2065,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 restart_coordinator # noqa: E402 # #658 MCP restart coordinator/impact import incident_bridge # noqa: E402 @@ -2310,7 +2311,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, @@ -2625,7 +2625,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, @@ -2636,6 +2641,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, } @@ -4406,6 +4420,135 @@ def gitea_lock_issue( @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): + 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_recover_dirty_orphaned_issue_worktree( issue_number: int, @@ -4679,6 +4822,7 @@ def gitea_recover_dirty_orphaned_issue_worktree( ) return result + @mcp.tool() def gitea_rebind_dirty_same_claimant_author_session( issue_number: int, @@ -4935,8 +5079,6 @@ def gitea_rebind_dirty_same_claimant_author_session( } return result - return result - @mcp.tool() def gitea_assess_work_issue_duplicate( diff --git a/issue_lock_store.py b/issue_lock_store.py index d011c1c..f216656 100644 --- a/issue_lock_store.py +++ b/issue_lock_store.py @@ -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._+-]+") @@ -257,6 +269,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)) @@ -340,6 +677,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 @@ -360,26 +707,113 @@ 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: @@ -397,56 +831,103 @@ def assess_lock_freshness( pid_int = None pid_alive = is_process_alive(pid_int) if pid_int 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, - "pid_missing": pid_missing, - } + lifecycle = lease_lifecycle_version(lock_data) + legacy = lifecycle != lease_policy.LIFECYCLE_HEARTBEAT_V1 + policy = lease_policy.policy_for(lease_task_class(lock_data)) - # #860: a PID-less lock must never be considered live merely because - # expiration / heartbeat fields are absent. Missing PID is insufficient - # evidence of a live owner; treat as malformed/stale so recovery routes - # can evaluate corroborating pins instead of blocking on a false live flag. - if pid_missing: - return { - "status": "malformed", - "live": False, - "stale": True, - "reason": ( - "lock has no usable session pid; cannot prove live ownership " - "(PID-less locks are never live by missing expiry alone)" - ), - "pid_alive": False, - "pid_missing": True, - "heartbeat_at": heartbeat_at.isoformat() if heartbeat_at else None, - "expires_at": expires_at.isoformat() if expires_at else None, - } - - if pid_int is not None and not pid_alive: - return { - "status": "stale", - "live": False, - "stale": True, - "reason": f"owner pid {pid_int} is not alive", - "pid_alive": False, - "pid_missing": False, - } - - return { - "status": "live", - "live": True, - "stale": False, - "reason": "lock heartbeat and lease are fresh", + evidence: dict[str, Any] = { "pid_alive": pid_alive, - "pid_missing": False, + "pid_missing": pid_missing, + "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: @@ -483,6 +964,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") @@ -557,7 +1059,28 @@ def assess_same_issue_lease_conflict( ) if recovery_sanctioned and existing_issue == issue_number and existing_branch == branch_name: return None - if is_lease_expired(existing_lock, now=now): + expired = is_lease_expired(existing_lock, now=now) + # #790 review #502/#516: a heartbeat-lifecycle lease that is non-live but + # whose absolute expires_at is still in the future must enter the same + # reclaim/renewal disposition as an expired lease. This happens for + # stale_absolute_cap (a session that keeps heartbeating past the 8h cap has + # expires_at = last_heartbeat + TTL in the future) and for + # stale_missed_heartbeat under an independent TTL>grace policy. Keying this + # gate on is_lease_expired alone left assess_expired_lock_reclaim — which + # already permits exactly those bands — unreachable from the acquisition + # path, so the load-bearing heartbeat was not load-bearing for foreign + # reclaim: the precise abandonment scenario #790 exists to fix. + # assess_foreign_lock_overwrite already keys on is_lease_live; this makes the + # same-issue path consistent with it. Legacy leases keep their absolute-expiry + # clock (AC-N8): is_lease_live already decides a legacy lease from expires_at + # / dead-PID alone and is_legacy_lease excludes it here, so this widening is a + # no-op for every pre-lifecycle lock. + heartbeat_non_live_reclaimable = ( + not expired + and not is_legacy_lease(existing_lock) + and not is_lease_live(existing_lock, now=now) + ) + if expired or heartbeat_non_live_reclaimable: # #760 AC1/AC2: exact-owner renewal is a different disposition from # foreign takeover and is evaluated first. Before this, both branches # below returned unconditionally, so the same_owner allowance further @@ -571,10 +1094,12 @@ def assess_same_issue_lease_conflict( reclaim = assess_expired_lock_reclaim(existing_lock, now=now) if reclaim.get("reclaim_allowed"): # #601: expired + dead pid / missing worktree may be reclaimed - # through the normal lock path (sanctioned overwrite). + # through the normal lock path (sanctioned overwrite). #790: the new + # stale heartbeat bands are reclaimable here on the same evidence. return None + descriptor = "expired" if expired else "non-live" return ( - f"Issue #{issue_number} has an expired {operation_type} lease on " + f"Issue #{issue_number} has an {descriptor} {operation_type} lease on " f"branch '{existing_branch}' from worktree '{existing_worktree}'. " "Recovery review is required before takeover (fail closed)" ) diff --git a/lease_policy.py b/lease_policy.py new file mode 100644 index 0000000..e778cd3 --- /dev/null +++ b/lease_policy.py @@ -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, + } diff --git a/task_capability_map.py b/task_capability_map.py index 9875a2d..878cf0a 100644 --- a/task_capability_map.py +++ b/task_capability_map.py @@ -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", + }, # #860: dirty orphaned same-claimant worktree recovery (explicit operation). "recover_dirty_orphaned_issue_worktree": { "permission": "gitea.issue.comment", @@ -520,6 +529,15 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = { "permission": "gitea.issue.comment", "role": "author", }, + + # #651 console analytics ingest — control-plane DB write, not a Gitea API + # call. Authority comes from console RBAC (operator+) plus phase gating; + # permission string is a non-Gitea runtime capability so no Gitea profile + # can satisfy it by accident. + "record_analytics_usage": { + "permission": "runtime.record_analytics_usage", + "role": "author", + }, } diff --git a/tests/test_control_plane_db.py b/tests/test_control_plane_db.py index 8f251eb..1752299 100644 --- a/tests/test_control_plane_db.py +++ b/tests/test_control_plane_db.py @@ -36,7 +36,7 @@ class ControlPlaneDBTest(unittest.TestCase): rows = dict(conn.execute("SELECT key, value FROM schema_meta").fetchall()) finally: conn.close() - self.assertEqual(rows["schema_version"], "4") + self.assertEqual(rows["schema_version"], "5") self.assertIn("DB coordinates", rows["architecture"]) self.assertIn("bridge", rows["architecture"].lower()) diff --git a/tests/test_issue_628_orchestration.py b/tests/test_issue_628_orchestration.py new file mode 100644 index 0000000..91e9437 --- /dev/null +++ b/tests/test_issue_628_orchestration.py @@ -0,0 +1,199 @@ +"""Regression tests for #628 building blocks (child scope only). + +Honest scope: unit coverage of pre-existing APIs used by umbrella #628. +This module does **not** implement or verify all 21 umbrella acceptance +criteria, automatic handoff store/retrieve, multi-worker product wiring, +or end-to-end orchestration. + +Covered building blocks: +- CTH format / parse / assess (`format_cth_body`, `parse_cth_comment`, + `assess_cth_comment`) +- Exclusive-ownership skip classification (`classify_skip` with + OWNERSHIP_FOREIGN vs OWNERSHIP_OWN) +- Durable dependency edges (`upsert_dependency_edge` / list) and skip + when dependency_unmet +- Edge state transition UNMET -> MET + +Parent umbrella remains #628; this slice is a scoped child issue only. +""" + +import unittest +from unittest.mock import MagicMock, patch +import os +import json +import tempfile + +from canonical_thread_handoff import ( + format_cth_body, + parse_cth_comment, + assess_cth_comment, +) +import dependency_graph +from control_plane_db import ControlPlaneDB +from allocator_service import ( + WorkCandidate, + classify_skip, + ROLE_AUTHOR, + ROLE_REVIEWER, + ROLE_MERGER, + ROLE_RECONCILER, + OWNERSHIP_OWN, + OWNERSHIP_FOREIGN, +) + + +class TestIssue628Orchestration(unittest.TestCase): + def setUp(self): + self._tmp = tempfile.TemporaryDirectory() + self.db_path = os.path.join(self._tmp.name, "cp.sqlite3") + self.db = ControlPlaneDB(self.db_path) + + def tearDown(self): + self._tmp.cleanup() + + def test_canonical_handoff_serialization_and_retrieval(self): + """Building block: format/parse/assess a CTH body (not full AC1/AC2 product path).""" + handoff = format_cth_body( + cth_type="Author Handoff", + status="completed", + next_owner="reviewer", + current_blocker="none", + decision="Implementation complete, tests passing", + proof="pytest tests/test_issue_628_orchestration.py passed", + next_action="Review PR and run reviewer pre-flight", + ready_to_paste_prompt="Review PR for child issue #878 (parent #628)", + ) + self.assertIn("CTH: Author Handoff", handoff) + + parsed = parse_cth_comment(handoff) + self.assertIsNotNone(parsed) + self.assertEqual(parsed["cth_type"], "Author Handoff") + + assessment = assess_cth_comment(handoff) + self.assertFalse(assessment["block"]) + + def test_exclusive_task_unit_single_owner(self): + """Building block: classify_skip foreign vs own ownership.""" + candidate = WorkCandidate( + kind="issue", + number=878, + title="Child #878 ownership classify candidate", + state="open", + labels=["status:in-progress"], + blocked=False, + dependency_unmet=False, + ) + # Foreign ownership MUST be skipped + skip_foreign = classify_skip( + c=candidate, + role=ROLE_AUTHOR, + terminal_pr=None, + claim_ownership=OWNERSHIP_FOREIGN, + ) + self.assertIsNotNone(skip_foreign) + self.assertIn("active lease", skip_foreign) + + # Own/Self claim remains selectable for session resumption + skip_self = classify_skip( + c=candidate, + role=ROLE_AUTHOR, + terminal_pr=None, + claim_ownership=OWNERSHIP_OWN, + ) + self.assertIsNone(skip_self) + + def test_durable_dependency_graph_blocking(self): + """Building block: unmet dependency edge + classify_skip on dependency_unmet.""" + # Upsert a blocking dependency edge between issue 878 and blocker 601 + self.db.upsert_dependency_edge( + remote="prgs", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + source_kind="issue", + source_number=878, + target_kind="issue", + target_number=601, + edge_type=dependency_graph.EDGE_ISSUE_BLOCKED_BY_ISSUE, + state=dependency_graph.STATE_UNMET, + blocking_condition="Target issue #601 is not closed", + completion_condition="Target issue #601 is closed", + evidence={"source": "unit_test"}, + ) + + edges = self.db.list_dependency_edges( + remote="prgs", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + source_kind="issue", + source_number=878, + ) + self.assertEqual(len(edges), 1) + self.assertEqual(edges[0]["state"], "unmet") + self.assertEqual(edges[0]["target_number"], 601) + + # When dependency is unmet, candidate is blocked from selection + candidate = WorkCandidate( + kind="issue", + number=878, + title="Blocked candidate", + state="open", + labels=[], + blocked=False, + dependency_unmet=True, + dependency_reason="issue#878 is blocked by unmet dependency issue#601", + ) + skip_reason = classify_skip( + c=candidate, + role=ROLE_AUTHOR, + terminal_pr=None, + claim_ownership=OWNERSHIP_OWN, + ) + self.assertIsNotNone(skip_reason) + self.assertIn("issue#601", skip_reason) + + def test_dependency_completion_reevaluation(self): + """Building block: dependency edge state can transition UNMET -> MET.""" + self.db.upsert_dependency_edge( + remote="prgs", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + source_kind="issue", + source_number=878, + target_kind="issue", + target_number=601, + edge_type=dependency_graph.EDGE_ISSUE_BLOCKED_BY_ISSUE, + state=dependency_graph.STATE_UNMET, + blocking_condition="Target issue #601 is open", + completion_condition="Target issue #601 is closed", + evidence={"source": "unit_test"}, + ) + + # Mark edge as met upon target issue closure + self.db.upsert_dependency_edge( + remote="prgs", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + source_kind="issue", + source_number=878, + target_kind="issue", + target_number=601, + edge_type=dependency_graph.EDGE_ISSUE_BLOCKED_BY_ISSUE, + state=dependency_graph.STATE_MET, + blocking_condition="Target issue #601 is open", + completion_condition="Target issue #601 is closed", + evidence={"source": "target_closed_event"}, + ) + + edges = self.db.list_dependency_edges( + remote="prgs", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + source_kind="issue", + source_number=878, + ) + self.assertEqual(len(edges), 1) + self.assertEqual(edges[0]["state"], "met") + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_issue_790_heartbeat_mcp_path.py b/tests/test_issue_790_heartbeat_mcp_path.py new file mode 100644 index 0000000..57e0d5b --- /dev/null +++ b/tests/test_issue_790_heartbeat_mcp_path.py @@ -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", "test@example.com") + 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() diff --git a/tests/test_issue_790_lease_policy.py b/tests/test_issue_790_lease_policy.py new file mode 100644 index 0000000..34d52d6 --- /dev/null +++ b/tests/test_issue_790_lease_policy.py @@ -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() diff --git a/tests/test_issue_794_conflict_gate_reclaim.py b/tests/test_issue_794_conflict_gate_reclaim.py new file mode 100644 index 0000000..16f000c --- /dev/null +++ b/tests/test_issue_794_conflict_gate_reclaim.py @@ -0,0 +1,193 @@ +"""Conflict gate honors heartbeat-lifecycle non-live reclaim bands (#790 review #502/#516). + +``assess_expired_lock_reclaim`` already permits reclaim for the heartbeat-lifecycle +stale bands (``stale_missed_heartbeat``, ``stale_absolute_cap``) without a dead PID. +But ``assess_same_issue_lease_conflict`` used to enter that reclaim branch only under +``is_lease_expired`` (``expires_at <= now``). For a heartbeat-lifecycle lease that is +non-live yet whose ``expires_at`` is still in the future, the acquisition gate fell +through to the foreign "already has an active lease" block and never consulted the +reclaim assessor — so the load-bearing heartbeat was not load-bearing for foreign +reclaim, the exact abandonment scenario #790 exists to fix. + +Two future-``expires_at`` non-live shapes are reachable: + +* ``stale_absolute_cap`` — a session that keeps heartbeating past the 8h absolute cap + has ``expires_at = last_heartbeat + TTL`` in the future (default policy). +* ``stale_missed_heartbeat`` — under an independent TTL>grace policy the heartbeat + grace lapses while ``expires_at`` is still ahead. + +These tests pin: both reclaim from a foreign acquirer; a live lease still blocks a +foreign acquirer; the same owner may reclaim its own abandoned heartbeat lease; and a +legacy (pre-lifecycle) lock keeps its absolute-``expires_at`` clock — a non-expired +legacy lock with a dead PID is *not* reclaimable through this path. +""" + +from __future__ import annotations + +import os +import sys +import unittest +from datetime import datetime, timedelta, timezone +from unittest import mock + +sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) + +import issue_lock_store as ils # noqa: E402 +import lease_policy # noqa: E402 + +ISSUE = 790 +OWNER_BRANCH = f"fix/issue-{ISSUE}-slice-a-heartbeat-policy" +OWNER_WORKTREE = "/tmp/wt-790-owner" +FOREIGN_BRANCH = f"fix/issue-{ISSUE}-foreign-attempt" +FOREIGN_WORKTREE = "/tmp/wt-790-foreign" + + +def _ts(moment: datetime) -> str: + return ( + moment.astimezone(timezone.utc) + .replace(microsecond=0) + .isoformat() + .replace("+00:00", "Z") + ) + + +class _ConflictGateBase(unittest.TestCase): + def setUp(self): + self.now = datetime(2026, 7, 23, 12, 0, 0, tzinfo=timezone.utc) + + def _lock( + self, + *, + lifecycle: str | None, + created_ago: timedelta, + heartbeat_ago: timedelta, + expires_in: timedelta, + pid: int, + ) -> dict: + lease: dict = { + "operation_type": ils.AUTHOR_ISSUE_WORK_LEASE, + "issue_number": ISSUE, + "branch": OWNER_BRANCH, + "worktree_path": OWNER_WORKTREE, + "created_at": _ts(self.now - created_ago), + "last_heartbeat_at": _ts(self.now - heartbeat_ago), + "expires_at": _ts(self.now + expires_in), + } + if lifecycle is not None: + lease["lifecycle_version"] = lifecycle + lease["task_session_id"] = "author_issue_work-deadbeefdeadbeef" + return { + "issue_number": ISSUE, + "branch_name": OWNER_BRANCH, + "remote": "prgs", + "org": "Scaled-Tech-Consulting", + "repo": "Gitea-Tools", + "worktree_path": OWNER_WORKTREE, + "session_pid": pid, + "pid": pid, + "work_lease": lease, + } + + def _foreign_conflict(self, existing: dict) -> str | None: + return ils.assess_same_issue_lease_conflict( + existing, + issue_number=ISSUE, + branch_name=FOREIGN_BRANCH, + worktree_path=FOREIGN_WORKTREE, + now=self.now, + ) + + def _same_owner_conflict(self, existing: dict) -> str | None: + return ils.assess_same_issue_lease_conflict( + existing, + issue_number=ISSUE, + branch_name=OWNER_BRANCH, + worktree_path=OWNER_WORKTREE, + now=self.now, + ) + + +class TestHeartbeatNonLiveFutureExpiresReclaim(_ConflictGateBase): + def test_stale_absolute_cap_future_expires_allows_foreign_reclaim(self): + # Created >8h ago, heartbeated one minute ago, expires 9 min in the FUTURE, + # owner PID alive: freshness = stale_absolute_cap, live=False, not expired. + existing = self._lock( + lifecycle=lease_policy.LIFECYCLE_HEARTBEAT_V1, + created_ago=timedelta(hours=9), + heartbeat_ago=timedelta(minutes=1), + expires_in=timedelta(minutes=9), + pid=os.getpid(), + ) + freshness = ils.assess_lock_freshness(existing, now=self.now) + self.assertEqual(freshness["status"], ils.STATUS_STALE_ABSOLUTE_CAP) + self.assertFalse(freshness["live"]) + self.assertFalse(ils.is_lease_expired(existing, now=self.now)) + with mock.patch.object(ils, "is_process_alive", return_value=True): + self.assertIsNone(self._foreign_conflict(existing)) + + def test_missed_heartbeat_ttl_gt_grace_future_expires_allows_foreign_reclaim(self): + # Heartbeat grace (default 10 min) lapsed 5 min ago, but a TTL>grace policy + # leaves expires_at 20 min in the FUTURE: stale_missed_heartbeat, not expired. + existing = self._lock( + lifecycle=lease_policy.LIFECYCLE_HEARTBEAT_V1, + created_ago=timedelta(minutes=30), + heartbeat_ago=timedelta(minutes=15), + expires_in=timedelta(minutes=20), + pid=os.getpid(), + ) + freshness = ils.assess_lock_freshness(existing, now=self.now) + self.assertEqual(freshness["status"], ils.STATUS_STALE_MISSED_HEARTBEAT) + self.assertFalse(freshness["live"]) + self.assertFalse(ils.is_lease_expired(existing, now=self.now)) + with mock.patch.object(ils, "is_process_alive", return_value=True): + self.assertIsNone(self._foreign_conflict(existing)) + + def test_same_owner_may_reclaim_its_own_abandoned_heartbeat_lease(self): + existing = self._lock( + lifecycle=lease_policy.LIFECYCLE_HEARTBEAT_V1, + created_ago=timedelta(hours=9), + heartbeat_ago=timedelta(minutes=1), + expires_in=timedelta(minutes=9), + pid=os.getpid(), + ) + with mock.patch.object(ils, "is_process_alive", return_value=True): + self.assertIsNone(self._same_owner_conflict(existing)) + + +class TestLiveAndLegacyStillBlockForeign(_ConflictGateBase): + def test_live_heartbeat_lease_still_blocks_foreign(self): + existing = self._lock( + lifecycle=lease_policy.LIFECYCLE_HEARTBEAT_V1, + created_ago=timedelta(minutes=5), + heartbeat_ago=timedelta(minutes=1), + expires_in=timedelta(minutes=9), + pid=os.getpid(), + ) + freshness = ils.assess_lock_freshness(existing, now=self.now) + self.assertTrue(freshness["live"]) + with mock.patch.object(ils, "is_process_alive", return_value=True): + block = self._foreign_conflict(existing) + self.assertIn("already has an active", block or "") + + def test_legacy_lease_future_expires_dead_pid_is_not_reclaimed_here(self): + # AC-N8: a legacy lock keeps its absolute expires_at clock. Non-expired + + # dead PID is non-live, but the widened band excludes legacy, so the + # foreign acquirer is still blocked rather than silently reclaiming. + existing = self._lock( + lifecycle=None, + created_ago=timedelta(hours=9), + heartbeat_ago=timedelta(hours=9), + expires_in=timedelta(hours=2), + pid=4_194_304, # far above any live pid on a test host + ) + with mock.patch.object(ils, "is_process_alive", return_value=False): + freshness = ils.assess_lock_freshness(existing, now=self.now) + self.assertTrue(freshness["legacy_lease"]) + self.assertFalse(freshness["live"]) + self.assertFalse(ils.is_lease_expired(existing, now=self.now)) + block = self._foreign_conflict(existing) + self.assertIn("already has an active", block or "") + + +if __name__ == "__main__": + unittest.main() diff --git a/tests/test_webui_analytics.py b/tests/test_webui_analytics.py new file mode 100644 index 0000000..bf1a419 --- /dev/null +++ b/tests/test_webui_analytics.py @@ -0,0 +1,252 @@ +"""Unit and integration tests for Model Usage & Performance Analytics (#651).""" + +from __future__ import annotations + +import os +import tempfile +import unittest +from starlette.testclient import TestClient + +import control_plane_db +from webui.analytics_loader import ( + ANALYTICS_SCHEMA_VERSION, + compute_percentile, + load_analytics, + record_usage, +) +from webui.app import create_app +from webui import console_redaction + + +class AnalyticsLoaderTest(unittest.TestCase): + + def setUp(self) -> None: + self.temp_dir = tempfile.TemporaryDirectory() + self.db_path = os.path.join(self.temp_dir.name, "test_control_plane.sqlite3") + os.environ["GITEA_CONTROL_PLANE_DB"] = self.db_path + self.db = control_plane_db.ControlPlaneDB(db_path=self.db_path) + + def tearDown(self) -> None: + self.temp_dir.cleanup() + + def test_compute_percentile(self) -> None: + self.assertIsNone(compute_percentile([], 50.0)) + self.assertEqual(compute_percentile([100], 50.0), 100.0) + + # 2 elements: [100, 200] + self.assertEqual(compute_percentile([100, 200], 50.0), 150.0) + + # 100 elements: 1..100 + vals = list(range(1, 101)) + self.assertEqual(compute_percentile(vals, 50.0), 50.5) + self.assertAlmostEqual(compute_percentile(vals, 90.0), 90.1) + + def test_record_and_aggregate_usage(self) -> None: + # Record event 1 (complete data) + u1 = record_usage( + db_path=self.db_path, + remote="dadeschools", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + role="author", + model="gemini-3.6-flash", + issue_number=651, + stage="implementation", + input_tokens=1000, + output_tokens=500, + estimated_cost_usd=0.0015, + latency_ms=200, + duration_ms=3000, + metadata={"secret_key": "secret123", "note": "token=secret123"}, + ) + self.assertGreater(u1, 0) + + # Record event 2 (missing tokens and cost -> unknown) + u2 = record_usage( + db_path=self.db_path, + remote="dadeschools", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + role="reviewer", + model="claude-3-5-sonnet", + pr_number=846, + stage="review", + latency_ms=500, + duration_ms=6000, + ) + self.assertGreater(u2, u1) + + snapshot = load_analytics( + db_path=self.db_path, + remote="dadeschools", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + ) + + self.assertTrue(snapshot.ok) + self.assertEqual(snapshot.schema_version, ANALYTICS_SCHEMA_VERSION) + self.assertEqual(snapshot.total_events, 2) + + # Verify overall summary + summary = snapshot.overall_summary + self.assertEqual(summary.total_events, 2) + self.assertEqual(summary.events_with_tokens, 1) + self.assertEqual(summary.total_tokens, 1500) + self.assertEqual(summary.events_with_cost, 1) + self.assertEqual(summary.estimated_cost_usd, 0.0015) + self.assertEqual(summary.events_with_latency, 2) + self.assertEqual(summary.latency_p50_ms, 350.0) + + # Verify missing data handling (AC 3: not zero-fabricated) + reviewer_model = snapshot.by_model.get("claude-3-5-sonnet") + self.assertIsNotNone(reviewer_model) + self.assertEqual(reviewer_model.total_events, 1) + self.assertEqual(reviewer_model.events_with_tokens, 0) + self.assertIsNone(reviewer_model.total_tokens) + self.assertEqual(reviewer_model.display_tokens, "Unknown") + self.assertEqual(reviewer_model.events_with_cost, 0) + self.assertIsNone(reviewer_model.estimated_cost_usd) + self.assertEqual(reviewer_model.display_cost, "Unknown") + + # Verify redaction (AC 4) + e1 = [e for e in snapshot.events if e.usage_id == u1][0] + self.assertIsNotNone(e1.metadata) + self.assertNotIn("secret123", e1.metadata) + self.assertIn("[REDACTED]", e1.metadata) + + def test_missing_db_fail_soft(self) -> None: + invalid_path = "/nonexistent_path_dir/db.sqlite3" + snapshot = load_analytics(db_path=invalid_path) + self.assertFalse(snapshot.ok) + self.assertIn("control_plane_db_unavailable", snapshot.reason) + self.assertEqual(snapshot.overall_summary.display_tokens, "Unknown") + + +class AnalyticsWebUITest(unittest.TestCase): + + def setUp(self) -> None: + self.temp_dir = tempfile.TemporaryDirectory() + self.db_path = os.path.join(self.temp_dir.name, "test_webui.sqlite3") + os.environ["GITEA_CONTROL_PLANE_DB"] = self.db_path + self.app = create_app() + self.client = TestClient(self.app) + + record_usage( + db_path=self.db_path, + remote="dadeschools", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + role="author", + model="gemini-3.6-flash", + issue_number=651, + stage="implementation", + input_tokens=2000, + output_tokens=1000, + estimated_cost_usd=0.003, + latency_ms=150, + duration_ms=2500, + ) + + def tearDown(self) -> None: + self.temp_dir.cleanup() + + def test_analytics_html_route(self) -> None: + response = self.client.get("/analytics") + self.assertEqual(response.status_code, 200) + self.assertIn("Model Usage & Performance Analytics", response.text) + self.assertIn("gemini-3.6-flash", response.text) + self.assertIn("3,000", response.text) + + def test_analytics_api_route(self) -> None: + response = self.client.get("/api/v1/analytics") + self.assertEqual(response.status_code, 200) + data = response.json() + self.assertTrue(data["ok"]) + self.assertEqual(data["total_events"], 1) + self.assertIn("gemini-3.6-flash", data["by_model"]) + + def test_analytics_ingest_unauthorized_denied(self) -> None: + """F2: unauthenticated POST must not write the control-plane DB.""" + payload = { + "remote": "dadeschools", + "org": "Scaled-Tech-Consulting", + "repo": "Gitea-Tools", + "role": "reviewer", + "model": "claude-3-5-sonnet", + "pr_number": 846, + "stage": "review", + "input_tokens": 500, + "output_tokens": 100, + "latency_ms": 400, + "metadata": "Review note token=secret456", + } + response = self.client.post("/api/v1/analytics/usage", json=payload) + self.assertEqual(response.status_code, 403) + res_json = response.json() + self.assertFalse(res_json.get("ok", True)) + self.assertEqual(res_json.get("error"), "unauthorized") + authorization = res_json.get("authorization") or {} + self.assertFalse(authorization.get("allowed")) + self.assertFalse(authorization.get("execution_enabled")) + + # No new row written + res2 = self.client.get("/api/v1/analytics") + self.assertEqual(res2.status_code, 200) + self.assertEqual(res2.json()["total_events"], 1) + + def test_html_escapes_script_bearing_model_role_stage(self) -> None: + """F1: stored XSS — dynamic model/role/stage must render escaped.""" + xss = '' + record_usage( + db_path=self.db_path, + remote="dadeschools", + org="Scaled-Tech-Consulting", + repo="Gitea-Tools", + role=xss, + model=xss, + stage=xss, + issue_number=999, + status="success", + ) + response = self.client.get("/analytics") + self.assertEqual(response.status_code, 200) + # Raw tag must not appear; escaped form must. + self.assertNotIn("", response.text) + self.assertIn("<script>alert(1)</script>", response.text) + + def test_load_analytics_coerces_none_scope(self) -> None: + """F4: None remote/org/repo become empty strings, never None.""" + snapshot = load_analytics(db_path=self.db_path, remote=None, org=None, repo=None) + self.assertIsInstance(snapshot.remote, str) + self.assertIsInstance(snapshot.org, str) + self.assertIsInstance(snapshot.repo, str) + self.assertEqual(snapshot.remote, "") + self.assertEqual(snapshot.org, "") + self.assertEqual(snapshot.repo, "") + + def test_usage_events_retention_max_rows(self) -> None: + """F3: record_usage_event enforces USAGE_EVENTS_MAX_ROWS.""" + db = control_plane_db.ControlPlaneDB(db_path=self.db_path) + original_max = db.USAGE_EVENTS_MAX_ROWS + try: + db.USAGE_EVENTS_MAX_ROWS = 3 + for i in range(5): + db.record_usage_event( + remote="dadeschools", + org="org", + repo="repo", + role="author", + model=f"model-{i}", + stage="test", + ) + rows = db.query_usage_events(limit=100) + self.assertLessEqual(len(rows), 3) + # Newest three retained + models = {r["model"] for r in rows} + self.assertEqual(models, {"model-2", "model-3", "model-4"}) + finally: + db.USAGE_EVENTS_MAX_ROWS = original_max + + +if __name__ == "__main__": + unittest.main() diff --git a/webui/analytics_loader.py b/webui/analytics_loader.py new file mode 100644 index 0000000..15a6138 --- /dev/null +++ b/webui/analytics_loader.py @@ -0,0 +1,433 @@ +"""Model usage, token cost, latency, and workflow-performance analytics (#651, Phase 4). + +Ingests session instrumentation metrics, aggregates usage/cost/latency percentiles +by project, role, model, issue/PR, and stage, enforcing secret redaction and +explicitly rendering missing metrics as "Unknown" without zero-fabrication. +""" + +from __future__ import annotations + +import math +from dataclasses import asdict, dataclass +from typing import Any, Sequence + +import control_plane_db +from webui import console_redaction + +ANALYTICS_SCHEMA_VERSION = 1 + + +@dataclass(frozen=True) +class UsageEvent: + usage_id: int + session_id: str | None + remote: str + org: str + repo: str + project_id: str | None + role: str + model: str + issue_number: int | None + pr_number: int | None + stage: str + input_tokens: int | None + output_tokens: int | None + total_tokens: int | None + estimated_cost_usd: float | None + latency_ms: int | None + duration_ms: int | None + status: str + metadata: str | None + created_at: str + + def to_dict(self) -> dict[str, Any]: + d = asdict(self) + if d["metadata"]: + d["metadata"] = console_redaction.redact_text(d["metadata"]) + return d + + +@dataclass(frozen=True) +class GroupMetrics: + name: str + total_events: int + events_with_tokens: int + input_tokens: int | None + output_tokens: int | None + total_tokens: int | None + events_with_cost: int + estimated_cost_usd: float | None + events_with_latency: int + latency_p50_ms: float | None + latency_p90_ms: float | None + latency_p95_ms: float | None + latency_p99_ms: float | None + latency_avg_ms: float | None + events_with_duration: int + duration_avg_ms: float | None + display_tokens: str + display_cost: str + display_latency_p50: str + display_latency_p90: str + display_duration_avg: str + + def to_dict(self) -> dict[str, Any]: + return asdict(self) + + +@dataclass(frozen=True) +class AnalyticsSnapshot: + ok: bool + reason: str + schema_version: int + remote: str + org: str + repo: str + total_events: int + overall_summary: GroupMetrics + by_project: dict[str, GroupMetrics] + by_role: dict[str, GroupMetrics] + by_model: dict[str, GroupMetrics] + by_work_item: dict[str, GroupMetrics] + by_stage: dict[str, GroupMetrics] + events: tuple[UsageEvent, ...] + + def to_dict(self) -> dict[str, Any]: + return { + "ok": self.ok, + "reason": self.reason, + "schema_version": self.schema_version, + "remote": self.remote, + "org": self.org, + "repo": self.repo, + "total_events": self.total_events, + "overall_summary": self.overall_summary.to_dict(), + "by_project": {k: v.to_dict() for k, v in self.by_project.items()}, + "by_role": {k: v.to_dict() for k, v in self.by_role.items()}, + "by_model": {k: v.to_dict() for k, v in self.by_model.items()}, + "by_work_item": {k: v.to_dict() for k, v in self.by_work_item.items()}, + "by_stage": {k: v.to_dict() for k, v in self.by_stage.items()}, + "events": [e.to_dict() for e in self.events], + } + + +def compute_percentile(values: Sequence[float | int], percentile: float) -> float | None: + if not values: + return None + sorted_vals = sorted(values) + n = len(sorted_vals) + if n == 1: + return float(sorted_vals[0]) + k = (n - 1) * (percentile / 100.0) + f = math.floor(k) + c = math.ceil(k) + if f == c: + return float(sorted_vals[int(f)]) + d0 = sorted_vals[int(f)] * (c - k) + d1 = sorted_vals[int(c)] * (k - f) + return float(d0 + d1) + + +def aggregate_events(group_name: str, events: Sequence[UsageEvent]) -> GroupMetrics: + total_events = len(events) + if total_events == 0: + return GroupMetrics( + name=group_name, + total_events=0, + events_with_tokens=0, + input_tokens=None, + output_tokens=None, + total_tokens=None, + events_with_cost=0, + estimated_cost_usd=None, + events_with_latency=0, + latency_p50_ms=None, + latency_p90_ms=None, + latency_p95_ms=None, + latency_p99_ms=None, + latency_avg_ms=None, + events_with_duration=0, + duration_avg_ms=None, + display_tokens="Unknown", + display_cost="Unknown", + display_latency_p50="Unknown", + display_latency_p90="Unknown", + display_duration_avg="Unknown", + ) + + token_events = [ + e for e in events + if e.total_tokens is not None or e.input_tokens is not None or e.output_tokens is not None + ] + events_with_tokens = len(token_events) + if events_with_tokens > 0: + input_tokens = sum(e.input_tokens or 0 for e in token_events) + output_tokens = sum(e.output_tokens or 0 for e in token_events) + total_tokens = sum( + e.total_tokens if e.total_tokens is not None else ((e.input_tokens or 0) + (e.output_tokens or 0)) + for e in token_events + ) + display_tokens = f"{total_tokens:,}" + else: + input_tokens = None + output_tokens = None + total_tokens = None + display_tokens = "Unknown" + + cost_events = [e for e in events if e.estimated_cost_usd is not None] + events_with_cost = len(cost_events) + if events_with_cost > 0: + estimated_cost_usd = round(sum(e.estimated_cost_usd for e in cost_events), 6) + display_cost = f"${estimated_cost_usd:.4f}" + else: + estimated_cost_usd = None + display_cost = "Unknown" + + latency_vals = [e.latency_ms for e in events if e.latency_ms is not None] + events_with_latency = len(latency_vals) + if events_with_latency > 0: + latency_p50_ms = compute_percentile(latency_vals, 50.0) + latency_p90_ms = compute_percentile(latency_vals, 90.0) + latency_p95_ms = compute_percentile(latency_vals, 95.0) + latency_p99_ms = compute_percentile(latency_vals, 99.0) + latency_avg_ms = round(sum(latency_vals) / events_with_latency, 2) + display_latency_p50 = f"{round(latency_p50_ms, 1)} ms" if latency_p50_ms is not None else "Unknown" + display_latency_p90 = f"{round(latency_p90_ms, 1)} ms" if latency_p90_ms is not None else "Unknown" + else: + latency_p50_ms = None + latency_p90_ms = None + latency_p95_ms = None + latency_p99_ms = None + latency_avg_ms = None + display_latency_p50 = "Unknown" + display_latency_p90 = "Unknown" + + duration_vals = [e.duration_ms for e in events if e.duration_ms is not None] + events_with_duration = len(duration_vals) + if events_with_duration > 0: + duration_avg_ms = round(sum(duration_vals) / events_with_duration, 2) + display_duration_avg = f"{round(duration_avg_ms / 1000.0, 2)} s" if duration_avg_ms >= 1000 else f"{round(duration_avg_ms, 1)} ms" + else: + duration_avg_ms = None + display_duration_avg = "Unknown" + + return GroupMetrics( + name=group_name, + total_events=total_events, + events_with_tokens=events_with_tokens, + input_tokens=input_tokens, + output_tokens=output_tokens, + total_tokens=total_tokens, + events_with_cost=events_with_cost, + estimated_cost_usd=estimated_cost_usd, + events_with_latency=events_with_latency, + latency_p50_ms=latency_p50_ms, + latency_p90_ms=latency_p90_ms, + latency_p95_ms=latency_p95_ms, + latency_p99_ms=latency_p99_ms, + latency_avg_ms=latency_avg_ms, + events_with_duration=events_with_duration, + duration_avg_ms=duration_avg_ms, + display_tokens=display_tokens, + display_cost=display_cost, + display_latency_p50=display_latency_p50, + display_latency_p90=display_latency_p90, + display_duration_avg=display_duration_avg, + ) + + +def record_usage( + *, + db_path: str | None = None, + session_id: str | None = None, + remote: str = "dadeschools", + org: str = "", + repo: str = "", + project_id: str | None = None, + role: str = "unknown", + model: str = "unknown", + issue_number: int | None = None, + pr_number: int | None = None, + stage: str = "unknown", + input_tokens: int | None = None, + output_tokens: int | None = None, + total_tokens: int | None = None, + estimated_cost_usd: float | None = None, + latency_ms: int | None = None, + duration_ms: int | None = None, + status: str = "success", + metadata: str | dict[str, Any] | None = None, + created_at: str | None = None, +) -> int: + """Ingest/record a single usage event with optional metrics.""" + db = control_plane_db.ControlPlaneDB(db_path=db_path) + return db.record_usage_event( + session_id=session_id, + remote=remote, + org=org, + repo=repo, + project_id=project_id, + role=role, + model=model, + issue_number=issue_number, + pr_number=pr_number, + stage=stage, + input_tokens=input_tokens, + output_tokens=output_tokens, + total_tokens=total_tokens, + estimated_cost_usd=estimated_cost_usd, + latency_ms=latency_ms, + duration_ms=duration_ms, + status=status, + metadata=metadata, + created_at=created_at, + ) + + +def load_analytics( + *, + db_path: str | None = None, + remote: str | None = None, + org: str | None = None, + repo: str | None = None, + project_id: str | None = None, + role: str | None = None, + model: str | None = None, + stage: str | None = None, + issue_number: int | None = None, + pr_number: int | None = None, + limit: int = 500, +) -> AnalyticsSnapshot: + """Load analytics snapshot aggregated by project, role, model, issue/PR, and stage.""" + remote_filter = (remote or "").strip() or None + org_filter = (org or "").strip() or None + repo_filter = (repo or "").strip() or None + role_filter = (role or "").strip() or None + model_filter = (model or "").strip() or None + stage_filter = (stage or "").strip() or None + + # F4: coerce optional scope filters to str so AnalyticsSnapshot never holds None. + scope_remote = (remote or "").strip() + scope_org = (org or "").strip() + scope_repo = (repo or "").strip() + + try: + db = control_plane_db.ControlPlaneDB(db_path=db_path) + rows = db.query_usage_events( + remote=remote_filter, + org=org_filter, + repo=repo_filter, + project_id=project_id, + role=role_filter, + model=model_filter, + stage=stage_filter, + issue_number=issue_number, + pr_number=pr_number, + limit=limit, + ) + except Exception as exc: + empty_summary = aggregate_events("Overall", []) + return AnalyticsSnapshot( + ok=False, + reason=f"control_plane_db_unavailable: {exc}", + schema_version=ANALYTICS_SCHEMA_VERSION, + remote=scope_remote, + org=scope_org, + repo=scope_repo, + total_events=0, + overall_summary=empty_summary, + by_project={}, + by_role={}, + by_model={}, + by_work_item={}, + by_stage={}, + events=(), + ) + + parsed_events: list[UsageEvent] = [] + for r in rows: + meta = console_redaction.redact_text(r.get("metadata")) if r.get("metadata") else None + parsed_events.append( + UsageEvent( + usage_id=r["usage_id"], + session_id=r.get("session_id"), + remote=r.get("remote") or scope_remote, + org=r.get("org") or scope_org, + repo=r.get("repo") or scope_repo, + project_id=r.get("project_id"), + role=r.get("role") or "unknown", + model=r.get("model") or "unknown", + issue_number=r.get("issue_number"), + pr_number=r.get("pr_number"), + stage=r.get("stage") or "unknown", + input_tokens=r.get("input_tokens"), + output_tokens=r.get("output_tokens"), + total_tokens=r.get("total_tokens"), + estimated_cost_usd=r.get("estimated_cost_usd"), + latency_ms=r.get("latency_ms"), + duration_ms=r.get("duration_ms"), + status=r.get("status") or "success", + metadata=meta, + created_at=r.get("created_at") or "", + ) + ) + + overall_summary = aggregate_events("Overall", parsed_events) + + # Group by project + groups_by_project: dict[str, list[UsageEvent]] = {} + for e in parsed_events: + key = e.project_id or (f"{e.org}/{e.repo}" if e.org and e.repo else "default") + groups_by_project.setdefault(key, []).append(e) + by_project = {k: aggregate_events(k, v) for k, v in groups_by_project.items()} + + # Group by role + groups_by_role: dict[str, list[UsageEvent]] = {} + for e in parsed_events: + groups_by_role.setdefault(e.role, []).append(e) + by_role = {k: aggregate_events(k, v) for k, v in groups_by_role.items()} + + # Group by model + groups_by_model: dict[str, list[UsageEvent]] = {} + for e in parsed_events: + groups_by_model.setdefault(e.model, []).append(e) + by_model = {k: aggregate_events(k, v) for k, v in groups_by_model.items()} + + # Group by work item + groups_by_work_item: dict[str, list[UsageEvent]] = {} + for e in parsed_events: + if e.issue_number: + key = f"issue #{e.issue_number}" + elif e.pr_number: + key = f"pr #{e.pr_number}" + else: + key = "unlinked" + groups_by_work_item.setdefault(key, []).append(e) + by_work_item = {k: aggregate_events(k, v) for k, v in groups_by_work_item.items()} + + # Group by stage + groups_by_stage: dict[str, list[UsageEvent]] = {} + for e in parsed_events: + groups_by_stage.setdefault(e.stage, []).append(e) + by_stage = {k: aggregate_events(k, v) for k, v in groups_by_stage.items()} + + return AnalyticsSnapshot( + ok=True, + reason="ok", + schema_version=ANALYTICS_SCHEMA_VERSION, + remote=scope_remote, + org=scope_org, + repo=scope_repo, + total_events=len(parsed_events), + overall_summary=overall_summary, + by_project=by_project, + by_role=by_role, + by_model=by_model, + by_work_item=by_work_item, + by_stage=by_stage, + events=tuple(parsed_events), + ) + + +def snapshot_to_dict(snapshot: AnalyticsSnapshot) -> dict[str, Any]: + return snapshot.to_dict() diff --git a/webui/analytics_views.py b/webui/analytics_views.py new file mode 100644 index 0000000..506aa96 --- /dev/null +++ b/webui/analytics_views.py @@ -0,0 +1,248 @@ +"""HTML views for the Model Usage & Performance Analytics console (#651).""" + +from __future__ import annotations + +import html + +from webui.analytics_loader import AnalyticsSnapshot, GroupMetrics, UsageEvent +from webui.layout import render_page + + +def _escape(text: object) -> str: + """HTML-escape dynamic analytics fields (mirrors audit_views / project_views).""" + return html.escape(str(text), quote=True) + + +def _render_badge(text: str, badge_type: str = "muted") -> str: + return f'{_escape(text)}' + + +def _render_group_table(title: str, groups: dict[str, GroupMetrics], key_header: str = "Group") -> str: + if not groups: + return ( + f"

{_escape(title)}

" + '

No telemetry events recorded for this dimension.

' + ) + + rows = [] + for key, g in sorted(groups.items(), key=lambda x: x[1].total_events, reverse=True): + cost_cell = ( + f'{_escape(g.display_cost)}' + if g.events_with_cost > 0 + else _render_badge("Unknown") + ) + tokens_cell = ( + _escape(g.display_tokens) + if g.events_with_tokens > 0 + else _render_badge("Unknown") + ) + lat_p50 = ( + _escape(g.display_latency_p50) + if g.events_with_latency > 0 + else _render_badge("Unknown") + ) + lat_p90 = ( + _escape(g.display_latency_p90) + if g.events_with_latency > 0 + else _render_badge("Unknown") + ) + dur_avg = ( + _escape(g.display_duration_avg) + if g.events_with_duration > 0 + else _render_badge("Unknown") + ) + + rows.append( + "" + f"{_escape(key)}" + f"{g.total_events}" + f"{tokens_cell}" + f"{cost_cell}" + f"{lat_p50}" + f"{lat_p90}" + f"{dur_avg}" + "" + ) + + rows_html = "".join(rows) + return f""" +

{_escape(title)}

+
+ + + + + + + + + + + + + + {rows_html} + +
{_escape(key_header)}EventsTotal TokensEst. CostLatency (p50)Latency (p90)Avg Stage Duration
+
+ """ + + +def _render_events_table(events: tuple[UsageEvent, ...]) -> str: + if not events: + return ( + "

Recent Usage & Instrumentation Events

" + '

No individual telemetry events recorded yet. Opt-in instrumentation via session logging or authorized POST /api/v1/analytics/usage.

' + ) + + rows = [] + for e in list(events)[-50:]: # Display latest 50 + if e.issue_number is not None: + work_item = f"issue #{e.issue_number}" + elif e.pr_number is not None: + work_item = f"pr #{e.pr_number}" + else: + work_item = "unlinked" + tokens = ( + _escape(f"{e.total_tokens:,}") + if e.total_tokens is not None + else _render_badge("Unknown") + ) + cost = ( + _escape(f"${e.estimated_cost_usd:.4f}") + if e.estimated_cost_usd is not None + else _render_badge("Unknown") + ) + latency = ( + _escape(f"{e.latency_ms} ms") + if e.latency_ms is not None + else _render_badge("Unknown") + ) + duration = ( + _escape(f"{e.duration_ms} ms") + if e.duration_ms is not None + else _render_badge("Unknown") + ) + status_badge = _render_badge( + e.status, "success" if e.status == "success" else "danger" + ) + + rows.append( + "" + f"#{e.usage_id}" + f"{_escape(e.created_at)}" + f"{_escape(e.role)}" + f"{_escape(e.model)}" + f"{_escape(e.stage)}" + f"{_escape(work_item)}" + f"{tokens}" + f"{cost}" + f"{latency}" + f"{duration}" + f"{status_badge}" + "" + ) + + rows_html = "".join(rows) + return f""" +

Recent Telemetry Events

+
+ + + + + + + + + + + + + + + + + + {rows_html} + +
IDTimestampRoleModelStageWork ItemTokensCostLatencyDurationStatus
+
+ """ + + +def render_analytics_page(snapshot: AnalyticsSnapshot) -> str: + """Render the main Model Usage & Performance Analytics console page.""" + summary = snapshot.overall_summary + + kpi_tokens = ( + _escape(summary.display_tokens) + if summary.events_with_tokens > 0 + else _render_badge("Unknown") + ) + kpi_cost = ( + _escape(summary.display_cost) + if summary.events_with_cost > 0 + else _render_badge("Unknown") + ) + kpi_lat_p50 = ( + _escape(summary.display_latency_p50) + if summary.events_with_latency > 0 + else _render_badge("Unknown") + ) + kpi_dur_avg = ( + _escape(summary.display_duration_avg) + if summary.events_with_duration > 0 + else _render_badge("Unknown") + ) + + status_notice = "" + if not snapshot.ok: + status_notice = ( + f'
Degraded Data Source: ' + f'{_escape(snapshot.reason)}
' + ) + + body_html = f""" +

Model Usage & Performance Analytics (Phase 4)

+

+ Durable console analytics for model usage, token cost, latency percentiles, and workflow-stage performance correlated to issues, PRs, and worker roles. +

+ + {status_notice} + +
+ Note on telemetry fidelity: Missing data or untracked metrics are explicitly labeled as Unknown. No token costs or latency metrics are zero-fabricated. +
+ +
+
+ Total Events +

{summary.total_events}

+
+
+ Total Tokens +

{kpi_tokens}

+
+
+ Est. Token Cost +

{kpi_cost}

+
+
+ Latency (p50) +

{kpi_lat_p50}

+
+
+ Avg Stage Duration +

{kpi_dur_avg}

+
+
+ + {_render_group_table("Usage & Cost by Model", snapshot.by_model, "Model")} + {_render_group_table("Performance by Workflow Stage", snapshot.by_stage, "Stage")} + {_render_group_table("Usage & Cost by Role", snapshot.by_role, "Role")} + {_render_group_table("Work Item Analytics", snapshot.by_work_item, "Work Item")} + {_render_events_table(snapshot.events)} + """ + + return render_page(title="Model Usage & Performance Analytics", body_html=body_html) diff --git a/webui/app.py b/webui/app.py index f5147b5..f7648eb 100644 --- a/webui/app.py +++ b/webui/app.py @@ -47,6 +47,12 @@ 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.analytics_loader import ( + load_analytics, + record_usage, + snapshot_to_dict as analytics_snapshot_to_dict, +) +from webui.analytics_views import render_analytics_page from webui.system_health import ( API_PATH as SYSTEM_HEALTH_API_PATH, load_system_health, @@ -567,6 +573,114 @@ async def api_v1_timeline(request: Request) -> JSONResponse: return JSONResponse(timeline_snapshot_to_dict(snapshot), status_code=status_code) +async def analytics(request: Request) -> HTMLResponse: + """Read-only model usage, token cost, latency, and performance analytics HTML view (#651).""" + snapshot = load_analytics( + remote=request.query_params.get("remote"), + org=request.query_params.get("org"), + repo=request.query_params.get("repo"), + role=request.query_params.get("role"), + model=request.query_params.get("model"), + stage=request.query_params.get("stage"), + issue_number=_query_int(request, "issue"), + pr_number=_query_int(request, "pr"), + limit=_query_int(request, "limit") or 200, + ) + return HTMLResponse(render_analytics_page(snapshot)) + + +async def api_v1_analytics(request: Request) -> JSONResponse: + """Read-only model usage, token cost, latency, and performance analytics API (#651).""" + snapshot = load_analytics( + remote=request.query_params.get("remote"), + org=request.query_params.get("org"), + repo=request.query_params.get("repo"), + role=request.query_params.get("role"), + model=request.query_params.get("model"), + stage=request.query_params.get("stage"), + issue_number=_query_int(request, "issue"), + pr_number=_query_int(request, "pr"), + limit=_query_int(request, "limit") or 500, + ) + status_code = 200 if snapshot.ok else 500 + return JSONResponse(analytics_snapshot_to_dict(snapshot), status_code=status_code) + + +async def api_v1_analytics_ingest(request: Request) -> JSONResponse: + """Optional session instrumentation ingestion endpoint (#651). + + Fail-closed write: every request is authorized through console_authz + (``record_analytics_usage``) before any control-plane DB mutation. Phase 1 + keeps ``execution_enabled=False`` and denies unauthenticated callers, so + this route cannot be used as an unauthenticated write or XSS injection + vector (PR #876 F2). + """ + try: + body = await request.json() + except Exception: + body = {} + if not isinstance(body, dict): + body = {} + + principal = resolve_principal(headers=dict(request.headers)) + decision = authorize( + "record_analytics_usage", principal, for_execution=True + ) + allowed = bool(decision.allowed and decision.execution_enabled) + console_audit.record_event( + action_id="record_analytics_usage", + result=( + console_audit.RESULT_ALLOWED + if allowed + else console_audit.RESULT_DENIED + ), + decision=decision, + principal=principal, + target=_audit_target("record_analytics_usage", body), + request_id=_request_id(), + detail=decision.detail, + ) + authorization = decision.to_dict() + if not allowed: + return JSONResponse( + { + "ok": False, + "error": "unauthorized", + "detail": ( + "POST /api/v1/analytics/usage requires an authenticated " + "principal with record_analytics_usage execution enabled" + ), + "authorization": authorization, + }, + status_code=403, + ) + + usage_id = record_usage( + session_id=body.get("session_id"), + remote=body.get("remote", "dadeschools"), + org=body.get("org", ""), + repo=body.get("repo", ""), + project_id=body.get("project_id"), + role=body.get("role", "unknown"), + model=body.get("model", "unknown"), + issue_number=body.get("issue_number") or body.get("issue"), + pr_number=body.get("pr_number") or body.get("pr"), + stage=body.get("stage", "unknown"), + input_tokens=body.get("input_tokens"), + output_tokens=body.get("output_tokens"), + total_tokens=body.get("total_tokens"), + estimated_cost_usd=body.get("estimated_cost_usd"), + latency_ms=body.get("latency_ms"), + duration_ms=body.get("duration_ms"), + status=body.get("status", "success"), + metadata=body.get("metadata"), + ) + return JSONResponse( + {"ok": True, "usage_id": usage_id, "authorization": authorization}, + status_code=201, + ) + + async def method_not_allowed(request: Request, _exc: Exception) -> Response: path = request.url.path if path in _AUDIT_MUTATION_PATHS and request.method == "POST": @@ -608,6 +722,10 @@ def create_app(*, bind_host: str | None = None) -> Starlette: Route("/runtime", runtime, methods=["GET"]), Route("/api/runtime", api_runtime, methods=["GET"]), Route("/api/v1/timeline", api_v1_timeline, methods=["GET"]), + Route("/analytics", analytics, methods=["GET"]), + Route("/api/analytics", api_v1_analytics, methods=["GET"]), + Route("/api/v1/analytics", api_v1_analytics, methods=["GET"]), + Route("/api/v1/analytics/usage", api_v1_analytics_ingest, methods=["POST"]), Route("/audit", audit, methods=["GET", "POST"]), Route("/api/audit", api_audit, methods=["GET", "POST"]), Route("/worktrees", worktrees, methods=["GET"]), diff --git a/webui/console_authz.py b/webui/console_authz.py index 06efe86..0c51037 100644 --- a/webui/console_authz.py +++ b/webui/console_authz.py @@ -236,6 +236,19 @@ _ACTION_SPECS: tuple[ConsoleAction, ...] = ( phase=3, summary="Remove a remote feature branch.", ), + # #651 analytics ingest: local control-plane write, not a Gitea mutation. + # Phase 2 gated write so Phase 1 (ACTIVE_PHASE=1) fails closed on execution. + ConsoleAction( + action_id="record_analytics_usage", + task_key="record_analytics_usage", + action_class=CLASS_WRITE, + minimum_role=OPERATOR, + requires_confirmation=True, + dual_control=False, + break_glass=False, + phase=2, + summary="Ingest a model-usage / latency analytics event into the control-plane DB.", + ), # #642: sanctioned daemon lifecycle. These exist so operators have an # audited path off `pkill -f mcp_server.py` (#630). Restart drops every # in-flight request on a namespace, so it carries the same dual-control and diff --git a/webui/nav.py b/webui/nav.py index c24643d..45d1212 100644 --- a/webui/nav.py +++ b/webui/nav.py @@ -65,6 +65,7 @@ NAV_GROUPS: tuple[NavGroup, ...] = ( )), NavGroup("Insights", ( NavItem("/insights", "Insights", "stub"), + NavItem("/analytics", "Analytics"), NavItem("/audit", "Audit"), )), )