Files
Gitea-Tools/allocator_service.py
T
sysadminandClaude Opus 4.8 d17f055e86 fix(allocator): pre-rank exclusions and candidates_json transport (#776)
Expose exclude_issue_numbers on gitea_allocate_next_work, remove excluded
numbers before ranking, normalize decoded-list and JSON-string
candidates_json fail-closed, and return candidate-set fingerprints for
dry-run/apply CAS. Same-owner leases on excluded issues surface a
structured resume/release blocker.

Closes #776

Co-Authored-By: Claude Opus 4.8 (1M context) <[email protected]>
2026-07-20 22:40:39 -04:00

1102 lines
42 KiB
Python

"""Controller-owned work allocator policy (#600).
Builds on the #613 control-plane DB substrate (``ControlPlaneDB.assign_and_lease``).
Workers must not self-select exclusive work under the standard multi-LLM
workflow. They call ``gitea_allocate_next_work`` which:
1. Inspects candidate Gitea issues/PRs (never raw monitoring incidents).
2. Applies ADR routing policy (role, terminal path, leases, blocked, deps).
3. Atomically assigns + leases the selected item in one DB transaction.
This module is pure selection + substrate orchestration. Gitea I/O for live
inventory lives in the MCP tool wrapper so tests can inject candidates.
#612 remains downstream: bridge-created Gitea issues become candidates only
after they exist as normal issues; this module never assigns incidents.
"""
from __future__ import annotations
import hashlib
import json
import os
import uuid
from dataclasses import dataclass, field
from typing import Any, Mapping, Sequence
from control_plane_db import (
ControlPlaneDB,
ControlPlaneError,
InvalidWorkKindError,
LeaseRequiredError,
WORK_KINDS,
)
# Outcomes required by #600 / ADR §5.
OUTCOME_ASSIGNED = "assigned_work"
OUTCOME_WAIT = "wait"
OUTCOME_BLOCKED_TERMINAL = "blocked_by_terminal_path"
OUTCOME_BLOCKED_LEASE = "blocked_by_active_lease"
OUTCOME_NEEDS_CONTROLLER = "needs_controller"
OUTCOME_NO_SAFE = "no_safe_work"
OUTCOME_ROLE_INELIGIBLE = "role_ineligible"
OUTCOME_PREVIEW = "preview" # dry-run only (apply=false)
# #765: ownership could not be established for every remaining candidate.
OUTCOME_OWNERSHIP_DEFECT = "allocator_ownership_defect"
# #776: excluded issue still carries a live same-owner lease — resume or release.
OUTCOME_BLOCKED_EXCLUDED_OWN_LEASE = "blocked_by_excluded_own_lease"
# #776: dry-run/apply candidate-set fingerprint mismatch (CAS drift).
OUTCOME_CANDIDATE_SET_DRIFT = "candidate_set_drift"
# #765 skip reason code for work already claimed by a different controller.
SKIP_CLAIMED_BY_OTHER_SESSION = "claimed_by_other_session"
# #776: controller-supplied pre-rank exclusion.
SKIP_EXCLUDED_BY_CONTROLLER = "excluded_by_controller"
# Ownership verdicts for a live claim on a candidate (#765).
OWNERSHIP_OWN = "own"
OWNERSHIP_FOREIGN = "foreign"
OWNERSHIP_UNKNOWN = "unknown"
# Human-readable statement of how a winner is chosen (#758 AC10). Reported
# alongside allocator results so the flat status:ready tier and its
# oldest-number tie-break are explicit rather than incidental.
SELECTION_POLICY = (
"rank complete inventory by (priority desc, PRs before issues, "
"number asc); status:ready issues share priority 20, so the oldest "
"eligible number wins ties; result limits never affect selection"
)
ROLE_AUTHOR = "author"
ROLE_REVIEWER = "reviewer"
ROLE_MERGER = "merger"
ROLE_RECONCILER = "reconciler"
ROLE_CONTROLLER = "controller"
VALID_ROLES = frozenset(
{ROLE_AUTHOR, ROLE_REVIEWER, ROLE_MERGER, ROLE_RECONCILER, ROLE_CONTROLLER}
)
# Default action matrices by role (mutation gate will re-check).
ROLE_ACTIONS: dict[str, tuple[tuple[str, ...], tuple[str, ...]]] = {
ROLE_AUTHOR: (
("implement", "comment", "push", "create_pr"),
("approve", "merge", "request_changes", "self_select_without_assignment"),
),
ROLE_REVIEWER: (
("review", "comment", "approve", "request_changes"),
("merge", "push", "create_pr", "self_select_without_assignment"),
),
ROLE_MERGER: (
("merge", "comment"),
("approve", "request_changes", "push", "create_pr", "self_select_without_assignment"),
),
ROLE_RECONCILER: (
("comment", "diagnose", "cleanup"),
("approve", "merge", "push", "create_pr", "self_select_without_assignment"),
),
ROLE_CONTROLLER: (
("comment", "diagnose", "allocate"),
("approve", "merge", "push", "create_pr"),
),
}
@dataclass
class WorkCandidate:
"""One assignable Gitea issue or PR presented to the allocator."""
kind: str # issue | pr
number: int
state: str = "open"
labels: tuple[str, ...] = ()
title: str = ""
priority: int = 0
head_sha: str | None = None
# Routing signals (callers derive from Gitea / review feedback).
request_changes_current_head: bool = False
approval_on_current_head: bool = False
approval_stale: bool = False
approval_contaminated: bool = False
mergeable: bool = False
blocked: bool = False
dependency_unmet: bool = False
dependency_reason: str | None = None
already_claimed_elsewhere: bool = False
def __post_init__(self) -> None:
self.kind = (self.kind or "").strip().lower()
self.state = (self.state or "open").strip().lower()
self.labels = tuple(
str(x).strip().lower() for x in (self.labels or ()) if str(x).strip()
)
if self.kind not in WORK_KINDS:
raise InvalidWorkKindError(
f"candidate kind '{self.kind}' is not assignable; only "
f"{sorted(WORK_KINDS)} (never raw incidents)"
)
def as_dict(self) -> dict[str, Any]:
return {
"kind": self.kind,
"number": self.number,
"state": self.state,
"labels": list(self.labels),
"title": self.title,
"priority": self.priority,
"head_sha": self.head_sha,
"request_changes_current_head": self.request_changes_current_head,
"approval_on_current_head": self.approval_on_current_head,
"approval_stale": self.approval_stale,
"approval_contaminated": self.approval_contaminated,
"mergeable": self.mergeable,
"blocked": self.blocked,
"dependency_unmet": self.dependency_unmet,
"dependency_reason": self.dependency_reason,
}
@dataclass
class SkipRecord:
kind: str
number: int
reason: str
reason_code: str | None = None
def as_dict(self) -> dict[str, Any]:
return {
"kind": self.kind,
"number": self.number,
"reason": self.reason,
"reason_code": self.reason_code,
}
CONTROLLER_INSTANCE_ENV = "GITEA_CONTROLLER_INSTANCE_ID"
def resolve_controller_instance_id(
env: Mapping[str, str] | None = None,
) -> str | None:
"""Return this controller's stable identity, or ``None`` if undeclared.
Deliberately has no derived fallback. The obvious candidates are unsafe:
``session_id`` is regenerated per invocation, and the MCP process pid is
shared by every controller attached to the same daemon — two independent
controllers really do report the same pid and profile. Guessing from either
would let one controller adopt another's lease, which is the failure #765
exists to prevent. When this returns ``None``, live claims are treated as
unidentified: they are excluded from selection and reported as ownership
defects rather than adopted.
"""
source = env if env is not None else os.environ
return (source.get(CONTROLLER_INSTANCE_ENV) or "").strip() or None
def classify_claim_ownership(
claim: dict[str, Any] | None,
*,
session_id: str | None,
controller_instance_id: str | None,
) -> str | None:
"""Classify a live claim as own / foreign / unknown ownership (#765).
Returns ``None`` when the candidate carries no live claim.
Session ids are regenerated per allocator invocation, so they only prove
ownership positively (an exact match is certainly this session). The
durable signal is ``controller_instance_id``. When either side lacks one,
ownership is *unknown*: the allocator must not assume that a lease sharing
the same profile belongs to this controller, so unknown is treated as
not-ours for selection purposes and reported as an ownership defect.
"""
if not claim:
return None
claim_session = str(claim.get("session_id") or "").strip()
claim_instance = str(claim.get("controller_instance_id") or "").strip()
own_session = str(session_id or "").strip()
own_instance = str(controller_instance_id or "").strip()
if claim_session and own_session and claim_session == own_session:
return OWNERSHIP_OWN
if claim_instance and own_instance:
return (
OWNERSHIP_OWN if claim_instance == own_instance else OWNERSHIP_FOREIGN
)
if not claim_instance and not own_instance:
# Neither side declares a controller identity. The session ids differ
# (an exact match returned OWN above), so this is simply someone
# else's lease: foreign, and we wait rather than adopt.
return OWNERSHIP_FOREIGN
# Exactly one side is identified, so the two cannot be compared: this may
# or may not be our own task under a different session id. Never guess.
return OWNERSHIP_UNKNOWN
def normalize_role(role: str | None, *, profile_name: str | None = None) -> str:
"""Map profile/role strings to a canonical allocator role."""
raw = (role or "").strip().lower()
if raw in VALID_ROLES:
return raw
prof = (profile_name or "").strip().lower()
for token in VALID_ROLES:
if token in prof or prof.endswith(f"-{token}"):
return token
if "author" in raw:
return ROLE_AUTHOR
if "review" in raw:
return ROLE_REVIEWER
if "merg" in raw:
return ROLE_MERGER
if "reconcil" in raw:
return ROLE_RECONCILER
if "control" in raw:
return ROLE_CONTROLLER
raise ControlPlaneError(
f"unknown allocator role '{role}' (profile={profile_name!r}); "
f"expected one of {sorted(VALID_ROLES)}"
)
def expected_role_for_candidate(c: WorkCandidate) -> str:
"""ADR §5.3 routing: which role should take this work next."""
if c.kind == "pr":
if c.approval_contaminated:
return ROLE_RECONCILER
if c.request_changes_current_head:
return ROLE_AUTHOR
if c.approval_stale:
return ROLE_REVIEWER
if c.approval_on_current_head and c.mergeable:
return ROLE_MERGER
# Open PR without terminal verdict → reviewer
return ROLE_REVIEWER
# Issues: ready work → author by default; blocked stays controller/none
labels = set(c.labels)
if "status:blocked" in labels or c.blocked:
return ROLE_CONTROLLER
return ROLE_AUTHOR
def role_actions(role: str) -> tuple[tuple[str, ...], tuple[str, ...]]:
return ROLE_ACTIONS.get(role, ROLE_ACTIONS[ROLE_AUTHOR])
def classify_skip(
c: WorkCandidate,
*,
role: str,
terminal_pr: int | None,
claim_ownership: str | None = None,
) -> str | None:
"""Return skip reason, or None if candidate is selectable for *role*.
*claim_ownership* (#765) is the verdict from
:func:`classify_claim_ownership` for this candidate's live claim. Foreign
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.
"""
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:
return f"{c.kind}#{c.number} is blocked"
if c.dependency_unmet:
return (
c.dependency_reason
or f"{c.kind}#{c.number} has unmet dependencies"
)
if claim_ownership in (OWNERSHIP_FOREIGN, OWNERSHIP_UNKNOWN):
detail = (
"owned by another controller instance"
if claim_ownership == OWNERSHIP_FOREIGN
else "owner could not be identified; never adopt on a guess"
)
return (
f"{c.kind}#{c.number} {SKIP_CLAIMED_BY_OTHER_SESSION}: "
f"active lease {detail}"
)
if c.already_claimed_elsewhere:
return f"{c.kind}#{c.number} already claimed elsewhere"
if c.kind == "pr" and not (c.head_sha or "").strip():
return f"pr#{c.number} missing head_sha pin"
# Terminal path first: when an active terminal PR exists, only that PR
# (or controller diagnosis) is assignable for review-path roles.
if terminal_pr is not None and c.kind == "pr" and c.number != terminal_pr:
if role in (ROLE_REVIEWER, ROLE_MERGER):
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 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} expects role '{expected}', active role is '{role}'"
)
# Ready-gate for issues: prefer status:ready when labels present.
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)})"
return None
def sort_candidates(candidates: Sequence[WorkCandidate]) -> list[WorkCandidate]:
"""Rank candidates deterministically (#758 AC10).
Ordering key, in precedence order:
1. ``priority`` descending — the loader scores ``status:ready`` issues at
20 and everything else at 1, so the ready queue ties at a single value
by design;
2. PRs before issues — in-flight review work drains before new authoring;
3. ``number`` ascending — oldest first, which is what actually breaks the
flat ``status:ready`` tie.
Because the ready tier is intentionally flat, rule 3 decides most real
selections. That is only safe when ranking sees the *complete* candidate
inventory: truncating before this call silently redefines "oldest" as
"oldest among whatever survived the slice", which is the defect #758
fixed. Callers must rank everything and bound reporting afterwards.
"""
return sorted(
candidates,
key=lambda c: (-int(c.priority), c.kind != "pr", int(c.number)),
)
def _require_strict_int(value: Any, *, field: str) -> int:
"""Parse an issue number; reject bools and non-integers (#776)."""
if isinstance(value, bool) or not isinstance(value, int):
raise ValueError(
f"{field} must be an integer (booleans and non-integers rejected; "
f"got {type(value).__name__})"
)
return int(value)
def normalize_exclude_issue_numbers(
exclude_issue_numbers: Any = None,
) -> list[int]:
"""Normalize controller-supplied exclusions to a sorted unique int list (#776).
``None`` / omitted → empty list (existing behavior). Accepts a list/tuple of
integers. Rejects scalars, bools-as-ints, nested structures, and strings.
"""
if exclude_issue_numbers is None:
return []
if isinstance(exclude_issue_numbers, (str, bytes)) or not isinstance(
exclude_issue_numbers, (list, tuple)
):
raise ValueError(
"exclude_issue_numbers must be a list of integers "
f"(got {type(exclude_issue_numbers).__name__})"
)
out: list[int] = []
seen: set[int] = set()
for idx, raw in enumerate(exclude_issue_numbers):
num = _require_strict_int(raw, field=f"exclude_issue_numbers[{idx}]")
if num not in seen:
seen.add(num)
out.append(num)
return sorted(out)
def candidate_set_fingerprint(
candidates: Sequence[WorkCandidate],
*,
exclude_issue_numbers: Sequence[int] | None = None,
) -> str:
"""Stable CAS fingerprint of normalized candidate set + exclusions (#776 AC4)."""
payload = {
"candidates": sorted(
({"kind": c.kind, "number": int(c.number)} for c in candidates),
key=lambda x: (x["kind"], x["number"]),
),
"exclude_issue_numbers": list(
normalize_exclude_issue_numbers(exclude_issue_numbers)
),
}
blob = json.dumps(payload, sort_keys=True, separators=(",", ":"))
return hashlib.sha256(blob.encode("utf-8")).hexdigest()
def normalize_candidates_payload(raw: Any) -> list[WorkCandidate]:
"""Decode MCP ``candidates_json`` from list or JSON string (#776 AC3).
Accepts:
* an already-decoded ``list`` of candidate dicts (native MCP transport);
* a valid JSON string that decodes to such a list (backward compatible).
Rejects malformed JSON, scalars, non-list containers, invalid records,
booleans-as-integers, and unsupported types with fail-closed ``ValueError``.
"""
if raw is None:
raise ValueError("candidates_json is empty")
if isinstance(raw, (bytes, bytearray)):
try:
raw = raw.decode("utf-8")
except Exception as exc: # noqa: BLE001
raise ValueError(
f"candidates_json bytes are not valid utf-8: {exc}"
) from exc
if isinstance(raw, str):
text = raw.strip()
if not text:
raise ValueError("candidates_json string is empty")
try:
decoded = json.loads(text)
except json.JSONDecodeError as exc:
raise ValueError(
f"malformed candidates_json JSON: {exc.msg} at pos {exc.pos}"
) from exc
raw = decoded
if not isinstance(raw, list):
raise ValueError(
"candidates_json must be a JSON list (or already-decoded list); "
f"got {type(raw).__name__}"
)
candidates: list[WorkCandidate] = []
for idx, item in enumerate(raw):
if not isinstance(item, dict):
raise ValueError(
f"invalid candidate record at index {idx}: expected object, "
f"got {type(item).__name__}"
)
try:
candidates.append(candidate_from_dict(item))
except (KeyError, TypeError, ValueError, InvalidWorkKindError) as exc:
raise ValueError(
f"invalid candidate record at index {idx}: {exc}"
) from exc
return candidates
def allocate_next_work(
db: ControlPlaneDB,
*,
session_id: str,
role: str,
remote: str,
org: str,
repo: str,
candidates: Sequence[WorkCandidate],
apply: bool = False,
profile_name: str | None = None,
username: str | None = None,
lease_ttl_seconds: int | None = None,
controller_instance_id: str | None = None,
claims: Mapping[tuple[str, int], dict[str, Any]] | None = None,
exclude_issue_numbers: Sequence[int] | None = None,
expected_candidate_set_fingerprint: 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.
*exclude_issue_numbers* (#776): numbers removed before ranking. Omitted /
empty preserves prior behavior.
*expected_candidate_set_fingerprint* (#776 AC4): when set on apply, rejects
material candidate-set drift vs a prior dry-run.
Never uses file locks or comment-only leases as the assignment source.
"""
if db is None:
return {
"success": False,
"outcome": OUTCOME_NO_SAFE,
"reasons": [
"control-plane DB substrate unavailable (fail closed, #600/#613)"
],
"skipped": [],
"assignment": None,
"substrate": "control_plane_db",
"file_lock_only": False,
"comment_lease_only": False,
}
try:
role_norm = normalize_role(role, profile_name=profile_name)
except ControlPlaneError as exc:
return {
"success": False,
"outcome": OUTCOME_ROLE_INELIGIBLE,
"reasons": [str(exc)],
"skipped": [],
"assignment": None,
"substrate": "control_plane_db",
}
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
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:
terminal = db.get_active_terminal_lock(remote=remote, org=org, repo=repo)
except Exception as exc: # noqa: BLE001
return {
"success": False,
"outcome": OUTCOME_NO_SAFE,
"reasons": [f"terminal lock lookup failed: {exc} (fail closed)"],
"skipped": [],
"assignment": None,
"substrate": "control_plane_db",
}
terminal_pr = int(terminal["terminal_pr"]) if terminal else None
# #765: live claims exclude work owned by a *different* controller before
# ranking, so one session's in-progress task cannot blockade the queue.
# #776 AC5: load live claims in this call path immediately before selection
# (and before apply reserve) so ownership is never stale within the
# allocation attempt. Test callers may inject *claims* explicitly.
if claims is None:
try:
claims = db.list_active_claims(remote=remote, org=org, repo=repo)
except Exception as exc: # noqa: BLE001
return {
"success": False,
"outcome": OUTCOME_NO_SAFE,
"reasons": [
f"active claim lookup failed: {exc} (fail closed, #765)"
],
"skipped": [],
"assignment": None,
"substrate": "control_plane_db",
}
try:
exclude_nums = normalize_exclude_issue_numbers(exclude_issue_numbers)
except ValueError as exc:
return {
"success": False,
"outcome": OUTCOME_NO_SAFE,
"apply": bool(apply),
"reasons": [
f"invalid exclude_issue_numbers: {exc} (fail closed, #776)"
],
"skipped": [],
"assignment": None,
"substrate": "control_plane_db",
}
exclude_set = set(exclude_nums)
cas_fp = candidate_set_fingerprint(
candidates, exclude_issue_numbers=exclude_nums
)
expected_fp = (expected_candidate_set_fingerprint or "").strip() or None
if expected_fp and expected_fp != cas_fp:
return {
"success": False,
"outcome": OUTCOME_CANDIDATE_SET_DRIFT,
"apply": bool(apply),
"reasons": [
"candidate-set fingerprint drift: apply rejected rather than "
"silently leasing a different candidate (#776 AC4)"
],
"candidate_set_fingerprint": cas_fp,
"expected_candidate_set_fingerprint": expected_fp,
"exclude_issue_numbers": list(exclude_nums),
"skipped": [],
"assignment": None,
"substrate": "control_plane_db",
"file_lock_only": False,
"comment_lease_only": False,
}
skipped: list[SkipRecord] = []
claims_excluded: list[dict[str, Any]] = []
ownership_defects: list[dict[str, Any]] = []
controller_excluded: list[dict[str, Any]] = []
# #776 AC2: remove excluded numbers *before* ranking / selection / lease.
rankable: list[WorkCandidate] = []
for c in candidates:
if int(c.number) in exclude_set:
reason = (
f"{c.kind}#{c.number} {SKIP_EXCLUDED_BY_CONTROLLER}: "
"controller pre-rank exclusion"
)
skipped.append(
SkipRecord(
c.kind,
c.number,
reason,
SKIP_EXCLUDED_BY_CONTROLLER,
)
)
controller_excluded.append(
{
"kind": c.kind,
"number": c.number,
"reason_code": SKIP_EXCLUDED_BY_CONTROLLER,
}
)
# #776 AC5: same-owner live lease on an excluded issue is a
# structured resume/release blocker, never a silent strand.
claim = claims.get((c.kind, int(c.number))) if claims else None
ownership = classify_claim_ownership(
claim,
session_id=session_id,
controller_instance_id=controller_instance_id,
)
if ownership == OWNERSHIP_OWN and claim:
return {
"success": True,
"outcome": OUTCOME_BLOCKED_EXCLUDED_OWN_LEASE,
"apply": bool(apply),
"role": role_norm,
"profile_name": profile_name,
"username": username,
"session_id": session_id,
"remote": remote,
"org": org,
"repo": repo,
"selected": None,
"expected_role_next": None,
"reasons": [
f"{c.kind}#{c.number} is excluded_by_controller but "
"carries a live same-owner lease; resume or release "
"that lease before allocating other work (#776 AC5)"
],
"skipped": [s.as_dict() for s in skipped],
"terminal_pr": terminal_pr,
"assignment": None,
"substrate": "control_plane_db",
"file_lock_only": False,
"comment_lease_only": False,
"controller_instance_id": controller_instance_id,
"claims_excluded": list(claims_excluded),
"ownership_defects": list(ownership_defects),
"controller_excluded": list(controller_excluded),
"exclude_issue_numbers": list(exclude_nums),
"candidate_set_fingerprint": cas_fp,
"blocked_lease": {
"kind": c.kind,
"number": c.number,
"lease_id": claim.get("lease_id"),
"owner_session_id": claim.get("session_id"),
"owner_controller_instance_id": claim.get(
"controller_instance_id"
),
"expires_at": claim.get("expires_at"),
"safe_next_action": (
"resume the same-owner lease or release it, then "
"re-run allocation without stranding the excluded "
"issue"
),
},
}
continue
rankable.append(c)
ordered = sort_candidates(rankable)
selected: WorkCandidate | None = None
for c in ordered:
claim = claims.get((c.kind, int(c.number))) if claims else None
ownership = classify_claim_ownership(
claim,
session_id=session_id,
controller_instance_id=controller_instance_id,
)
reason = classify_skip(
c,
role=role_norm,
terminal_pr=terminal_pr,
claim_ownership=ownership,
)
if reason:
is_claim_skip = SKIP_CLAIMED_BY_OTHER_SESSION in reason
skipped.append(
SkipRecord(
c.kind,
c.number,
reason,
SKIP_CLAIMED_BY_OTHER_SESSION if is_claim_skip else None,
)
)
if is_claim_skip and claim:
record = {
"kind": c.kind,
"number": c.number,
"ownership": ownership,
"lease_id": claim.get("lease_id"),
"owner_session_id": claim.get("session_id"),
"owner_controller_instance_id": claim.get(
"controller_instance_id"
),
"expires_at": claim.get("expires_at"),
}
claims_excluded.append(record)
if ownership == OWNERSHIP_UNKNOWN:
ownership_defects.append(record)
continue
selected = c
break
if selected is None:
# If terminal lock blocks all review work, surface that explicitly.
owner_session_id: str | None = None
if terminal_pr is not None and role_norm in (ROLE_REVIEWER, ROLE_MERGER):
outcome = OUTCOME_BLOCKED_TERMINAL
reasons = [
f"no safe work for role '{role_norm}': active terminal-review "
f"lock on PR #{terminal_pr} (resolve terminal path first, #332/#600)"
]
elif ownership_defects:
# #765: every remaining candidate is claimed and at least one owner
# could not be identified. Report the defect; never adopt.
outcome = OUTCOME_OWNERSHIP_DEFECT
reasons = [
f"no safe assignable work for role '{role_norm}': "
f"{len(ownership_defects)} candidate(s) carry an active lease "
"whose controller ownership could not be established. Record a "
"controller_instance_id on those sessions; the allocator will "
"not assume a shared profile means shared ownership (#765)."
]
elif claims_excluded:
outcome = OUTCOME_WAIT
reasons = [
f"no unclaimed work for role '{role_norm}': "
f"{len(claims_excluded)} candidate(s) are actively claimed by "
"another controller. Waiting; their leases are not adopted (#765)."
]
# Preserve the pre-#765 wait contract: name the blocking owner.
owner_session_id = claims_excluded[0].get("owner_session_id")
elif controller_excluded and not ordered:
# #776 AC7: every candidate was controller-excluded → wait / no lease.
outcome = OUTCOME_WAIT
reasons = [
f"no assignable work for role '{role_norm}': all "
f"{len(controller_excluded)} candidate(s) were removed by "
f"{SKIP_EXCLUDED_BY_CONTROLLER} before ranking; no assignment "
"or lease created (#776 AC7)"
]
else:
outcome = OUTCOME_NO_SAFE
reasons = [
f"no safe assignable work for role '{role_norm}' "
f"among {len(ordered)} rankable candidates "
f"({len(controller_excluded)} controller-excluded)"
]
return {
"success": True,
"outcome": outcome,
"apply": bool(apply),
"role": role_norm,
"profile_name": profile_name,
"username": username,
"session_id": session_id,
"remote": remote,
"org": org,
"repo": repo,
"selected": None,
"expected_role_next": None,
"reasons": reasons,
"skipped": [s.as_dict() for s in skipped],
"terminal_pr": terminal_pr,
"assignment": None,
"substrate": "control_plane_db",
"file_lock_only": False,
"comment_lease_only": False,
"controller_instance_id": controller_instance_id,
"claims_excluded": list(claims_excluded),
"ownership_defects": list(ownership_defects),
"controller_excluded": list(controller_excluded),
"exclude_issue_numbers": list(exclude_nums),
"candidate_set_fingerprint": cas_fp,
"owner_session_id": owner_session_id,
"downstream_note": (
"#612 incident bridge remains downstream of #600; "
"allocator never assigns raw monitoring incidents"
),
}
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})"
),
}
if not apply:
return {
"success": True,
"outcome": OUTCOME_PREVIEW,
"apply": False,
"role": role_norm,
"profile_name": profile_name,
"username": username,
"session_id": session_id,
"remote": remote,
"org": org,
"repo": repo,
"selected": selection,
"expected_role_next": expected_role,
"reasons": [
"dry-run only (apply=false); no assignment/lease created — "
"call again with apply=true to reserve via control-plane DB"
],
"skipped": [s.as_dict() for s in skipped],
"terminal_pr": terminal_pr,
"assignment": None,
"substrate": "control_plane_db",
"file_lock_only": False,
"comment_lease_only": False,
"controller_instance_id": controller_instance_id,
"claims_excluded": list(claims_excluded),
"ownership_defects": list(ownership_defects),
"controller_excluded": list(controller_excluded),
"exclude_issue_numbers": list(exclude_nums),
"candidate_set_fingerprint": cas_fp,
"downstream_note": (
"#612 incident bridge remains downstream of #600; "
"allocator never assigns raw monitoring incidents"
),
}
# Atomic reserve via #613 substrate.
ttl = lease_ttl_seconds if lease_ttl_seconds is not None else None
try:
kwargs: dict[str, Any] = {
"session_id": session_id,
"role": role_norm,
"remote": remote,
"org": org,
"repo": repo,
"kind": selected.kind,
"number": selected.number,
"expected_head_sha": selected.head_sha,
"allowed_actions": allowed,
"forbidden_actions": forbidden,
"phase": "allocated",
}
if ttl is not None:
kwargs["lease_ttl_seconds"] = int(ttl)
result = db.assign_and_lease(**kwargs)
except (InvalidWorkKindError, LeaseRequiredError, ControlPlaneError) as exc:
return {
"success": False,
"outcome": OUTCOME_NO_SAFE,
"apply": True,
"role": role_norm,
"session_id": session_id,
"remote": remote,
"org": org,
"repo": repo,
"selected": selection,
"expected_role_next": expected_role,
"reasons": [f"atomic assign+lease failed: {exc} (fail closed, #613)"],
"skipped": [s.as_dict() for s in skipped],
"terminal_pr": terminal_pr,
"assignment": None,
"substrate": "control_plane_db",
"file_lock_only": False,
"comment_lease_only": False,
}
if result.outcome == "wait":
return {
"success": True,
"outcome": OUTCOME_WAIT,
"apply": True,
"role": role_norm,
"session_id": session_id,
"remote": remote,
"org": org,
"repo": repo,
"selected": selection,
"expected_role_next": expected_role,
"reasons": [result.reason or "foreign active lease"],
"skipped": [s.as_dict() for s in skipped],
"terminal_pr": terminal_pr,
"assignment": result.as_dict(),
"owner_session_id": result.owner_session_id,
"substrate": "control_plane_db",
"file_lock_only": False,
"comment_lease_only": False,
}
if result.outcome == "no_safe_work":
return {
"success": True,
"outcome": OUTCOME_NO_SAFE,
"apply": True,
"role": role_norm,
"session_id": session_id,
"remote": remote,
"org": org,
"repo": repo,
"selected": selection,
"expected_role_next": expected_role,
"reasons": [result.reason or "no_safe_work"],
"skipped": [s.as_dict() for s in skipped],
"terminal_pr": terminal_pr,
"assignment": result.as_dict(),
"substrate": "control_plane_db",
"file_lock_only": False,
"comment_lease_only": False,
}
# assigned
return {
"success": True,
"outcome": OUTCOME_ASSIGNED,
"apply": True,
"role": role_norm,
"profile_name": profile_name,
"username": username,
"session_id": session_id,
"remote": remote,
"org": org,
"repo": repo,
"selected": selection,
"expected_role_next": expected_role,
"reasons": [
selection["reason_selected"],
result.reason or "atomic assign+lease created",
],
"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",
},
"next_valid_command": _next_command(role_norm, selected),
"substrate": "control_plane_db",
"file_lock_only": False,
"comment_lease_only": False,
"controller_instance_id": controller_instance_id,
"claims_excluded": list(claims_excluded),
"ownership_defects": list(ownership_defects),
"controller_excluded": list(controller_excluded),
"exclude_issue_numbers": list(exclude_nums),
"candidate_set_fingerprint": cas_fp,
"downstream_note": (
"#612 incident bridge remains downstream of #600; "
"allocator never assigns raw monitoring incidents"
),
}
def _next_command(role: str, c: WorkCandidate) -> str:
if role == ROLE_AUTHOR and c.kind == "issue":
return f"implement issue #{c.number} under a branches/ worktree; open PR when ready"
if role == ROLE_AUTHOR and c.kind == "pr":
return (
f"address REQUEST_CHANGES on PR #{c.number} at head "
f"{(c.head_sha or '')[:12]} and push fixes"
)
if role == ROLE_REVIEWER:
return (
f"review PR #{c.number} pinned at head {(c.head_sha or '')[:12]} "
"via full reviewer workflow"
)
if role == ROLE_MERGER:
return (
f"merge PR #{c.number} only with explicit operator MERGE "
f"authorization at head {(c.head_sha or '')[:12]}"
)
if role == ROLE_RECONCILER:
return f"diagnose contested state for {c.kind}#{c.number}"
return f"proceed on {c.kind}#{c.number} under role {role}"
def candidate_from_dict(data: dict[str, Any]) -> WorkCandidate:
"""Build a WorkCandidate from a plain dict (tests / MCP inventory).
#776: reject booleans-as-integers and non-int numbers fail-closed.
"""
if "number" not in data:
raise KeyError("number")
number = _require_strict_int(data["number"], field="number")
priority_raw = data.get("priority") or 0
if isinstance(priority_raw, bool) or not isinstance(priority_raw, (int, float)):
# Allow numeric strings only for priority? Keep strict for bools.
if isinstance(priority_raw, str) and priority_raw.strip().lstrip("-").isdigit():
priority = int(priority_raw)
else:
raise ValueError(
f"priority must be numeric (booleans rejected; "
f"got {type(priority_raw).__name__})"
)
else:
priority = int(priority_raw)
return WorkCandidate(
kind=str(data.get("kind") or "issue"),
number=number,
state=str(data.get("state") or "open"),
labels=tuple(data.get("labels") or ()),
title=str(data.get("title") or ""),
priority=priority,
head_sha=data.get("head_sha"),
request_changes_current_head=bool(data.get("request_changes_current_head")),
approval_on_current_head=bool(data.get("approval_on_current_head")),
approval_stale=bool(data.get("approval_stale")),
approval_contaminated=bool(data.get("approval_contaminated")),
mergeable=bool(data.get("mergeable")),
blocked=bool(data.get("blocked")),
dependency_unmet=bool(data.get("dependency_unmet")),
dependency_reason=data.get("dependency_reason"),
already_claimed_elsewhere=bool(data.get("already_claimed_elsewhere")),
)