fix: authoritative cross-role generic queue allocation (Closes #840)

Add controller-owned cross_role allocation mode that inspects the full
queue and returns one selection with required role/profile/action and
lease evidence. Document process_work_queue routing, normalize
controller role metadata, and keep the dashboard explanatory only.

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
This commit is contained in:
2026-07-23 00:05:35 -04:00
co-authored by Claude Opus 4.8
parent 689c60fc7c
commit 648d9464ba
7 changed files with 992 additions and 56 deletions
+294 -40
View File
@@ -78,6 +78,33 @@ VALID_ROLES = frozenset(
{ROLE_AUTHOR, ROLE_REVIEWER, ROLE_MERGER, ROLE_RECONCILER, ROLE_CONTROLLER}
)
# Allocation modes (#840).
# role_scoped: only candidates whose expected role matches the caller role.
# cross_role: controller-owned generic queue selection — inspect full queue,
# rank/eligibility canonically, return one selection naming the required
# downstream role/profile. Controller routes; it does not perform mutations.
ALLOCATION_MODE_ROLE_SCOPED = "role_scoped"
ALLOCATION_MODE_CROSS_ROLE = "cross_role"
VALID_ALLOCATION_MODES = frozenset(
{ALLOCATION_MODE_ROLE_SCOPED, ALLOCATION_MODE_CROSS_ROLE}
)
# Default execution-profile / MCP-namespace names for each role.
DEFAULT_ROLE_PROFILES: dict[str, str] = {
ROLE_AUTHOR: "prgs-author",
ROLE_REVIEWER: "prgs-reviewer",
ROLE_MERGER: "prgs-merger",
ROLE_RECONCILER: "prgs-reconciler",
ROLE_CONTROLLER: "prgs-controller",
}
DEFAULT_ROLE_NAMESPACES: dict[str, str] = {
ROLE_AUTHOR: "gitea-author",
ROLE_REVIEWER: "gitea-reviewer",
ROLE_MERGER: "gitea-merger",
ROLE_RECONCILER: "gitea-reconciler",
ROLE_CONTROLLER: "gitea-controller",
}
# Default action matrices by role (mutation gate will re-check).
ROLE_ACTIONS: dict[str, tuple[tuple[str, ...], tuple[str, ...]]] = {
ROLE_AUTHOR: (
@@ -259,6 +286,126 @@ def normalize_role(role: str | None, *, profile_name: str | None = None) -> str:
)
def resolve_allocation_mode(
role: str,
allocation_mode: str | None = None,
) -> str:
"""Resolve allocation mode; controller defaults to cross_role (#840)."""
raw = (allocation_mode or "").strip().lower()
if raw:
if raw not in VALID_ALLOCATION_MODES:
raise ControlPlaneError(
f"unknown allocation_mode {allocation_mode!r}; expected one of "
f"{sorted(VALID_ALLOCATION_MODES)}"
)
return raw
if role == ROLE_CONTROLLER:
return ALLOCATION_MODE_CROSS_ROLE
return ALLOCATION_MODE_ROLE_SCOPED
def required_profile_for_role(
role: str,
*,
profile_name: str | None = None,
) -> str:
"""Map a required role to the canonical execution profile name."""
role_norm = (role or "").strip().lower()
# Preserve remote/env prefix from the active profile when present
# (e.g. dadeschools-author → dadeschools-reviewer).
active = (profile_name or "").strip()
if active:
lower = active.lower()
for token in ("author", "reviewer", "merger", "reconciler", "controller"):
if lower.endswith(f"-{token}") or lower == token:
prefix = active[: -len(token)].rstrip("-")
if prefix:
return f"{prefix}-{role_norm}"
return role_norm
return DEFAULT_ROLE_PROFILES.get(role_norm, f"prgs-{role_norm}")
def required_namespace_for_role(
role: str,
*,
profile_name: str | None = None,
) -> str:
"""Map a required role to the canonical MCP namespace name."""
role_norm = (role or "").strip().lower()
profile = required_profile_for_role(role_norm, profile_name=profile_name)
# Namespace is typically gitea-<role>; keep stable mapping when profile is
# non-prgs (still gitea-<role> for isolation).
return DEFAULT_ROLE_NAMESPACES.get(role_norm, f"gitea-{role_norm}")
def selected_action_for_candidate(c: WorkCandidate, required_role: str) -> str:
"""Canonical next action for the selected work under *required_role*."""
role = (required_role or "").strip().lower()
if role == ROLE_AUTHOR:
if c.kind == "pr" and c.request_changes_current_head:
return "address_pr_change_requests"
if c.kind == "pr":
return "update_pr"
return "implement"
if role == ROLE_REVIEWER:
if c.approval_stale:
return "re_review"
return "review"
if role == ROLE_MERGER:
return "merge"
if role == ROLE_RECONCILER:
if c.approval_contaminated:
return "reconcile_contaminated_approval"
return "reconcile"
if role == ROLE_CONTROLLER:
return "diagnose"
return "process"
def build_selection_dict(
selected: WorkCandidate,
*,
active_role: str,
required_role: str,
profile_name: str | None = None,
allocation_mode: str,
) -> dict[str, Any]:
"""Authoritative single selection payload for allocator results (#840)."""
action = selected_action_for_candidate(selected, required_role)
req_profile = required_profile_for_role(
required_role, profile_name=profile_name
)
req_ns = required_namespace_for_role(
required_role, profile_name=profile_name
)
return {
"kind": selected.kind,
"number": selected.number,
"title": selected.title,
"labels": list(selected.labels),
"head_sha": selected.head_sha,
"priority": selected.priority,
"expected_role_next": required_role,
"required_role": required_role,
"selected_action": action,
"action": action,
"required_profile": req_profile,
"required_namespace": req_ns,
"pinned": {
"kind": selected.kind,
"number": selected.number,
"head_sha": selected.head_sha,
"issue_number": selected.number if selected.kind == "issue" else None,
"pr_number": selected.number if selected.kind == "pr" else None,
},
"reason_selected": (
f"highest-priority eligible candidate under allocation_mode="
f"'{allocation_mode}' (active_role={active_role}, "
f"required_role={required_role}, action={action})"
),
}
def expected_role_for_candidate(c: WorkCandidate) -> str:
"""ADR §5.3 routing: which role should take this work next."""
if c.kind == "pr":
@@ -289,6 +436,7 @@ def classify_skip(
role: str,
terminal_pr: int | None,
claim_ownership: str | None = None,
allocation_mode: str | None = None,
) -> str | None:
"""Return skip reason, or None if candidate is selectable for *role*.
@@ -297,7 +445,12 @@ def classify_skip(
and unknown claims are excluded so one session's in-progress task can never
blockade the queue for a different controller; ``own`` stays selectable so
a controller can resume its own work.
*allocation_mode* (#840): ``cross_role`` (controller default) ranks the full
queue and selects the highest-priority eligible item for any downstream
role. ``role_scoped`` retains prior role-match filtering.
"""
mode = resolve_allocation_mode(role, allocation_mode)
if c.state in ("merged", "closed"):
return f"{c.kind}#{c.number} is {c.state}; never assign"
if c.blocked or "status:blocked" in c.labels:
@@ -322,34 +475,58 @@ def classify_skip(
if c.kind == "pr" and not (c.head_sha or "").strip():
return f"pr#{c.number} missing head_sha pin"
expected = expected_role_for_candidate(c)
# Terminal path first: when an active terminal PR exists, only that PR
# (or controller diagnosis) is assignable for review-path roles.
# is assignable for review-path roles (or for work whose expected role is
# review/merge under cross_role selection).
if terminal_pr is not None and c.kind == "pr" and c.number != terminal_pr:
if role in (ROLE_REVIEWER, ROLE_MERGER):
terminal_roles = (ROLE_REVIEWER, ROLE_MERGER)
if mode == ALLOCATION_MODE_CROSS_ROLE:
if expected in terminal_roles:
return (
f"pr#{c.number} skipped: active terminal-review lock on "
f"PR #{terminal_pr} must be resolved first"
)
elif role in terminal_roles:
return (
f"pr#{c.number} skipped: active terminal-review lock on "
f"PR #{terminal_pr} must be resolved first"
)
expected = expected_role_for_candidate(c)
if role == ROLE_CONTROLLER:
# Controller may inspect anything but only assigns diagnosis targets
# when contaminated / blocked.
if mode == ALLOCATION_MODE_CROSS_ROLE:
# Cross-role controller selection: eligibility only — no active-role
# match filter. The selection payload names required_role.
pass
elif role == ROLE_CONTROLLER:
# Legacy diagnosis-only controller path (role_scoped): only reconciler-
# needed targets. Prefer cross_role for generic queue allocation.
if expected == ROLE_RECONCILER or c.blocked:
return None
return f"{c.kind}#{c.number} does not require controller (expected {expected})"
if role != expected:
return (
f"{c.kind}#{c.number} does not require controller "
f"(expected {expected})"
)
elif role != expected:
return (
f"{c.kind}#{c.number} expects role '{expected}', active role is '{role}'"
)
# Ready-gate for issues: prefer status:ready when labels present.
# Applies for author-bound work in both modes (cross_role only gates
# author-expected issues so reconciler/reviewer PRs stay selectable).
if c.kind == "issue" and c.labels:
if "status:ready" not in c.labels and "status:in-progress" not in c.labels:
# Allow unlabeled open issues; only skip explicit non-ready states.
if any(l.startswith("status:") for l in c.labels):
return f"issue#{c.number} not status:ready ({','.join(c.labels)})"
gate_role = expected if mode == ALLOCATION_MODE_CROSS_ROLE else role
if gate_role in (ROLE_AUTHOR, ROLE_CONTROLLER):
if (
"status:ready" not in c.labels
and "status:in-progress" not in c.labels
):
if any(l.startswith("status:") for l in c.labels):
return (
f"issue#{c.number} not status:ready "
f"({','.join(c.labels)})"
)
return None
@@ -501,12 +678,19 @@ def allocate_next_work(
claims: Mapping[tuple[str, int], dict[str, Any]] | None = None,
exclude_issue_numbers: Sequence[int] | None = None,
expected_candidate_set_fingerprint: str | None = None,
allocation_mode: str | None = None,
) -> 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.
*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
per-role filtering. Controller routes only — never grants author/reviewer/
merger/reconciler mutation rights to the controller session.
*exclude_issue_numbers* (#776): numbers removed before ranking. Omitted /
empty preserves prior behavior.
@@ -541,6 +725,19 @@ def allocate_next_work(
"substrate": "control_plane_db",
}
try:
mode = resolve_allocation_mode(role_norm, allocation_mode)
except ControlPlaneError as exc:
return {
"success": False,
"outcome": OUTCOME_ROLE_INELIGIBLE,
"reasons": [str(exc)],
"skipped": [],
"assignment": None,
"substrate": "control_plane_db",
"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(
@@ -748,6 +945,7 @@ def allocate_next_work(
role=role_norm,
terminal_pr=terminal_pr,
claim_ownership=ownership,
allocation_mode=mode,
)
if reason:
is_claim_skip = SKIP_CLAIMED_BY_OTHER_SESSION in reason
@@ -828,6 +1026,10 @@ def allocate_next_work(
"outcome": outcome,
"apply": bool(apply),
"role": role_norm,
"allocation_mode": mode,
"routing_role": role_norm,
"required_role": None,
"selected_action": None,
"profile_name": profile_name,
"username": username,
"session_id": session_id,
@@ -840,6 +1042,12 @@ def allocate_next_work(
"skipped": [s.as_dict() for s in skipped],
"terminal_pr": terminal_pr,
"assignment": None,
"allocation_evidence": {
"mode": "empty",
"allocation_mode": mode,
"lease_created": False,
"selection_policy": SELECTION_POLICY,
},
"substrate": "control_plane_db",
"file_lock_only": False,
"comment_lease_only": False,
@@ -852,25 +1060,37 @@ def allocate_next_work(
"owner_session_id": owner_session_id,
"downstream_note": (
"#612 incident bridge remains downstream of #600; "
"allocator never assigns raw monitoring incidents"
"allocator never assigns raw monitoring incidents; "
"controller routes only under cross_role (#840)"
),
}
expected_role = expected_role_for_candidate(selected)
allowed, forbidden = role_actions(role_norm)
selection = {
"kind": selected.kind,
"number": selected.number,
"title": selected.title,
"labels": list(selected.labels),
"head_sha": selected.head_sha,
"priority": selected.priority,
"expected_role_next": expected_role,
"reason_selected": (
f"highest-priority candidate for role '{role_norm}' "
f"(expected_role={expected_role})"
),
}
# Cross-role: lease/action matrix follows the required downstream role so
# evidence names the worker that must act. Controller session still owns
# the routing decision; mutation isolation is enforced by role gates on
# mutation tools (controller profile lacks author/review/merge ops).
lease_role = (
expected_role if mode == ALLOCATION_MODE_CROSS_ROLE else role_norm
)
allowed, forbidden = role_actions(lease_role)
# Controller must never receive mutation-class rights via cross-role apply.
if role_norm == ROLE_CONTROLLER:
ctrl_allowed, ctrl_forbidden = role_actions(ROLE_CONTROLLER)
# Keep controller session capability evidence separate from lease_role.
controller_allowed_actions = ctrl_allowed
controller_forbidden_actions = ctrl_forbidden
else:
controller_allowed_actions = allowed
controller_forbidden_actions = forbidden
selection = build_selection_dict(
selected,
active_role=role_norm,
required_role=expected_role,
profile_name=profile_name,
allocation_mode=mode,
)
if not apply:
return {
@@ -878,6 +1098,12 @@ def allocate_next_work(
"outcome": OUTCOME_PREVIEW,
"apply": False,
"role": role_norm,
"allocation_mode": mode,
"routing_role": role_norm,
"required_role": expected_role,
"selected_action": selection["selected_action"],
"required_profile": selection["required_profile"],
"required_namespace": selection["required_namespace"],
"profile_name": profile_name,
"username": username,
"session_id": session_id,
@@ -893,6 +1119,12 @@ def allocate_next_work(
"skipped": [s.as_dict() for s in skipped],
"terminal_pr": terminal_pr,
"assignment": None,
"allocation_evidence": {
"mode": "preview",
"allocation_mode": mode,
"lease_created": False,
"selection_policy": SELECTION_POLICY,
},
"substrate": "control_plane_db",
"file_lock_only": False,
"comment_lease_only": False,
@@ -902,9 +1134,12 @@ def allocate_next_work(
"controller_excluded": list(controller_excluded),
"exclude_issue_numbers": list(exclude_nums),
"candidate_set_fingerprint": cas_fp,
"controller_allowed_actions": list(controller_allowed_actions),
"controller_forbidden_actions": list(controller_forbidden_actions),
"downstream_note": (
"#612 incident bridge remains downstream of #600; "
"allocator never assigns raw monitoring incidents"
"allocator never assigns raw monitoring incidents; "
"controller routes only under cross_role (#840)"
),
}
@@ -913,7 +1148,7 @@ def allocate_next_work(
try:
kwargs: dict[str, Any] = {
"session_id": session_id,
"role": role_norm,
"role": lease_role,
"remote": remote,
"org": org,
"repo": repo,
@@ -992,11 +1227,27 @@ def allocate_next_work(
}
# assigned
lease_proof = {
"assignment_id": result.assignment_id,
"lease_id": result.lease_id,
"expires_at": result.expires_at,
"expected_head_sha": result.expected_head_sha,
"allowed_actions": list(result.allowed_actions),
"forbidden_actions": list(result.forbidden_actions),
"lease_role": lease_role,
"source": "control_plane_db.assign_and_lease",
}
return {
"success": True,
"outcome": OUTCOME_ASSIGNED,
"apply": True,
"role": role_norm,
"allocation_mode": mode,
"routing_role": role_norm,
"required_role": expected_role,
"selected_action": selection["selected_action"],
"required_profile": selection["required_profile"],
"required_namespace": selection["required_namespace"],
"profile_name": profile_name,
"username": username,
"session_id": session_id,
@@ -1012,16 +1263,16 @@ def allocate_next_work(
"skipped": [s.as_dict() for s in skipped],
"terminal_pr": terminal_pr,
"assignment": result.as_dict(),
"lease_proof": {
"assignment_id": result.assignment_id,
"lease_id": result.lease_id,
"expires_at": result.expires_at,
"expected_head_sha": result.expected_head_sha,
"allowed_actions": list(result.allowed_actions),
"forbidden_actions": list(result.forbidden_actions),
"source": "control_plane_db.assign_and_lease",
"lease_proof": lease_proof,
"allocation_evidence": {
"mode": "assigned",
"allocation_mode": mode,
"lease_created": True,
"lease_role": lease_role,
"lease_proof": lease_proof,
"selection_policy": SELECTION_POLICY,
},
"next_valid_command": _next_command(role_norm, selected),
"next_valid_command": _next_command(lease_role, selected),
"substrate": "control_plane_db",
"file_lock_only": False,
"comment_lease_only": False,
@@ -1031,9 +1282,12 @@ def allocate_next_work(
"controller_excluded": list(controller_excluded),
"exclude_issue_numbers": list(exclude_nums),
"candidate_set_fingerprint": cas_fp,
"controller_allowed_actions": list(controller_allowed_actions),
"controller_forbidden_actions": list(controller_forbidden_actions),
"downstream_note": (
"#612 incident bridge remains downstream of #600; "
"allocator never assigns raw monitoring incidents"
"allocator never assigns raw monitoring incidents; "
"controller routes only under cross_role (#840)"
),
}