From 433f66add864062df7c211d9b52fc74fcfccfb2f Mon Sep 17 00:00:00 2001 From: Jason Walker <913443@dadeschools.net> Date: Sat, 25 Jul 2026 01:47:38 -0400 Subject: [PATCH 1/2] feat(webui): request preview, authorization, and workflow initiation (Closes #643) Operators had to paste a role prompt into a terminal to start work, and nothing enforced that the allocator had been consulted first, so two sessions could reach for the same issue and each believe it was theirs. This adds a request surface: a desired role, an issue or PR, and a stated intent, answered by an authorization decision and - on confirmation - an exclusive assignment from the allocator. Preview (POST /api/v1/requests/preview, and the /requests form) runs five checks and reports authorize/deny with a reason for each: console authorization, capability resolution for the desired role, lease availability, whether the allocator would independently select this work unit, and head pinning for PR work. It is read-only - it calls the allocator with apply=false and writes only an audit line. An unauthorized principal never reaches the allocator or the control-plane DB, so a denial cannot enumerate the queue. Initiation (POST /api/v1/requests/apply) never assigns the requested item directly. It runs a dry-run first and proceeds only when the allocator would independently pick that exact work unit, carrying the dry-run's candidate_set_fingerprint as a CAS pin; otherwise it returns wait or blocked and mutates nothing. An active claim on the work unit rejects a duplicate assign before one is attempted. A returned assignment carries a handoff block naming the required profile, namespace, and the actions that stay forbidden. Authorization reuses the #633 model rather than adding a second one. The new initiate_workflow action is operator-class because its outcome is a claim, not a Gitea verdict: requesting reviewer or merger work reserves that work but grants no right to approve or merge. Execution is gated by a new per-action execution_env_flag (WEBUI_REQUESTS_EXECUTION), deliberately in place of raising ACTIVE_PHASE - a phase bump would enable execution for every phase-2 action at once, including ones whose execution path is not implemented. Actions that declare no flag are unchanged and still report execution_enabled false. Every preview and apply emits a console audit record correlated to the resulting assignment by correlation.request_id. Fail-closed throughout: an unreadable control-plane DB, an incomplete queue inventory (#758), an allocator that raises, an unpinned PR head, a moved PR head, and an unconfirmed apply all deny without mutating. Files: - webui/request_service.py (new) - request model, preview, initiation - webui/request_views.py (new) - form and preview rendering, escaped - tests/test_webui_request_initiation.py (new) - 52 tests - webui/console_authz.py - initiate_workflow action, execution_wired() - webui/app.py - /requests, /api/v1/requests/preview, /api/v1/requests/apply - webui/nav.py - Requests nav entry - webui/traffic_loader.py - public candidates_from_queue_snapshot alias - docs/webui-requests.md (new), docs/webui-authz-audit.md Validation: full suite on this branch 5242 passed, 6 skipped, 899 subtests, 23 failed. Clean master baseline at 2f4dec83 in an equivalent branches/ worktree: 5190 passed, 6 skipped, 867 subtests, the same 23 tests failed. The branch adds 52 passing tests and introduces no new full-suite failure signature. Closes #643 Co-Authored-By: Claude Opus 4.8 (1M context) --- docs/webui-authz-audit.md | 50 +- docs/webui-requests.md | 160 ++++ tests/test_webui_request_initiation.py | 917 +++++++++++++++++++++++ webui/app.py | 116 +++ webui/console_authz.py | 79 +- webui/nav.py | 1 + webui/request_service.py | 967 +++++++++++++++++++++++++ webui/request_views.py | 164 +++++ webui/traffic_loader.py | 10 + 9 files changed, 2447 insertions(+), 17 deletions(-) create mode 100644 docs/webui-requests.md create mode 100644 tests/test_webui_request_initiation.py create mode 100644 webui/request_service.py create mode 100644 webui/request_views.py diff --git a/docs/webui-authz-audit.md b/docs/webui-authz-audit.md index 2dcffac..ff9cc38 100644 --- a/docs/webui-authz-audit.md +++ b/docs/webui-authz-audit.md @@ -94,6 +94,7 @@ already define, and a regression test asserts each mapping matches. | `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 | +| `initiate_workflow` | operator | gated_write | `gitea.read` | Yes | No | No | 2 | **Dual control** means the acting principal may not be the sole authority: a second distinct principal must confirm. **Break-glass** means the action is @@ -112,6 +113,12 @@ by the console — both hand off to a host supervisor, and neither exposes a raw process kill. See [`sanctioned-restart-controls.md`](sanctioned-restart-controls.md) (#642). +`initiate_workflow` (#643) is operator-class because its outcome is a *claim*, +not a Gitea verdict. Requesting reviewer or merger work reserves that work +through the allocator; it does not grant the right to approve or merge, which +stays with the MCP role profile and its own capability gates. See +[`webui-requests.md`](webui-requests.md). + ### Authorization decision `authorize(action_id, principal, for_execution=False)` returns a decision @@ -126,9 +133,24 @@ record and **denies by default**. The deny reasons are closed and enumerated: | `phase_not_active` | Execution requested for an action whose phase is not open. | | `allowed_preview_only` | Authorized — preview only, execution still disabled. | -There is no implicit allow branch. Even the allow result reports -`execution_enabled: false` while the console is in Phase 1, so no caller can -read an allow as permission to mutate. +There is no implicit allow branch. + +`execution_enabled` on the decision reports whether the action has a live +execution path at all, and is computed by `execution_wired(action)`. There are +exactly two ways to be wired: + +1. the action's `phase` is at or below `ACTIVE_PHASE`; or +2. the action declares an `execution_env_flag` **and** that variable is set. + +Every action that declares no flag therefore reports `execution_enabled: false` +while the console is in Phase 1, so no caller can read an allow as permission +to mutate. The per-action flag exists because raising `ACTIVE_PHASE` would +enable execution for every action of that phase at once, including ones whose +execution path is not implemented. One implemented action goes live on its own +flag instead of dragging its unimplemented phase-mates with it. + +`initiate_workflow` is the only action that currently declares a flag +(`WEBUI_REQUESTS_EXECUTION`), and it stays denied until an operator sets it. ## Secret redaction @@ -235,13 +257,22 @@ second one. The integration points are already wired and observable: instead of adding a parallel check. - **`GET /api/console/security-model`** publishes the RBAC matrix, redaction policy, and audit policy as JSON for operators and tests. +- **`POST /api/v1/requests/preview` and `.../apply`** (#643) are the first + actions to use this model for a real execution path. Preview always returns a + decision and an audited `previewed` record; apply requires `confirm=true`, + emits `succeeded` or `denied`, and reserves work only through the allocator. + See [`webui-requests.md`](webui-requests.md). -To open Phase 2, a child issue must: raise `ACTIVE_PHASE`, implement the -confirmation and dual-control flow the matrix already declares, emit a -`succeeded` or `failed` record alongside the `gitea_audit` mutation record, and -keep `viewer` unable to reach any of it. Turning on execution without the -confirmation flow contradicts a declared requirement and is a review failure, -not a shortcut. +A Phase 2 action must: use `execution_wired` rather than a private enable flag, +implement the confirmation and dual-control flow the matrix already declares, +emit a `succeeded` or `failed` record alongside the `gitea_audit` mutation +record, and keep `viewer` unable to reach any of it. Turning on execution +without the confirmation flow contradicts a declared requirement and is a +review failure, not a shortcut. + +Raising `ACTIVE_PHASE` remains the way to open a whole phase at once, and is +deliberately *not* what #643 did: an action-scoped opt-in cannot enable an +action whose execution path nobody wrote. ## Local-dev mode @@ -294,6 +325,7 @@ Until Phase 2 wires it, probe protection rests on network placement alone, as | `WEBUI_ROLE_MAP` | unset | JSON subject → role map | | `WEBUI_REQUIRE_PROBE_AUTH` | unset | Require auth for non-public probes | | `WEBUI_CONSOLE_AUDIT_LOG` | unset | Append-only audit sink path | +| `WEBUI_REQUESTS_EXECUTION` | unset | Opt in to `initiate_workflow` execution (#643) | All are read server-side only. None is ever rendered into a page or returned by an API. diff --git a/docs/webui-requests.md b/docs/webui-requests.md new file mode 100644 index 0000000..b8fbdb2 --- /dev/null +++ b/docs/webui-requests.md @@ -0,0 +1,160 @@ +# Web console requests: intent preview and workflow initiation (#643) + +**Phase 2. Preview is always live and always read-only. Initiation is wired but +denied until an operator opts in.** + +Before this surface, starting role work meant pasting a prompt into a terminal +and trusting the operator to have checked the allocator first. Nothing enforced +that check, so two sessions could reach for the same issue and each believe it +was theirs. This page replaces the paste with a *request*: a desired role, an +issue or PR, and a stated intent, answered by an authorization decision and — +on confirmation — an exclusive assignment from the allocator. + +| Concern | Module | +|---------|--------| +| Request model, preview, initiation | `webui/request_service.py` | +| Form and preview rendering | `webui/request_views.py` | +| Authorization | `webui/console_authz.py` (`initiate_workflow`) | +| Audit | `webui/console_audit.py` | +| Ownership substrate | `allocator_service.py` + `control_plane_db.py` | + +## Surfaces + +| Path | Method | Purpose | +|------|--------|---------| +| `/requests` | GET | Request form | +| `/requests` | POST | Render an intent preview. **Never assigns.** | +| `/api/v1/requests/preview` | POST | Intent preview as JSON | +| `/api/v1/requests/apply` | POST | Initiate — confirmed, audited, allocator-owned | + +The HTML form has no initiate button on purpose. Initiating requires a +confirmed POST to `/api/v1/requests/apply`, so a stray form submission cannot +reserve work as a side effect. + +## The request + +```json +{ + "desired_role": "author", + "work_kind": "issue", + "work_number": 643, + "intent_summary": "implement request preview and initiation", + "remote": "prgs", + "org": "Scaled-Tech-Consulting", + "repo": "Gitea-Tools", + "expected_head_sha": null +} +``` + +`desired_role` is one of `author`, `reviewer`, `merger`, `reconciler`, +`controller`. `work_kind` is `issue` or `pr`. `remote`/`org`/`repo` default to +the first project in the registry when omitted; when neither the request nor +the registry resolves them, the request is rejected rather than pointed at some +other repository. `intent_summary` is required — it is what the audit record +states as the reason — and is truncated to 500 characters. + +Parsing rejects rather than corrects. An unknown role, an unknown work kind, a +non-positive number, or a missing intent each return `400` with a `reason_code` +and the offending `field`. + +## Preview + +Five checks, each with its own verdict, reason code, and detail: + +| Check | Passes when | +|-------|-------------| +| `authorization` | The console principal holds `operator` or above | +| `capability` | The desired role maps to a declared profile and MCP namespace | +| `lease_availability` | No active claim holds the work unit | +| `next_safe_action` | The allocator would independently select this exact work unit | +| `head_pin` | PR work resolves to a head SHA, and a supplied SHA still matches | + +A preview also returns the role's `allowed_actions` and `prohibited_actions` +(from `allocator_service.ROLE_ACTIONS`), the `required_profile` and +`required_namespace` the work must run under, and a `correlation_id` that ties +the preview to its audit record and to any assignment that follows. + +Preview is read-only in the strict sense: it calls the allocator with +`apply=false` and writes nothing but an audit line. An unauthorized principal +never reaches the allocator or the control-plane DB at all, so a denial cannot +be used to enumerate the queue. + +## Initiation + +`POST /api/v1/requests/apply` refuses in this order, and every refusal returns +before any assignment is attempted: + +| Condition | Outcome | Status | +|-----------|---------|--------| +| Unparseable request | `invalid_request` | 400 | +| Not authorized, or execution not wired | `denied` | 403 | +| `confirm` not set | `denied` / `confirmation_required` | 409 | +| Work unit already claimed | `blocked` / `duplicate_assignment` | 409 | +| Allocator would select other work | `wait` / `not_next_safe_work` | 409 | +| Allocator declines on apply | `blocked` or `wait` | 409 | +| Evidence unavailable | `wait` / `evidence_unavailable` | 503 | +| Assigned | `assigned_work` | 201 | + +A success returns the assignment plus a `handoff` block naming the profile, the +namespace, and the actions that stay forbidden — enough for the operator to +continue in the right MCP namespace without guessing. + +### Why apply runs the allocator twice + +The allocator is the only source of exclusive ownership (#600 / #613), and it +selects work; it does not take orders. So `apply` runs a dry-run first and +proceeds only when the allocator would independently pick the requested work +unit. If it would not, the request reports `wait` and mutates nothing. + +A request is therefore a *confirmation* of the allocator's decision, never an +override of it. The apply call carries the dry-run's +`candidate_set_fingerprint` as a CAS pin (#776), so a queue that changed +between the two calls fails closed rather than assigning against a stale view. +The result is checked again on the way out: an assignment naming a different +work unit is not read as success. + +### Fail-closed defaults + +- An unreadable control-plane DB denies. It is never treated as "nothing holds + this work unit". +- An incomplete queue inventory denies (#758). Ranking a partial candidate set + can select the wrong work. +- An allocator that raises denies. +- PR work with no resolvable head SHA denies; a supplied SHA that no longer + matches denies with `head_moved`. + +## Enabling initiation + +Execution is wired off. Set `WEBUI_REQUESTS_EXECUTION=1` to enable it for the +`initiate_workflow` action only — see +[`webui-authz-audit.md`](webui-authz-audit.md) for why this is an +action-scoped flag rather than a phase bump. With the variable unset, `apply` +returns `403` with `reason_code: unauthorized` no matter who asks. + +Enabling execution does **not** enable approvals or merges. Those are phase 3 +console actions and remain forbidden in every path here; the console reserves +work and hands off, and the MCP role profile enforces what that role may then +do. + +## Audit + +Every preview and every apply emits a console audit record (schema in +[`webui-authz-audit.md`](webui-authz-audit.md)): + +| Event | `result` | +|-------|----------| +| Preview | `previewed` | +| Refusal at any stage | `denied` | +| Assignment created | `succeeded` | + +`correlation.request_id` carries the request's `correlation_id`, and a +successful record's `metadata` carries `assignment_id` and `lease_id`, so an +assignment can be traced back to the intent that produced it. The operator's +`intent_summary` travels in `metadata` and passes through the standard +redaction pass before persistence like every other field. + +## Non-goals + +- No browser-initiated approve or merge, in this phase or any other. +- No bypass of allocator exclusive ownership; no self-selection of work. +- No auto-start from raw monitoring incidents (#612 stays downstream). diff --git a/tests/test_webui_request_initiation.py b/tests/test_webui_request_initiation.py new file mode 100644 index 0000000..c733898 --- /dev/null +++ b/tests/test_webui_request_initiation.py @@ -0,0 +1,917 @@ +"""Request preview, authorization, and workflow initiation tests (#643). + +Covers each acceptance criterion: + +* AC1 — preview shows authorize/deny with reasons. +* AC2 — apply creates an exclusive assignment or returns wait/blocked. +* AC3 — duplicate assign rejected. +* AC4 — preview / apply / deny / collision are all exercised. +* AC5 — the UI never renders a secret, and messaging stays brief. + +Required tests named in the issue: allocator integration with fakes, and +gated-action tests. The allocator is injected as a fake throughout so no test +touches Gitea or reserves real work; one class asserts the *real* default +allocator refuses an incomplete inventory rather than ranking a partial set. +""" + +from __future__ import annotations + +import json +import os +import pathlib +import sys +import tempfile +import unittest +from typing import Any +from unittest import mock + +from tests.webui_testclient import TestClient + +sys.path.insert(0, str(pathlib.Path(__file__).resolve().parents[1])) + +import allocator_service # noqa: E402 +from webui import console_audit, console_authz, request_service # noqa: E402 +from webui.app import create_app # noqa: E402 +from webui.console_redaction import scan_for_secrets # noqa: E402 +from webui.request_views import render_requests_page # noqa: E402 + +EXEC_FLAG = "WEBUI_REQUESTS_EXECUTION" + +SCOPE = { + "remote": "prgs", + "org": "Scaled-Tech-Consulting", + "repo": "Gitea-Tools", +} + + +def _principal(role: str) -> console_authz.Principal: + return console_authz.Principal( + subject=f"{role}@example.com", + role=role, + identity_source=console_authz.IDENTITY_ACCESS_PROXY, + authenticated=True, + ) + + +def _request( + *, + role: str = "author", + kind: str = "issue", + number: int = 643, + intent: str = "implement request preview and initiation", + head: str | None = None, +) -> request_service.WorkRequest: + parsed, error = request_service.parse_request( + { + "desired_role": role, + "work_kind": kind, + "work_number": number, + "intent_summary": intent, + "expected_head_sha": head, + **SCOPE, + } + ) + assert error is None, error + assert parsed is not None + return parsed + + +def _selection( + *, kind: str = "issue", number: int = 643, head_sha: str | None = None +) -> dict[str, Any]: + return { + "kind": kind, + "number": number, + "title": "Web Console: Requests, intent preview, authorization", + "head_sha": head_sha, + "selected_action": "implement", + "expected_role_next": "author", + } + + +def _fake_allocator( + *, + selection: dict[str, Any] | None = None, + preview_outcome: str = allocator_service.OUTCOME_PREVIEW, + apply_outcome: str = allocator_service.OUTCOME_ASSIGNED, + assignment: dict[str, Any] | None = None, + calls: list[dict[str, Any]] | None = None, +): + """Build an allocator double that records how it was called.""" + chosen = selection if selection is not None else _selection() + made = ( + assignment + if assignment is not None + else { + "assignment_id": "asn-test-0001", + "lease_id": "lease-test-0001", + "session_id": "webui-request-test", + "expected_head_sha": chosen.get("head_sha"), + } + ) + + def _allocator(*, request, apply, expected_candidate_set_fingerprint=None): + if calls is not None: + calls.append( + { + "apply": apply, + "role": request.desired_role, + "fingerprint": expected_candidate_set_fingerprint, + } + ) + return { + "outcome": apply_outcome if apply else preview_outcome, + "selected": dict(chosen), + "reasons": ["fake allocator"], + "candidate_set_fingerprint": "fp-test", + "candidate_count": 3, + "inventory_complete": True, + "selection_policy": allocator_service.SELECTION_POLICY, + "substrate": "control_plane_db", + "assignment": dict(made) if apply else None, + } + + return _allocator + + +def _no_claims(_request): + return {} + + +def _claimed(role: str = "author"): + def _source(request): + return { + request.work_key: { + "lease_id": "lease-foreign-9999", + "session_id": "prgs-author-999-foreign", + "role": role, + "expires_at": "2026-07-25T09:10:37Z", + } + } + + return _source + + +class TestRequestParsing(unittest.TestCase): + """The request model rejects rather than guesses.""" + + def test_valid_request_round_trips(self): + req = _request() + self.assertEqual(req.work_key, ("issue", 643)) + self.assertEqual(req.display_ref, "#643") + self.assertEqual(req.to_dict()["desired_role"], "author") + + def test_unknown_role_rejected(self): + parsed, error = request_service.parse_request( + { + "desired_role": "admin", + "work_kind": "issue", + "work_number": 1, + "intent_summary": "x", + **SCOPE, + } + ) + self.assertIsNone(parsed) + self.assertEqual(error.reason_code, "unknown_role") + self.assertEqual(error.field_name, "desired_role") + + def test_unknown_work_kind_rejected(self): + parsed, error = request_service.parse_request( + { + "desired_role": "author", + "work_kind": "branch", + "work_number": 1, + "intent_summary": "x", + **SCOPE, + } + ) + self.assertIsNone(parsed) + self.assertEqual(error.reason_code, "unknown_work_kind") + + def test_non_positive_number_rejected(self): + for value in (0, -3): + with self.subTest(value=value): + parsed, error = request_service.parse_request( + { + "desired_role": "author", + "work_kind": "issue", + "work_number": value, + "intent_summary": "x", + **SCOPE, + } + ) + self.assertIsNone(parsed) + self.assertEqual(error.reason_code, "invalid_work_number") + + def test_missing_intent_rejected(self): + parsed, error = request_service.parse_request( + { + "desired_role": "author", + "work_kind": "issue", + "work_number": 1, + **SCOPE, + } + ) + self.assertIsNone(parsed) + self.assertEqual(error.reason_code, "missing_intent") + + def test_intent_is_bounded(self): + req = _request(intent="x" * 5000) + self.assertEqual(len(req.intent_summary), request_service.MAX_INTENT_CHARS) + + def test_unresolved_scope_rejected(self): + parsed, error = request_service.parse_request( + { + "desired_role": "author", + "work_kind": "issue", + "work_number": 1, + "intent_summary": "x", + } + ) + self.assertIsNone(parsed) + self.assertEqual(error.reason_code, "scope_unresolved") + + def test_default_scope_fills_missing_fields(self): + parsed, error = request_service.parse_request( + { + "desired_role": "author", + "work_kind": "issue", + "work_number": 7, + "intent_summary": "x", + }, + default_scope=SCOPE, + ) + self.assertIsNone(error) + self.assertEqual(parsed.repo, "Gitea-Tools") + + +class TestPreviewAuthorizeDeny(unittest.TestCase): + """AC1 — preview shows authorize/deny with reasons.""" + + def test_authorized_preview_names_every_check(self): + preview = request_service.preview_request( + _request(), + principal=_principal(console_authz.OPERATOR), + allocator=_fake_allocator(), + claims_source=_no_claims, + audit=False, + ) + self.assertTrue(preview.authorized) + self.assertEqual( + {c.name for c in preview.checks}, + { + request_service.CHECK_AUTHORIZATION, + request_service.CHECK_CAPABILITY, + request_service.CHECK_LEASE_AVAILABILITY, + request_service.CHECK_NEXT_SAFE_ACTION, + request_service.CHECK_HEAD_PIN, + }, + ) + self.assertEqual(preview.required_profile, "prgs-author") + self.assertEqual(preview.required_namespace, "gitea-author") + + def test_every_check_carries_a_reason(self): + preview = request_service.preview_request( + _request(), + principal=_principal(console_authz.OPERATOR), + allocator=_fake_allocator(), + claims_source=_no_claims, + audit=False, + ) + for check in preview.checks: + with self.subTest(check=check.name): + self.assertTrue(check.reason_code.strip()) + self.assertTrue(check.detail.strip()) + + def test_anonymous_preview_denied_with_reason(self): + preview = request_service.preview_request( + _request(), + allocator=_fake_allocator(), + claims_source=_no_claims, + audit=False, + ) + self.assertFalse(preview.authorized) + self.assertEqual(preview.reason_code, console_authz.DENY_UNAUTHENTICATED) + + def test_viewer_preview_denied_for_insufficient_role(self): + preview = request_service.preview_request( + _request(), + principal=_principal(console_authz.VIEWER), + allocator=_fake_allocator(), + claims_source=_no_claims, + audit=False, + ) + self.assertFalse(preview.authorized) + self.assertEqual(preview.reason_code, console_authz.DENY_INSUFFICIENT_ROLE) + + def test_denied_preview_never_reaches_the_allocator(self): + """A denial must not double as a queue oracle.""" + calls: list[dict[str, Any]] = [] + request_service.preview_request( + _request(), + principal=_principal(console_authz.VIEWER), + allocator=_fake_allocator(calls=calls), + claims_source=_no_claims, + audit=False, + ) + self.assertEqual(calls, []) + + def test_preview_lists_prohibited_actions(self): + preview = request_service.preview_request( + _request(role="author"), + principal=_principal(console_authz.OPERATOR), + allocator=_fake_allocator(), + claims_source=_no_claims, + audit=False, + ) + self.assertIn("merge", preview.prohibited_actions) + self.assertIn("approve", preview.prohibited_actions) + + def test_preview_reports_next_safe_action(self): + preview = request_service.preview_request( + _request(), + principal=_principal(console_authz.OPERATOR), + allocator=_fake_allocator(), + claims_source=_no_claims, + audit=False, + ) + self.assertIn("issue #643", preview.next_safe_action) + + def test_preview_never_mutates(self): + calls: list[dict[str, Any]] = [] + request_service.preview_request( + _request(), + principal=_principal(console_authz.OPERATOR), + allocator=_fake_allocator(calls=calls), + claims_source=_no_claims, + audit=False, + ) + self.assertEqual([c["apply"] for c in calls], [False]) + + +class TestPreviewFailClosed(unittest.TestCase): + """Missing evidence denies; it never reads as an absence of obstacles.""" + + def test_unreadable_claim_inventory_denies(self): + def _boom(_request): + raise RuntimeError("db unavailable") + + preview = request_service.preview_request( + _request(), + principal=_principal(console_authz.OPERATOR), + allocator=_fake_allocator(), + claims_source=_boom, + audit=False, + ) + self.assertFalse(preview.authorized) + self.assertEqual( + preview.reason_code, request_service.REASON_EVIDENCE_UNAVAILABLE + ) + + def test_allocator_failure_denies(self): + def _boom(**_kwargs): + raise RuntimeError("allocator exploded") + + preview = request_service.preview_request( + _request(), + principal=_principal(console_authz.OPERATOR), + allocator=_boom, + claims_source=_no_claims, + audit=False, + ) + self.assertFalse(preview.authorized) + self.assertEqual( + preview.reason_code, request_service.REASON_EVIDENCE_UNAVAILABLE + ) + + def test_allocator_selecting_other_work_denies(self): + preview = request_service.preview_request( + _request(number=643), + principal=_principal(console_authz.OPERATOR), + allocator=_fake_allocator(selection=_selection(number=999)), + claims_source=_no_claims, + audit=False, + ) + self.assertFalse(preview.authorized) + self.assertEqual(preview.reason_code, request_service.REASON_NOT_NEXT_SAFE) + self.assertIn("#999", preview.detail) + + def test_pr_without_head_sha_denies(self): + preview = request_service.preview_request( + _request(role="reviewer", kind="pr", number=898), + principal=_principal(console_authz.OPERATOR), + allocator=_fake_allocator( + selection=_selection(kind="pr", number=898, head_sha=None) + ), + claims_source=_no_claims, + audit=False, + ) + self.assertFalse(preview.authorized) + self.assertEqual( + preview.reason_code, request_service.REASON_EVIDENCE_UNAVAILABLE + ) + + def test_pr_head_moved_denies(self): + preview = request_service.preview_request( + _request(role="reviewer", kind="pr", number=898, head="a" * 40), + principal=_principal(console_authz.OPERATOR), + allocator=_fake_allocator( + selection=_selection(kind="pr", number=898, head_sha="b" * 40) + ), + claims_source=_no_claims, + audit=False, + ) + self.assertFalse(preview.authorized) + self.assertEqual(preview.reason_code, "head_moved") + + def test_pr_head_matching_passes(self): + preview = request_service.preview_request( + _request(role="reviewer", kind="pr", number=898, head="b" * 40), + principal=_principal(console_authz.OPERATOR), + allocator=_fake_allocator( + selection=_selection(kind="pr", number=898, head_sha="b" * 40) + ), + claims_source=_no_claims, + audit=False, + ) + self.assertTrue(preview.authorized) + + +class TestApplyExecutionGate(unittest.TestCase): + """Execution stays wired off unless an operator opts in explicitly.""" + + def test_action_is_registered_and_unwired_by_default(self): + action = console_authz.get_action(request_service.ACTION_ID) + self.assertIsNotNone(action) + self.assertEqual(action.phase, 2) + self.assertEqual(action.minimum_role, console_authz.OPERATOR) + self.assertTrue(action.requires_confirmation) + self.assertFalse(console_authz.execution_wired(action, env={})) + + def test_flag_named_but_unset_does_not_wire(self): + action = console_authz.get_action(request_service.ACTION_ID) + self.assertFalse(console_authz.execution_wired(action, env={EXEC_FLAG: "no"})) + self.assertTrue(console_authz.execution_wired(action, env={EXEC_FLAG: "1"})) + + def test_opting_in_wires_only_this_action(self): + env = {EXEC_FLAG: "1"} + for action_id, action in console_authz.ACTIONS.items(): + with self.subTest(action=action_id): + self.assertEqual( + console_authz.execution_wired(action, env=env), + action_id == request_service.ACTION_ID, + ) + + def test_apply_denied_while_unwired(self): + with mock.patch.dict(os.environ, {EXEC_FLAG: ""}): + result = request_service.apply_request( + _request(), + principal=_principal(console_authz.OPERATOR), + confirm=True, + allocator=_fake_allocator(), + claims_source=_no_claims, + ) + self.assertFalse(result["ok"]) + self.assertEqual(result["outcome"], request_service.OUTCOME_DENIED) + self.assertEqual(result["reason_code"], request_service.REASON_UNAUTHORIZED) + self.assertFalse(result["mutation_performed"]) + + +class TestApplyOutcomes(unittest.TestCase): + """AC2/AC3/AC4 — assignment, wait, blocked, and duplicate rejection.""" + + def setUp(self): + patcher = mock.patch.dict(os.environ, {EXEC_FLAG: "1"}) + patcher.start() + self.addCleanup(patcher.stop) + + def test_apply_creates_exclusive_assignment(self): + calls: list[dict[str, Any]] = [] + result = request_service.apply_request( + _request(), + principal=_principal(console_authz.OPERATOR), + confirm=True, + allocator=_fake_allocator(calls=calls), + claims_source=_no_claims, + ) + self.assertTrue(result["ok"]) + self.assertEqual(result["outcome"], allocator_service.OUTCOME_ASSIGNED) + self.assertEqual(result["assignment"]["assignment_id"], "asn-test-0001") + self.assertTrue(result["mutation_performed"]) + self.assertEqual(result["status_code"], 201) + # Dry-run first, then apply — never apply alone. + self.assertEqual([c["apply"] for c in calls], [False, True]) + # The apply call carries the fingerprint the dry-run produced. + self.assertEqual(calls[1]["fingerprint"], "fp-test") + + def test_assignment_returns_a_role_handoff(self): + result = request_service.apply_request( + _request(), + principal=_principal(console_authz.OPERATOR), + confirm=True, + allocator=_fake_allocator(), + claims_source=_no_claims, + ) + handoff = result["handoff"] + self.assertEqual(handoff["required_profile"], "prgs-author") + self.assertEqual(handoff["required_namespace"], "gitea-author") + self.assertEqual(handoff["assignment_id"], "asn-test-0001") + self.assertIn("merge", handoff["forbidden_actions"]) + + def test_unconfirmed_apply_refuses_before_the_allocator(self): + calls: list[dict[str, Any]] = [] + result = request_service.apply_request( + _request(), + principal=_principal(console_authz.OPERATOR), + confirm=False, + allocator=_fake_allocator(calls=calls), + claims_source=_no_claims, + ) + self.assertFalse(result["ok"]) + self.assertEqual( + result["reason_code"], request_service.REASON_CONFIRMATION_REQUIRED + ) + self.assertEqual(calls, []) + + def test_duplicate_assignment_rejected(self): + """AC3 — an active lease on the work unit blocks a second assign.""" + calls: list[dict[str, Any]] = [] + result = request_service.apply_request( + _request(), + principal=_principal(console_authz.OPERATOR), + confirm=True, + allocator=_fake_allocator(calls=calls), + claims_source=_claimed(), + ) + self.assertFalse(result["ok"]) + self.assertEqual(result["outcome"], request_service.OUTCOME_BLOCKED) + self.assertEqual( + result["reason_code"], request_service.REASON_DUPLICATE_ASSIGNMENT + ) + self.assertFalse(result["mutation_performed"]) + # The dry-run ran; the apply never did. + self.assertEqual([c["apply"] for c in calls], [False]) + + def test_not_next_safe_work_returns_wait_without_applying(self): + calls: list[dict[str, Any]] = [] + result = request_service.apply_request( + _request(number=643), + principal=_principal(console_authz.OPERATOR), + confirm=True, + allocator=_fake_allocator( + selection=_selection(number=999), calls=calls + ), + claims_source=_no_claims, + ) + self.assertFalse(result["ok"]) + self.assertEqual(result["outcome"], request_service.OUTCOME_WAIT) + self.assertEqual(result["reason_code"], request_service.REASON_NOT_NEXT_SAFE) + self.assertEqual([c["apply"] for c in calls], [False]) + + def test_allocator_declining_on_apply_returns_blocked(self): + result = request_service.apply_request( + _request(), + principal=_principal(console_authz.OPERATOR), + confirm=True, + allocator=_fake_allocator( + apply_outcome=allocator_service.OUTCOME_BLOCKED_LEASE, + assignment={}, + ), + claims_source=_no_claims, + ) + self.assertFalse(result["ok"]) + self.assertEqual(result["outcome"], request_service.OUTCOME_BLOCKED) + self.assertEqual( + result["reason_code"], request_service.REASON_ALLOCATOR_OUTCOME + ) + self.assertFalse(result["mutation_performed"]) + + def test_allocator_drift_on_apply_is_not_read_as_an_assignment(self): + """The apply call must return *this* work unit, not a substitute.""" + + def _drifting(*, request, apply, expected_candidate_set_fingerprint=None): + return { + "outcome": ( + allocator_service.OUTCOME_ASSIGNED + if apply + else allocator_service.OUTCOME_PREVIEW + ), + "selected": _selection(number=999 if apply else 643), + "assignment": {"assignment_id": "asn-wrong"} if apply else None, + "candidate_set_fingerprint": "fp-test", + } + + result = request_service.apply_request( + _request(number=643), + principal=_principal(console_authz.OPERATOR), + confirm=True, + allocator=_drifting, + claims_source=_no_claims, + ) + self.assertFalse(result["ok"]) + self.assertIsNone(result["assignment"]) + self.assertFalse(result["mutation_performed"]) + + def test_viewer_cannot_apply(self): + result = request_service.apply_request( + _request(), + principal=_principal(console_authz.VIEWER), + confirm=True, + allocator=_fake_allocator(), + claims_source=_no_claims, + ) + self.assertFalse(result["ok"]) + self.assertEqual(result["reason_code"], request_service.REASON_UNAUTHORIZED) + + def test_anonymous_cannot_apply(self): + result = request_service.apply_request( + _request(), + confirm=True, + allocator=_fake_allocator(), + claims_source=_no_claims, + ) + self.assertFalse(result["ok"]) + self.assertFalse(result["mutation_performed"]) + + +class TestAllocatorIntegrationFakes(unittest.TestCase): + """The real default allocator refuses a partial inventory (#758).""" + + def test_incomplete_inventory_returns_none(self): + from webui.queue_loader import PaginationMeta, QueueSnapshot + + snapshot = QueueSnapshot( + project_id="p", + repo_label="r", + prs=(), + issues=(), + pr_pagination=PaginationMeta( + page=1, + per_page=50, + returned_count=50, + has_more=True, + is_final_page=False, + inventory_complete=False, + pages_fetched=1, + ), + issue_pagination=None, + ) + with mock.patch( + "webui.queue_loader.load_queue_snapshot", return_value=snapshot + ): + result = request_service.default_allocator( + request=_request(), apply=False + ) + self.assertIsNone(result) + + def test_fetch_error_returns_none(self): + from webui.queue_loader import QueueSnapshot + + snapshot = QueueSnapshot( + project_id="p", + repo_label="r", + prs=(), + issues=(), + pr_pagination=None, + issue_pagination=None, + fetch_error="no credentials", + ) + with mock.patch( + "webui.queue_loader.load_queue_snapshot", return_value=snapshot + ): + result = request_service.default_allocator( + request=_request(), apply=True + ) + self.assertIsNone(result) + + +class TestAuditRecords(unittest.TestCase): + """Every preview and apply is auditable, correlated, and redacted.""" + + def setUp(self): + handle = tempfile.NamedTemporaryFile( + mode="w", suffix=".jsonl", delete=False + ) + handle.close() + self.sink = handle.name + self.addCleanup( + lambda: os.path.exists(self.sink) and os.remove(self.sink) + ) + patcher = mock.patch.dict( + os.environ, {console_audit.AUDIT_LOG_ENV: self.sink} + ) + patcher.start() + self.addCleanup(patcher.stop) + + def _records(self) -> list[dict[str, Any]]: + with open(self.sink, encoding="utf-8") as handle: + return [json.loads(line) for line in handle if line.strip()] + + def test_preview_is_audited_with_a_correlation_id(self): + preview = request_service.preview_request( + _request(), + principal=_principal(console_authz.OPERATOR), + allocator=_fake_allocator(), + claims_source=_no_claims, + ) + records = self._records() + self.assertEqual(len(records), 1) + record = records[0] + self.assertEqual(record["action"], request_service.ACTION_ID) + self.assertEqual(record["result"], console_audit.RESULT_PREVIEWED) + self.assertEqual( + record["correlation"]["request_id"], preview.correlation_id + ) + self.assertEqual(record["target"]["ref"], "#643") + + def test_denied_apply_is_audited(self): + with mock.patch.dict(os.environ, {EXEC_FLAG: ""}): + request_service.apply_request( + _request(), + principal=_principal(console_authz.VIEWER), + confirm=True, + allocator=_fake_allocator(), + claims_source=_no_claims, + ) + self.assertEqual(self._records()[-1]["result"], console_audit.RESULT_DENIED) + + def test_assignment_is_audited_and_correlated(self): + with mock.patch.dict(os.environ, {EXEC_FLAG: "1"}): + result = request_service.apply_request( + _request(), + principal=_principal(console_authz.OPERATOR), + confirm=True, + allocator=_fake_allocator(), + claims_source=_no_claims, + ) + record = self._records()[-1] + self.assertEqual(record["result"], console_audit.RESULT_SUCCEEDED) + self.assertEqual( + record["correlation"]["request_id"], result["correlation_id"] + ) + self.assertEqual( + record["metadata"]["assignment_id"], + result["assignment"]["assignment_id"], + ) + + def test_intent_bearing_a_secret_is_not_persisted_raw(self): + request_service.preview_request( + _request(intent="use token=ghp_ABCDEFGHIJKLMNOPQRSTUVWXYZ0123456789"), + principal=_principal(console_authz.OPERATOR), + allocator=_fake_allocator(), + claims_source=_no_claims, + ) + for record in self._records(): + with self.subTest(event=record.get("event_id")): + self.assertFalse(scan_for_secrets(record)) + + +class TestRequestRoutes(unittest.TestCase): + """The HTTP surface: form page, preview API, apply API.""" + + def setUp(self): + self.client = TestClient(create_app()) + + def test_requests_page_renders_form(self): + response = self.client.get("/requests") + self.assertEqual(response.status_code, 200) + body = response.text + self.assertIn("Requests", body) + self.assertIn("desired_role", body) + self.assertIn("intent_summary", body) + + def test_requests_page_is_linked_from_nav(self): + from webui.nav import nav_hrefs + + self.assertIn("/requests", nav_hrefs()) + + def test_preview_api_rejects_an_invalid_request(self): + response = self.client.post( + "/api/v1/requests/preview", + json={ + "desired_role": "wizard", + "work_kind": "issue", + "work_number": 1, + "intent_summary": "x", + **SCOPE, + }, + ) + self.assertEqual(response.status_code, 400) + self.assertEqual(response.json()["reason_code"], "unknown_role") + + def test_preview_api_denies_anonymous(self): + response = self.client.post( + "/api/v1/requests/preview", + json={ + "desired_role": "author", + "work_kind": "issue", + "work_number": 643, + "intent_summary": "x", + **SCOPE, + }, + ) + self.assertEqual(response.status_code, 403) + payload = response.json() + self.assertFalse(payload["authorized"]) + self.assertFalse(payload["mutation_performed"]) + + def test_apply_api_denies_anonymous(self): + response = self.client.post( + "/api/v1/requests/apply", + json={ + "desired_role": "author", + "work_kind": "issue", + "work_number": 643, + "intent_summary": "x", + "confirm": True, + **SCOPE, + }, + ) + self.assertEqual(response.status_code, 403) + payload = response.json() + self.assertFalse(payload["ok"]) + self.assertFalse(payload["mutation_performed"]) + self.assertIsNone(payload["assignment"]) + + def test_apply_api_rejects_an_invalid_request(self): + response = self.client.post( + "/api/v1/requests/apply", + json={ + "desired_role": "author", + "work_kind": "issue", + "work_number": -1, + "intent_summary": "x", + **SCOPE, + }, + ) + self.assertEqual(response.status_code, 400) + + def test_request_apis_are_post_only(self): + """GET is not a way in. The app's 405 handler renders a read-only + method against a write route as 404, so that is what is asserted.""" + for path in ("/api/v1/requests/preview", "/api/v1/requests/apply"): + with self.subTest(path=path): + self.assertEqual(self.client.get(path).status_code, 404) + + def test_form_post_previews_and_never_assigns(self): + response = self.client.post( + "/requests", + data={ + "desired_role": "author", + "work_kind": "issue", + "work_number": "643", + "intent_summary": "implement the request surface", + "remote": "prgs", + "org": "Scaled-Tech-Consulting", + "repo": "Gitea-Tools", + }, + ) + self.assertEqual(response.status_code, 200) + self.assertIn("Intent preview", response.text) + + +class TestRenderingSafety(unittest.TestCase): + """AC5 — the page escapes hostile input and shows no secret.""" + + def test_intent_is_escaped(self): + preview = request_service.preview_request( + _request(intent=""), + principal=_principal(console_authz.OPERATOR), + allocator=_fake_allocator(), + claims_source=_no_claims, + audit=False, + ) + html = render_requests_page(preview=preview) + self.assertNotIn("", html) + self.assertIn("<script>", html) + + def test_page_renders_a_denial_without_a_preview(self): + _, error = request_service.parse_request( + { + "desired_role": "wizard", + "work_kind": "issue", + "work_number": 1, + "intent_summary": "x", + **SCOPE, + } + ) + html = render_requests_page(error=error) + self.assertIn("Request rejected", html) + self.assertIn("unknown_role", html) + + def test_page_shows_no_credential_material(self): + preview = request_service.preview_request( + _request(), + principal=_principal(console_authz.OPERATOR), + allocator=_fake_allocator(), + claims_source=_no_claims, + audit=False, + ) + html = render_requests_page(preview=preview) + for needle in ("token=", "Bearer ", "password"): + with self.subTest(needle=needle): + self.assertNotIn(needle, html) + + +if __name__ == "__main__": # pragma: no cover + unittest.main() diff --git a/webui/app.py b/webui/app.py index cedcef1..9526f49 100644 --- a/webui/app.py +++ b/webui/app.py @@ -67,6 +67,8 @@ from webui.system_health import ( snapshot_to_dict as system_health_to_dict, ) from webui.system_health_views import render_system_health_page +from webui import request_service +from webui.request_views import render_requests_page _READ_ONLY_METHODS = frozenset({"GET", "HEAD", "OPTIONS"}) _AUDIT_MUTATION_PATHS = frozenset({"/audit", "/api/audit"}) @@ -722,6 +724,109 @@ async def api_v1_analytics_ingest(request: Request) -> JSONResponse: ) +def _default_request_scope() -> dict[str, str]: + """Resolve remote/org/repo from the project registry for request forms. + + Returns an empty mapping when the registry cannot be read, which makes + ``parse_request`` reject a request that did not name its own scope rather + than letting it default to some other repository. + """ + from webui.queue_loader import _host_from_url # host normalisation helper + + registry, error = _load_project_registry() + if error is not None or not registry.projects: + return {} + project = registry.projects[0] + host = _host_from_url(project.remote_host) + return { + "remote": _derive_remote(host), + "org": project.gitea_owner or "", + "repo": project.repo_name or "", + } + + +async def _request_payload(request: Request) -> dict[str, object]: + """Read a request body as JSON or form-encoded. Never raises.""" + content_type = (request.headers.get("content-type") or "").lower() + if "application/json" in content_type: + try: + body = await request.json() + except Exception: + return {} + return dict(body) if isinstance(body, dict) else {} + try: + form = await request.form() + except Exception: + return {} + return {key: form[key] for key in form} + + +async def requests_page(request: Request) -> HTMLResponse: + """Operator request form and intent preview (#643). + + POST here only ever *previews*. Initiation is a separate confirmed call to + ``/api/v1/requests/apply`` so that submitting this form cannot reserve + work as a side effect. + """ + submitted: dict[str, object] = {} + preview = None + error = None + if request.method == "POST": + submitted = await _request_payload(request) + work_request, error = request_service.parse_request( + submitted, default_scope=_default_request_scope() + ) + if work_request is not None: + preview = request_service.preview_request( + work_request, + principal=resolve_principal(headers=dict(request.headers)), + ) + return HTMLResponse( + render_requests_page( + preview=preview, error=error, submitted=submitted + ) + ) + + +async def api_v1_request_preview(request: Request) -> JSONResponse: + """Dry-run authorization and intent preview for a work request (#643).""" + payload = await _request_payload(request) + work_request, error = request_service.parse_request( + payload, default_scope=_default_request_scope() + ) + if work_request is None: + return JSONResponse(error.to_dict(), status_code=400) + preview = request_service.preview_request( + work_request, + principal=resolve_principal(headers=dict(request.headers)), + ) + return JSONResponse( + preview.to_dict(), status_code=200 if preview.authorized else 403 + ) + + +async def api_v1_request_apply(request: Request) -> JSONResponse: + """Initiate a previewed work request through the allocator (#643). + + Fail-closed at every step: unauthorized, unconfirmed, not-next-safe, and + already-claimed all return without attempting an assignment. + """ + payload = await _request_payload(request) + work_request, error = request_service.parse_request( + payload, default_scope=_default_request_scope() + ) + if work_request is None: + return JSONResponse(error.to_dict(), status_code=400) + confirm = _truthy_flag(str(payload.get("confirm") or "")) + result = request_service.apply_request( + work_request, + principal=resolve_principal(headers=dict(request.headers)), + confirm=confirm, + ) + status = int(result.pop("status_code", 403)) + return JSONResponse(result, status_code=status) + + async def method_not_allowed(request: Request, _exc: Exception) -> Response: path = request.url.path if path in _AUDIT_MUTATION_PATHS and request.method == "POST": @@ -786,6 +891,17 @@ def create_app(*, bind_host: str | None = None) -> Starlette: api_action_attempt, methods=["POST"], ), + Route("/requests", requests_page, methods=["GET", "POST"]), + Route( + "/api/v1/requests/preview", + api_v1_request_preview, + methods=["POST"], + ), + Route( + "/api/v1/requests/apply", + api_v1_request_apply, + methods=["POST"], + ), Route("/api/leases", api_leases, methods=["GET"]), Route("/api/v1/inventory", api_inventory, methods=["GET"]), Route( diff --git a/webui/console_authz.py b/webui/console_authz.py index 0c51037..282d831 100644 --- a/webui/console_authz.py +++ b/webui/console_authz.py @@ -115,6 +115,12 @@ class ConsoleAction: break_glass: bool phase: int summary: str + # Opt-in switch for an action whose execution path is genuinely wired + # ahead of its phase becoming globally active (#643). Naming a variable + # here enables nothing on its own: the variable must also be set in the + # environment. An action that leaves this ``None`` can only execute once + # ACTIVE_PHASE reaches its phase, exactly as before. + execution_env_flag: str | None = None @property def mcp_permission(self) -> str: @@ -277,6 +283,27 @@ _ACTION_SPECS: tuple[ConsoleAction, ...] = ( phase=2, summary="Restart one MCP namespace via the host supervisor.", ), + # #643: submit a work request — desired role, issue/PR, intent — and let + # the allocator reserve it. This is the one Phase 2 action whose execution + # path is actually implemented (``webui.request_service``), so it carries + # the opt-in flag; it stays denied until an operator sets that variable. + # Authority is operator-class because the outcome is a claim, not a Gitea + # verdict: initiating reviewer or merger *work* does not grant the right + # to approve or merge, which stays with the MCP role profile. + ConsoleAction( + action_id="initiate_workflow", + task_key="allocate_next_work", + action_class=CLASS_WRITE, + minimum_role=OPERATOR, + requires_confirmation=True, + dual_control=False, + break_glass=False, + phase=2, + summary=( + "Preview and initiate allocator-owned workflow work for a role." + ), + execution_env_flag="WEBUI_REQUESTS_EXECUTION", + ), ) ACTIONS: dict[str, ConsoleAction] = {a.action_id: a for a in _ACTION_SPECS} @@ -430,6 +457,33 @@ ALLOW_PREVIEW = "allowed_preview_only" # gated on this model landing; nothing here enables it. ACTIVE_PHASE = 1 +_TRUTHY = frozenset({"1", "true", "yes", "on"}) + + +def execution_wired( + action: ConsoleAction | None, env: dict[str, str] | None = None +) -> bool: + """Whether *action* has a live execution path right now. + + Two ways to be wired, and only two. The action's phase is active, or the + action declares an opt-in environment variable *and* that variable is set. + Everything else — including every action that never declares a flag — is + unwired, so the default across the registry stays deny. + + Bumping ``ACTIVE_PHASE`` would enable execution for every action of that + phase at once. The per-action flag exists so a single implemented action + can go live without dragging its unimplemented phase-mates with it. + """ + if action is None: + return False + if action.phase <= ACTIVE_PHASE: + return True + flag = (action.execution_env_flag or "").strip() + if not flag: + return False + source = env if env is not None else os.environ + return (source.get(flag) or "").strip().lower() in _TRUTHY + @dataclass(frozen=True) class AuthorizationDecision: @@ -469,16 +523,19 @@ def authorize( principal: Principal | None = None, *, for_execution: bool = False, + env: dict[str, str] | None = None, ) -> AuthorizationDecision: """Decide whether *principal* may invoke *action_id*. Deny by default. ``for_execution`` distinguishes a read-only preview from a real invocation. - Even an allowed decision reports ``execution_enabled=False`` while the - console is in Phase 1, so no caller can read an allow as permission to - mutate. + ``execution_enabled`` reports whether the action has a live execution path + at all (:func:`execution_wired`) — for every action without an explicit + opt-in flag that stays ``False`` while the console is in Phase 1, so no + caller can read an allow as permission to mutate. """ who = principal if principal is not None else ANONYMOUS action = get_action(action_id) + wired = execution_wired(action, env) if action is None: return AuthorizationDecision( @@ -497,7 +554,7 @@ def authorize( "requires_confirmation": action.requires_confirmation, "dual_control": action.dual_control, "break_glass": action.break_glass, - "execution_enabled": False, + "execution_enabled": wired, } if not who.authenticated: @@ -530,13 +587,19 @@ def authorize( **base, ) - if for_execution and action.phase > ACTIVE_PHASE: + if for_execution and not wired: return AuthorizationDecision( allowed=False, reason_code=DENY_PHASE_NOT_ACTIVE, detail=( f"Action {action_id!r} belongs to phase {action.phase}; the " - f"console is in phase {ACTIVE_PHASE}. Execution is not wired." + f"console is in phase {ACTIVE_PHASE}" + + ( + f" and {action.execution_env_flag} is not set" + if action.execution_env_flag + else "" + ) + + ". Execution is not wired." ), **base, ) @@ -545,8 +608,8 @@ def authorize( allowed=True, reason_code=ALLOW_PREVIEW, detail=( - "Principal holds the required role. Preview only — execution " - "remains disabled until the Phase 2 action framework ships." + "Principal holds the required role. Execution proceeds only for an " + "action with a wired execution path; everything else is preview." ), **base, ) diff --git a/webui/nav.py b/webui/nav.py index 129b2a4..dbde5e1 100644 --- a/webui/nav.py +++ b/webui/nav.py @@ -45,6 +45,7 @@ NAV_GROUPS: tuple[NavGroup, ...] = ( NavItem("/queue", "Queue"), NavItem("/leases", "Leases"), NavItem("/actions", "Actions"), + NavItem("/requests", "Requests"), )), NavGroup("Runtime/Sessions", ( NavItem("/runtime", "Runtime health"), diff --git a/webui/request_service.py b/webui/request_service.py new file mode 100644 index 0000000..6941027 --- /dev/null +++ b/webui/request_service.py @@ -0,0 +1,967 @@ +"""Operator work-request preview and initiation (#643, Phase 2). + +An operator's alternative to pasting a role prompt into a terminal. A +*request* names three things — the role to run as, the issue or PR to run +against, and what the operator intends — and this module answers two questions +about it: + +* **Preview** (:func:`preview_request`) — would that request be authorized, + is the work unit actually free, is it the next safe thing that role should + touch, and which actions stay prohibited? Read-only, always. It creates no + assignment and never mutates. +* **Initiate** (:func:`apply_request`) — turn an authorized request into an + *exclusive assignment*, and only ever through the allocator. + +Three invariants hold and are the reason this module exists rather than a +direct call to :func:`allocator_service.allocate_next_work` from a route: + +1. **The allocator remains the only source of exclusive ownership** (#600 / + #613). ``apply`` never assigns the requested item directly. It runs a + dry-run first and proceeds only when the allocator would independently pick + that exact item; otherwise it reports ``wait`` and mutates nothing. A + request is therefore a *confirmation* of the allocator's decision, never an + override of it. +2. **Duplicate assignment is rejected before it is attempted.** An active + claim on the work unit — held by any session, this one included — blocks. +3. **Fail closed at every unknown.** An unparseable request, an unavailable + control-plane DB, an incomplete queue inventory, or an unresolved + authorization all deny. There is no branch that proceeds on missing + evidence. + +Authorization comes from :mod:`webui.console_authz` (``initiate_workflow``) +and every outcome is audited through :mod:`webui.console_audit`, correlated to +the resulting assignment by ``correlation_id``. +""" + +from __future__ import annotations + +import uuid +from dataclasses import dataclass, field +from typing import Any, Callable, Mapping, Sequence + +import allocator_service +from task_capability_map import required_permission, required_role +from webui import console_audit, console_authz + +# The console action this module is gated by. Registered in console_authz. +ACTION_ID = "initiate_workflow" + +KIND_ISSUE = "issue" +KIND_PR = "pr" +WORK_KINDS: tuple[str, ...] = (KIND_ISSUE, KIND_PR) + +REQUESTABLE_ROLES: tuple[str, ...] = ( + allocator_service.ROLE_AUTHOR, + allocator_service.ROLE_REVIEWER, + allocator_service.ROLE_MERGER, + allocator_service.ROLE_RECONCILER, + allocator_service.ROLE_CONTROLLER, +) + +# Intent is operator prose echoed back into an audit record. Bounded so a +# pasted transcript cannot bloat the append-only log. +MAX_INTENT_CHARS = 500 + +# --- Outcomes --------------------------------------------------------------- +OUTCOME_ASSIGNED = allocator_service.OUTCOME_ASSIGNED +OUTCOME_WAIT = allocator_service.OUTCOME_WAIT +OUTCOME_BLOCKED = "blocked" +OUTCOME_DENIED = "denied" +OUTCOME_INVALID = "invalid_request" +OUTCOME_PREVIEW = allocator_service.OUTCOME_PREVIEW + +# --- Reason codes ----------------------------------------------------------- +REASON_AUTHORIZED = "request_authorized" +REASON_PREVIEW_OK = "preview_authorized" +REASON_UNAUTHORIZED = "unauthorized" +REASON_NOT_NEXT_SAFE = "not_next_safe_work" +REASON_DUPLICATE_ASSIGNMENT = "duplicate_assignment" +REASON_CONFIRMATION_REQUIRED = "confirmation_required" +REASON_EVIDENCE_UNAVAILABLE = "evidence_unavailable" +REASON_ALLOCATOR_OUTCOME = "allocator_declined" + +# --- Check names ------------------------------------------------------------ +CHECK_AUTHORIZATION = "authorization" +CHECK_CAPABILITY = "capability" +CHECK_LEASE_AVAILABILITY = "lease_availability" +CHECK_NEXT_SAFE_ACTION = "next_safe_action" +CHECK_HEAD_PIN = "head_pin" + + +# --- Request model ---------------------------------------------------------- + + +@dataclass(frozen=True) +class WorkRequest: + """One operator request: a role, a work unit, and a stated intent.""" + + desired_role: str + work_kind: str + work_number: int + intent_summary: str + remote: str + org: str + repo: str + expected_head_sha: str | None = None + + @property + def work_key(self) -> tuple[str, int]: + return (self.work_kind, self.work_number) + + @property + def display_ref(self) -> str: + return f"#{self.work_number}" + + def to_dict(self) -> dict[str, Any]: + return { + "desired_role": self.desired_role, + "work_kind": self.work_kind, + "work_number": self.work_number, + "intent_summary": self.intent_summary, + "remote": self.remote, + "org": self.org, + "repo": self.repo, + "expected_head_sha": self.expected_head_sha, + } + + +@dataclass(frozen=True) +class RequestError: + """A rejected request, with the field that caused the rejection.""" + + reason_code: str + detail: str + field_name: str | None = None + + def to_dict(self) -> dict[str, Any]: + return { + "ok": False, + "outcome": OUTCOME_INVALID, + "reason_code": self.reason_code, + "detail": self.detail, + "field": self.field_name, + } + + +def _clean(value: Any) -> str: + return str(value or "").strip() + + +def parse_request( + payload: Mapping[str, Any] | None, + *, + default_scope: Mapping[str, str] | None = None, +) -> tuple[WorkRequest | None, RequestError | None]: + """Validate an operator payload into a :class:`WorkRequest`. + + Returns ``(request, None)`` or ``(None, error)``. Never raises and never + guesses: an unknown role, an unknown work kind, or a non-positive number is + an error rather than a silently corrected value. + """ + body = dict(payload or {}) + scope = dict(default_scope or {}) + + role = _clean(body.get("desired_role") or body.get("role")).lower() + if role not in REQUESTABLE_ROLES: + return None, RequestError( + reason_code="unknown_role", + detail=( + f"desired_role must be one of {', '.join(REQUESTABLE_ROLES)}; " + f"got {role or '(empty)'!r}." + ), + field_name="desired_role", + ) + + kind = _clean(body.get("work_kind") or body.get("kind")).lower() + if kind not in WORK_KINDS: + return None, RequestError( + reason_code="unknown_work_kind", + detail=( + f"work_kind must be 'issue' or 'pr'; got {kind or '(empty)'!r}." + ), + field_name="work_kind", + ) + + raw_number = body.get("work_number") + if raw_number is None: + raw_number = ( + body.get("pr_number") if kind == KIND_PR else body.get("issue_number") + ) + if raw_number is None: + raw_number = body.get("number") + try: + number = int(str(raw_number).strip()) + except (TypeError, ValueError): + return None, RequestError( + reason_code="invalid_work_number", + detail=f"work_number must be an integer; got {raw_number!r}.", + field_name="work_number", + ) + if number <= 0: + return None, RequestError( + reason_code="invalid_work_number", + detail="work_number must be a positive issue or PR number.", + field_name="work_number", + ) + + intent = _clean(body.get("intent_summary") or body.get("intent")) + if not intent: + return None, RequestError( + reason_code="missing_intent", + detail="intent_summary is required so the audit record states why.", + field_name="intent_summary", + ) + intent = intent[:MAX_INTENT_CHARS] + + remote = _clean(body.get("remote")) or _clean(scope.get("remote")) + org = _clean(body.get("org")) or _clean(scope.get("org")) + repo = _clean(body.get("repo")) or _clean(scope.get("repo")) + if not (remote and org and repo): + return None, RequestError( + reason_code="scope_unresolved", + detail=( + "remote, org, and repo could not be resolved from the request " + "or the project registry." + ), + field_name="repo", + ) + + head = _clean(body.get("expected_head_sha")) or None + + return ( + WorkRequest( + desired_role=role, + work_kind=kind, + work_number=number, + intent_summary=intent, + remote=remote, + org=org, + repo=repo, + expected_head_sha=head, + ), + None, + ) + + +# --- Preview ---------------------------------------------------------------- + + +@dataclass(frozen=True) +class RequestCheck: + """One named precondition and its verdict.""" + + name: str + ok: bool + reason_code: str + detail: str + evidence: dict[str, Any] = field(default_factory=dict) + + def to_dict(self) -> dict[str, Any]: + return { + "name": self.name, + "ok": self.ok, + "reason_code": self.reason_code, + "detail": self.detail, + "evidence": dict(self.evidence), + } + + +@dataclass(frozen=True) +class RequestPreview: + """The full intent preview for one request. Read-only in every field.""" + + request: WorkRequest + authorized: bool + reason_code: str + detail: str + authorization: dict[str, Any] + checks: tuple[RequestCheck, ...] + prohibited_actions: tuple[str, ...] + allowed_actions: tuple[str, ...] + next_safe_action: str + required_profile: str + required_namespace: str + required_permission: str + correlation_id: str + allocator_evidence: dict[str, Any] = field(default_factory=dict) + + @property + def failed_checks(self) -> tuple[RequestCheck, ...]: + return tuple(c for c in self.checks if not c.ok) + + def to_dict(self) -> dict[str, Any]: + return { + "ok": self.authorized, + "outcome": OUTCOME_PREVIEW, + "dry_run": True, + "mutation_performed": False, + "authorized": self.authorized, + "reason_code": self.reason_code, + "detail": self.detail, + "request": self.request.to_dict(), + "authorization": dict(self.authorization), + "checks": [c.to_dict() for c in self.checks], + "failed_checks": [c.name for c in self.failed_checks], + "prohibited_actions": list(self.prohibited_actions), + "allowed_actions": list(self.allowed_actions), + "next_safe_action": self.next_safe_action, + "required_profile": self.required_profile, + "required_namespace": self.required_namespace, + "required_permission": self.required_permission, + "correlation_id": self.correlation_id, + "allocator_evidence": dict(self.allocator_evidence), + } + + +AllocatorFn = Callable[..., dict[str, Any] | None] +ClaimsFn = Callable[["WorkRequest"], Mapping[tuple[str, int], dict[str, Any]]] + + +def _correlation_id() -> str: + return f"req-{uuid.uuid4().hex}" + + +def _selection_matches( + selection: Mapping[str, Any] | None, request: WorkRequest +) -> bool: + if not selection: + return False + kind = _clean(selection.get("kind")).lower() + try: + number_int = int(selection.get("number")) + except (TypeError, ValueError): + return False + return (kind, number_int) == request.work_key + + +def _authorization_check( + decision: console_authz.AuthorizationDecision, +) -> RequestCheck: + return RequestCheck( + name=CHECK_AUTHORIZATION, + ok=bool(decision.allowed), + reason_code=decision.reason_code, + detail=decision.detail, + evidence={ + "subject": decision.principal.subject, + "role": decision.principal.role, + "required_role": decision.required_role, + "identity_source": decision.principal.identity_source, + }, + ) + + +def _capability_check(request: WorkRequest) -> RequestCheck: + """Whether the requested role maps to a declared MCP capability. + + The console never invents an authority: the permission and role come from + ``task_capability_map`` via the same ``allocate_next_work`` task the MCP + allocator gates on. + """ + # The remote-prefixed hint keeps a dadeschools request from being told to + # run under a prgs profile; ``required_profile_for_role`` preserves the + # prefix when one is present and falls back to its own default otherwise. + profile_hint = f"{request.remote}-{request.desired_role}" + try: + profile = allocator_service.required_profile_for_role( + request.desired_role, profile_name=profile_hint + ) + namespace = allocator_service.required_namespace_for_role( + request.desired_role, profile_name=profile_hint + ) + except Exception as exc: # noqa: BLE001 — an unresolved role is a denial + return RequestCheck( + name=CHECK_CAPABILITY, + ok=False, + reason_code="capability_unresolved", + detail=( + f"no profile/namespace maps to role {request.desired_role!r}: " + f"{exc}" + ), + ) + resolved = bool(profile and namespace) + return RequestCheck( + name=CHECK_CAPABILITY, + ok=resolved, + reason_code="capability_resolved" if resolved else "capability_unresolved", + detail=( + f"role {request.desired_role!r} runs under profile {profile!r} in " + f"MCP namespace {namespace!r}." + ), + evidence={ + "required_profile": profile, + "required_namespace": namespace, + "required_permission": required_permission("allocate_next_work"), + "capability_role": required_role("allocate_next_work"), + }, + ) + + +def _lease_check( + request: WorkRequest, + claims: Mapping[tuple[str, int], dict[str, Any]] | None, +) -> RequestCheck: + """Whether the work unit is free of an active claim. + + ``claims is None`` means the control-plane DB could not be read. That is a + failure, not an absence of claims: an unreadable substrate must never read + as "nothing holds this". + """ + if claims is None: + return RequestCheck( + name=CHECK_LEASE_AVAILABILITY, + ok=False, + reason_code=REASON_EVIDENCE_UNAVAILABLE, + detail=( + "active-claim inventory is unavailable; refusing to treat an " + "unreadable control-plane DB as an unclaimed work unit." + ), + ) + claim = claims.get(request.work_key) + if claim: + return RequestCheck( + name=CHECK_LEASE_AVAILABILITY, + ok=False, + reason_code=REASON_DUPLICATE_ASSIGNMENT, + detail=( + f"{request.work_kind} {request.display_ref} already carries an " + f"active {claim.get('role') or 'unknown'} lease." + ), + evidence={ + "lease_id": claim.get("lease_id"), + "session_id": claim.get("session_id"), + "role": claim.get("role"), + "expires_at": claim.get("expires_at"), + }, + ) + return RequestCheck( + name=CHECK_LEASE_AVAILABILITY, + ok=True, + reason_code="lease_available", + detail=f"no active lease holds {request.work_kind} {request.display_ref}.", + ) + + +def _next_safe_action_check( + request: WorkRequest, allocation: Mapping[str, Any] | None +) -> RequestCheck: + """Whether the allocator would independently select this exact work unit.""" + if not allocation: + return RequestCheck( + name=CHECK_NEXT_SAFE_ACTION, + ok=False, + reason_code=REASON_EVIDENCE_UNAVAILABLE, + detail="allocator dry-run produced no result; refusing to proceed.", + ) + selection = allocation.get("selected") or {} + outcome = _clean(allocation.get("outcome")) + if not _selection_matches(selection, request): + chosen = ( + f"{_clean(selection.get('kind')) or 'unknown'} #{selection.get('number')}" + if selection + else "nothing" + ) + return RequestCheck( + name=CHECK_NEXT_SAFE_ACTION, + ok=False, + reason_code=REASON_NOT_NEXT_SAFE, + detail=( + f"the allocator would select {chosen} for role " + f"{request.desired_role!r}, not {request.work_kind} " + f"{request.display_ref}. Requests confirm the allocator's " + "decision; they never override it." + ), + evidence={ + "allocator_outcome": outcome, + "allocator_selection": dict(selection), + "reasons": list(allocation.get("reasons") or ()), + }, + ) + return RequestCheck( + name=CHECK_NEXT_SAFE_ACTION, + ok=True, + reason_code="next_safe_work", + detail=( + f"the allocator selects {request.work_kind} {request.display_ref} " + f"for role {request.desired_role!r}." + ), + evidence={ + "allocator_outcome": outcome, + "selected_action": _clean(selection.get("selected_action")), + "expected_role_next": _clean(selection.get("expected_role_next")), + }, + ) + + +def _head_pin_check( + request: WorkRequest, allocation: Mapping[str, Any] | None +) -> RequestCheck: + """PR work must be pinned to a head SHA; issue work has nothing to pin.""" + if request.work_kind != KIND_PR: + return RequestCheck( + name=CHECK_HEAD_PIN, + ok=True, + reason_code="head_pin_not_applicable", + detail="issue work carries no head SHA to pin.", + ) + selection = (allocation or {}).get("selected") or {} + allocator_head = _clean(selection.get("head_sha")) or None + if not allocator_head: + return RequestCheck( + name=CHECK_HEAD_PIN, + ok=False, + reason_code=REASON_EVIDENCE_UNAVAILABLE, + detail=( + "the allocator reported no head SHA for this PR; PR work " + "cannot be initiated unpinned." + ), + ) + if request.expected_head_sha and request.expected_head_sha != allocator_head: + return RequestCheck( + name=CHECK_HEAD_PIN, + ok=False, + reason_code="head_moved", + detail=( + "the requested head SHA does not match the PR's current head; " + "re-preview against the live head before initiating." + ), + evidence={ + "requested_head_sha": request.expected_head_sha, + "current_head_sha": allocator_head, + }, + ) + return RequestCheck( + name=CHECK_HEAD_PIN, + ok=True, + reason_code="head_pinned", + detail=f"PR {request.display_ref} is pinned at {allocator_head}.", + evidence={"head_sha": allocator_head}, + ) + + +def _next_safe_action_text( + request: WorkRequest, checks: Sequence[RequestCheck], authorized: bool +) -> str: + if authorized: + return ( + f"Confirm and initiate {request.desired_role} work on " + f"{request.work_kind} {request.display_ref} via the allocator." + ) + for check in checks: + if not check.ok: + return f"Resolve {check.name}: {check.detail}" + return "No safe action; the request is not authorized." + + +def preview_request( + request: WorkRequest, + *, + principal: console_authz.Principal | None = None, + allocator: AllocatorFn | None = None, + claims_source: ClaimsFn | None = None, + correlation_id: str | None = None, + audit: bool = True, +) -> RequestPreview: + """Build the read-only intent preview for *request*. Never mutates.""" + who = principal or console_authz.ANONYMOUS + corr = correlation_id or _correlation_id() + decision = console_authz.authorize(ACTION_ID, who, for_execution=False) + + allocation: dict[str, Any] | None = None + claims: Mapping[tuple[str, int], dict[str, Any]] | None = None + checks: list[RequestCheck] = [_authorization_check(decision)] + if decision.allowed: + # An unauthorized principal never reaches the allocator or the + # control-plane DB: a denial must not double as a queue oracle. + allocation = _run_allocator(request, allocator, apply=False) + claims = _load_claims(request, claims_source) + checks.append(_capability_check(request)) + checks.append(_lease_check(request, claims)) + checks.append(_next_safe_action_check(request, allocation)) + checks.append(_head_pin_check(request, allocation)) + + authorized = all(c.ok for c in checks) + allowed_actions, prohibited_actions = allocator_service.role_actions( + request.desired_role + ) + capability = next((c for c in checks if c.name == CHECK_CAPABILITY), None) + evidence = capability.evidence if capability else {} + + if authorized: + reason_code = REASON_PREVIEW_OK + detail = ( + "Request is authorized. Preview only — nothing has been assigned." + ) + else: + first_failure = next(c for c in checks if not c.ok) + reason_code, detail = first_failure.reason_code, first_failure.detail + + preview = RequestPreview( + request=request, + authorized=authorized, + reason_code=reason_code, + detail=detail, + authorization=decision.to_dict(), + checks=tuple(checks), + prohibited_actions=tuple(prohibited_actions), + allowed_actions=tuple(allowed_actions), + next_safe_action=_next_safe_action_text(request, checks, authorized), + required_profile=str(evidence.get("required_profile") or ""), + required_namespace=str(evidence.get("required_namespace") or ""), + required_permission=str(evidence.get("required_permission") or ""), + correlation_id=corr, + allocator_evidence=_allocator_evidence(allocation), + ) + + if audit: + _audit( + request, + result=console_audit.RESULT_PREVIEWED, + decision=decision, + principal=who, + reason_code=reason_code, + detail=detail, + correlation_id=corr, + metadata={ + "intent_summary": request.intent_summary, + "desired_role": request.desired_role, + "authorized": authorized, + "failed_checks": [c.name for c in preview.failed_checks], + "phase": "preview", + }, + ) + return preview + + +# --- Initiation ------------------------------------------------------------- + + +def apply_request( + request: WorkRequest, + *, + principal: console_authz.Principal | None = None, + confirm: bool = False, + allocator: AllocatorFn | None = None, + claims_source: ClaimsFn | None = None, + correlation_id: str | None = None, +) -> dict[str, Any]: + """Initiate *request* as an exclusive assignment, or refuse. + + The only path to an assignment is the allocator agreeing, on a dry-run, + that this work unit is what the requested role should take next. Every + refusal returns before any mutation is attempted. + """ + who = principal or console_authz.ANONYMOUS + corr = correlation_id or _correlation_id() + + execution_decision = console_authz.authorize(ACTION_ID, who, for_execution=True) + authorization = execution_decision.to_dict() + + def _refuse( + outcome: str, + reason_code: str, + detail: str, + *, + status: int, + extra: dict[str, Any] | None = None, + ) -> dict[str, Any]: + _audit( + request, + result=console_audit.RESULT_DENIED, + decision=execution_decision, + principal=who, + reason_code=reason_code, + detail=detail, + correlation_id=corr, + metadata={ + "intent_summary": request.intent_summary, + "desired_role": request.desired_role, + "phase": "apply", + "outcome": outcome, + }, + ) + payload: dict[str, Any] = { + "ok": False, + "outcome": outcome, + "reason_code": reason_code, + "detail": detail, + "request": request.to_dict(), + "authorization": authorization, + "assignment": None, + "correlation_id": corr, + "mutation_performed": False, + "status_code": status, + } + payload.update(extra or {}) + return payload + + if not (execution_decision.allowed and execution_decision.execution_enabled): + return _refuse( + OUTCOME_DENIED, + REASON_UNAUTHORIZED, + execution_decision.detail, + status=403, + ) + + # Confirmation is a property of the action in the RBAC model, so it is read + # from there rather than assumed here. + action = console_authz.get_action(ACTION_ID) + if action is not None and action.requires_confirmation and not confirm: + return _refuse( + OUTCOME_DENIED, + REASON_CONFIRMATION_REQUIRED, + ( + "This action requires explicit confirmation. Re-submit with " + "confirm=true after reviewing the preview." + ), + status=409, + ) + + preview = preview_request( + request, + principal=who, + allocator=allocator, + claims_source=claims_source, + correlation_id=corr, + audit=False, + ) + if not preview.authorized: + outcome = ( + OUTCOME_BLOCKED + if preview.reason_code == REASON_DUPLICATE_ASSIGNMENT + else OUTCOME_WAIT + ) + return _refuse( + outcome, + preview.reason_code, + preview.detail, + status=409, + extra={"preview": preview.to_dict()}, + ) + + fingerprint = ( + _clean(preview.allocator_evidence.get("candidate_set_fingerprint")) or None + ) + allocation = _run_allocator( + request, + allocator, + apply=True, + expected_candidate_set_fingerprint=fingerprint, + ) + if not allocation: + return _refuse( + OUTCOME_WAIT, + REASON_EVIDENCE_UNAVAILABLE, + "the allocator returned no result; nothing was assigned.", + status=503, + ) + + assignment = allocation.get("assignment") or None + outcome = _clean(allocation.get("outcome")) + assigned = bool( + outcome == allocator_service.OUTCOME_ASSIGNED + and assignment + and _selection_matches(allocation.get("selected"), request) + ) + if not assigned: + blocked = outcome in { + allocator_service.OUTCOME_BLOCKED_LEASE, + allocator_service.OUTCOME_BLOCKED_TERMINAL, + allocator_service.OUTCOME_BLOCKED_EXCLUDED_OWN_LEASE, + } + return _refuse( + OUTCOME_BLOCKED if blocked else OUTCOME_WAIT, + REASON_ALLOCATOR_OUTCOME, + ( + f"the allocator returned {outcome or 'no outcome'} rather than " + "an assignment for this work unit; nothing was assigned." + ), + status=409, + extra={"allocator_evidence": _allocator_evidence(allocation)}, + ) + + _audit( + request, + result=console_audit.RESULT_SUCCEEDED, + decision=execution_decision, + principal=who, + reason_code=REASON_AUTHORIZED, + detail=( + f"assigned {request.work_kind} {request.display_ref} to role " + f"{request.desired_role}." + ), + correlation_id=corr, + metadata={ + "intent_summary": request.intent_summary, + "desired_role": request.desired_role, + "phase": "apply", + "outcome": OUTCOME_ASSIGNED, + "assignment_id": assignment.get("assignment_id"), + "lease_id": assignment.get("lease_id"), + }, + ) + return { + "ok": True, + "outcome": OUTCOME_ASSIGNED, + "reason_code": REASON_AUTHORIZED, + "detail": ( + "Exclusive assignment created via the allocator. Continue in the " + f"{preview.required_namespace or 'assigned'} MCP namespace." + ), + "request": request.to_dict(), + "authorization": authorization, + "assignment": dict(assignment), + "handoff": { + "assignment_id": assignment.get("assignment_id"), + "lease_id": assignment.get("lease_id"), + "session_id": assignment.get("session_id"), + "required_profile": preview.required_profile, + "required_namespace": preview.required_namespace, + "allowed_actions": list(preview.allowed_actions), + "forbidden_actions": list(preview.prohibited_actions), + "expected_head_sha": assignment.get("expected_head_sha"), + }, + "correlation_id": corr, + "mutation_performed": True, + "status_code": 201, + "allocator_evidence": _allocator_evidence(allocation), + } + + +# --- Adapters --------------------------------------------------------------- + + +def _allocator_evidence(allocation: Mapping[str, Any] | None) -> dict[str, Any]: + """Reduce an allocator result to the non-secret fields worth surfacing.""" + if not allocation: + return {} + return { + "outcome": allocation.get("outcome"), + "selected": allocation.get("selected"), + "reasons": list(allocation.get("reasons") or ()), + "candidate_set_fingerprint": allocation.get("candidate_set_fingerprint"), + "candidate_count": allocation.get("candidate_count"), + "inventory_complete": allocation.get("inventory_complete"), + "selection_policy": allocation.get("selection_policy"), + "substrate": allocation.get("substrate"), + } + + +def _run_allocator( + request: WorkRequest, + allocator: AllocatorFn | None, + *, + apply: bool, + expected_candidate_set_fingerprint: str | None = None, +) -> dict[str, Any] | None: + fn = allocator or default_allocator + try: + result = fn( + request=request, + apply=apply, + expected_candidate_set_fingerprint=expected_candidate_set_fingerprint, + ) + except Exception: # noqa: BLE001 — an allocator failure denies, never proceeds + return None + return result if isinstance(result, dict) else None + + +def _load_claims( + request: WorkRequest, claims_source: ClaimsFn | None +) -> Mapping[tuple[str, int], dict[str, Any]] | None: + fn = claims_source or default_claims_source + try: + claims = fn(request) + except Exception: # noqa: BLE001 — an unreadable substrate is a denial + return None + return claims if isinstance(claims, Mapping) else None + + +def default_claims_source( + request: WorkRequest, +) -> Mapping[tuple[str, int], dict[str, Any]]: + """Live active-claim inventory from the #613 control-plane DB.""" + import control_plane_db + + db = control_plane_db.ControlPlaneDB() + return db.list_active_claims( + remote=request.remote, org=request.org, repo=request.repo + ) + + +def default_allocator( + *, + request: WorkRequest, + apply: bool, + expected_candidate_set_fingerprint: str | None = None, +) -> dict[str, Any] | None: + """Run the real allocator over the live queue for *request*'s scope. + + An incomplete candidate inventory returns ``None`` rather than a ranking + over a partial set (#758): selecting from a short list can pick the wrong + work unit, so the request denies instead. + """ + import control_plane_db + + from webui.queue_loader import load_queue_snapshot + from webui.traffic_loader import candidates_from_queue_snapshot + + snapshot = load_queue_snapshot() + if snapshot.fetch_error: + return None + for pagination in (snapshot.pr_pagination, snapshot.issue_pagination): + if pagination is not None and not pagination.inventory_complete: + return None + + candidates = candidates_from_queue_snapshot(snapshot) + db = control_plane_db.ControlPlaneDB() + result = allocator_service.allocate_next_work( + db, + session_id=f"webui-request-{uuid.uuid4().hex[:12]}", + role=request.desired_role, + remote=request.remote, + org=request.org, + repo=request.repo, + candidates=candidates, + apply=bool(apply), + allocation_mode="role_scoped", + expected_candidate_set_fingerprint=expected_candidate_set_fingerprint, + ) + if isinstance(result, dict): + result.setdefault("candidate_count", len(candidates)) + result.setdefault("inventory_complete", True) + result.setdefault("selection_policy", allocator_service.SELECTION_POLICY) + return result + + +# --- Audit ------------------------------------------------------------------ + + +def _audit( + request: WorkRequest, + *, + result: str, + decision: console_authz.AuthorizationDecision, + principal: console_authz.Principal, + reason_code: str, + detail: str, + correlation_id: str, + metadata: dict[str, Any], +) -> dict[str, Any]: + return console_audit.record_event( + action_id=ACTION_ID, + result=result, + decision=decision, + principal=principal, + target={ + "kind": request.work_kind, + "ref": request.display_ref, + "remote": request.remote, + "org": request.org, + "repo": request.repo, + }, + reason_code=reason_code, + request_id=correlation_id, + detail=detail, + metadata=metadata, + ) diff --git a/webui/request_views.py b/webui/request_views.py new file mode 100644 index 0000000..fe85679 --- /dev/null +++ b/webui/request_views.py @@ -0,0 +1,164 @@ +"""HTML views for the operator request surface (#643). + +The form is deliberately a *preview* form. It has no initiate button, because +initiating requires a confirmed POST to ``/api/v1/requests/apply`` and a stray +form submission must not be able to produce one by accident. + +Nothing rendered here is trusted input: every interpolated value is escaped, +and the page renders only values the service already produced rather than +echoing a raw request body back. +""" + +from __future__ import annotations + +import html +import json +from typing import Any + +from webui.layout import render_page +from webui.request_service import ( + REQUESTABLE_ROLES, + WORK_KINDS, + RequestError, + RequestPreview, +) + +REQUESTS_PATH = "/requests" +PREVIEW_API_PATH = "/api/v1/requests/preview" +APPLY_API_PATH = "/api/v1/requests/apply" + + +def _escape(text: Any) -> str: + return html.escape(str(text if text is not None else ""), quote=True) + + +REQUEST_PAGE_STYLES = """ + +""" + + +def _options(values: tuple[str, ...], selected: Any) -> str: + return "".join( + f"" + for value in values + ) + + +def _form(values: dict[str, Any] | None = None) -> str: + current = dict(values or {}) + number = current.get("work_number") + return ( + f"
" + "" + "" + "" + "" + "" + "" + "

Preview is read-only and creates no assignment. " + f"Initiating requires a confirmed POST to {APPLY_API_PATH}." + "

" + "
" + ) + + +def _checks_block(preview: RequestPreview) -> str: + rows = [] + for check in preview.checks: + verdict = "PASS" if check.ok else "FAIL" + css = "verdict-ok" if check.ok else "verdict-fail" + rows.append( + "
  • " + f"{verdict} " + f"{_escape(check.name)} — {_escape(check.detail)} " + f"({_escape(check.reason_code)})" + "
  • " + ) + return "" + + +def _preview_block(preview: RequestPreview) -> str: + verdict = "AUTHORIZED" if preview.authorized else "DENIED" + prohibited = "".join( + f"{_escape(action)}" for action in preview.prohibited_actions + ) + request = preview.request + evidence = json.dumps(preview.allocator_evidence, indent=2, default=str) + return ( + "

    Intent preview

    " + f"

    {verdict} — {_escape(preview.detail)}

    " + "

    " + f"Role {_escape(request.desired_role)} · " + f"{_escape(request.work_kind)} {_escape(request.display_ref)}" + f" · profile {_escape(preview.required_profile)} · " + f"namespace {_escape(preview.required_namespace)} · " + f"permission {_escape(preview.required_permission)}" + "

    " + f"

    Intent: {_escape(request.intent_summary)}

    " + f"{_checks_block(preview)}" + f"

    Next safe action: " + f"{_escape(preview.next_safe_action)}

    " + "

    Prohibited for this role: " + + (prohibited or "none declared") + + "

    " + "

    Correlation id " + f"{_escape(preview.correlation_id)}

    " + "
    Allocator evidence" + f"
    {_escape(evidence)}
    " + "
    " + ) + + +def _error_block(error: RequestError) -> str: + field = ( + f"

    Field: {_escape(error.field_name)}

    " + if error.field_name + else "" + ) + return ( + "

    Request rejected

    " + f"

    {_escape(error.reason_code)} — " + f"{_escape(error.detail)}

    {field}" + ) + + +def render_requests_page( + *, + preview: RequestPreview | None = None, + error: RequestError | None = None, + submitted: dict[str, Any] | None = None, +) -> str: + """Render the request form, plus a preview or rejection when one exists.""" + body = ( + "

    Requests

    " + "

    Submit a work request — desired role, issue or PR, and intent — " + "and see whether it would be authorized before anything is reserved. " + "Initiation goes through the allocator (#600/#613); this console never " + "self-selects work, never approves, and never merges.

    " + + _form(submitted) + + (_error_block(error) if error is not None else "") + + (_preview_block(preview) if preview is not None else "") + + f"

    Preview API · " + "RBAC model

    " + + REQUEST_PAGE_STYLES + ) + return render_page(title="Requests", body_html=body) diff --git a/webui/traffic_loader.py b/webui/traffic_loader.py index 912d10e..63993ff 100644 --- a/webui/traffic_loader.py +++ b/webui/traffic_loader.py @@ -201,6 +201,16 @@ def _candidates_from_queue_snapshot(q_snap: QueueSnapshot) -> list[WorkCandidate return candidates +def candidates_from_queue_snapshot(q_snap: QueueSnapshot) -> list[WorkCandidate]: + """Public alias for :func:`_candidates_from_queue_snapshot` (#643). + + The request-initiation service ranks the same candidate set this view + renders, so both must agree on how a queue row becomes a candidate. One + construction, two callers — not two that can drift apart. + """ + return _candidates_from_queue_snapshot(q_snap) + + def _claim_lease_records(inventory: dict[str, Any] | None) -> list[dict[str, Any]]: """Normalize ``build_claim_inventory`` entries into lease records. -- 2.43.7 From 53ce1b1a5ec6476bbe8a153e58f9e13f2bbf1ff9 Mon Sep 17 00:00:00 2001 From: Jason Walker <913443@dadeschools.net> Date: Sat, 25 Jul 2026 03:32:24 -0400 Subject: [PATCH 2/2] fix(webui): compensate stray allocations and keep preview side-effect free MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Addresses both blockers from the PR #902 review (review 589) for issue #643. B1 — an allocator-created assignment could be orphaned and reported as no mutation. apply_request re-previews, then re-runs the allocator with apply=True. The CAS fingerprint hashes only {kind, number} plus exclusions, so a competing lease taken on the requested unit inside the window leaves the fingerprint identical: the pin passes, the selection loop skips the now-claimed unit and commits an assignment on the *next* one, and _selection_matches then fails on egress. The old code returned mutation_performed False with that lease still committed and owned by a synthetic session nothing heartbeats. The window contains a second full load_queue_snapshot(), so it is seconds wide, and foreign sessions acting on this repo concurrently are an observed condition. The egress mismatch now releases the assignment the allocator created before refusing. When the release succeeds the refusal reports mutation_performed False and a compensation record; when it fails the response carries mutation_performed True, an explicit orphaned_assignment, and a gitea_release_workflow_lease reclaim action, because claiming nothing changed while a lease is live is the defect rather than a report of it. The same state was reachable through _run_allocator's bare except Exception: allocate_next_work only catches InvalidWorkKindError, LeaseRequiredError and ControlPlaneError, so anything raised after assign_and_lease committed arrived as "no result" with a durable lease. A None result on the apply path now sweeps and releases whatever this flow's session owns. That sweep is only possible because of the B2 fix below — the session id is now stable across the flow, so the lease is findable. B2 — the "read-only" preview wrote to the control-plane DB. allocate_next_work called db.upsert_session and db.expire_stale_leases unconditionally, before the apply branch was consulted, and default_allocator minted a fresh webui-request- per call. Every preview therefore appended a never-reused session row and mutated global lease state while the payload said dry_run True / mutation_performed False, driven by an operator refreshing a form. One apply wrote two rows and bound the lease to the second, which is why B1's orphan had no reclaimable owner. allocate_next_work gains a keyword-only side_effect_free flag, default False so every existing caller is byte-for-byte unchanged. Under the flag both writes are suppressed and expired leases are instead filtered out of the claim map in memory, which reaches the same selection the sweep would have produced without persisting anything; a claim whose expiry cannot be parsed is kept, since an unreadable expiry is not evidence that work is free. side_effect_free with apply=True fails closed rather than silently reserving. request_service mints one session id per request flow and threads it through both the dry-run and the apply, and the dry-run now routes through the side-effect-free path. Coverage. test_allocator_drift_on_apply_is_not_read_as_an_assignment asserted the defect — it built the orphan state and then required mutation_performed to be False, which a leak satisfies. It now requires the compensating release. Added: release-failure surfacing a reclaim action, an assignment with no lease id, the post-commit exception route, and session-id identity across the flow. The two areas the review named as having zero coverage now have it: default_allocator past its two fail-closed early returns (side_effect_free routing, session-id pass-through and minting, scope and fingerprint propagation) and default_claims_source (scoped read, and an unreadable substrate denying rather than reading as "nothing claimed"). allocator_service gains side-effect-free tests against a real temp-file DB plus the expiry-filter unit tests. Every new guard was mutation-tested by reverting it one at a time: B1 compensation removed → 3 failures; post-commit sweep removed → 1; session id re-minted per call → 2; side_effect_free ignored → 3; in-memory expiry filter removed → 1. Verification: WEBUI_TEST_OFFLINE=1 python -m pytest tests/ -q from this branches/ worktree gives 23 failed, 5263 passed, 6 skipped, 899 subtests. The sorted FAILED set is identical to the reviewer's clean-master baseline at 2f4dec83 (23 failed, 5190 passed) — no new, changed, or disappeared failure — and 21 tests were added over the reviewed head's 5242. Zero conflict markers; py_compile passes; git diff --check clean. Refs #643, PR #902 Co-Authored-By: Claude Opus 5 (1M context) Claude-Session: https://claude.ai/code/session_01V6xFqovhbArPv61j9KCGkL --- allocator_service.py | 122 +++++++-- tests/test_allocator_service.py | 158 ++++++++++++ tests/test_webui_request_initiation.py | 327 ++++++++++++++++++++++++- webui/request_service.py | 282 +++++++++++++++++++-- 4 files changed, 839 insertions(+), 50 deletions(-) diff --git a/allocator_service.py b/allocator_service.py index 9f32fd9..a085d25 100644 --- a/allocator_service.py +++ b/allocator_service.py @@ -23,6 +23,7 @@ import json import os import uuid from dataclasses import dataclass, field +from datetime import datetime, timezone from typing import Any, Mapping, Sequence from control_plane_db import ( @@ -738,6 +739,46 @@ def normalize_exclude_issue_numbers( return sorted(out) +def _claim_expires_at(claim: Any) -> datetime | None: + """Parse a claim's ``expires_at``, or ``None`` when it is absent/malformed.""" + if not isinstance(claim, Mapping): + return None + text = str(claim.get("expires_at") or "").strip() + if not text: + return None + if text.endswith("Z"): + text = text[:-1] + "+00:00" + try: + parsed = datetime.fromisoformat(text) + except ValueError: + return None + if parsed.tzinfo is None: + parsed = parsed.replace(tzinfo=timezone.utc) + return parsed.astimezone(timezone.utc) + + +def _drop_expired_claims( + claims: Mapping[tuple[str, int], dict[str, Any]], + *, + now: datetime | None = None, +) -> dict[tuple[str, int], dict[str, Any]]: + """Claims minus those whose lease has already expired (#643). + + The read-only mirror of ``expire_stale_leases``: the sweep marks such rows + ``expired`` so they stop being returned as claims, and this reaches the same + view without writing. A claim with no parseable ``expires_at`` is **kept** — + an unreadable expiry is not evidence that work is free. + """ + moment = now or datetime.now(timezone.utc) + kept: dict[tuple[str, int], dict[str, Any]] = {} + for key, claim in (claims or {}).items(): + expires_at = _claim_expires_at(claim) + if expires_at is not None and expires_at <= moment: + continue + kept[key] = claim + return kept + + def candidate_set_fingerprint( candidates: Sequence[WorkCandidate], *, @@ -826,12 +867,22 @@ def allocate_next_work( exclude_issue_numbers: Sequence[int] | None = None, expected_candidate_set_fingerprint: str | None = None, allocation_mode: str | None = None, + side_effect_free: bool = False, ) -> dict[str, Any]: """Select and optionally reserve the next work unit via control-plane DB. *apply=False* (default): dry-run selection only — no lease/assignment. *apply=True*: atomic ``assign_and_lease`` for the selected candidate. + *side_effect_free* (#643): a dry run that writes **nothing** to the + control-plane DB. A plain ``apply=False`` still registered a session row and + swept stale leases globally, so a caller advertising a read-only preview was + mutating on every call. Under this flag both writes are suppressed and stale + leases are instead filtered out of the claim map in memory, which yields the + same selection the sweep would have produced without persisting anything. + Incompatible with *apply* — the combination fails closed rather than + silently reserving. + *allocation_mode* (#840): ``cross_role`` (default for controller) inspects the complete queue and returns one authoritative selection naming the required downstream role/profile/action. ``role_scoped`` keeps prior @@ -885,40 +936,57 @@ def allocate_next_work( "allocation_mode": (allocation_mode or "").strip() or None, } - session_id = (session_id or "").strip() or f"alloc-{uuid.uuid4().hex[:12]}" - try: - db.upsert_session( - session_id=session_id, - role=role_norm, - profile=profile_name, - pid=os.getpid(), - controller_instance_id=controller_instance_id, - ) - except Exception as exc: # noqa: BLE001 — surface structured + # A side-effect-free run may never reserve: reserving is a write, and the + # flag is the caller's assertion that this call writes nothing (#643). + if side_effect_free and apply: return { "success": False, "outcome": OUTCOME_NO_SAFE, + "apply": True, "reasons": [ - f"failed to register session in control-plane DB: {exc} " - "(fail closed, #613)" + "side_effect_free is incompatible with apply=True; an " + "assignment is a write (fail closed, #643)" ], "skipped": [], "assignment": None, "substrate": "control_plane_db", } - # Expire stale leases globally before selection. - try: - db.expire_stale_leases() - except Exception as exc: # noqa: BLE001 - return { - "success": False, - "outcome": OUTCOME_NO_SAFE, - "reasons": [f"lease expiry failed: {exc} (fail closed)"], - "skipped": [], - "assignment": None, - "substrate": "control_plane_db", - } + session_id = (session_id or "").strip() or f"alloc-{uuid.uuid4().hex[:12]}" + if not side_effect_free: + try: + db.upsert_session( + session_id=session_id, + role=role_norm, + profile=profile_name, + pid=os.getpid(), + controller_instance_id=controller_instance_id, + ) + except Exception as exc: # noqa: BLE001 — surface structured + return { + "success": False, + "outcome": OUTCOME_NO_SAFE, + "reasons": [ + f"failed to register session in control-plane DB: {exc} " + "(fail closed, #613)" + ], + "skipped": [], + "assignment": None, + "substrate": "control_plane_db", + } + + # Expire stale leases globally before selection. + try: + db.expire_stale_leases() + except Exception as exc: # noqa: BLE001 + return { + "success": False, + "outcome": OUTCOME_NO_SAFE, + "reasons": [f"lease expiry failed: {exc} (fail closed)"], + "skipped": [], + "assignment": None, + "substrate": "control_plane_db", + } terminal = None try: @@ -953,6 +1021,12 @@ def allocate_next_work( "assignment": None, "substrate": "control_plane_db", } + if side_effect_free: + # ``list_active_claims`` filters on status alone, so without the + # global sweep an already-expired lease would still read as a live + # claim and the preview would report work as taken that is free. + # Drop those in memory: same view the sweep produces, no write. + claims = _drop_expired_claims(claims) try: exclude_nums = normalize_exclude_issue_numbers(exclude_issue_numbers) diff --git a/tests/test_allocator_service.py b/tests/test_allocator_service.py index 3e252a1..5c37d0d 100644 --- a/tests/test_allocator_service.py +++ b/tests/test_allocator_service.py @@ -7,6 +7,7 @@ import tempfile import threading import unittest from concurrent.futures import ThreadPoolExecutor, as_completed +from datetime import datetime, timezone from allocator_service import ( OUTCOME_ASSIGNED, @@ -15,6 +16,7 @@ from allocator_service import ( OUTCOME_PREVIEW, OUTCOME_WAIT, WorkCandidate, + _drop_expired_claims, allocate_next_work, candidate_from_dict, classify_skip, @@ -362,5 +364,161 @@ class AllocatorServiceTest(unittest.TestCase): self.assertIn("unavailable", res["reasons"][0].lower()) +class SideEffectFreeAllocationTest(unittest.TestCase): + """``side_effect_free`` dry runs write nothing to the control plane (#643). + + A plain ``apply=False`` still called ``upsert_session`` and + ``expire_stale_leases`` before the apply branch was consulted, so a caller + advertising a read-only preview mutated on every call — one unreferenced + session row per preview, plus a global lease sweep. + """ + + def setUp(self) -> None: + self._tmp = tempfile.TemporaryDirectory() + self.db = ControlPlaneDB(os.path.join(self._tmp.name, "cp.sqlite3")) + + def tearDown(self) -> None: + self._tmp.cleanup() + + def _alloc(self, **kwargs): + defaults = dict( + db=self.db, + session_id="s-preview", + role="author", + remote="prgs", + org="org", + repo="repo", + candidates=[ + WorkCandidate(kind="issue", number=643, labels=("status:ready",)) + ], + apply=False, + profile_name="prgs-author", + username="jcwalker3", + ) + defaults.update(kwargs) + return allocate_next_work(**defaults) + + def _session_ids(self) -> set[str]: + return {str(r.get("session_id")) for r in self.db.list_sessions()} + + def test_side_effect_free_preview_writes_no_session_row(self): + before = self._session_ids() + result = self._alloc(side_effect_free=True) + self.assertEqual(result["outcome"], OUTCOME_PREVIEW) + self.assertEqual(self._session_ids(), before) + self.assertNotIn("s-preview", self._session_ids()) + + def test_plain_dry_run_still_registers_a_session(self): + # The default is unchanged for every existing caller. + self._alloc() + self.assertIn("s-preview", self._session_ids()) + + def test_repeated_previews_do_not_accumulate_rows(self): + for index in range(5): + self._alloc(side_effect_free=True, session_id=f"s-{index}") + self.assertEqual(self._session_ids(), set()) + + def test_side_effect_free_does_not_sweep_stale_leases(self): + self.db.upsert_session(session_id="owner", role="author", pid=1) + assigned = self.db.assign_and_lease( + session_id="owner", + role="author", + remote="prgs", + org="org", + repo="repo", + kind="issue", + number=999, + lease_ttl_seconds=-60, # already expired + ) + self.assertEqual(assigned.outcome, "assigned") + + self._alloc(side_effect_free=True) + + # The expired row is still 'active' in the DB: nothing swept it. + statuses = { + r["lease_id"]: r["status"] + for r in self.db.list_leases( + remote="prgs", org="org", repo="repo", + statuses=("active", "expired"), + ) + } + self.assertEqual(statuses.get(assigned.lease_id), "active") + + def test_expired_claims_are_filtered_in_memory_so_work_stays_selectable(self): + """The read-only mirror of the sweep: expired claims must not block.""" + self.db.upsert_session(session_id="owner", role="author", pid=1) + self.db.assign_and_lease( + session_id="owner", + role="author", + remote="prgs", + org="org", + repo="repo", + kind="issue", + number=643, + lease_ttl_seconds=-60, # expired: must not withhold #643 + ) + result = self._alloc(side_effect_free=True) + self.assertEqual(result["outcome"], OUTCOME_PREVIEW) + self.assertEqual(result["selected"]["number"], 643) + + def test_a_live_claim_still_withholds_the_work(self): + self.db.upsert_session(session_id="owner", role="author", pid=1) + self.db.assign_and_lease( + session_id="owner", + role="author", + remote="prgs", + org="org", + repo="repo", + kind="issue", + number=643, + lease_ttl_seconds=3600, + ) + result = self._alloc(side_effect_free=True) + self.assertNotEqual(result["outcome"], OUTCOME_ASSIGNED) + self.assertNotEqual((result.get("selected") or {}).get("number"), 643) + + def test_side_effect_free_with_apply_fails_closed(self): + result = self._alloc(side_effect_free=True, apply=True) + self.assertFalse(result["success"]) + self.assertEqual(result["outcome"], OUTCOME_NO_SAFE) + self.assertIsNone(result["assignment"]) + self.assertIn("incompatible with apply", result["reasons"][0]) + # And it reserved nothing. + self.assertEqual( + self.db.list_leases(remote="prgs", org="org", repo="repo"), [] + ) + + +class DropExpiredClaimsTest(unittest.TestCase): + """The in-memory expiry filter behind side-effect-free previews (#643).""" + + def test_unparseable_expiry_is_kept_rather_than_assumed_free(self): + claims = { + ("issue", 1): {"lease_id": "l1", "expires_at": "not-a-date"}, + ("issue", 2): {"lease_id": "l2"}, + ("issue", 3): {"lease_id": "l3", "expires_at": None}, + } + self.assertEqual(_drop_expired_claims(claims), claims) + + def test_expired_dropped_and_future_kept(self): + now = datetime(2026, 7, 25, 12, 0, tzinfo=timezone.utc) + claims = { + ("issue", 1): {"expires_at": "2026-07-25T11:59:59+00:00"}, + ("issue", 2): {"expires_at": "2026-07-25T12:00:01+00:00"}, + ("issue", 3): {"expires_at": "2026-07-25T12:00:00+00:00"}, # boundary + } + kept = _drop_expired_claims(claims, now=now) + self.assertEqual(set(kept), {("issue", 2)}) + + def test_naive_and_zulu_timestamps_are_treated_as_utc(self): + now = datetime(2026, 7, 25, 12, 0, tzinfo=timezone.utc) + claims = { + ("issue", 1): {"expires_at": "2026-07-25T11:00:00"}, # naive, past + ("issue", 2): {"expires_at": "2026-07-25T13:00:00Z"}, # zulu, future + } + kept = _drop_expired_claims(claims, now=now) + self.assertEqual(set(kept), {("issue", 2)}) + + if __name__ == "__main__": unittest.main() diff --git a/tests/test_webui_request_initiation.py b/tests/test_webui_request_initiation.py index c733898..e05978c 100644 --- a/tests/test_webui_request_initiation.py +++ b/tests/test_webui_request_initiation.py @@ -16,6 +16,7 @@ allocator refuses an incomplete inventory rather than ranking a partial set. from __future__ import annotations +import contextlib import json import os import pathlib @@ -89,6 +90,63 @@ def _selection( } +@contextlib.contextmanager +def _patched_control_plane( + *, + release_sink: list[tuple[str, str]] | None = None, + release_error: Exception | None = None, + leases=None, +): + """Stand in for ``control_plane_db`` so compensation paths are observable. + + The service imports the module inside the function, so patching + ``sys.modules`` is what intercepts it. No real DB is opened. + """ + module = mock.MagicMock() + db = mock.MagicMock() + + def _release(lease_id, *, session_id): + if release_error is not None: + raise release_error + if release_sink is not None: + release_sink.append((lease_id, session_id)) + + db.release_lease.side_effect = _release + db.list_leases.side_effect = leases or (lambda **_kwargs: []) + module.ControlPlaneDB.return_value = db + with mock.patch.dict(sys.modules, {"control_plane_db": module}): + yield db + + +def _drifting_allocator(*, lease_id: str | None = "lease-wrong"): + """An allocator that previews the requested unit but assigns another. + + This is the #643 B1 race in miniature: the CAS fingerprint hashes only + ``{kind, number}``, so a lease taken on the requested unit inside the window + leaves the fingerprint identical while the selection moves on. + """ + assignment: dict[str, Any] = {"assignment_id": "asn-wrong"} + if lease_id: + assignment["lease_id"] = lease_id + + def _drifting( + *, request, apply, expected_candidate_set_fingerprint=None, session_id=None + ): + return { + "outcome": ( + allocator_service.OUTCOME_ASSIGNED + if apply + else allocator_service.OUTCOME_PREVIEW + ), + "selected": _selection(number=999 if apply else 643), + "assignment": dict(assignment) if apply else None, + "candidate_set_fingerprint": "fp-test", + "session_id": session_id, + } + + return _drifting + + def _fake_allocator( *, selection: dict[str, Any] | None = None, @@ -587,30 +645,170 @@ class TestApplyOutcomes(unittest.TestCase): self.assertFalse(result["mutation_performed"]) def test_allocator_drift_on_apply_is_not_read_as_an_assignment(self): - """The apply call must return *this* work unit, not a substitute.""" + """The apply call must return *this* work unit, not a substitute. - def _drifting(*, request, apply, expected_candidate_set_fingerprint=None): + Drift is not merely refused: the allocator has already committed the + substitute assignment by the time egress rejects it, so the refusal must + also release it. Asserting only ``mutation_performed is False`` would + pass just as well against a leak. + """ + released: list[tuple[str, str]] = [] + + with _patched_control_plane(release_sink=released): + result = request_service.apply_request( + _request(number=643), + principal=_principal(console_authz.OPERATOR), + confirm=True, + allocator=_drifting_allocator(), + claims_source=_no_claims, + ) + + self.assertFalse(result["ok"]) + self.assertIsNone(result["assignment"]) + # The substitute assignment was released, so nothing durable survives. + self.assertEqual(len(released), 1) + self.assertEqual(released[0][0], "lease-wrong") + compensation = result["compensation"] + self.assertTrue(compensation["released"]) + self.assertEqual(compensation["lease_id"], "lease-wrong") + self.assertEqual(compensation["assignment_id"], "asn-wrong") + self.assertEqual(compensation["selected"]["number"], 999) + self.assertFalse(result["mutation_performed"]) + self.assertNotIn("orphaned_assignment", result) + + def test_drift_whose_release_fails_reports_the_orphan_and_a_reclaim(self): + """A release that fails must surface the leak, never swallow it.""" + with _patched_control_plane(release_error=RuntimeError("db is read-only")): + result = request_service.apply_request( + _request(number=643), + principal=_principal(console_authz.OPERATOR), + confirm=True, + allocator=_drifting_allocator(), + claims_source=_no_claims, + ) + + self.assertFalse(result["ok"]) + self.assertIsNone(result["assignment"]) + # A lease really is out there; saying "nothing changed" would be a lie. + self.assertTrue(result["mutation_performed"]) + orphan = result["orphaned_assignment"] + self.assertFalse(orphan["released"]) + self.assertTrue(orphan["attempted"]) + self.assertIn("db is read-only", orphan["error"]) + reclaim = orphan["reclaim_action"] + self.assertEqual(reclaim["tool"], "gitea_release_workflow_lease") + self.assertEqual(reclaim["lease_id"], "lease-wrong") + + def test_drift_without_a_lease_id_still_surfaces_a_reclaim(self): + """An assignment with no lease id cannot be released — say so.""" + result = request_service.apply_request( + _request(number=643), + principal=_principal(console_authz.OPERATOR), + confirm=True, + allocator=_drifting_allocator(lease_id=None), + claims_source=_no_claims, + ) + self.assertFalse(result["ok"]) + self.assertTrue(result["mutation_performed"]) + orphan = result["orphaned_assignment"] + self.assertFalse(orphan["attempted"]) + self.assertEqual(orphan["assignment_id"], "asn-wrong") + self.assertIn("reclaim_action", orphan) + + def test_exception_after_commit_releases_the_session_lease(self): + """A post-commit exception surfaces as no result — with a live lease. + + ``allocate_next_work`` catches only three exception types, so anything + else raised after ``assign_and_lease`` committed reaches the caller as + ``None`` while the lease is durable. The stable per-flow session id is + what makes that lease findable. + """ + released: list[tuple[str, str]] = [] + seen: list[str | None] = [] + + def _explodes_after_commit( + *, request, apply, expected_candidate_set_fingerprint=None, session_id=None + ): + seen.append(session_id) + if not apply: + return { + "outcome": allocator_service.OUTCOME_PREVIEW, + "selected": _selection(number=643), + "assignment": None, + "candidate_set_fingerprint": "fp-test", + } + raise KeyError("selection['required_profile']") + + def _leases(**_kwargs): + return [ + { + "lease_id": "lease-committed", + "session_id": seen[-1], + "work_kind": "issue", + "work_number": 643, + } + ] + + with _patched_control_plane(release_sink=released, leases=_leases): + result = request_service.apply_request( + _request(number=643), + principal=_principal(console_authz.OPERATOR), + confirm=True, + allocator=_explodes_after_commit, + claims_source=_no_claims, + ) + + self.assertFalse(result["ok"]) + self.assertEqual( + result["reason_code"], request_service.REASON_EVIDENCE_UNAVAILABLE + ) + self.assertEqual(len(released), 1) + self.assertEqual(released[0][0], "lease-committed") + self.assertEqual( + result["compensation"]["released"][0]["lease_id"], "lease-committed" + ) + self.assertFalse(result["mutation_performed"]) + + def test_one_session_id_spans_the_dry_run_and_the_apply(self): + """Both halves of an apply share one control-plane identity.""" + seen: list[str | None] = [] + + def _recording( + *, request, apply, expected_candidate_set_fingerprint=None, session_id=None + ): + seen.append(session_id) return { "outcome": ( allocator_service.OUTCOME_ASSIGNED if apply else allocator_service.OUTCOME_PREVIEW ), - "selected": _selection(number=999 if apply else 643), - "assignment": {"assignment_id": "asn-wrong"} if apply else None, + "selected": _selection(number=643), + "assignment": ( + { + "assignment_id": "asn-ok", + "lease_id": "lease-ok", + "expected_head_sha": None, + } + if apply + else None + ), "candidate_set_fingerprint": "fp-test", + "session_id": session_id, } result = request_service.apply_request( _request(number=643), principal=_principal(console_authz.OPERATOR), confirm=True, - allocator=_drifting, + allocator=_recording, claims_source=_no_claims, ) - self.assertFalse(result["ok"]) - self.assertIsNone(result["assignment"]) - self.assertFalse(result["mutation_performed"]) + self.assertTrue(result["ok"]) + self.assertEqual(len(seen), 2) + self.assertTrue(all(s for s in seen)) + self.assertEqual(seen[0], seen[1], "dry-run and apply must share one id") + self.assertTrue(seen[0].startswith("webui-request-")) def test_viewer_cannot_apply(self): result = request_service.apply_request( @@ -634,6 +832,119 @@ class TestApplyOutcomes(unittest.TestCase): self.assertFalse(result["mutation_performed"]) +class TestDefaultAllocatorPastTheEarlyReturns(unittest.TestCase): + """``default_allocator`` beyond its two fail-closed guards (#643). + + Both prior tests returned before ``ControlPlaneDB`` was ever constructed, so + the session id, the ``side_effect_free`` routing and the CAS round-trip had + no coverage at all — which is how a preview that writes session rows shipped. + """ + + def setUp(self): + from webui.queue_loader import QueueSnapshot + + self.snapshot = QueueSnapshot( + project_id="p", + repo_label="r", + prs=(), + issues=(), + pr_pagination=None, + issue_pagination=None, + ) + + @contextlib.contextmanager + def _harness(self): + """Run the real ``default_allocator`` against a recorded allocator call.""" + calls: dict[str, Any] = {} + + def _allocate(db, **kwargs): + calls.update(kwargs) + return {"outcome": allocator_service.OUTCOME_PREVIEW} + + module = mock.MagicMock() + with mock.patch.dict(sys.modules, {"control_plane_db": module}), \ + mock.patch( + "webui.queue_loader.load_queue_snapshot", + return_value=self.snapshot, + ), \ + mock.patch( + "webui.traffic_loader.candidates_from_queue_snapshot", + return_value=[], + ), \ + mock.patch.object( + allocator_service, "allocate_next_work", side_effect=_allocate + ): + yield calls + + def test_preview_runs_side_effect_free_and_never_applies(self): + with self._harness() as calls: + result = request_service.default_allocator( + request=_request(), apply=False + ) + self.assertIsNotNone(result) + self.assertTrue(calls["side_effect_free"]) + self.assertFalse(calls["apply"]) + + def test_apply_is_not_side_effect_free(self): + with self._harness() as calls: + request_service.default_allocator(request=_request(), apply=True) + self.assertFalse(calls["side_effect_free"]) + self.assertTrue(calls["apply"]) + + def test_caller_session_id_is_passed_through_verbatim(self): + with self._harness() as calls: + request_service.default_allocator( + request=_request(), apply=True, session_id="webui-request-fixed" + ) + self.assertEqual(calls["session_id"], "webui-request-fixed") + + def test_absent_session_id_is_minted_with_the_expected_shape(self): + with self._harness() as calls: + request_service.default_allocator(request=_request(), apply=False) + self.assertTrue(str(calls["session_id"]).startswith("webui-request-")) + + def test_scope_and_fingerprint_reach_the_allocator(self): + with self._harness() as calls: + request_service.default_allocator( + request=_request(number=664), + apply=False, + expected_candidate_set_fingerprint="fp-pinned", + ) + self.assertEqual(calls["expected_candidate_set_fingerprint"], "fp-pinned") + self.assertEqual(calls["remote"], SCOPE["remote"]) + self.assertEqual(calls["org"], SCOPE["org"]) + self.assertEqual(calls["repo"], SCOPE["repo"]) + self.assertEqual(calls["allocation_mode"], "role_scoped") + + +class TestDefaultClaimsSource(unittest.TestCase): + """``default_claims_source`` had no test at all (#643).""" + + def test_claims_are_read_for_the_request_scope(self): + db = mock.MagicMock() + db.list_active_claims.return_value = {("issue", 643): {"lease_id": "l1"}} + module = mock.MagicMock() + module.ControlPlaneDB.return_value = db + with mock.patch.dict(sys.modules, {"control_plane_db": module}): + claims = request_service.default_claims_source(_request()) + self.assertEqual(claims, {("issue", 643): {"lease_id": "l1"}}) + db.list_active_claims.assert_called_once_with( + remote=SCOPE["remote"], org=SCOPE["org"], repo=SCOPE["repo"] + ) + + def test_an_unreadable_substrate_denies_rather_than_returning_empty(self): + module = mock.MagicMock() + module.ControlPlaneDB.side_effect = RuntimeError("no db") + with mock.patch.dict(sys.modules, {"control_plane_db": module}): + # _load_claims converts the failure into None, which fails the + # lease check closed; an empty mapping would read as "nothing + # claimed" and wrongly authorize. + claims = request_service._load_claims( + _request(), request_service.default_claims_source + ) + self.assertIsNone(claims) + + class TestAllocatorIntegrationFakes(unittest.TestCase): """The real default allocator refuses a partial inventory (#758).""" diff --git a/webui/request_service.py b/webui/request_service.py index 6941027..5926971 100644 --- a/webui/request_service.py +++ b/webui/request_service.py @@ -35,6 +35,7 @@ the resulting assignment by ``correlation_id``. from __future__ import annotations +import inspect import uuid from dataclasses import dataclass, field from typing import Any, Callable, Mapping, Sequence @@ -321,6 +322,17 @@ def _correlation_id() -> str: return f"req-{uuid.uuid4().hex}" +def _new_session_id() -> str: + """Mint the control-plane session id for one request flow (#643). + + Minted once per flow and reused by the dry-run and the apply, so a lease the + apply creates is owned by an id the caller still holds and can release. When + each call minted its own id, an apply wrote two session rows and bound the + lease to the second — leaving it with no owner able to reclaim it. + """ + return f"webui-request-{uuid.uuid4().hex[:12]}" + + def _selection_matches( selection: Mapping[str, Any] | None, request: WorkRequest ) -> bool: @@ -561,8 +573,15 @@ def preview_request( claims_source: ClaimsFn | None = None, correlation_id: str | None = None, audit: bool = True, + session_id: str | None = None, ) -> RequestPreview: - """Build the read-only intent preview for *request*. Never mutates.""" + """Build the read-only intent preview for *request*. Never mutates. + + "Never mutates" is now literal on the default path: the allocator runs + ``side_effect_free``, so no control-plane session row is written and no lease + sweep runs. *session_id* is threaded from an enclosing apply so both halves + of that flow share one control-plane identity (#643). + """ who = principal or console_authz.ANONYMOUS corr = correlation_id or _correlation_id() decision = console_authz.authorize(ACTION_ID, who, for_execution=False) @@ -573,7 +592,9 @@ def preview_request( if decision.allowed: # An unauthorized principal never reaches the allocator or the # control-plane DB: a denial must not double as a queue oracle. - allocation = _run_allocator(request, allocator, apply=False) + allocation = _run_allocator( + request, allocator, apply=False, session_id=session_id + ) claims = _load_claims(request, claims_source) checks.append(_capability_check(request)) checks.append(_lease_check(request, claims)) @@ -650,9 +671,20 @@ def apply_request( The only path to an assignment is the allocator agreeing, on a dry-run, that this work unit is what the requested role should take next. Every refusal returns before any mutation is attempted. + + One exception is unavoidable and is compensated rather than prevented: the + allocator can commit an assignment and only then reveal, on egress, that it + selected a different work unit than the one requested. That assignment is + released before refusing, and ``mutation_performed`` reports whether any + durable control-plane state survives this call — ``False`` once the stray + assignment is gone, ``True`` with an explicit ``orphaned_assignment`` and + reclaim action when the release did not succeed (#643). """ who = principal or console_authz.ANONYMOUS corr = correlation_id or _correlation_id() + # One identity for the whole flow, so a lease the apply creates is owned by + # an id this function still holds and can release. + session_id = _new_session_id() execution_decision = console_authz.authorize(ACTION_ID, who, for_execution=True) authorization = execution_decision.to_dict() @@ -664,6 +696,7 @@ def apply_request( *, status: int, extra: dict[str, Any] | None = None, + mutation_performed: bool = False, ) -> dict[str, Any]: _audit( request, @@ -689,7 +722,9 @@ def apply_request( "authorization": authorization, "assignment": None, "correlation_id": corr, - "mutation_performed": False, + # True only when durable control-plane state survives this refusal — + # an assignment the allocator created that could not be released. + "mutation_performed": bool(mutation_performed), "status_code": status, } payload.update(extra or {}) @@ -724,6 +759,7 @@ def apply_request( claims_source=claims_source, correlation_id=corr, audit=False, + session_id=session_id, ) if not preview.authorized: outcome = ( @@ -747,21 +783,37 @@ def apply_request( allocator, apply=True, expected_candidate_set_fingerprint=fingerprint, + session_id=session_id, ) if not allocation: + # "No result" is not proof of "no write": allocate_next_work only catches + # three exception types, so anything raised after assign_and_lease + # committed lands here with a lease already durable. Release whatever + # this flow's session owns before refusing. + sweep = _sweep_session_leases(request, session_id=session_id) + extra: dict[str, Any] = {} + detail = "the allocator returned no result; nothing was assigned." + if sweep.get("released"): + extra["compensation"] = sweep + detail = ( + "the allocator returned no result after creating an assignment; " + "the assignment was released and nothing remains assigned." + ) + elif sweep.get("error"): + extra["compensation"] = sweep return _refuse( OUTCOME_WAIT, REASON_EVIDENCE_UNAVAILABLE, - "the allocator returned no result; nothing was assigned.", + detail, status=503, + extra=extra or None, ) assignment = allocation.get("assignment") or None outcome = _clean(allocation.get("outcome")) + committed = bool(outcome == allocator_service.OUTCOME_ASSIGNED and assignment) assigned = bool( - outcome == allocator_service.OUTCOME_ASSIGNED - and assignment - and _selection_matches(allocation.get("selected"), request) + committed and _selection_matches(allocation.get("selected"), request) ) if not assigned: blocked = outcome in { @@ -769,15 +821,45 @@ def apply_request( allocator_service.OUTCOME_BLOCKED_TERMINAL, allocator_service.OUTCOME_BLOCKED_EXCLUDED_OWN_LEASE, } + extra = {"allocator_evidence": _allocator_evidence(allocation)} + detail = ( + f"the allocator returned {outcome or 'no outcome'} rather than " + "an assignment for this work unit; nothing was assigned." + ) + if committed: + # The allocator assigned a *different* unit than the one requested. + # That assignment is real and durable; releasing it is the whole + # difference between a refusal and a silent leak. + compensation = _release_stray_assignment( + request, allocation, session_id=session_id + ) + extra["compensation"] = compensation + if compensation.get("released"): + detail = ( + "the allocator assigned a different work unit than the one " + "requested; that assignment was released and nothing " + "remains assigned." + ) + else: + extra["orphaned_assignment"] = compensation + return _refuse( + OUTCOME_BLOCKED if blocked else OUTCOME_WAIT, + REASON_ALLOCATOR_OUTCOME, + ( + "the allocator assigned a different work unit than the " + "one requested and it could not be released; it must be " + "reclaimed explicitly." + ), + status=409, + extra=extra, + mutation_performed=True, + ) return _refuse( OUTCOME_BLOCKED if blocked else OUTCOME_WAIT, REASON_ALLOCATOR_OUTCOME, - ( - f"the allocator returned {outcome or 'no outcome'} rather than " - "an assignment for this work unit; nothing was assigned." - ), + detail, status=409, - extra={"allocator_evidence": _allocator_evidence(allocation)}, + extra=extra, ) _audit( @@ -853,19 +935,175 @@ def _run_allocator( *, apply: bool, expected_candidate_set_fingerprint: str | None = None, + session_id: str | None = None, ) -> dict[str, Any] | None: fn = allocator or default_allocator + kwargs: dict[str, Any] = { + "request": request, + "apply": apply, + "expected_candidate_set_fingerprint": expected_candidate_set_fingerprint, + } + if session_id: + # Injected allocators in tests predate this argument; only pass it to + # callables that accept it so a fake signature is never broken. + if _accepts_session_id(fn): + kwargs["session_id"] = session_id try: - result = fn( - request=request, - apply=apply, - expected_candidate_set_fingerprint=expected_candidate_set_fingerprint, - ) + result = fn(**kwargs) except Exception: # noqa: BLE001 — an allocator failure denies, never proceeds return None return result if isinstance(result, dict) else None +def _accepts_session_id(fn: AllocatorFn) -> bool: + try: + params = inspect.signature(fn).parameters + except (TypeError, ValueError): # builtins / C callables + return False + if "session_id" in params: + return True + return any(p.kind is inspect.Parameter.VAR_KEYWORD for p in params.values()) + + +def _release_stray_assignment( + request: WorkRequest, + allocation: Mapping[str, Any], + *, + session_id: str | None, +) -> dict[str, Any]: + """Release an assignment the allocator committed for the wrong work unit. + + The apply path can commit an assignment and then discover, on egress, that + the allocator selected a *different* unit than the one requested — the CAS + fingerprint hashes only ``{kind, number}`` plus exclusions, so a lease taken + on the requested unit inside the window leaves the fingerprint identical and + the pin passes while the selection moves on. Without compensation that lease + is orphaned: owned by a session nothing heartbeats, and invisible because the + response says nothing was mutated. + + Returns a record describing what was attempted so the caller can report it, + including a reclaim action when the release itself did not succeed. + """ + assignment = allocation.get("assignment") or {} + lease_id = _clean(assignment.get("lease_id")) or None + assignment_id = _clean(assignment.get("assignment_id")) or None + owner = _clean(allocation.get("session_id")) or session_id or None + selected = allocation.get("selected") or {} + record: dict[str, Any] = { + "attempted": False, + "released": False, + "lease_id": lease_id, + "assignment_id": assignment_id, + "owner_session_id": owner, + "selected": { + "kind": _clean(selected.get("kind")).lower() or None, + "number": selected.get("number"), + }, + } + if not lease_id: + # Nothing durable to release (or the allocator reported no lease id); + # still surface the assignment id so a leak is never silent. + record["detail"] = ( + "the allocator reported an assignment with no lease id; nothing " + "could be released automatically." + ) + if assignment_id: + record["reclaim_action"] = _reclaim_action(record) + return record + if not owner: + record["detail"] = ( + "the owning session id is unknown; the assignment cannot be " + "released automatically." + ) + record["reclaim_action"] = _reclaim_action(record) + return record + + record["attempted"] = True + try: + import control_plane_db + + control_plane_db.ControlPlaneDB().release_lease(lease_id, session_id=owner) + except Exception as exc: # noqa: BLE001 — a failed release must be reported + record["error"] = str(exc) + record["detail"] = ( + "the allocator created an assignment for a different work unit and " + "releasing it failed; it must be reclaimed explicitly." + ) + record["reclaim_action"] = _reclaim_action(record) + return record + + record["released"] = True + record["detail"] = ( + "the allocator created an assignment for a different work unit; it was " + "released, so no assignment persists from this request." + ) + return record + + +def _reclaim_action(record: Mapping[str, Any]) -> dict[str, Any]: + """The explicit operator action for an assignment that outlived its request.""" + lease_id = record.get("lease_id") + return { + "tool": "gitea_release_workflow_lease", + "lease_id": lease_id, + "session_id": record.get("owner_session_id"), + "assignment_id": record.get("assignment_id"), + "instructions": ( + "A control-plane assignment was created for a work unit other than " + "the requested one and could not be released automatically. Call " + f"gitea_release_workflow_lease(lease_id={lease_id!r}, " + f"session_id={record.get('owner_session_id')!r}) to reclaim it, or " + "wait for the lease TTL to expire." + ), + } + + +def _sweep_session_leases( + request: WorkRequest, *, session_id: str | None +) -> dict[str, Any]: + """Release any lease this flow's session owns after an indeterminate apply. + + ``_run_allocator`` returns ``None`` for *any* exception, and + ``allocate_next_work`` only catches ``InvalidWorkKindError``, + ``LeaseRequiredError`` and ``ControlPlaneError`` — so an unexpected error + raised *after* ``assign_and_lease`` committed surfaces as "no result" with a + lease already written. Because the session id is now stable across the flow, + that lease is findable: anything this session owns after a failed apply is by + definition unclaimed by anyone, so it is released. + """ + record: dict[str, Any] = {"attempted": False, "released": [], "session_id": session_id} + if not session_id: + return record + record["attempted"] = True + try: + import control_plane_db + + db = control_plane_db.ControlPlaneDB() + leases = db.list_leases( + remote=request.remote, + org=request.org, + repo=request.repo, + statuses=("active",), + ) + for row in leases: + if _clean(row.get("session_id")) != session_id: + continue + lease_id = _clean(row.get("lease_id")) or None + if not lease_id: + continue + db.release_lease(lease_id, session_id=session_id) + record["released"].append( + { + "lease_id": lease_id, + "work_kind": row.get("work_kind"), + "work_number": row.get("work_number"), + } + ) + except Exception as exc: # noqa: BLE001 — best-effort; never mask the refusal + record["error"] = str(exc) + return record + + def _load_claims( request: WorkRequest, claims_source: ClaimsFn | None ) -> Mapping[tuple[str, int], dict[str, Any]] | None: @@ -894,12 +1132,19 @@ def default_allocator( request: WorkRequest, apply: bool, expected_candidate_set_fingerprint: str | None = None, + session_id: str | None = None, ) -> dict[str, Any] | None: """Run the real allocator over the live queue for *request*'s scope. An incomplete candidate inventory returns ``None`` rather than a ranking over a partial set (#758): selecting from a short list can pick the wrong work unit, so the request denies instead. + + *session_id* is supplied by the caller so the dry-run and the apply share one + control-plane identity; the lease an apply creates is then owned by an id the + request flow still holds. A dry run additionally goes through the allocator's + ``side_effect_free`` path, so a preview writes no session row and sweeps no + leases (#643). """ import control_plane_db @@ -917,7 +1162,7 @@ def default_allocator( db = control_plane_db.ControlPlaneDB() result = allocator_service.allocate_next_work( db, - session_id=f"webui-request-{uuid.uuid4().hex[:12]}", + session_id=session_id or _new_session_id(), role=request.desired_role, remote=request.remote, org=request.org, @@ -926,6 +1171,7 @@ def default_allocator( apply=bool(apply), allocation_mode="role_scoped", expected_candidate_set_fingerprint=expected_candidate_set_fingerprint, + side_effect_free=not apply, ) if isinstance(result, dict): result.setdefault("candidate_count", len(candidates)) -- 2.43.7