From df3417e845ed95764c182c9f309f1a4a46a5c46a Mon Sep 17 00:00:00 2001 From: Jory Irving <46251616+joryirving@users.noreply.github.com> Date: Fri, 18 Sep 2026 11:52:26 -0600 Subject: [PATCH] fix(bridge): cap every coder-creating path by CODER_AGENT_SLOTS (no unbounded coder concurrency) Three coder-creating paths were uncapped, so a batch of work could recreate coders concurrently and thrash the single-slot local 27B (observed: 6 at once, pipeline to a crawl). #344 capped the issue-claim loop, reconcile_failures, and drain_pr_fixes; this completes the set so NO path can spawn a coder past the cap: - reconcile_pr_fixes same-tier retry (retry / retry-progress): gate the recreate on the current coder's free_slots; defer (leave the tombstone) when full. - reconcile_pr_fixes escalation: gate the next-tier recreate on THAT tier's free_slots (coder-frontier has its own cap), so a burst of escalations cannot exceed the stronger coder either. - redrive_infra: gate the infra-recovery recreate on free_slots; return False to defer (reconcile_infra_parked keeps the marker and retries next tick). coder_load is threaded from run_tick through reconcile_failures -> infra redrive -> reconcile_pr_fixes -> drain, mutated in place, so all paths draw from one shared per-tick pool. Every defer is transient (retry next tick), never a silent BLOCK. Escalation to frontier is still reached -- gated by frontier's own slots, not the local coder's. --- bridge/main.py | 21 ++++++++++--- bridge/prfix.py | 31 ++++++++++++++++++- tests/test_prfix.py | 74 +++++++++++++++++++++++++++++++++++++++++++++ 3 files changed, 120 insertions(+), 6 deletions(-) diff --git a/bridge/main.py b/bridge/main.py index b2c7378..ad6ff53 100644 --- a/bridge/main.py +++ b/bridge/main.py @@ -19,6 +19,7 @@ coder_agent_for, coder_candidates, coders_saturated, + free_slots, revision_coder_agent_for, gate_profile_for, parse_gate_profiles, @@ -1156,17 +1157,26 @@ def redrive_infra(issue: dict, model: str) -> bool: item = replace(item, lane=cfg.lanes[0] if cfg.lanes else "local") language = cfg.gate_profiles.get(item.repo, {}).get("language") branch = _branch_name(item) + redrive_coder = coder_agent_for( + item.lane, language, cfg.lane_coder_agents, cfg.base_coder_agents, + repo=item.repo, repo_coder_agents=cfg.repo_coder_agents, + issue_number=item.issue_number, + ) + # Gate the infra redrive by coder capacity: a batch of infra-failed + # issues going healthy at once must not recreate their coders past the + # slot cap. Returning False defers -- reconcile_infra_parked keeps the + # infra marker and retries next tick. + if cfg.coder_slots and free_slots(redrive_coder, coder_load, cfg.coder_slots) <= 0: + return False manifest = build_workload( item, cfg.namespace, gate_profile_for(item.repo, cfg.gate_profiles), cfg.agent_name, attempt=1, - coder_agent=coder_agent_for( - item.lane, language, cfg.lane_coder_agents, cfg.base_coder_agents, - repo=item.repo, repo_coder_agents=cfg.repo_coder_agents, - issue_number=item.issue_number, - ), verify_enabled=cfg.verify_enabled, self_go=cfg.self_go, + coder_agent=redrive_coder, + verify_enabled=cfg.verify_enabled, self_go=cfg.self_go, revise_from_branch=branch, ) create_workload(manifest) + coder_load[redrive_coder] = coder_load.get(redrive_coder, 0) + 1 payload = { "issueId": item.issue_id, "repoFullName": item.repo, @@ -1278,6 +1288,7 @@ def get_pr_fix_signature(repo, pr) -> str: lane_agents=cfg.pr_fix_lane_agents, get_pr_fix_signature=get_pr_fix_signature, update_pr_branch=update_pr_branch, + agent_load=coder_load, agent_slots=cfg.coder_slots, ): logger.info(line) diff --git a/bridge/prfix.py b/bridge/prfix.py index 1376305..1899cb8 100644 --- a/bridge/prfix.py +++ b/bridge/prfix.py @@ -337,7 +337,8 @@ def reconcile_pr_fixes(list_prfix_workloads, delete_workload, create_workload, lane_agents=None, get_pr_fix_signature=lambda repo, pr: "", update_pr_branch=lambda repo, pr: False, - progress_max_attempts=PR_FIX_PROGRESS_MAX_ATTEMPTS) -> list: + progress_max_attempts=PR_FIX_PROGRESS_MAX_ATTEMPTS, + agent_load=None, agent_slots=None) -> list: """Settle prior fix Workloads: Succeeded -> verify the PR is actually mergeable (pr_is_mergeable) before marking FIXED, delete only if the mark succeeded (else leave the tombstone so the next tick retries the mark); @@ -351,6 +352,10 @@ def reconcile_pr_fixes(list_prfix_workloads, delete_workload, create_workload, untouched. Per-Workload isolation so one wedged delete/create/mark cannot abort the pass or the drain that follows.""" results = [] + slots = agent_slots or {} + # Mutated in place as same-tier retries recreate, so the shared coder_load + # carries this pass's draws into the pr-fix drain later in the same tick. + load = agent_load if agent_load is not None else {} for wl in list_prfix_workloads(): meta = wl.get("metadata") or {} name = meta.get("name") or "?" @@ -433,6 +438,20 @@ def reconcile_pr_fixes(list_prfix_workloads, delete_workload, create_workload, results.append(f"{name}:checks-pending:{attempt}/{max_attempts}") # Mark failed, still failing check, or Failed phase -> retry or BLOCKED elif attempt < max_attempts: + # Cap the same-tier retry by coder capacity. Recreating this + # Workload onto a coder already at its CODER_AGENT_SLOTS capacity + # stacks another generation on a single-slot local model. The + # drain and issue-retry paths gate this (#344); this pr-fix + # retry/recreate path did not, so a batch of failing pr-fixes + # recreated concurrently and thrashed the 27B. Defer -- leave the + # tombstone (do NOT delete) so it retries once a slot frees; this + # is a transient defer, not a BLOCK. Escalation (the else branch + # below) targets a different tier/coder and is intentionally not + # gated here. + retry_coder = _prfix_current_coder(wl) + if slots and retry_coder and free_slots(retry_coder, load, slots) <= 0: + results.append(f"{name}:retry-deferred:coder-busy:{retry_coder}") + continue # Signature-aware budgeting: charge the attempt budget by # *failure signature*, not by attempt count (#133). A retry # against the same wall still ticks attempt++; a retry against @@ -474,6 +493,7 @@ def reconcile_pr_fixes(list_prfix_workloads, delete_workload, create_workload, ) manifest["metadata"]["annotations"][PROGRESS_ANNOTATION] = str(next_progress) create_workload(manifest) + load[retry_coder] = load.get(retry_coder, 0) + 1 tag = "not-mergeable-retry-progress" if merge_status == "checks_failed" else "retry-progress" results.append( f"{name}:{tag}:{next_progress}/{progress_max_attempts}" @@ -484,6 +504,7 @@ def reconcile_pr_fixes(list_prfix_workloads, delete_workload, create_workload, wl, attempt + 1, signature=new_sig if new_sig else "", )) + load[retry_coder] = load.get(retry_coder, 0) + 1 tag = "not-mergeable-retry" if merge_status == "checks_failed" else "retry" results.append(f"{name}:{tag}:{attempt + 1}/{max_attempts}") else: @@ -496,8 +517,16 @@ def reconcile_pr_fixes(list_prfix_workloads, delete_workload, create_workload, nxt = next_prfix_lane(current_lane) next_coder = pr_fix_coder_for(nxt, lane_agents or {}) if nxt else None if nxt and next_coder and next_coder != _prfix_current_coder(wl): + # Escalation creates a coder on the next tier; gate it by that + # tier's own CODER_AGENT_SLOTS so a burst of escalations cannot + # exceed the stronger coder's capacity either. Defer (leave the + # tombstone) when that tier is full. + if slots and free_slots(next_coder, load, slots) <= 0: + results.append(f"{name}:escalate-deferred:coder-busy:{next_coder}") + continue delete_workload(name) create_workload(escalate_prfix_manifest(wl, nxt, next_coder)) + load[next_coder] = load.get(next_coder, 0) + 1 results.append(f"{name}:escalate:{current_lane or 'NORMAL'}->{nxt}") else: if repo and pr is not None: diff --git a/tests/test_prfix.py b/tests/test_prfix.py index 5676cb5..2487023 100644 --- a/tests/test_prfix.py +++ b/tests/test_prfix.py @@ -354,6 +354,80 @@ def test_reconcile_failed_under_max_deletes_and_recreates(): assert out == ["prfix-o-r-5:retry:2/3"] +def test_reconcile_retry_deferred_when_coder_at_capacity(): + # Same-tier retry onto a coder already at CODER_AGENT_SLOTS capacity would + # stack another generation on the single-slot 27B. Defer: leave the tombstone + # (no delete, no recreate), retry next tick once a slot frees. Not a BLOCK. + created, deleted = [], [] + out = reconcile_pr_fixes( + list_prfix_workloads=lambda: [_wl(5, "Failed", attempt=1, coder="coder")], + delete_workload=deleted.append, create_workload=created.append, + mark_pr_fix=lambda *a: (_ for _ in ()).throw(AssertionError("no mark")), + max_attempts=3, + agent_load={"coder": 1}, agent_slots={"coder": 1}, + ) + assert created == [] and deleted == [] + assert out == ["prfix-o-r-5:retry-deferred:coder-busy:coder"] + + +def test_reconcile_retry_proceeds_and_draws_load_when_slot_free(): + load = {} + created, deleted = [], [] + out = reconcile_pr_fixes( + list_prfix_workloads=lambda: [_wl(5, "Failed", attempt=1, coder="coder")], + delete_workload=deleted.append, create_workload=created.append, + mark_pr_fix=lambda *a: (_ for _ in ()).throw(AssertionError("no mark")), + max_attempts=3, + agent_load=load, agent_slots={"coder": 1}, + ) + assert out == ["prfix-o-r-5:retry:2/3"] and len(created) == 1 + assert load == {"coder": 1} # drawn down in place for the drain that follows + + +def test_reconcile_escalate_not_gated_by_busy_qwen(): + # At the cap, escalation targets coder-frontier (a different tier), so a busy + # qwen coder must NOT block it -- only same-tier retries are gated. + created, marks = [], [] + out = reconcile_pr_fixes( + list_prfix_workloads=lambda: [_wl(5, "Failed", attempt=3, lane="NORMAL", coder="coder")], + delete_workload=lambda n: None, create_workload=created.append, + mark_pr_fix=lambda *a: marks.append(a), + max_attempts=3, lane_agents=DEFAULT_PRFIX_LANE_AGENTS, + agent_load={"coder": 1}, agent_slots={"coder": 1}, + ) + assert marks == [] # not blocked + assert created and created[0]["metadata"]["labels"]["lane"] == "ESCALATED" + assert any(":escalate:" in x for x in out) + + +def test_reconcile_retry_uncapped_when_no_slots_configured(): + created = [] + out = reconcile_pr_fixes( + list_prfix_workloads=lambda: [_wl(5, "Failed", attempt=1, coder="coder")], + delete_workload=lambda n: None, create_workload=created.append, + mark_pr_fix=lambda *a: (_ for _ in ()).throw(AssertionError("no mark")), + max_attempts=3, + agent_load={"coder": 9}, agent_slots={}, + ) + assert out == ["prfix-o-r-5:retry:2/3"] and len(created) == 1 + + +def test_reconcile_escalate_deferred_when_frontier_at_capacity(): + # Escalation targets coder-frontier; if THAT tier is itself at capacity, + # defer (leave the tombstone) rather than exceed its slots too. + created, marks = [], [] + out = reconcile_pr_fixes( + list_prfix_workloads=lambda: [_wl(5, "Failed", attempt=3, lane="NORMAL", coder="coder")], + delete_workload=lambda n: (_ for _ in ()).throw(AssertionError("no delete on defer")), + create_workload=created.append, + mark_pr_fix=lambda *a: marks.append(a), + max_attempts=3, lane_agents=DEFAULT_PRFIX_LANE_AGENTS, + agent_load={"coder-frontier": 4}, agent_slots={"coder": 1, "coder-frontier": 4}, + ) + assert created == [] and marks == [] + assert out == ["prfix-o-r-5:escalate-deferred:coder-busy:coder-frontier"] + + def test_reconcile_normal_at_max_escalates_to_frontier(): # NORMAL tier exhausted -> escalate to ESCALATED (coder-frontier) with a fresh # attempt budget, NOT straight to BLOCKED/needs-human.