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
21 changes: 16 additions & 5 deletions bridge/main.py
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
coder_agent_for,
coder_candidates,
coders_saturated,
free_slots,
revision_coder_agent_for,
gate_profile_for,
parse_gate_profiles,
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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)

Expand Down
31 changes: 30 additions & 1 deletion bridge/prfix.py
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand All @@ -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 "?"
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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}"
Expand All @@ -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:
Expand All @@ -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:
Expand Down
74 changes: 74 additions & 0 deletions tests/test_prfix.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
Loading