diff --git a/bridge/main.py b/bridge/main.py index 6bbf041..b2c7378 100644 --- a/bridge/main.py +++ b/bridge/main.py @@ -1084,11 +1084,18 @@ def queue_for(lane: str) -> list: current_lane_for = {} logger.warning("lane-index-failed", extra={"error": _redact_token(repr(e))}) + # Shared coder-load accounting for the whole tick, computed before the retry + # pass so reconcile_failures, the issue-claim loop and the pr-fix drain draw + # from one pool -- a retried coder, a fresh issue coder and a pr-fix coder + # cannot each independently fill the same single slot. reconcile_failures + # mutates it in place as retries recreate. + coder_load = load_by_coder_agent() if cfg.coder_slots else {} # Retry failed workloads first (so a re-run this tick uses the current config), # then claim new work. for line in reconcile_failures( cfg.agent_name, list_failed_workloads, create_workload, delete_workload, cfg.namespace, cfg.gate_profiles, cfg.max_attempts, + agent_load=coder_load, agent_slots=cfg.coder_slots, escalate=escalate if cfg.escalation_lane else None, escalation_lane=cfg.escalation_lane, lane_coder_agents=cfg.lane_coder_agents, @@ -1184,7 +1191,6 @@ def redrive_infra(issue: dict, model: str) -> bool: # Cap concurrent in-progress work so the pipeline drains a bounded set # instead of claiming the whole backlog at once (0 = uncapped). active = count_active_workloads() if cfg.max_in_progress else 0 - coder_load = load_by_coder_agent() if cfg.coder_slots else {} # Fix-first work-stealing (issue #134): named agents are removed from the # issue rotation while they still hold fix work or have a full slot, so # the fix lane's single slot stays uncontended. A JSON list ["coder"] or diff --git a/bridge/retry.py b/bridge/retry.py index 436a8ee..b5336e8 100644 --- a/bridge/retry.py +++ b/bridge/retry.py @@ -7,6 +7,7 @@ from bridge.workload import ( build_workload, coder_agent_for, + free_slots, gate_profile_for, _branch_name, ATTEMPT_ANNOTATION, @@ -640,6 +641,8 @@ def reconcile_failures( park_infra: Optional[Callable[[ClaimedItem, str, int], bool]] = None, needs_human_for: Optional[NeedsHumanFor] = None, ensure_human_label: Optional[EnsureHumanLabel] = None, + agent_load: Optional[dict] = None, + agent_slots: Optional[dict] = None, ) -> list: """Retry Failed bridge Workloads, bounded by max_attempts. @@ -679,6 +682,10 @@ def reconcile_failures( """ lane_coder_agents = lane_coder_agents or {} base_coder_agents = base_coder_agents or {} + slots = agent_slots or {} + # Mutated in place as retries recreate, so the caller's shared coder_load + # carries this pass's draws into the issue-claim and pr-fix drain passes. + load = agent_load if agent_load is not None else {} repo_coder_agents = repo_coder_agents or {} results = [] _park_exhausted = _park_exhausted_factory( @@ -1065,9 +1072,23 @@ def reconcile_failures( # deletion never completes, LLMKube#949) or a create that races must # not abort the rest of the reconcile pass and the claim pass — one # bad Workload previously crashed the whole bridge run every tick. + language = gate_profiles.get(item.repo, {}).get("language") + coder_agent = coder_agent_for( + item.lane, language, lane_coder_agents, base_coder_agents, + repo=item.repo, repo_coder_agents=repo_coder_agents, + issue_number=item.issue_number, + ) + # Cap the retry by coder capacity. A retried Workload recreated onto a + # coder already at its CODER_AGENT_SLOTS capacity would put a second + # coder on a single-slot local model — the fresh-dispatch paths (claim + # loop, pr-fix drain) are capped, but this recreation path was not. + # Defer to a later tick: leave the Failed tombstone (do NOT delete) so + # list_failed() re-offers it once a slot frees. + if slots and free_slots(coder_agent, load, slots) <= 0: + results.append(f"{name}:retry-deferred:coder-busy:{coder_agent}") + continue try: delete_workload(name) - language = gate_profiles.get(item.repo, {}).get("language") # Infra errors are not real rejections — the request never reached # the agent — so a retry against the same backend must not spend # the verdict budget. Real verdicts increment as before. Infra @@ -1086,17 +1107,14 @@ def reconcile_failures( agent_name, next_attempt, infra_attempt=next_infra_attempt, - coder_agent=coder_agent_for( - item.lane, language, lane_coder_agents, base_coder_agents, - repo=item.repo, repo_coder_agents=repo_coder_agents, - issue_number=item.issue_number, - ), + coder_agent=coder_agent, feedback=feedback, verify_enabled=verify_enabled, self_go=self_go, revise_from_branch=revise_from, ) create_workload(manifest) + load[coder_agent] = load.get(coder_agent, 0) + 1 except Exception as e: results.append(f"{name}:retry-error:{e}") continue diff --git a/tests/test_retry.py b/tests/test_retry.py index 5538779..29b9409 100644 --- a/tests/test_retry.py +++ b/tests/test_retry.py @@ -69,6 +69,50 @@ def test_reconcile_retries_below_max_deletes_and_recreates_at_next_attempt(): assert m["spec"]["gateProfile"] == {"language": "generic"} +def test_reconcile_retry_deferred_when_coder_at_capacity(): + # A retry recreated onto a coder already at CODER_AGENT_SLOTS capacity would + # put a second coder on a single-slot local model. Defer: leave the Failed + # tombstone (do not delete) so it re-offers once a slot frees. + r = _Recorder([_failed_wl("wl-misospace-dispatch-7", attempt=1, lane="local")]) + out = reconcile_failures( + "foreman-coder", r.list_failed, r.create, r.delete, + namespace="llm", gate_profiles={}, max_attempts=3, + lane_coder_agents={"local": "coder"}, + agent_load={"coder": 1}, agent_slots={"coder": 1}, + ) + assert out == ["wl-misospace-dispatch-7:retry-deferred:coder-busy:coder"] + assert r.deleted == [] # tombstone left for the next tick + assert r.created == [] + + +def test_reconcile_retry_proceeds_and_draws_shared_load_when_slot_free(): + # With a free slot the retry recreates normally and draws the coder slot + # down in place, so the same-tick issue-claim / pr-fix passes see it. + load = {} + r = _Recorder([_failed_wl("wl-misospace-dispatch-7", attempt=1, lane="local")]) + out = reconcile_failures( + "foreman-coder", r.list_failed, r.create, r.delete, + namespace="llm", gate_profiles={}, max_attempts=3, + lane_coder_agents={"local": "coder"}, + agent_load=load, agent_slots={"coder": 1}, + ) + assert out == ["wl-misospace-dispatch-7:retry:2/3"] + assert len(r.created) == 1 + assert load == {"coder": 1} # drawn down in place for downstream passes + + +def test_reconcile_retry_uncapped_when_no_slots_configured(): + # No agent_slots -> legacy behavior, retry always recreates. + r = _Recorder([_failed_wl("wl-misospace-dispatch-7", attempt=1, lane="local")]) + out = reconcile_failures( + "foreman-coder", r.list_failed, r.create, r.delete, + namespace="llm", gate_profiles={}, max_attempts=3, + agent_load={"coder": 5}, agent_slots={}, + ) + assert out == ["wl-misospace-dispatch-7:retry:2/3"] + assert len(r.created) == 1 + + def test_reconcile_gives_up_at_max_without_touching_the_workload(): r = _Recorder([_failed_wl("wl-misospace-dispatch-7", attempt=3)]) out = reconcile_failures("foreman-coder", r.list_failed, r.create, r.delete,