Skip to content
This repository was archived by the owner on Sep 20, 2026. It is now read-only.
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
28 changes: 27 additions & 1 deletion bridge/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -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).

Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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)

Expand Down
33 changes: 25 additions & 8 deletions bridge/prfix.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand All @@ -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}")
Expand Down
16 changes: 16 additions & 0 deletions tests/test_bridge_runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
66 changes: 66 additions & 0 deletions tests/test_prfix.py
Original file line number Diff line number Diff line change
Expand Up @@ -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 = []
Expand Down