diff --git a/bridge/main.py b/bridge/main.py index 55b14d5..6bbf041 100644 --- a/bridge/main.py +++ b/bridge/main.py @@ -636,6 +636,30 @@ def _active_workloads( return active +def _active_coder_pool_workloads( + api: client.CustomObjectsApi, namespace: str +) -> List[Dict[str, Any]]: + """Active (non-terminal) Workloads whose coder draws on the shared pool. + + Both issue (``created-by=dispatch-bridge``) and pr-fix + (``created-by=dispatch-bridge-prfix``) Workloads, in one set-based list call. + A pr-fix coder shares the same coder Agent (and the same GPU) as an issue + coder, so it must count toward ``_load_by_coder_agent`` or a pr-fix coder and + an issue coder both land on the same single-slot local model. pr-fix + Workloads are pipeline-shaped (the coder is named on the ``issue-fix`` step, + resolved by ``_coder_agent_name``'s pipeline fallback), not + ``spec.coderAgentRef``. + """ + active: List[Dict[str, Any]] = [] + for workload in _list_workloads_by_label( + api, namespace, f"created-by in (dispatch-bridge,{PRFIX_CREATED_BY})" + ): + phase = (workload.get("status") or {}).get("phase") + if phase not in _TERMINAL_PHASES: + active.append(workload) + return active + + def _coder_still_busy(tasks: List[Dict[str, Any]]) -> bool: """Fail-closed busy check for one Workload's tasks (#180). @@ -690,7 +714,8 @@ def _refs(workloads: List[Dict[str, Any]]) -> Dict[str, int]: load[ref] = load.get(ref, 0) + 1 return load - workloads = _active_workloads(api, namespace) if workloads is None else workloads + if workloads is None: + workloads = _active_coder_pool_workloads(api, namespace) tasks_by_workload = _list_agentic_tasks_by_workload(api, namespace) if tasks_by_workload is None: return _refs(workloads) @@ -1259,6 +1284,7 @@ def get_pr_fix_signature(repo, pr) -> str: cfg.gate_profiles, cfg.pr_fix_lane_agents, cfg.agent_name, cfg.namespace, verify_enabled=cfg.verify_enabled, self_go=cfg.self_go, + agent_load=coder_load, agent_slots=cfg.coder_slots, ): logger.info(line) diff --git a/bridge/prfix.py b/bridge/prfix.py index 13bc6e9..1376305 100644 --- a/bridge/prfix.py +++ b/bridge/prfix.py @@ -3,7 +3,7 @@ from bridge.workload import ( CODER_AGENT, VERIFIER_AGENT, ATTEMPT_ANNOTATION, SIGNATURE_ANNOTATION, - PROGRESS_ANNOTATION, LANE_CODER_WILDCARD, gate_profile_for, + PROGRESS_ANNOTATION, LANE_CODER_WILDCARD, gate_profile_for, free_slots, ) # Lane values dispatch assigns to a PR-fix item. NEEDS_HUMAN is never actioned @@ -194,13 +194,25 @@ def build_fix_workload(item, namespace, gate_profile, agent_name, coder_agent, a def drain_pr_fixes(list_queued, existing_prfix_names, create_workload, gate_profiles, lane_agents, agent_name, namespace, - verify_enabled: bool = True, self_go: list[str] | None = None) -> list: - """Create a fix Workload per newly-QUEUED item. list_queued returns raw - dicts already filtered to actionable lanes by the API query. An item is - skipped when it has no branch (nothing to amend) or already has an - in-flight prfix Workload (reconcile owns it; the item stays QUEUED). One - bad item never aborts the pass.""" + verify_enabled: bool = True, self_go: list[str] | None = None, + agent_load: Optional[dict] = None, + agent_slots: Optional[dict] = None) -> list: + """Create a fix Workload per newly-QUEUED item, bounded by coder capacity. + + list_queued returns raw dicts already filtered to actionable lanes by the + API query. An item is skipped when it has no branch (nothing to amend), when + it already has an in-flight prfix Workload (reconcile owns it; the item stays + QUEUED), or when its coder Agent is already at its ``CODER_AGENT_SLOTS`` + capacity. Without the capacity guard every queued item dispatched at once: + fine when they route to a cloud API, but a pileup of concurrent coders on a + single-slot local model thrashes it. ``agent_load`` starts from the coders + already in flight (issue + pr-fix, shared pool) and is drawn down as this + pass dispatches; a saturated coder leaves its item QUEUED for the next tick. + When ``agent_slots`` is empty the guard is inert and the legacy dispatch-all + behavior stands. One bad item never aborts the pass.""" lane_agents = lane_agents or {} + slots = agent_slots or {} + load = dict(agent_load or {}) results = [] for raw in list_queued(): item = parse_pr_fix_item(raw) @@ -215,15 +227,20 @@ def drain_pr_fixes(list_queued, existing_prfix_names, create_workload, if name in existing_prfix_names: results.append(f"{tag}:skip:in-flight") continue + coder_agent = pr_fix_coder_for(item.lane, lane_agents) + if slots and free_slots(coder_agent, load, slots) <= 0: + results.append(f"{tag}:skip:coder-busy:{coder_agent}") + continue try: manifest = build_fix_workload( item, namespace, gate_profile_for(item.repo, gate_profiles), - agent_name, pr_fix_coder_for(item.lane, lane_agents), attempt=1, + agent_name, coder_agent, attempt=1, verify_enabled=verify_enabled, signature=failure_signature(item), self_go=self_go, ) create_workload(manifest) + load[coder_agent] = load.get(coder_agent, 0) + 1 results.append(f"{tag}:created:{name}") except Exception as e: results.append(f"{tag}:error:{e}") diff --git a/tests/test_bridge_runtime.py b/tests/test_bridge_runtime.py index ee812ea..1c87f00 100644 --- a/tests/test_bridge_runtime.py +++ b/tests/test_bridge_runtime.py @@ -333,6 +333,22 @@ def test_counts_pipeline_shaped_coder(self) -> None: ) assert _load_by_coder_agent(api, "ns") == {"coder-py": 1} + def test_counts_prfix_workload_toward_shared_pool(self) -> None: + # A pr-fix Workload (created-by=dispatch-bridge-prfix, pipeline-shaped) + # shares the coder pool with issue coders, so it must count toward load, + # or a pr-fix coder and an issue coder both land on one single-slot model. + wl = _pipeline_workload("p", "Running", coder="coder-py") + wl["metadata"]["labels"]["created-by"] = "dispatch-bridge-prfix" + api = FakeAPI( + responses={ + "list_namespaced_custom_object": [ + {"items": [wl]}, # combined issue+prfix pool list + {"items": [_task("t1", "p", phase="Running")]}, + ] + } + ) + assert _load_by_coder_agent(api, "ns") == {"coder-py": 1} + def test_terminal_issue_fix_with_running_review_frees_slot(self) -> None: # A terminal issue-fix task plus a running review does NOT hold the # coder busy: only issue-fix tasks count, and only while non-terminal. diff --git a/tests/test_prfix.py b/tests/test_prfix.py index 1b29aff..5676cb5 100644 --- a/tests/test_prfix.py +++ b/tests/test_prfix.py @@ -193,6 +193,72 @@ def create(m): assert any("o/r#5:error:" in line for line in out) +def test_drain_caps_concurrent_coders_at_slot_capacity(): + # Two queued NORMAL items resolve to the same coder; coder:1 means the first + # dispatches and the second is left QUEUED (skip:coder-busy) rather than + # piling two concurrent coders onto a single-slot local model. + created = [] + out = drain_pr_fixes( + list_queued=lambda: [_raw(pr=5), _raw(pr=6)], + existing_prfix_names=set(), + create_workload=created.append, + gate_profiles={}, lane_agents={"NORMAL": "coder"}, agent_name="a", + namespace="llm", + agent_slots={"coder": 1}, + ) + assert [m["metadata"]["name"] for m in created] == ["prfix-o-r-5"] + assert "o/r#5:created:prfix-o-r-5" in out + assert any("o/r#6:skip:coder-busy:coder" in line for line in out) + + +def test_drain_counts_inflight_load_against_cap(): + # A coder already in flight (agent_load) occupies the only slot, so a fresh + # queued item is held instead of dispatched alongside it. + created = [] + out = drain_pr_fixes( + list_queued=lambda: [_raw(pr=7)], + existing_prfix_names=set(), + create_workload=created.append, + gate_profiles={}, lane_agents={"NORMAL": "coder"}, agent_name="a", + namespace="llm", + agent_load={"coder": 1}, agent_slots={"coder": 1}, + ) + assert created == [] + assert any("o/r#7:skip:coder-busy:coder" in line for line in out) + + +def test_drain_uncapped_when_no_slots_configured(): + # Empty agent_slots keeps the legacy dispatch-all behavior. + created = [] + drain_pr_fixes( + list_queued=lambda: [_raw(pr=5), _raw(pr=6)], + existing_prfix_names=set(), + create_workload=created.append, + gate_profiles={}, lane_agents={"NORMAL": "coder"}, agent_name="a", + namespace="llm", + agent_slots={}, + ) + assert [m["metadata"]["name"] for m in created] == [ + "prfix-o-r-5", "prfix-o-r-6" + ] + + +def test_drain_higher_capacity_lane_dispatches_more(): + # A lane whose coder has capacity 4 dispatches multiple in one pass. + created = [] + drain_pr_fixes( + list_queued=lambda: [_raw(pr=5), _raw(pr=6), _raw(pr=7)], + existing_prfix_names=set(), + create_workload=created.append, + gate_profiles={}, lane_agents={"NORMAL": "coder-frontier"}, agent_name="a", + namespace="llm", + agent_slots={"coder": 1, "coder-frontier": 4}, + ) + assert [m["metadata"]["name"] for m in created] == [ + "prfix-o-r-5", "prfix-o-r-6", "prfix-o-r-7" + ] + + def test_drain_gateless_creates_issue_fix_only_no_verify(): """verify_enabled=False drain creates a Workload with issue-fix only, no verify.""" created = []