diff --git a/allocator_service.py b/allocator_service.py index 40787ce..84b5175 100644 --- a/allocator_service.py +++ b/allocator_service.py @@ -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-; keep stable mapping when profile is + # non-prgs (still gitea- 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)" ), } diff --git a/gitea_mcp_server.py b/gitea_mcp_server.py index e8d6ab2..11acf51 100644 --- a/gitea_mcp_server.py +++ b/gitea_mcp_server.py @@ -234,12 +234,25 @@ def _effective_workspace_role() -> str: def _profile_role_kind(profile: dict) -> str: - """Resolve a profile's declared role before inferring from permissions.""" - role = (profile.get("role") or profile.get("role_kind") or "").strip() + """Resolve a profile's declared role before inferring from permissions. + + Declared ``role`` / ``role_kind`` always wins so a controller profile is + never reclassified as reconciler from permission inference (#840). + """ + role = (profile.get("role") or profile.get("role_kind") or "").strip().lower() if role: + # Normalize aliases / case. + if "control" in role: + return "controller" return role profile_name = (profile.get("profile_name") or "").strip().lower() - for candidate in ("reconciler", "merger", "reviewer", "author"): + for candidate in ( + "controller", + "reconciler", + "merger", + "reviewer", + "author", + ): if candidate in profile_name: return candidate return _role_kind( @@ -15684,7 +15697,8 @@ def mcp_get_control_plane_guide( profile = get_profile() allowed = profile["allowed_operations"] forbidden = profile["forbidden_operations"] - role = _role_kind(allowed, forbidden) + # Prefer declared profile role so controller is not mislabeled reconciler (#840). + role = _profile_role_kind(profile) username = _authenticated_username(h) identity = { @@ -15743,6 +15757,16 @@ def mcp_get_control_plane_guide( "user, and merging requires explicit operator authorization plus the " "'MERGE PR ' confirmation. " "Review and merge are separate workflow roles. A reviewer approval is not merge authorization.") + elif role == "controller": + guidance.append( + "Controller profile: route work via " + "gitea_route_task_session(task_type='process_work_queue') then " + "gitea_allocate_next_work (allocation_mode=cross_role by default). " + "The allocator returns exactly one authoritative selection with " + "required_role / required_profile / selected_action. Do not " + "implement, review, approve, or merge in this session — schedule " + "the matching role namespace instead. Dashboard output is " + "explanatory only and never replaces allocator selection.") elif role == "mixed": guidance.append( "WARNING: this profile allows both authoring and " @@ -15952,7 +15976,8 @@ def gitea_whoami( "environment": profile.get("environment"), "service": profile.get("service"), "identity": profile.get("identity"), - "role": profile.get("role"), + "role": profile.get("role") or _profile_role_kind(profile), + "role_kind": _profile_role_kind(profile), "profile_address": profile.get("profile_path"), "execution_profile": profile.get("execution_profile"), "audit_label": profile.get("audit_label"), @@ -20736,9 +20761,10 @@ def gitea_allocate_next_work( candidates_json: Any = None, exclude_issue_numbers: list[int] | None = None, expected_candidate_set_fingerprint: str | None = None, + allocation_mode: str | None = None, limit: int = 50, ) -> dict: - """Controller-owned next-work allocator using the #613 control-plane DB (#600). + """Controller-owned next-work allocator using the #613 control-plane DB (#600/#840). Workers must not self-select exclusive work under the standard multi-LLM workflow. Call this tool instead. @@ -20748,6 +20774,14 @@ def gitea_allocate_next_work( ``ControlPlaneDB.assign_and_lease`` (never file locks or comment-only leases as the coordination source). + *allocation_mode* (#840): when the active role is controller (or mode is + ``cross_role``), inspect the complete queue and return exactly one + authoritative selection with selected item, action, required_role, + required_profile/namespace, pins, and allocation/lease evidence. + Role-scoped workers pass ``role=author|reviewer|merger|reconciler`` (or + omit for profile role) for single-role filtering. Controller routes only + and does not perform downstream mutations. + Outcomes include: ``assigned_work``, ``preview``, ``wait``, ``blocked_by_terminal_path``, ``no_safe_work``, ``role_ineligible``, ``blocked_by_excluded_own_lease``, ``candidate_set_drift``. @@ -20882,6 +20916,7 @@ def gitea_allocate_next_work( controller_instance_id=allocator_service.resolve_controller_instance_id(), exclude_issue_numbers=exclude_issue_numbers, expected_candidate_set_fingerprint=expected_candidate_set_fingerprint, + allocation_mode=allocation_mode, ) except ValueError as exc: return { diff --git a/namespace_workspace_binding.py b/namespace_workspace_binding.py index fa006b6..be2ca82 100644 --- a/namespace_workspace_binding.py +++ b/namespace_workspace_binding.py @@ -24,7 +24,12 @@ ROLE_WORKTREE_ENVS: dict[str, str] = { "reconciler": RECONCILER_WORKTREE_ENV, } -NON_AUTHOR_ROLES = frozenset({"reviewer", "merger", "reconciler"}) +# Controller has no task worktree env — it routes only (#840). +KNOWN_ROLE_KINDS = frozenset( + {"author", "reviewer", "merger", "reconciler", "controller"} +) + +NON_AUTHOR_ROLES = frozenset({"reviewer", "merger", "reconciler", "controller"}) def normalize_role_kind( @@ -37,8 +42,12 @@ def normalize_role_kind( profile = (profile_name or "").strip().lower() if role == "reviewer" and "merger" in profile: return "merger" + if "controller" in profile or role == "controller": + return "controller" if role in ROLE_WORKTREE_ENVS: return role + if role in KNOWN_ROLE_KINDS: + return role return "author" @@ -80,7 +89,7 @@ def resolve_namespace_workspace( """ env_map = env if env is not None else os.environ role = normalize_role_kind(role_kind, profile_name=profile_name) - role_env_key = ROLE_WORKTREE_ENVS[role] + role_env_key = ROLE_WORKTREE_ENVS.get(role) # #618: durable author resolution — no silent control/master fallback. if role == "author" and verify_paths: @@ -108,13 +117,17 @@ def resolve_namespace_workspace( ) return workspace, source + role_env_candidate = ( + (_env_value(env_map, role_env_key), f"{role_env_key} environment variable", True) + if role_env_key + else (None, "no role worktree env", True) + ) for candidate, source, env_sourced in ( (worktree_path, "worktree_path argument", False), (worktree, "worktree argument", False), (_env_value(env_map, ACTIVE_WORKTREE_ENV), f"{ACTIVE_WORKTREE_ENV} environment variable", True), - (_env_value(env_map, role_env_key), - f"{role_env_key} environment variable", True), + role_env_candidate, (session_lease_worktree if role in {"reviewer", "merger"} else None, "reviewer PR lease worktree", False), # Author lock derivation is handled by the durable path above when @@ -433,7 +446,8 @@ def assess_namespace_mutation_workspace( reasons.append( f"{role} mutation blocked: workspace is the stable control checkout; " f"create or reconnect to a session-owned worktree under branches/ " - f"or set {ROLE_WORKTREE_ENVS[role]} / {ACTIVE_WORKTREE_ENV}" + f"or set {ROLE_WORKTREE_ENVS.get(role, ACTIVE_WORKTREE_ENV)} / " + f"{ACTIVE_WORKTREE_ENV}" ) elif ( role in {"reviewer", "merger"} diff --git a/role_session_router.py b/role_session_router.py index 4f331b8..f756c71 100644 --- a/role_session_router.py +++ b/role_session_router.py @@ -81,6 +81,12 @@ AUTHOR_TASKS = frozenset({ "reconcile_landed_pr", }) +CONTROLLER_TASKS = frozenset({ + "process_work_queue", + "process-work-queue", + "cross_role_allocate", +}) + RECONCILER_TASKS = frozenset({ "cleanup_merged_pr_branch", # #729: delete_branch is reconciler-owned (gitea.branch.delete is granted @@ -132,6 +138,10 @@ TASK_REQUIRED_ROLE = { "reconcile_close_superseded_pr": "reconciler", "reconcile_close_satisfied_issue": "reconciler", "reconcile_create_followup_issue": "reconciler", + # #840: controller-owned generic queue allocation / routing. + "process_work_queue": "controller", + "process-work-queue": "controller", + "cross_role_allocate": "controller", } WRONG_ROLE_REVIEWER_MSG = ( @@ -147,6 +157,10 @@ WRONG_ROLE_MERGER_MSG = ( "Wrong role/session for merger task. Launch merger MCP namespace." ) +WRONG_ROLE_CONTROLLER_MSG = ( + "Wrong role/session for controller task. Launch controller MCP namespace." +) + _session_last_route: dict | None = None @@ -281,6 +295,27 @@ def route_task_session( _record_route(result) return result + if required_role == "controller": + result = { + "task_type": task_type, + "required_role": required_role, + "active_role": active_role_kind, + "active_profile": active_profile, + "route_result": ROUTE_WRONG_ROLE, + "downstream_allowed": False, + "reasons": [ + WRONG_ROLE_CONTROLLER_MSG, + "Controller tasks (process_work_queue / cross-role allocate) " + "cannot run in author, reviewer, merger, or reconciler " + "worker sessions.", + ], + "message": WRONG_ROLE_CONTROLLER_MSG, + "runtime_switching_supported": runtime_switching_supported, + "profile_switch_blocked": not runtime_switching_supported, + } + _record_route(result) + return result + if required_role == "author": route = ROUTE_TO_AUTHOR message = ( diff --git a/task_capability_map.py b/task_capability_map.py index 628c444..8c060a8 100644 --- a/task_capability_map.py +++ b/task_capability_map.py @@ -318,8 +318,10 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = { "permission": "gitea.pr.create", "role": "author", }, - # #600: controller-owned allocator — any authenticated profile may call; - # routing enforces role match to selected work. Uses control-plane DB (#613). + # #600: workers and controller may call with gitea.read; role-scoped workers + # pass role=author|reviewer|merger|reconciler. Cross-role routing is the + # controller default (#840). The canonical generic queue *task type* is + # process_work_queue (controller-only below). "allocate_next_work": { "permission": "gitea.read", "role": "author", @@ -328,6 +330,19 @@ TASK_CAPABILITY_MAP: dict[str, dict[str, str]] = { "permission": "gitea.read", "role": "author", }, + # #840: documented generic queue task — controller routes only. + "process_work_queue": { + "permission": "gitea.read", + "role": "controller", + }, + "process-work-queue": { + "permission": "gitea.read", + "role": "controller", + }, + "cross_role_allocate": { + "permission": "gitea.read", + "role": "controller", + }, # #601 first-class lease lifecycle — inspect/list need read; mutations gate on # ownership in the control-plane DB (not a separate Gitea write permission). diff --git a/tests/test_cross_role_queue_allocation.py b/tests/test_cross_role_queue_allocation.py new file mode 100644 index 0000000..4126d30 --- /dev/null +++ b/tests/test_cross_role_queue_allocation.py @@ -0,0 +1,581 @@ +"""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() diff --git a/workflow_dashboard.py b/workflow_dashboard.py index d2b10ee..e47b2dc 100644 --- a/workflow_dashboard.py +++ b/workflow_dashboard.py @@ -65,9 +65,11 @@ PROMPT_RECONCILER = ( "reconciliation (already-landed / post-merge cleanup). Do not approve or merge." ) PROMPT_CONTROLLER = ( - "CONTROLLER session: inspect gitea_workflow_dashboard + control-plane leases, " - "diagnose blocked/terminal-locked items for {remote}/{org}/{repo}, and schedule " - "exactly one fresh role-scoped cycle. Do not implement, review, or merge in-band." + "CONTROLLER session: call gitea_route_task_session(task_type='process_work_queue') " + "then gitea_allocate_next_work (cross_role default) for {remote}/{org}/{repo}; " + "use the returned required_role/profile/action to schedule exactly one downstream " + "role cycle. Dashboard is explanatory only and never replaces allocator selection. " + "Do not implement, review, approve, or merge in-band." ) PROMPT_IDLE = ( "IDLE: no safe assignable work for role '{role}' on {remote}/{org}/{repo}. "