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) <[email protected]>
460 lines
16 KiB
Python
460 lines
16 KiB
Python
"""Traffic-control view loader for Phase 1 operator web console (#640).
|
|
|
|
Combines queue snapshots, inventory leases, dependency graph classifications,
|
|
and workflow dashboard rules to deliver full traffic-control visibility:
|
|
runnable, leased (in-progress), blocked (dependency/lock), needs-controller,
|
|
and terminal-complete candidates.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
from dataclasses import dataclass
|
|
from typing import Any, Callable, Sequence
|
|
|
|
from webui.project_registry import find_project, load_registry
|
|
from webui.queue_loader import load_queue_snapshot, QueueSnapshot
|
|
from webui.lease_loader import load_lease_snapshot, LeaseSnapshot
|
|
from workflow_dashboard import (
|
|
DashboardSnapshot,
|
|
QueueEntry,
|
|
RoleNextAction,
|
|
build_workflow_dashboard,
|
|
DASHBOARD_ROLES,
|
|
)
|
|
from allocator_service import WorkCandidate
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class TrafficItem:
|
|
kind: str # "issue" or "pr"
|
|
number: int
|
|
title: str
|
|
traffic_state: str # "runnable", "leased", "blocked", "needs_controller", "terminal_complete"
|
|
expected_role: str
|
|
safe_for_roles: tuple[str, ...]
|
|
badges: tuple[str, ...]
|
|
block_reason: str | None = None
|
|
lease_info: dict[str, Any] | None = None
|
|
head_sha: str | None = None
|
|
|
|
@property
|
|
def is_safe(self) -> bool:
|
|
return self.block_reason is None and bool(self.safe_for_roles)
|
|
|
|
def as_dict(self) -> dict[str, Any]:
|
|
return {
|
|
"kind": self.kind,
|
|
"number": self.number,
|
|
"title": self.title,
|
|
"traffic_state": self.traffic_state,
|
|
"expected_role": self.expected_role,
|
|
"safe_for_roles": list(self.safe_for_roles),
|
|
"badges": list(self.badges),
|
|
"block_reason": self.block_reason,
|
|
"lease_info": self.lease_info,
|
|
"head_sha": self.head_sha,
|
|
"is_safe": self.is_safe,
|
|
}
|
|
|
|
|
|
@dataclass(frozen=True)
|
|
class TrafficSnapshot:
|
|
project_id: str
|
|
repo_label: str
|
|
runnable: tuple[TrafficItem, ...]
|
|
leased: tuple[TrafficItem, ...]
|
|
blocked: tuple[TrafficItem, ...]
|
|
needs_controller: tuple[TrafficItem, ...]
|
|
terminal_complete: tuple[TrafficItem, ...]
|
|
next_roles: tuple[dict[str, Any], ...]
|
|
fetch_error: str | None = None
|
|
inventory_complete: bool = True
|
|
|
|
def as_dict(self) -> dict[str, Any]:
|
|
return {
|
|
"project_id": self.project_id,
|
|
"repo_label": self.repo_label,
|
|
"runnable": [i.as_dict() for i in self.runnable],
|
|
"leased": [i.as_dict() for i in self.leased],
|
|
"blocked": [i.as_dict() for i in self.blocked],
|
|
"needs_controller": [i.as_dict() for i in self.needs_controller],
|
|
"terminal_complete": [i.as_dict() for i in self.terminal_complete],
|
|
"next_roles": list(self.next_roles),
|
|
"fetch_error": self.fetch_error,
|
|
"inventory_complete": self.inventory_complete,
|
|
}
|
|
|
|
|
|
def _classify_traffic_item(
|
|
entry: QueueEntry,
|
|
*,
|
|
lease_info: dict[str, Any] | None = None,
|
|
) -> TrafficItem:
|
|
"""Classify a QueueEntry into a TrafficItem with explicit traffic state."""
|
|
badges = list(entry.badges)
|
|
block_reason = entry.block_reason
|
|
expected_role = entry.expected_role
|
|
|
|
entry_is_safe = entry.block_reason is None and bool(entry.safe_for_roles)
|
|
# Lease state is checked first: an item that is both leased and blocked is
|
|
# reported as leased. That is safe by construction — a leased item is never
|
|
# placed in the runnable lane — and it keeps the operator's attention on the
|
|
# session that currently owns the work. The blocker text still renders.
|
|
if lease_info is not None or "in-progress" in badges or "claimed" in badges:
|
|
state = "leased"
|
|
elif expected_role == "reconciler" or "terminal-lock" in badges:
|
|
state = "terminal_complete"
|
|
elif expected_role == "controller" or "contaminated" in badges or "needs-controller" in badges:
|
|
state = "needs_controller"
|
|
elif (
|
|
block_reason is not None
|
|
or "blocked" in badges
|
|
or "dependency-unmet" in badges
|
|
or "blocked-by-terminal" in badges
|
|
or "status:blocked" in badges
|
|
):
|
|
state = "blocked"
|
|
elif entry_is_safe:
|
|
state = "runnable"
|
|
else:
|
|
state = "needs_controller"
|
|
|
|
return TrafficItem(
|
|
kind=entry.kind,
|
|
number=entry.number,
|
|
title=entry.title,
|
|
traffic_state=state,
|
|
expected_role=expected_role,
|
|
safe_for_roles=entry.safe_for_roles,
|
|
badges=tuple(badges),
|
|
block_reason=block_reason,
|
|
lease_info=lease_info,
|
|
head_sha=entry.head_sha,
|
|
)
|
|
|
|
|
|
# Claim statuses from ``issue_claim_heartbeat.build_claim_inventory`` that mean
|
|
# a live worker currently holds the issue. Everything else (``stale``,
|
|
# ``phantom``, ``reclaimable``, ``not_claimed``) is reported through the
|
|
# dashboard's stale-lease channel and is never rendered as an active lease.
|
|
_ACTIVE_CLAIM_STATUSES = frozenset({"active", "awaiting_review"})
|
|
|
|
# Statuses that positively mean "not an active lease" for any lease record.
|
|
_INACTIVE_LEASE_STATUSES = frozenset(
|
|
{"expired", "stale", "released", "moot", "reclaimable", "phantom", "not_claimed"}
|
|
)
|
|
|
|
|
|
def _candidates_from_queue_snapshot(q_snap: QueueSnapshot) -> list[WorkCandidate]:
|
|
"""Build allocator candidates from the queue loader's authoritative signals.
|
|
|
|
Display badges (``blocked``/``claimed``/``duplicate``/``stale``/
|
|
``in-review``/``open``) are rendering hints, not routing state, so nothing
|
|
here branches on them. Every routing field comes from
|
|
``QueueItem.signals`` — the raw Gitea payload values.
|
|
|
|
The queue loader reads ``/pulls`` and ``/issues`` only; it never fetches
|
|
review verdicts. ``request_changes_current_head`` / ``approval_on_current_head``
|
|
are therefore left at their fail-safe ``False`` rather than being guessed
|
|
from badges: an unproven approval must never route a PR to the merger.
|
|
"""
|
|
candidates: list[WorkCandidate] = []
|
|
|
|
for pr in q_snap.prs:
|
|
signals = pr.signals or {}
|
|
head_sha = str(signals.get("head_sha") or "").strip()
|
|
mergeable = signals.get("mergeable")
|
|
labels = tuple(str(x) for x in (signals.get("labels") or ()))
|
|
candidates.append(
|
|
WorkCandidate(
|
|
kind="pr",
|
|
number=pr.number,
|
|
state="open",
|
|
labels=labels,
|
|
title=pr.title,
|
|
# Full 40-char SHA from head.sha — never the 12-char display value.
|
|
head_sha=head_sha or None,
|
|
priority=5,
|
|
mergeable=mergeable is True,
|
|
blocked=mergeable is False or "status:blocked" in labels,
|
|
)
|
|
)
|
|
|
|
for issue in q_snap.issues:
|
|
signals = issue.signals or {}
|
|
labels = tuple(str(x) for x in (signals.get("labels") or ()))
|
|
lowered = {label.lower() for label in labels}
|
|
candidates.append(
|
|
WorkCandidate(
|
|
kind="issue",
|
|
number=issue.number,
|
|
state="open",
|
|
labels=labels,
|
|
title=issue.title,
|
|
priority=20 if "status:ready" in lowered else 10,
|
|
blocked="status:blocked" in lowered,
|
|
# A live claim by another session is not this session's work.
|
|
already_claimed_elsewhere="status:in-progress" in lowered,
|
|
)
|
|
)
|
|
|
|
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.
|
|
|
|
The inventory contract is ``{"entries", "counts", "heartbeat_lease_minutes",
|
|
"reclaim_after_minutes", "in_progress_total"}``. Each entry is keyed by
|
|
``issue_number``; the subject kind is therefore always ``issue``.
|
|
"""
|
|
entries = (inventory or {}).get("entries") or ()
|
|
records: list[dict[str, Any]] = []
|
|
for entry in entries:
|
|
if not isinstance(entry, dict):
|
|
continue
|
|
number = entry.get("issue_number")
|
|
if number is None:
|
|
continue
|
|
try:
|
|
number_int = int(number)
|
|
except (TypeError, ValueError):
|
|
continue
|
|
heartbeat = entry.get("latest_heartbeat") or {}
|
|
record = dict(entry)
|
|
record.update(
|
|
{
|
|
"kind": "issue",
|
|
"number": number_int,
|
|
"role": "author",
|
|
"lease_source": "issue-claim-heartbeat",
|
|
}
|
|
)
|
|
if isinstance(heartbeat, dict):
|
|
if heartbeat.get("session_id") and not record.get("session_id"):
|
|
record["session_id"] = heartbeat.get("session_id")
|
|
if heartbeat.get("author") and not record.get("author"):
|
|
record["author"] = heartbeat.get("author")
|
|
records.append(record)
|
|
return records
|
|
|
|
|
|
def _lease_subject(lease: dict[str, Any]) -> tuple[str, int] | None:
|
|
"""Return the ``(kind, number)`` a lease record actually covers.
|
|
|
|
Fails closed: a record that does not identify exactly one subject is
|
|
dropped rather than attributed to a guessed work item (#640 — never invent
|
|
a lease, and never attach a PR lease to a same-numbered issue).
|
|
"""
|
|
kind = str(lease.get("kind") or lease.get("work_kind") or "").strip().lower()
|
|
pr_number = lease.get("pr_number")
|
|
issue_number = lease.get("issue_number")
|
|
|
|
if kind not in ("pr", "issue"):
|
|
if pr_number is not None and issue_number is None:
|
|
kind = "pr"
|
|
elif issue_number is not None and pr_number is None:
|
|
kind = "issue"
|
|
else:
|
|
return None
|
|
|
|
number = lease.get("number")
|
|
if number is None:
|
|
number = lease.get("work_number")
|
|
if number is None:
|
|
number = pr_number if kind == "pr" else issue_number
|
|
if number is None:
|
|
return None
|
|
try:
|
|
return kind, int(number)
|
|
except (TypeError, ValueError):
|
|
return None
|
|
|
|
|
|
def _is_active_lease(lease: dict[str, Any]) -> bool:
|
|
"""True when the record proves a worker currently holds the item."""
|
|
if lease.get("stale") or lease.get("expired"):
|
|
return False
|
|
status = str(lease.get("status") or lease.get("lease_status") or "").strip().lower()
|
|
if status in _INACTIVE_LEASE_STATUSES:
|
|
return False
|
|
if lease.get("lease_source") == "issue-claim-heartbeat":
|
|
return status in _ACTIVE_CLAIM_STATUSES
|
|
return True
|
|
|
|
|
|
def load_traffic_snapshot(
|
|
*,
|
|
candidates: Sequence[WorkCandidate] | None = None,
|
|
leases: Sequence[dict[str, Any]] | None = None,
|
|
terminal_pr: int | None = None,
|
|
fetch_queue_snapshot: Callable[[], QueueSnapshot] | None = None,
|
|
fetch_lease_snapshot: Callable[[], LeaseSnapshot] | None = None,
|
|
project_id: str = "gitea-tools",
|
|
) -> TrafficSnapshot:
|
|
"""Load and compute the traffic-control snapshot."""
|
|
try:
|
|
reg = load_registry()
|
|
proj = find_project(reg, project_id)
|
|
repo_label = proj.remote_repo if proj else "Scaled-Tech-Consulting/Gitea-Tools"
|
|
except Exception:
|
|
repo_label = "Scaled-Tech-Consulting/Gitea-Tools"
|
|
|
|
# Injected candidates path (pure unit testing)
|
|
if candidates is not None:
|
|
dashboard = build_workflow_dashboard(
|
|
candidates=candidates,
|
|
leases=leases,
|
|
terminal_pr=terminal_pr,
|
|
inventory_complete=True,
|
|
)
|
|
return _build_traffic_snapshot_from_dashboard(
|
|
project_id=project_id,
|
|
repo_label=repo_label,
|
|
dashboard=dashboard,
|
|
leases=leases or (),
|
|
)
|
|
|
|
# Live snapshot loading
|
|
q_loader = fetch_queue_snapshot or load_queue_snapshot
|
|
l_loader = fetch_lease_snapshot or load_lease_snapshot
|
|
|
|
try:
|
|
q_snap = q_loader()
|
|
l_snap = l_loader()
|
|
except Exception as exc: # noqa: BLE001
|
|
return TrafficSnapshot(
|
|
project_id=project_id,
|
|
repo_label=repo_label,
|
|
runnable=(),
|
|
leased=(),
|
|
blocked=(),
|
|
needs_controller=(),
|
|
terminal_complete=(),
|
|
next_roles=(),
|
|
fetch_error=f"Failed to load traffic state: {exc}",
|
|
inventory_complete=False,
|
|
)
|
|
|
|
if q_snap.fetch_error or l_snap.fetch_error:
|
|
err = q_snap.fetch_error or l_snap.fetch_error
|
|
return TrafficSnapshot(
|
|
project_id=project_id,
|
|
repo_label=repo_label,
|
|
runnable=(),
|
|
leased=(),
|
|
blocked=(),
|
|
needs_controller=(),
|
|
terminal_complete=(),
|
|
next_roles=(),
|
|
fetch_error=err,
|
|
inventory_complete=False,
|
|
)
|
|
|
|
candidate_list = _candidates_from_queue_snapshot(q_snap)
|
|
|
|
raw_leases: list[dict[str, Any]] = _claim_lease_records(l_snap.claim_inventory)
|
|
for r_lease in l_snap.reviewer_leases or ():
|
|
if not isinstance(r_lease, dict):
|
|
continue
|
|
# Always pin reviewer leases to the PR subject, even if a linked
|
|
# issue_number is present on the marker (#640 B2).
|
|
normalized = dict(r_lease)
|
|
subject = normalized.get("pr_number") or normalized.get("number")
|
|
if subject is None:
|
|
continue
|
|
try:
|
|
pr_num = int(subject)
|
|
except (TypeError, ValueError):
|
|
continue
|
|
normalized["kind"] = "pr"
|
|
normalized["number"] = pr_num
|
|
normalized["pr_number"] = pr_num
|
|
normalized.setdefault("role", "reviewer")
|
|
raw_leases.append(normalized)
|
|
|
|
dashboard = build_workflow_dashboard(
|
|
candidates=candidate_list,
|
|
leases=raw_leases,
|
|
inventory_complete=q_snap.pr_pagination.inventory_complete if q_snap.pr_pagination else True,
|
|
)
|
|
|
|
return _build_traffic_snapshot_from_dashboard(
|
|
project_id=project_id,
|
|
repo_label=repo_label,
|
|
dashboard=dashboard,
|
|
leases=raw_leases,
|
|
)
|
|
|
|
|
|
def _build_traffic_snapshot_from_dashboard(
|
|
*,
|
|
project_id: str,
|
|
repo_label: str,
|
|
dashboard: DashboardSnapshot,
|
|
leases: Sequence[dict[str, Any]],
|
|
) -> TrafficSnapshot:
|
|
"""Classify dashboard entries into the 5 traffic state buckets."""
|
|
all_entries = dashboard.open_prs + dashboard.open_issues
|
|
|
|
# Map each active lease onto the exact work item it covers. Records whose
|
|
# subject cannot be determined, and claims that are stale/phantom/
|
|
# reclaimable, are deliberately dropped instead of guessed.
|
|
lease_map: dict[tuple[str, int], dict[str, Any]] = {}
|
|
for lease in leases:
|
|
if not isinstance(lease, dict) or not _is_active_lease(lease):
|
|
continue
|
|
subject = _lease_subject(lease)
|
|
if subject is not None:
|
|
lease_map[subject] = lease
|
|
|
|
runnable: list[TrafficItem] = []
|
|
leased: list[TrafficItem] = []
|
|
blocked: list[TrafficItem] = []
|
|
needs_controller: list[TrafficItem] = []
|
|
terminal_complete: list[TrafficItem] = []
|
|
|
|
for entry in all_entries:
|
|
l_info = lease_map.get((entry.kind, entry.number))
|
|
item = _classify_traffic_item(entry, lease_info=l_info)
|
|
|
|
if item.traffic_state == "leased":
|
|
leased.append(item)
|
|
elif item.traffic_state == "terminal_complete":
|
|
terminal_complete.append(item)
|
|
elif item.traffic_state == "blocked":
|
|
blocked.append(item)
|
|
elif item.traffic_state == "needs_controller":
|
|
needs_controller.append(item)
|
|
else:
|
|
runnable.append(item)
|
|
|
|
next_roles = [dashboard.next_safe_by_role[r].as_dict() for r in DASHBOARD_ROLES if r in dashboard.next_safe_by_role]
|
|
|
|
return TrafficSnapshot(
|
|
project_id=project_id,
|
|
repo_label=repo_label,
|
|
runnable=tuple(runnable),
|
|
leased=tuple(leased),
|
|
blocked=tuple(blocked),
|
|
needs_controller=tuple(needs_controller),
|
|
terminal_complete=tuple(terminal_complete),
|
|
next_roles=tuple(next_roles),
|
|
fetch_error=None,
|
|
inventory_complete=dashboard.inventory_complete,
|
|
)
|
|
|
|
|
|
def snapshot_to_dict(snapshot: TrafficSnapshot) -> dict[str, Any]:
|
|
return snapshot.as_dict()
|