diff --git a/bridge/main.py b/bridge/main.py index 3423262..55b14d5 100644 --- a/bridge/main.py +++ b/bridge/main.py @@ -595,14 +595,33 @@ def _list_agentic_tasks_by_workload( def _coder_agent_name(workload: Dict[str, Any]) -> Optional[str]: - """Return a Workload's coder ref, or None when it cannot be resolved.""" + """Return a Workload's coder ref, or None when it cannot be resolved. + + Standard bridge Workloads name the coder at ``spec.coderAgentRef``. Pipeline- + shaped Workloads (``spec.pipeline`` — e.g. a retry re-dispatched with an + explicit pipeline) carry no top-level ``coderAgentRef`` and name the coder on + the ``issue-fix`` step's ``agentRef`` instead. Without the pipeline fallback + those Workloads resolve to None and contribute nothing to + ``_load_by_coder_agent``, so their coder never consumes a ``CODER_AGENT_SLOTS`` + slot and a second coder can be dispatched alongside them (observed: a + pipeline-shaped coder and a standard coder both running the ``coder`` agent at + once under a ``coder:1`` cap). + """ spec = workload.get("spec") or {} if not isinstance(spec, dict): return None - ref = spec.get("coderAgentRef") or {} - if not isinstance(ref, dict): - return None - return ref.get("name") + ref = spec.get("coderAgentRef") + if isinstance(ref, dict) and ref.get("name"): + return ref.get("name") + pipeline = spec.get("pipeline") + if isinstance(pipeline, list): + for step in pipeline: + if not isinstance(step, dict) or step.get("kind") != "issue-fix": + continue + step_ref = step.get("agentRef") + if isinstance(step_ref, dict) and step_ref.get("name"): + return step_ref.get("name") + return None def _active_workloads( diff --git a/tests/test_bridge_runtime.py b/tests/test_bridge_runtime.py index 298617b..ee812ea 100644 --- a/tests/test_bridge_runtime.py +++ b/tests/test_bridge_runtime.py @@ -115,6 +115,22 @@ def _task(name: str, workload: str, *, kind: str = "issue-fix", phase: str) -> d } +def _pipeline_workload(name: str, phase: str, *, coder: str) -> dict: + # Pipeline-shaped Workload: the coder is named on the issue-fix step's + # agentRef, with no top-level coderAgentRef (the shape a retry re-dispatch + # produces). See _coder_agent_name's pipeline fallback. + return { + "metadata": {"name": name, "labels": {"created-by": "dispatch-bridge"}}, + "spec": { + "pipeline": [ + {"name": "code", "kind": "issue-fix", "agentRef": {"name": coder}}, + {"name": "review", "kind": "review", "agentRef": {"name": "reviewer"}}, + ] + }, + "status": {"phase": phase}, + } + + def _prfix_workload(name: str, phase: str) -> dict: return { "metadata": {"name": name, "labels": {"created-by": "dispatch-bridge-prfix"}}, @@ -303,6 +319,20 @@ def test_counts_running_issue_fix(self) -> None: ) assert _load_by_coder_agent(api, "ns") == {"coder-py": 1} + def test_counts_pipeline_shaped_coder(self) -> None: + # A pipeline-shaped Workload names its coder on the issue-fix step, not at + # spec.coderAgentRef. It must still count toward coder load, or the coder + # slot leaks and a second coder runs alongside it under CODER_AGENT_SLOTS. + api = FakeAPI( + responses={ + "list_namespaced_custom_object": [ + {"items": [_pipeline_workload("a", "Running", coder="coder-py")]}, + {"items": [_task("t1", "a", 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.