fix(kanban): review reservation also skips per-profile-capped review rows
_any_spawnable_review now mirrors every gate _dispatch_lane_task applies to a review row this tick: profile exists, assignee under kanban.max_in_progress_per_profile, and not respawn-guarded (the guard half landed in the salvaged commit from #113600). per_profile_running is resolved before the reservation so the predicate can see it. Why: with spawn_budget == 1 a review row whose assignee is already at the per-profile cap reserved the only slot, the review loop then skipped it as capped, and the ready lane starved every tick — the same starvation shape as the guard case (#113598, last comment). The guard test is parametrized over both ways a review row can be refused (respawn guard, per-profile cap); the unguarded control keeps reserving.
This commit is contained in:
@@ -2084,23 +2084,32 @@ def _lane_rows(conn: sqlite3.Connection, status: str) -> list[sqlite3.Row]:
|
||||
|
||||
|
||||
def _any_spawnable_review(
|
||||
conn: sqlite3.Connection, review_rows: list[sqlite3.Row],
|
||||
conn: sqlite3.Connection,
|
||||
review_rows: list[sqlite3.Row],
|
||||
*,
|
||||
per_profile_cap: Optional[int] = None,
|
||||
per_profile_running: Optional[dict[str, int]] = None,
|
||||
) -> bool:
|
||||
"""Mirror review dispatch gates before reserving ready-lane capacity.
|
||||
|
||||
Unavailable profile metadata retains the historic fail-open behavior. A
|
||||
respawn-guarded review row cannot consume the reservation, however, so it
|
||||
must not withhold capacity from an otherwise ready task.
|
||||
review row that :func:`_dispatch_lane_task` would refuse this tick — its
|
||||
assignee already at the per-profile cap, or respawn-guarded — cannot
|
||||
consume the reservation, so it must not withhold capacity from an
|
||||
otherwise ready task (one such row would pin ``ready_budget`` to 0).
|
||||
"""
|
||||
if not review_rows:
|
||||
return False
|
||||
profile_exists = _profile_exists_fn()
|
||||
running = per_profile_running or {}
|
||||
for row in review_rows:
|
||||
assignee = row["assignee"]
|
||||
if not assignee:
|
||||
continue
|
||||
if profile_exists is not None and not profile_exists(assignee):
|
||||
continue
|
||||
if per_profile_cap is not None and running.get(assignee, 0) >= per_profile_cap:
|
||||
continue
|
||||
if check_respawn_guard(conn, row["id"], lane="review") is None:
|
||||
return True
|
||||
return False
|
||||
@@ -2159,15 +2168,10 @@ def _dispatch_once_locked(
|
||||
# Review rows are enumerated up front so the budget split can see whether
|
||||
# review work exists at all.
|
||||
review_rows = _lane_rows(conn, "review") if review_dispatch_enabled() else []
|
||||
# Review-lane reservation: the ready loop runs first and would otherwise
|
||||
# consume the ENTIRE shared budget, starving reviews under a sustained ready
|
||||
# backlog. When spawnable review work exists and there is any budget, hold
|
||||
# one slot back.
|
||||
ready_budget = spawn_budget
|
||||
if spawn_budget is not None and spawn_budget > 0 and _any_spawnable_review(conn, review_rows):
|
||||
ready_budget = max(spawn_budget - 1, 0)
|
||||
# Per-profile cap. Deferred tasks go to skipped_per_profile_capped, not
|
||||
# skipped_unassigned — "busy, retry later" differs from "needs routing".
|
||||
# Resolved BEFORE the review reservation so the reservation can see which
|
||||
# review rows the lane loop would refuse this tick.
|
||||
per_profile_cap = max_in_progress_per_profile if (
|
||||
# Per-profile concurrency cap (#21582): when set, track how many workers each assignee already has
|
||||
# in flight, and refuse to spawn when this would push that assignee past the cap. Prevents fan-out
|
||||
@@ -2184,6 +2188,16 @@ def _dispatch_once_locked(
|
||||
"GROUP BY assignee"
|
||||
):
|
||||
per_profile_running[prow["assignee"]] = int(prow["n"])
|
||||
# Review-lane reservation: the ready loop runs first and would otherwise
|
||||
# consume the ENTIRE shared budget, starving reviews under a sustained ready
|
||||
# backlog. When spawnable review work exists and there is any budget, hold
|
||||
# one slot back.
|
||||
ready_budget = spawn_budget
|
||||
if spawn_budget is not None and spawn_budget > 0 and _any_spawnable_review(
|
||||
conn, review_rows,
|
||||
per_profile_cap=per_profile_cap, per_profile_running=per_profile_running,
|
||||
):
|
||||
ready_budget = max(spawn_budget - 1, 0)
|
||||
lane_kwargs: dict[str, Any] = dict(
|
||||
dry_run=dry_run, ttl_seconds=ttl_seconds, board=board,
|
||||
failure_limit=failure_limit, spawn_fn=spawn_fn,
|
||||
|
||||
@@ -243,10 +243,34 @@ def test_review_lane_gets_reserved_slot_under_ready_backlog(
|
||||
assert review_id in spawned_ids
|
||||
|
||||
|
||||
def test_guarded_review_does_not_reserve_the_only_ready_slot(
|
||||
kanban_home, all_assignees_spawnable, monkeypatch,
|
||||
def _guard_review_row(conn: sqlite3.Connection, review_id: str) -> dict:
|
||||
"""Latest run ``rate_limited`` → ``check_respawn_guard`` returns a cooldown."""
|
||||
now = int(time.time())
|
||||
with kb.write_txn(conn):
|
||||
conn.execute(
|
||||
"INSERT INTO task_runs (task_id, profile, status, outcome, "
|
||||
"started_at, ended_at) VALUES (?, 'reviewer', 'rate_limited', "
|
||||
"'rate_limited', ?, ?)",
|
||||
(review_id, now, now),
|
||||
)
|
||||
assert kbd.check_respawn_guard(conn, review_id, lane="review") == "rate_limit_cooldown"
|
||||
return {"max_in_progress": 1}
|
||||
|
||||
|
||||
def _cap_review_row(conn: sqlite3.Connection, review_id: str) -> dict:
|
||||
"""``reviewer`` already has one running worker → the review row is per-profile capped."""
|
||||
busy_id = kb.create_task(conn, title="busy", assignee="reviewer")
|
||||
assert kb.claim_task(conn, busy_id) is not None
|
||||
return {"max_in_progress": 2, "max_in_progress_per_profile": 1}
|
||||
|
||||
|
||||
@pytest.mark.parametrize("make_unspawnable", [_guard_review_row, _cap_review_row])
|
||||
def test_unspawnable_review_does_not_reserve_the_only_ready_slot(
|
||||
kanban_home, all_assignees_spawnable, monkeypatch, make_unspawnable,
|
||||
):
|
||||
"""A review card in cooldown must not consume fairness reservation."""
|
||||
"""A review card the review loop would refuse this tick (respawn guard,
|
||||
per-profile cap) must not consume the fairness reservation — otherwise the
|
||||
ready lane starves every tick while the reserved slot goes unused."""
|
||||
import hermes_cli.config as cfgmod
|
||||
|
||||
monkeypatch.setattr(
|
||||
@@ -257,18 +281,9 @@ def test_guarded_review_does_not_reserve_the_only_ready_slot(
|
||||
spawns: list = []
|
||||
with kbc.connect() as conn:
|
||||
ready_id = kb.create_task(conn, title="ready-now", assignee="alice")
|
||||
review_id = _park_in_review(conn, "review-cooldown", "reviewer")
|
||||
now = int(time.time())
|
||||
with kb.write_txn(conn):
|
||||
conn.execute(
|
||||
"INSERT INTO task_runs (task_id, profile, status, outcome, "
|
||||
"started_at, ended_at) VALUES (?, 'reviewer', 'rate_limited', "
|
||||
"'rate_limited', ?, ?)",
|
||||
(review_id, now, now),
|
||||
)
|
||||
res = kbd.dispatch_once(
|
||||
conn, spawn_fn=_fake_spawn_factory(spawns), max_in_progress=1,
|
||||
)
|
||||
review_id = _park_in_review(conn, "review-unspawnable", "reviewer")
|
||||
caps = make_unspawnable(conn, review_id)
|
||||
res = kbd.dispatch_once(conn, spawn_fn=_fake_spawn_factory(spawns), **caps)
|
||||
|
||||
assert [task_id for task_id, *_ in res.spawned] == [ready_id]
|
||||
|
||||
|
||||
Reference in New Issue
Block a user