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
8 changes: 7 additions & 1 deletion bridge/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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
Expand Down
30 changes: 24 additions & 6 deletions bridge/retry.py
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
from bridge.workload import (
build_workload,
coder_agent_for,
free_slots,
gate_profile_for,
_branch_name,
ATTEMPT_ANNOTATION,
Expand Down Expand Up @@ -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.

Expand Down Expand Up @@ -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(
Expand Down Expand Up @@ -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
Expand All @@ -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
Expand Down
44 changes: 44 additions & 0 deletions tests/test_retry.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down