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]>
582 lines
20 KiB
Python
582 lines
20 KiB
Python
"""Authoritative controller cross-role generic queue allocation (#840)."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import os
|
|
import tempfile
|
|
import unittest
|
|
from unittest.mock import patch
|
|
|
|
from allocator_service import (
|
|
ALLOCATION_MODE_CROSS_ROLE,
|
|
ALLOCATION_MODE_ROLE_SCOPED,
|
|
OUTCOME_NO_SAFE,
|
|
OUTCOME_PREVIEW,
|
|
OUTCOME_WAIT,
|
|
ROLE_AUTHOR,
|
|
ROLE_CONTROLLER,
|
|
ROLE_MERGER,
|
|
ROLE_RECONCILER,
|
|
ROLE_REVIEWER,
|
|
WorkCandidate,
|
|
allocate_next_work,
|
|
build_selection_dict,
|
|
classify_skip,
|
|
required_namespace_for_role,
|
|
required_profile_for_role,
|
|
resolve_allocation_mode,
|
|
selected_action_for_candidate,
|
|
)
|
|
from control_plane_db import ControlPlaneDB
|
|
import role_session_router
|
|
from role_session_router import (
|
|
ROUTE_ALLOWED,
|
|
ROUTE_AMBIGUOUS,
|
|
ROUTE_WRONG_ROLE,
|
|
route_task_session,
|
|
)
|
|
import namespace_workspace_binding as nwb
|
|
import task_capability_map
|
|
|
|
|
|
class CrossRoleAllocationModeTest(unittest.TestCase):
|
|
def test_controller_defaults_to_cross_role(self) -> None:
|
|
self.assertEqual(
|
|
resolve_allocation_mode(ROLE_CONTROLLER),
|
|
ALLOCATION_MODE_CROSS_ROLE,
|
|
)
|
|
|
|
def test_worker_defaults_to_role_scoped(self) -> None:
|
|
for role in (ROLE_AUTHOR, ROLE_REVIEWER, ROLE_MERGER, ROLE_RECONCILER):
|
|
self.assertEqual(
|
|
resolve_allocation_mode(role),
|
|
ALLOCATION_MODE_ROLE_SCOPED,
|
|
)
|
|
|
|
def test_explicit_modes(self) -> None:
|
|
self.assertEqual(
|
|
resolve_allocation_mode(ROLE_CONTROLLER, "role_scoped"),
|
|
ALLOCATION_MODE_ROLE_SCOPED,
|
|
)
|
|
self.assertEqual(
|
|
resolve_allocation_mode(ROLE_AUTHOR, "cross_role"),
|
|
ALLOCATION_MODE_CROSS_ROLE,
|
|
)
|
|
|
|
|
|
class CrossRoleSelectionPayloadTest(unittest.TestCase):
|
|
def test_selection_contains_required_fields(self) -> None:
|
|
c = WorkCandidate(
|
|
kind="issue",
|
|
number=840,
|
|
labels=("status:ready",),
|
|
title="cross-role",
|
|
priority=20,
|
|
)
|
|
sel = build_selection_dict(
|
|
c,
|
|
active_role=ROLE_CONTROLLER,
|
|
required_role=ROLE_AUTHOR,
|
|
profile_name="prgs-controller",
|
|
allocation_mode=ALLOCATION_MODE_CROSS_ROLE,
|
|
)
|
|
self.assertEqual(sel["number"], 840)
|
|
self.assertEqual(sel["kind"], "issue")
|
|
self.assertEqual(sel["required_role"], ROLE_AUTHOR)
|
|
self.assertEqual(sel["selected_action"], "implement")
|
|
self.assertEqual(sel["action"], "implement")
|
|
self.assertEqual(sel["required_profile"], "prgs-author")
|
|
self.assertEqual(sel["required_namespace"], "gitea-author")
|
|
self.assertEqual(sel["pinned"]["number"], 840)
|
|
self.assertIsNone(sel["pinned"]["head_sha"])
|
|
|
|
def test_profile_prefix_preserved(self) -> None:
|
|
self.assertEqual(
|
|
required_profile_for_role(ROLE_REVIEWER, profile_name="dadeschools-controller"),
|
|
"dadeschools-reviewer",
|
|
)
|
|
self.assertEqual(
|
|
required_namespace_for_role(ROLE_MERGER),
|
|
"gitea-merger",
|
|
)
|
|
|
|
def test_selected_actions_per_role(self) -> None:
|
|
issue = WorkCandidate(kind="issue", number=1, labels=("status:ready",))
|
|
pr_review = WorkCandidate(kind="pr", number=2, head_sha="a" * 40)
|
|
pr_rc = WorkCandidate(
|
|
kind="pr",
|
|
number=3,
|
|
head_sha="b" * 40,
|
|
request_changes_current_head=True,
|
|
)
|
|
pr_merge = WorkCandidate(
|
|
kind="pr",
|
|
number=4,
|
|
head_sha="c" * 40,
|
|
approval_on_current_head=True,
|
|
mergeable=True,
|
|
)
|
|
pr_recon = WorkCandidate(
|
|
kind="pr",
|
|
number=5,
|
|
head_sha="d" * 40,
|
|
approval_contaminated=True,
|
|
)
|
|
self.assertEqual(selected_action_for_candidate(issue, ROLE_AUTHOR), "implement")
|
|
self.assertEqual(
|
|
selected_action_for_candidate(pr_rc, ROLE_AUTHOR),
|
|
"address_pr_change_requests",
|
|
)
|
|
self.assertEqual(
|
|
selected_action_for_candidate(pr_review, ROLE_REVIEWER), "review"
|
|
)
|
|
self.assertEqual(selected_action_for_candidate(pr_merge, ROLE_MERGER), "merge")
|
|
self.assertEqual(
|
|
selected_action_for_candidate(pr_recon, ROLE_RECONCILER),
|
|
"reconcile_contaminated_approval",
|
|
)
|
|
|
|
|
|
class CrossRoleAllocateServiceTest(unittest.TestCase):
|
|
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="ctrl-session",
|
|
role=ROLE_CONTROLLER,
|
|
remote="prgs",
|
|
org="org",
|
|
repo="repo",
|
|
candidates=[],
|
|
apply=False,
|
|
profile_name="prgs-controller",
|
|
username="controller-bot",
|
|
controller_instance_id="ctrl-1",
|
|
)
|
|
defaults.update(kwargs)
|
|
return allocate_next_work(**defaults)
|
|
|
|
def test_eligible_author_work(self) -> None:
|
|
cands = [
|
|
WorkCandidate(
|
|
kind="issue",
|
|
number=100,
|
|
labels=("status:ready",),
|
|
title="author work",
|
|
priority=20,
|
|
),
|
|
]
|
|
res = self._alloc(candidates=cands)
|
|
self.assertTrue(res["success"])
|
|
self.assertEqual(res["outcome"], OUTCOME_PREVIEW)
|
|
self.assertEqual(res["allocation_mode"], ALLOCATION_MODE_CROSS_ROLE)
|
|
self.assertIsNotNone(res["selected"])
|
|
self.assertEqual(res["selected"]["number"], 100)
|
|
self.assertEqual(res["required_role"], ROLE_AUTHOR)
|
|
self.assertEqual(res["selected_action"], "implement")
|
|
self.assertEqual(res["required_profile"], "prgs-author")
|
|
self.assertEqual(res["required_namespace"], "gitea-author")
|
|
self.assertIn("allocate", res["controller_allowed_actions"])
|
|
self.assertIn("merge", res["controller_forbidden_actions"])
|
|
self.assertFalse(res["allocation_evidence"]["lease_created"])
|
|
|
|
def test_eligible_reviewer_work(self) -> None:
|
|
cands = [
|
|
WorkCandidate(
|
|
kind="pr",
|
|
number=200,
|
|
head_sha="e" * 40,
|
|
title="needs review",
|
|
priority=30,
|
|
),
|
|
]
|
|
res = self._alloc(candidates=cands)
|
|
self.assertEqual(res["selected"]["number"], 200)
|
|
self.assertEqual(res["required_role"], ROLE_REVIEWER)
|
|
self.assertEqual(res["selected_action"], "review")
|
|
self.assertEqual(res["required_profile"], "prgs-reviewer")
|
|
self.assertEqual(res["selected"]["pinned"]["head_sha"], "e" * 40)
|
|
|
|
def test_eligible_merger_work(self) -> None:
|
|
cands = [
|
|
WorkCandidate(
|
|
kind="pr",
|
|
number=300,
|
|
head_sha="f" * 40,
|
|
approval_on_current_head=True,
|
|
mergeable=True,
|
|
priority=40,
|
|
),
|
|
]
|
|
res = self._alloc(candidates=cands)
|
|
self.assertEqual(res["selected"]["number"], 300)
|
|
self.assertEqual(res["required_role"], ROLE_MERGER)
|
|
self.assertEqual(res["selected_action"], "merge")
|
|
|
|
def test_eligible_reconciler_work(self) -> None:
|
|
cands = [
|
|
WorkCandidate(
|
|
kind="pr",
|
|
number=400,
|
|
head_sha="1" * 40,
|
|
approval_contaminated=True,
|
|
priority=50,
|
|
),
|
|
]
|
|
res = self._alloc(candidates=cands)
|
|
self.assertEqual(res["selected"]["number"], 400)
|
|
self.assertEqual(res["required_role"], ROLE_RECONCILER)
|
|
self.assertIn("reconcile", res["selected_action"])
|
|
|
|
def test_no_eligible_work(self) -> None:
|
|
cands = [
|
|
WorkCandidate(
|
|
kind="issue",
|
|
number=10,
|
|
labels=("status:blocked",),
|
|
blocked=True,
|
|
priority=99,
|
|
),
|
|
WorkCandidate(
|
|
kind="issue",
|
|
number=11,
|
|
labels=("status:ready",),
|
|
dependency_unmet=True,
|
|
dependency_reason="blocked by #10",
|
|
priority=98,
|
|
),
|
|
]
|
|
res = self._alloc(candidates=cands)
|
|
self.assertTrue(res["success"])
|
|
self.assertEqual(res["outcome"], OUTCOME_NO_SAFE)
|
|
self.assertIsNone(res["selected"])
|
|
self.assertEqual(res["allocation_mode"], ALLOCATION_MODE_CROSS_ROLE)
|
|
|
|
def test_leased_work_skipped(self) -> None:
|
|
cands = [
|
|
WorkCandidate(
|
|
kind="issue",
|
|
number=50,
|
|
labels=("status:ready",),
|
|
priority=20,
|
|
),
|
|
WorkCandidate(
|
|
kind="issue",
|
|
number=51,
|
|
labels=("status:ready",),
|
|
priority=10,
|
|
),
|
|
]
|
|
# Seed a foreign lease on issue 50 via assign_and_lease under another session.
|
|
other = allocate_next_work(
|
|
self.db,
|
|
session_id="other-worker",
|
|
role=ROLE_AUTHOR,
|
|
remote="prgs",
|
|
org="org",
|
|
repo="repo",
|
|
candidates=cands[:1],
|
|
apply=True,
|
|
profile_name="prgs-author",
|
|
controller_instance_id="other-ctrl",
|
|
)
|
|
self.assertEqual(other["outcome"], "assigned_work")
|
|
res = self._alloc(candidates=cands)
|
|
self.assertIsNotNone(res["selected"])
|
|
self.assertEqual(res["selected"]["number"], 51)
|
|
self.assertTrue(any(s["number"] == 50 for s in res["skipped"]))
|
|
self.assertTrue(res["claims_excluded"])
|
|
|
|
def test_dependencies_skipped(self) -> None:
|
|
cands = [
|
|
WorkCandidate(
|
|
kind="issue",
|
|
number=1,
|
|
labels=("status:ready",),
|
|
priority=99,
|
|
dependency_unmet=True,
|
|
dependency_reason="needs #2",
|
|
),
|
|
WorkCandidate(
|
|
kind="issue",
|
|
number=2,
|
|
labels=("status:ready",),
|
|
priority=1,
|
|
),
|
|
]
|
|
res = self._alloc(candidates=cands)
|
|
self.assertEqual(res["selected"]["number"], 2)
|
|
skipped = {s["number"]: s["reason"] for s in res["skipped"]}
|
|
self.assertIn(1, skipped)
|
|
self.assertIn("needs #2", skipped[1])
|
|
|
|
def test_pagination_limit_only_truncates_skip_report(self) -> None:
|
|
"""Ranking uses full inventory; reporting limit is MCP-layer only.
|
|
|
|
Service ranks all candidates; prove higher-priority eligible item
|
|
wins even when many skipped precede it.
|
|
"""
|
|
cands = []
|
|
for n in range(1, 30):
|
|
cands.append(
|
|
WorkCandidate(
|
|
kind="issue",
|
|
number=n,
|
|
labels=("status:ready",),
|
|
priority=100 - n,
|
|
dependency_unmet=True,
|
|
dependency_reason=f"dep {n}",
|
|
)
|
|
)
|
|
cands.append(
|
|
WorkCandidate(
|
|
kind="issue",
|
|
number=999,
|
|
labels=("status:ready",),
|
|
priority=1,
|
|
)
|
|
)
|
|
res = self._alloc(candidates=cands)
|
|
self.assertEqual(res["selected"]["number"], 999)
|
|
self.assertGreaterEqual(len(res["skipped"]), 29)
|
|
|
|
def test_role_scoped_controller_legacy_still_restricts(self) -> None:
|
|
"""role_scoped controller only takes reconciler-needed items."""
|
|
cands = [
|
|
WorkCandidate(
|
|
kind="issue",
|
|
number=1,
|
|
labels=("status:ready",),
|
|
priority=50,
|
|
),
|
|
WorkCandidate(
|
|
kind="pr",
|
|
number=2,
|
|
head_sha="a" * 40,
|
|
approval_contaminated=True,
|
|
priority=1,
|
|
),
|
|
]
|
|
res = self._alloc(
|
|
candidates=cands,
|
|
allocation_mode=ALLOCATION_MODE_ROLE_SCOPED,
|
|
)
|
|
self.assertEqual(res["allocation_mode"], ALLOCATION_MODE_ROLE_SCOPED)
|
|
self.assertEqual(res["selected"]["number"], 2)
|
|
self.assertEqual(res["required_role"], ROLE_RECONCILER)
|
|
|
|
def test_cross_role_prefers_highest_priority_across_roles(self) -> None:
|
|
cands = [
|
|
WorkCandidate(
|
|
kind="issue",
|
|
number=10,
|
|
labels=("status:ready",),
|
|
priority=10,
|
|
),
|
|
WorkCandidate(
|
|
kind="pr",
|
|
number=20,
|
|
head_sha="b" * 40,
|
|
priority=50,
|
|
),
|
|
WorkCandidate(
|
|
kind="pr",
|
|
number=30,
|
|
head_sha="c" * 40,
|
|
approval_on_current_head=True,
|
|
mergeable=True,
|
|
priority=20,
|
|
),
|
|
]
|
|
res = self._alloc(candidates=cands)
|
|
# PR #20 highest priority → reviewer
|
|
self.assertEqual(res["selected"]["number"], 20)
|
|
self.assertEqual(res["required_role"], ROLE_REVIEWER)
|
|
|
|
def test_apply_creates_lease_evidence_for_required_role(self) -> None:
|
|
cands = [
|
|
WorkCandidate(
|
|
kind="issue",
|
|
number=777,
|
|
labels=("status:ready",),
|
|
priority=20,
|
|
),
|
|
]
|
|
res = self._alloc(candidates=cands, apply=True)
|
|
self.assertEqual(res["outcome"], "assigned_work")
|
|
self.assertTrue(res["allocation_evidence"]["lease_created"])
|
|
self.assertEqual(res["allocation_evidence"]["lease_role"], ROLE_AUTHOR)
|
|
proof = res["lease_proof"]
|
|
self.assertIsNotNone(proof["lease_id"])
|
|
self.assertEqual(proof["lease_role"], ROLE_AUTHOR)
|
|
self.assertIn("implement", proof["allowed_actions"])
|
|
# Controller isolation: controller still forbids merge/push/create_pr
|
|
self.assertIn("merge", res["controller_forbidden_actions"])
|
|
self.assertIn("push", res["controller_forbidden_actions"])
|
|
|
|
def test_metadata_consistency_role_is_controller(self) -> None:
|
|
cands = [
|
|
WorkCandidate(
|
|
kind="issue",
|
|
number=1,
|
|
labels=("status:ready",),
|
|
),
|
|
]
|
|
res = self._alloc(candidates=cands)
|
|
self.assertEqual(res["role"], ROLE_CONTROLLER)
|
|
self.assertEqual(res["routing_role"], ROLE_CONTROLLER)
|
|
self.assertEqual(res["required_role"], ROLE_AUTHOR)
|
|
|
|
|
|
class ProcessWorkQueueRouterTest(unittest.TestCase):
|
|
def tearDown(self) -> None:
|
|
role_session_router.clear_route_state()
|
|
|
|
def test_process_work_queue_allowed_for_controller(self) -> None:
|
|
res = route_task_session(
|
|
"process_work_queue",
|
|
active_profile="prgs-controller",
|
|
active_role_kind="controller",
|
|
allowed_in_current_session=True,
|
|
)
|
|
self.assertEqual(res["route_result"], ROUTE_ALLOWED)
|
|
self.assertEqual(res["required_role"], "controller")
|
|
self.assertTrue(res["downstream_allowed"])
|
|
|
|
def test_process_work_queue_hyphen_alias(self) -> None:
|
|
res = route_task_session(
|
|
"process-work-queue",
|
|
active_profile="prgs-controller",
|
|
active_role_kind="controller",
|
|
allowed_in_current_session=True,
|
|
)
|
|
self.assertEqual(res["route_result"], ROUTE_ALLOWED)
|
|
|
|
def test_process_work_queue_wrong_role_for_author(self) -> None:
|
|
res = route_task_session(
|
|
"process_work_queue",
|
|
active_profile="prgs-author",
|
|
active_role_kind="author",
|
|
allowed_in_current_session=False,
|
|
)
|
|
self.assertEqual(res["route_result"], ROUTE_WRONG_ROLE)
|
|
self.assertEqual(res["required_role"], "controller")
|
|
self.assertFalse(res["downstream_allowed"])
|
|
|
|
def test_unknown_still_ambiguous(self) -> None:
|
|
res = route_task_session(
|
|
"not_a_real_task",
|
|
active_profile="prgs-controller",
|
|
active_role_kind="controller",
|
|
allowed_in_current_session=False,
|
|
)
|
|
self.assertEqual(res["route_result"], ROUTE_AMBIGUOUS)
|
|
|
|
def test_capability_map_process_work_queue_is_controller(self) -> None:
|
|
self.assertEqual(
|
|
task_capability_map.required_role("process_work_queue"),
|
|
"controller",
|
|
)
|
|
self.assertEqual(
|
|
task_capability_map.required_permission("process_work_queue"),
|
|
"gitea.read",
|
|
)
|
|
|
|
|
|
class ControllerRoleMetadataTest(unittest.TestCase):
|
|
def test_normalize_role_kind_controller(self) -> None:
|
|
self.assertEqual(
|
|
nwb.normalize_role_kind("controller"),
|
|
"controller",
|
|
)
|
|
self.assertEqual(
|
|
nwb.normalize_role_kind("author", profile_name="prgs-controller"),
|
|
"controller",
|
|
)
|
|
self.assertEqual(
|
|
nwb.normalize_role_kind("reconciler", profile_name="prgs-controller"),
|
|
"controller",
|
|
)
|
|
|
|
def test_profile_role_kind_prefers_declared_controller(self) -> None:
|
|
# Import from worktree package path via sys.path already set by pytest.
|
|
import gitea_mcp_server as mcp
|
|
|
|
profile = {
|
|
"profile_name": "prgs-controller",
|
|
"role": "controller",
|
|
"allowed_operations": [
|
|
"gitea.read",
|
|
"gitea.issue.comment",
|
|
"gitea.pr.close",
|
|
],
|
|
"forbidden_operations": [
|
|
"gitea.pr.approve",
|
|
"gitea.pr.merge",
|
|
"gitea.pr.create",
|
|
"gitea.branch.push",
|
|
],
|
|
}
|
|
# Declared role wins even if permissions look reconciler-like.
|
|
self.assertEqual(mcp._profile_role_kind(profile), "controller")
|
|
# Name-based fallback.
|
|
profile_no_role = dict(profile)
|
|
profile_no_role["role"] = None
|
|
profile_no_role["role_kind"] = None
|
|
self.assertEqual(mcp._profile_role_kind(profile_no_role), "controller")
|
|
|
|
def test_permission_inference_without_controller_name_stays_reconciler(self) -> None:
|
|
import gitea_mcp_server as mcp
|
|
|
|
# Pure permission inference still may return reconciler when no controller
|
|
# declaration exists — that is intentional for reconciler profiles.
|
|
role = mcp._role_kind(
|
|
["gitea.read", "gitea.pr.close", "gitea.issue.comment"],
|
|
["gitea.pr.approve", "gitea.pr.merge", "gitea.pr.create", "gitea.branch.push"],
|
|
)
|
|
self.assertEqual(role, "reconciler")
|
|
|
|
|
|
class DashboardRemainsExplanatoryTest(unittest.TestCase):
|
|
def test_dashboard_prompt_points_at_allocator_not_self_select(self) -> None:
|
|
import workflow_dashboard as wd
|
|
|
|
self.assertIn("gitea_allocate_next_work", wd.PROMPT_CONTROLLER)
|
|
self.assertIn("process_work_queue", wd.PROMPT_CONTROLLER)
|
|
self.assertIn("never replaces allocator", wd.PROMPT_CONTROLLER.lower())
|
|
self.assertNotIn("self-select", wd.PROMPT_CONTROLLER.lower())
|
|
|
|
|
|
class ClassifySkipCrossRoleTest(unittest.TestCase):
|
|
def test_controller_cross_role_accepts_author_issue(self) -> None:
|
|
c = WorkCandidate(kind="issue", number=1, labels=("status:ready",))
|
|
self.assertIsNone(
|
|
classify_skip(
|
|
c,
|
|
role=ROLE_CONTROLLER,
|
|
terminal_pr=None,
|
|
allocation_mode=ALLOCATION_MODE_CROSS_ROLE,
|
|
)
|
|
)
|
|
|
|
def test_legacy_controller_skips_author_issue(self) -> None:
|
|
c = WorkCandidate(kind="issue", number=1, labels=("status:ready",))
|
|
reason = classify_skip(
|
|
c,
|
|
role=ROLE_CONTROLLER,
|
|
terminal_pr=None,
|
|
allocation_mode=ALLOCATION_MODE_ROLE_SCOPED,
|
|
)
|
|
self.assertIsNotNone(reason)
|
|
self.assertIn("does not require controller", reason or "")
|
|
|
|
|
|
if __name__ == "__main__":
|
|
unittest.main()
|