fix(kanban): preserve review phase across retries
This commit is contained in:
parent
c230d1202f
commit
b90da8243b
|
|
@ -4196,6 +4196,33 @@ def _has_sticky_block(conn: sqlite3.Connection, task_id: str) -> bool:
|
|||
return bool(row) and row["kind"] == "blocked"
|
||||
|
||||
|
||||
def _resume_status_from_events(conn: sqlite3.Connection, task_id: str) -> str:
|
||||
"""Return the durable phase a blocked/dependency-wait task should resume.
|
||||
|
||||
Events written by review workers carry ``source_status``/``retry_status``;
|
||||
an explicit unblock that must wait for parents carries ``resume_status``.
|
||||
Legacy events omit these fields and therefore retain the historical
|
||||
``ready`` behavior.
|
||||
"""
|
||||
row = conn.execute(
|
||||
"SELECT payload FROM task_events "
|
||||
"WHERE task_id = ? AND kind IN ("
|
||||
"'blocked', 'block_loop_detected', 'dependency_wait', 'gave_up', "
|
||||
"'unblocked', 'changes_requested', 'review_reopened', 'reclaimed', "
|
||||
"'stale', 'timed_out', 'crashed', 'spawn_failed', 'rate_limited'"
|
||||
") ORDER BY id DESC LIMIT 1",
|
||||
(task_id,),
|
||||
).fetchone()
|
||||
try:
|
||||
payload = json.loads(row["payload"]) if row and row["payload"] else {}
|
||||
except (json.JSONDecodeError, TypeError):
|
||||
payload = {}
|
||||
for key in ("resume_status", "retry_status", "source_status"):
|
||||
if payload.get(key) == "review":
|
||||
return "review"
|
||||
return "ready"
|
||||
|
||||
|
||||
def recompute_ready(
|
||||
conn: sqlite3.Connection, failure_limit: int = None,
|
||||
) -> int:
|
||||
|
|
@ -4251,6 +4278,7 @@ def recompute_ready(
|
|||
(task_id,),
|
||||
).fetchall()
|
||||
if all(p["status"] in ("done", "archived") for p in parents):
|
||||
resume_status = _resume_status_from_events(conn, task_id)
|
||||
if cur_status == "blocked":
|
||||
# Don't auto-recover tasks that have hit the
|
||||
# circuit-breaker failure limit. Without this
|
||||
|
|
@ -4269,16 +4297,19 @@ def recompute_ready(
|
|||
if failures >= effective_limit:
|
||||
continue
|
||||
conn.execute(
|
||||
"UPDATE tasks SET status = 'ready' "
|
||||
"UPDATE tasks SET status = ? "
|
||||
"WHERE id = ? AND status = 'blocked'",
|
||||
(task_id,),
|
||||
(resume_status, task_id),
|
||||
)
|
||||
else:
|
||||
conn.execute(
|
||||
"UPDATE tasks SET status = 'ready' WHERE id = ? AND status = 'todo'",
|
||||
(task_id,),
|
||||
"UPDATE tasks SET status = ? WHERE id = ? AND status = 'todo'",
|
||||
(resume_status, task_id),
|
||||
)
|
||||
_append_event(conn, task_id, "promoted", None)
|
||||
_append_event(
|
||||
conn, task_id, "promoted",
|
||||
{"status": resume_status} if resume_status != "ready" else None,
|
||||
)
|
||||
promoted += 1
|
||||
return promoted
|
||||
|
||||
|
|
@ -4484,6 +4515,39 @@ def claim_review_task(
|
|||
return get_task(conn, task_id)
|
||||
|
||||
|
||||
def _retry_status_for_run(
|
||||
conn: sqlite3.Connection,
|
||||
task_id: str,
|
||||
run_id: Optional[int] = None,
|
||||
) -> str:
|
||||
"""Return the non-running phase an interrupted run must resume from.
|
||||
|
||||
Review claims record ``source_status=review`` on their claimed event. All
|
||||
other and legacy runs retry from ``ready``. Keeping this decision in one
|
||||
place prevents crash/timeout/reclaim paths from silently converting a
|
||||
reviewer run into an implementation run.
|
||||
"""
|
||||
if run_id is None:
|
||||
row = conn.execute(
|
||||
"SELECT current_run_id FROM tasks WHERE id = ?",
|
||||
(task_id,),
|
||||
).fetchone()
|
||||
run_id = row["current_run_id"] if row else None
|
||||
if run_id is None:
|
||||
return "ready"
|
||||
event = conn.execute(
|
||||
"SELECT payload FROM task_events "
|
||||
"WHERE task_id = ? AND run_id = ? AND kind = 'claimed' "
|
||||
"ORDER BY id DESC LIMIT 1",
|
||||
(task_id, int(run_id)),
|
||||
).fetchone()
|
||||
try:
|
||||
payload = json.loads(event["payload"]) if event and event["payload"] else {}
|
||||
except (json.JSONDecodeError, TypeError):
|
||||
payload = {}
|
||||
return "review" if payload.get("source_status") == "review" else "ready"
|
||||
|
||||
|
||||
def heartbeat_claim(
|
||||
conn: sqlite3.Connection,
|
||||
task_id: str,
|
||||
|
|
@ -4621,12 +4685,13 @@ def release_stale_claims(
|
|||
)
|
||||
continue
|
||||
with write_txn(conn):
|
||||
retry_status = _retry_status_for_run(conn, row["id"])
|
||||
cur = conn.execute(
|
||||
"UPDATE tasks SET status = 'ready', claim_lock = NULL, "
|
||||
"UPDATE tasks SET status = ?, claim_lock = NULL, "
|
||||
"claim_expires = NULL, worker_pid = NULL "
|
||||
"WHERE id = ? AND status = 'running' AND claim_lock IS ? "
|
||||
"AND claim_expires IS NOT NULL AND claim_expires < ?",
|
||||
(row["id"], row["claim_lock"], now),
|
||||
(retry_status, row["id"], row["claim_lock"], now),
|
||||
)
|
||||
if cur.rowcount != 1:
|
||||
continue
|
||||
|
|
@ -4650,6 +4715,7 @@ def release_stale_claims(
|
|||
"now": now,
|
||||
"host_local": host_local,
|
||||
"heartbeat_stale": bool(heartbeat_stale),
|
||||
"retry_status": retry_status,
|
||||
}
|
||||
payload.update(termination)
|
||||
_append_event(
|
||||
|
|
@ -4668,7 +4734,7 @@ def reclaim_task(
|
|||
reason: Optional[str] = None,
|
||||
signal_fn=None,
|
||||
) -> bool:
|
||||
"""Operator-driven reclaim: release the claim and reset to ``ready``.
|
||||
"""Operator-driven reclaim: release the claim and restore its source phase.
|
||||
|
||||
Unlike :func:`release_stale_claims` which only acts on tasks whose
|
||||
``claim_expires`` has passed, this function reclaims immediately
|
||||
|
|
@ -4693,12 +4759,13 @@ def reclaim_task(
|
|||
row["worker_pid"], prev_lock, signal_fn=signal_fn,
|
||||
)
|
||||
with write_txn(conn):
|
||||
retry_status = _retry_status_for_run(conn, task_id)
|
||||
cur = conn.execute(
|
||||
"UPDATE tasks SET status = 'ready', claim_lock = NULL, "
|
||||
"UPDATE tasks SET status = ?, claim_lock = NULL, "
|
||||
"claim_expires = NULL, worker_pid = NULL "
|
||||
"WHERE id = ? AND status IN ('running', 'ready', 'blocked') "
|
||||
"AND claim_lock IS ?",
|
||||
(task_id, prev_lock),
|
||||
(retry_status, task_id, prev_lock),
|
||||
)
|
||||
if cur.rowcount != 1:
|
||||
return False
|
||||
|
|
@ -4715,6 +4782,7 @@ def reclaim_task(
|
|||
"manual": True,
|
||||
"reason": reason,
|
||||
"prev_lock": prev_lock,
|
||||
"retry_status": retry_status,
|
||||
}
|
||||
payload.update(termination)
|
||||
_append_event(
|
||||
|
|
@ -5728,6 +5796,11 @@ def block_task(
|
|||
).fetchone()
|
||||
if cur_row is None:
|
||||
return False
|
||||
source_status = (
|
||||
_retry_status_for_run(conn, task_id)
|
||||
if cur_row["status"] == "running"
|
||||
else "ready"
|
||||
)
|
||||
prev_kind = cur_row["block_kind"] if "block_kind" in cur_row.keys() else None
|
||||
prev_recurrences = (
|
||||
int(cur_row["block_recurrences"])
|
||||
|
|
@ -5768,7 +5841,12 @@ def block_task(
|
|||
)
|
||||
_append_event(
|
||||
conn, task_id, "dependency_wait",
|
||||
{"reason": reason, "kind": kind}, run_id=run_id,
|
||||
{
|
||||
"reason": reason,
|
||||
"kind": kind,
|
||||
"source_status": source_status,
|
||||
},
|
||||
run_id=run_id,
|
||||
)
|
||||
_blocked_task = get_task(conn, task_id)
|
||||
_fire_kanban_lifecycle_hook(
|
||||
|
|
@ -5826,6 +5904,7 @@ def block_task(
|
|||
"kind": kind,
|
||||
"recurrences": recurrences,
|
||||
"limit": BLOCK_RECURRENCE_LIMIT,
|
||||
"source_status": source_status,
|
||||
},
|
||||
run_id=run_id,
|
||||
)
|
||||
|
|
@ -5878,7 +5957,12 @@ def block_task(
|
|||
)
|
||||
_append_event(
|
||||
conn, task_id, "blocked",
|
||||
{"reason": reason, "kind": kind, "recurrences": recurrences},
|
||||
{
|
||||
"reason": reason,
|
||||
"kind": kind,
|
||||
"recurrences": recurrences,
|
||||
"source_status": source_status,
|
||||
},
|
||||
run_id=run_id,
|
||||
)
|
||||
_blocked_task = get_task(conn, task_id)
|
||||
|
|
@ -6205,7 +6289,7 @@ def _landing_status_after_parents(conn: sqlite3.Connection, task_id: str) -> str
|
|||
|
||||
|
||||
def unblock_task(conn: sqlite3.Connection, task_id: str) -> bool:
|
||||
"""Transition ``blocked``/``scheduled`` -> ready or todo.
|
||||
"""Transition ``blocked``/``scheduled`` to its safe resumable phase.
|
||||
|
||||
Defensively closes any stale ``current_run_id`` pointer before flipping
|
||||
status. In the common path (``block_task`` closed the run already) this
|
||||
|
|
@ -6216,13 +6300,26 @@ def unblock_task(conn: sqlite3.Connection, task_id: str) -> bool:
|
|||
"""
|
||||
now = int(time.time())
|
||||
with write_txn(conn):
|
||||
current = conn.execute(
|
||||
"SELECT status FROM tasks WHERE id = ?",
|
||||
(task_id,),
|
||||
).fetchone()
|
||||
resume_status = (
|
||||
_resume_status_from_events(conn, task_id)
|
||||
if current and current["status"] == "blocked"
|
||||
else "ready"
|
||||
)
|
||||
_reclaim_dangling_run(
|
||||
conn, task_id, statuses=("blocked", "scheduled"), now=now,
|
||||
note="invariant recovery on unblock",
|
||||
)
|
||||
# Re-gate on parent completion before flipping 'blocked' back to
|
||||
# 'ready' (see :func:`_landing_status_after_parents`).
|
||||
new_status = _landing_status_after_parents(conn, task_id)
|
||||
# Re-gate on parent completion before restoring the source phase.
|
||||
landing_status = _landing_status_after_parents(conn, task_id)
|
||||
new_status = (
|
||||
"review"
|
||||
if landing_status == "ready" and resume_status == "review"
|
||||
else landing_status
|
||||
)
|
||||
# NOTE: deliberately does NOT touch ``block_recurrences`` or
|
||||
# ``block_kind``. Resetting the recurrence counter on unblock is exactly
|
||||
# the amnesia that let a cron unblock → worker re-block loop run
|
||||
|
|
@ -6243,7 +6340,11 @@ def unblock_task(conn: sqlite3.Connection, task_id: str) -> bool:
|
|||
return False
|
||||
_append_event(
|
||||
conn, task_id, "unblocked",
|
||||
{"status": new_status} if new_status != "ready" else None,
|
||||
(
|
||||
{"status": new_status, "resume_status": resume_status}
|
||||
if new_status != "ready" or resume_status != "ready"
|
||||
else None
|
||||
),
|
||||
)
|
||||
return True
|
||||
|
||||
|
|
@ -7546,8 +7647,8 @@ def enforce_max_runtime(
|
|||
"""Terminate workers whose per-task ``max_runtime_seconds`` has elapsed.
|
||||
|
||||
Sends SIGTERM, waits a short grace window, then SIGKILL. Emits a
|
||||
``timed_out`` event and drops the task back to ``ready`` so the next
|
||||
dispatcher tick re-spawns it — unless the spawn-failure circuit
|
||||
``timed_out`` event and restores the task's source phase so the next
|
||||
dispatcher tick re-spawns the same kind of worker — unless the circuit
|
||||
breaker has already given up, in which case the task stays blocked
|
||||
where ``_record_spawn_failure`` parked it.
|
||||
|
||||
|
|
@ -7610,13 +7711,14 @@ def enforce_max_runtime(
|
|||
pass
|
||||
|
||||
with write_txn(conn):
|
||||
retry_status = _retry_status_for_run(conn, tid)
|
||||
cur = conn.execute(
|
||||
"UPDATE tasks SET status = 'ready', claim_lock = NULL, "
|
||||
"UPDATE tasks SET status = ?, claim_lock = NULL, "
|
||||
"claim_expires = NULL, worker_pid = NULL, "
|
||||
"last_heartbeat_at = NULL "
|
||||
"WHERE id = ? AND status = 'running' "
|
||||
" AND worker_pid = ? AND claim_lock IS ?",
|
||||
(tid, pid, row["claim_lock"]),
|
||||
(retry_status, tid, pid, row["claim_lock"]),
|
||||
)
|
||||
if cur.rowcount == 1:
|
||||
payload = {
|
||||
|
|
@ -7624,6 +7726,7 @@ def enforce_max_runtime(
|
|||
"elapsed_seconds": int(elapsed),
|
||||
"limit_seconds": int(row["max_runtime_seconds"]),
|
||||
"sigkill": killed,
|
||||
"retry_status": retry_status,
|
||||
}
|
||||
run_id = _end_run(
|
||||
conn, tid,
|
||||
|
|
@ -7637,7 +7740,7 @@ def enforce_max_runtime(
|
|||
timed_out.append(tid)
|
||||
# Increment the unified failure counter. Outside the write_txn
|
||||
# above because ``_record_task_failure`` opens its own. If the
|
||||
# breaker trips, this flips the task ``ready → blocked`` and
|
||||
# breaker trips, this flips the retried task to ``blocked`` and
|
||||
# emits a ``gave_up`` event on top of the ``timed_out`` we
|
||||
# already emitted.
|
||||
if cur.rowcount == 1:
|
||||
|
|
@ -7647,7 +7750,11 @@ def enforce_max_runtime(
|
|||
outcome="timed_out",
|
||||
release_claim=False,
|
||||
end_run=False,
|
||||
event_payload_extra={"pid": pid, "sigkill": killed},
|
||||
event_payload_extra={
|
||||
"pid": pid,
|
||||
"sigkill": killed,
|
||||
"retry_status": retry_status,
|
||||
},
|
||||
)
|
||||
return timed_out
|
||||
|
||||
|
|
@ -7676,7 +7783,7 @@ def detect_stale_running(
|
|||
2. Its ``last_heartbeat_at`` is older than
|
||||
``_STALE_HEARTBEAT_GAP_SECONDS`` (or NULL — never sent a heartbeat).
|
||||
|
||||
On reclaim the task is reset to ``ready``, the run is closed with
|
||||
On reclaim the task is restored to its source phase, the run is closed with
|
||||
``outcome='stale'``, and the host-local worker (if still running) is
|
||||
terminated.
|
||||
|
||||
|
|
@ -7735,13 +7842,14 @@ def detect_stale_running(
|
|||
continue
|
||||
|
||||
with write_txn(conn):
|
||||
retry_status = _retry_status_for_run(conn, tid)
|
||||
cur = conn.execute(
|
||||
"UPDATE tasks SET status = 'ready', claim_lock = NULL, "
|
||||
"UPDATE tasks SET status = ?, claim_lock = NULL, "
|
||||
"claim_expires = NULL, worker_pid = NULL, "
|
||||
"last_heartbeat_at = NULL "
|
||||
"WHERE id = ? AND status = 'running' "
|
||||
" AND claim_lock IS ?",
|
||||
(tid, row["claim_lock"]),
|
||||
(retry_status, tid, row["claim_lock"]),
|
||||
)
|
||||
if cur.rowcount != 1:
|
||||
continue
|
||||
|
|
@ -7756,6 +7864,7 @@ def detect_stale_running(
|
|||
),
|
||||
"timeout_seconds": stale_timeout_seconds,
|
||||
"pid": int(pid) if pid else None,
|
||||
"retry_status": retry_status,
|
||||
}
|
||||
payload.update(termination)
|
||||
|
||||
|
|
@ -7776,7 +7885,7 @@ def detect_stale_running(
|
|||
|
||||
# Intentionally NOT calling _record_task_failure here. Stale reclaim
|
||||
# is dispatcher-side detection of an absent heartbeat; the task is
|
||||
# going straight back to ``ready`` for re-dispatch. Counting it as
|
||||
# going straight back to its source phase for re-dispatch. Counting it as
|
||||
# a worker failure would let two legitimately-long-running tasks
|
||||
# (>4h without explicit heartbeat) trip the circuit breaker and
|
||||
# auto-block, even though no worker actually failed. The 'stale'
|
||||
|
|
@ -7962,7 +8071,7 @@ def _protocol_violation_streak(conn: sqlite3.Connection, task_id: str) -> int:
|
|||
def detect_crashed_workers(conn: sqlite3.Connection) -> list[str]:
|
||||
"""Reclaim ``running`` tasks whose worker PID is no longer alive.
|
||||
|
||||
Appends a ``crashed`` event and drops the task back to ``ready``.
|
||||
Appends a ``crashed`` event and restores the task's source phase.
|
||||
Different from ``release_stale_claims``: this checks liveness
|
||||
immediately rather than waiting for the claim TTL.
|
||||
|
||||
|
|
@ -7981,7 +8090,7 @@ def detect_crashed_workers(conn: sqlite3.Connection) -> list[str]:
|
|||
When the reap registry shows the worker exited with the rate-limit
|
||||
sentinel (``KANBAN_RATE_LIMIT_EXIT_CODE``), the worker bailed on a
|
||||
provider quota wall, NOT a task failure. Such tasks are released back
|
||||
to ``ready`` WITHOUT counting a failure (so a long quota window can't
|
||||
to its source phase WITHOUT counting a failure (so a long quota window can't
|
||||
trip the breaker) and stamped with a quota-blocker error so
|
||||
``check_respawn_guard`` defers their respawn until the window clears.
|
||||
The ids are returned via the ``_last_rate_limited`` function attribute
|
||||
|
|
@ -8053,7 +8162,7 @@ def detect_crashed_workers(conn: sqlite3.Connection) -> list[str]:
|
|||
# Worker bailed because the provider rate-limited / exhausted
|
||||
# quota (EX_TEMPFAIL sentinel). This is NOT a task failure —
|
||||
# the task is fine, the account just hit a wall. Release it
|
||||
# back to ``ready`` so the respawn guard defers it until the
|
||||
# back to its source phase so the respawn guard defers it until the
|
||||
# quota window clears, and crucially do NOT count a failure
|
||||
# (skip ``_record_task_failure``) so a long quota window can't
|
||||
# trip the circuit breaker and permanently block the card.
|
||||
|
|
@ -8083,12 +8192,14 @@ def detect_crashed_workers(conn: sqlite3.Connection) -> list[str]:
|
|||
event_payload["exit_kind"] = kind
|
||||
event_payload["exit_code"] = code
|
||||
|
||||
retry_status = _retry_status_for_run(conn, row["id"])
|
||||
event_payload["retry_status"] = retry_status
|
||||
cur = conn.execute(
|
||||
"UPDATE tasks SET status = 'ready', claim_lock = NULL, "
|
||||
"UPDATE tasks SET status = ?, claim_lock = NULL, "
|
||||
"claim_expires = NULL, worker_pid = NULL "
|
||||
"WHERE id = ? AND status = 'running' "
|
||||
" AND worker_pid = ? AND claim_lock IS ?",
|
||||
(row["id"], pid, row["claim_lock"]),
|
||||
(retry_status, row["id"], pid, row["claim_lock"]),
|
||||
)
|
||||
if cur.rowcount == 1:
|
||||
# Rate-limited requeues are a clean release, not a crash —
|
||||
|
|
@ -8136,7 +8247,7 @@ def detect_crashed_workers(conn: sqlite3.Connection) -> list[str]:
|
|||
protocol_violation, error_text)
|
||||
)
|
||||
# Outside the main txn: account each crashed task and maybe trip the
|
||||
# breaker (the task transitions ready → blocked with a ``gave_up`` event
|
||||
# breaker (the retried task transitions to blocked with a ``gave_up`` event
|
||||
# on top of the event we already emitted).
|
||||
#
|
||||
# Protocol-violation crashes (clean exit, no terminal tool call) get a
|
||||
|
|
@ -8256,14 +8367,14 @@ def _record_task_failure(
|
|||
|
||||
* ``release_claim=True, end_run=True`` — spawn-failure path.
|
||||
Caller has a running task with an open run; this transitions
|
||||
it back to ``ready`` (or ``blocked`` when the breaker trips),
|
||||
it back to its source phase (or ``blocked`` when the breaker trips),
|
||||
releases the claim, and closes the run with ``outcome=<outcome>``.
|
||||
|
||||
* ``release_claim=False, end_run=False`` — timeout/crash path.
|
||||
Caller has ALREADY flipped the task to ``ready`` and closed the
|
||||
Caller has ALREADY restored the task's source phase and closed the
|
||||
run with the appropriate outcome. This just increments the
|
||||
counter; if the breaker trips, the task is re-transitioned
|
||||
``ready → blocked`` and a ``gave_up`` event is emitted.
|
||||
into ``blocked`` and a ``gave_up`` event is emitted.
|
||||
|
||||
``event_payload_extra`` merges into the ``gave_up`` event payload
|
||||
when the breaker trips, so callers can include outcome-specific
|
||||
|
|
@ -8289,11 +8400,16 @@ def _record_task_failure(
|
|||
blocked = False
|
||||
with write_txn(conn):
|
||||
row = conn.execute(
|
||||
"SELECT consecutive_failures, status, max_retries "
|
||||
"SELECT consecutive_failures, status, max_retries, current_run_id "
|
||||
"FROM tasks WHERE id = ?", (task_id,),
|
||||
).fetchone()
|
||||
if row is None:
|
||||
return False
|
||||
retry_status = (
|
||||
_retry_status_for_run(conn, task_id, row["current_run_id"])
|
||||
if release_claim
|
||||
else ("review" if row["status"] == "review" else "ready")
|
||||
)
|
||||
failures = int(row["consecutive_failures"]) + 1
|
||||
|
||||
# Per-task override wins over both caller-supplied and default
|
||||
|
|
@ -8316,17 +8432,17 @@ def _record_task_failure(
|
|||
"UPDATE tasks SET status = 'blocked', claim_lock = NULL, "
|
||||
"claim_expires = NULL, worker_pid = NULL, "
|
||||
"consecutive_failures = ?, last_failure_error = ? "
|
||||
"WHERE id = ? AND status IN ('running', 'ready')",
|
||||
"WHERE id = ? AND status IN ('running', 'ready', 'review')",
|
||||
(failures, error[:500], task_id),
|
||||
)
|
||||
else:
|
||||
# Timeout/crash path: task is already at ``ready``
|
||||
# with claim cleared; just flip to blocked + update
|
||||
# Timeout/crash path: source phase already restored with claim
|
||||
# cleared; just flip to blocked + update
|
||||
# counter fields.
|
||||
conn.execute(
|
||||
"UPDATE tasks SET status = 'blocked', "
|
||||
"consecutive_failures = ?, last_failure_error = ? "
|
||||
"WHERE id = ? AND status IN ('ready', 'running')",
|
||||
"WHERE id = ? AND status IN ('ready', 'review', 'running')",
|
||||
(failures, error[:500], task_id),
|
||||
)
|
||||
run_id = None
|
||||
|
|
@ -8341,6 +8457,7 @@ def _record_task_failure(
|
|||
"trigger_outcome": outcome,
|
||||
"effective_limit": effective_limit,
|
||||
"limit_source": limit_source,
|
||||
"retry_status": retry_status,
|
||||
},
|
||||
)
|
||||
payload = {
|
||||
|
|
@ -8349,6 +8466,7 @@ def _record_task_failure(
|
|||
"limit_source": limit_source,
|
||||
"error": error[:500],
|
||||
"trigger_outcome": outcome,
|
||||
"retry_status": retry_status,
|
||||
}
|
||||
if event_payload_extra:
|
||||
payload.update(event_payload_extra)
|
||||
|
|
@ -8359,17 +8477,16 @@ def _record_task_failure(
|
|||
else:
|
||||
# Below threshold.
|
||||
if release_claim:
|
||||
# Spawn path: transition running → ready + clear claim.
|
||||
# Spawn path: restore the claimed source phase + clear claim.
|
||||
conn.execute(
|
||||
"UPDATE tasks SET status = 'ready', claim_lock = NULL, "
|
||||
"UPDATE tasks SET status = ?, claim_lock = NULL, "
|
||||
"claim_expires = NULL, worker_pid = NULL, "
|
||||
"consecutive_failures = ?, last_failure_error = ? "
|
||||
"WHERE id = ? AND status = 'running'",
|
||||
(failures, error[:500], task_id),
|
||||
(retry_status, failures, error[:500], task_id),
|
||||
)
|
||||
else:
|
||||
# Timeout/crash path: task is already at ``ready`` via
|
||||
# its own UPDATE. Just bookkeep the counter + last error.
|
||||
# Timeout/crash path: caller already restored the source phase.
|
||||
conn.execute(
|
||||
"UPDATE tasks SET consecutive_failures = ?, "
|
||||
"last_failure_error = ? WHERE id = ?",
|
||||
|
|
@ -8381,11 +8498,18 @@ def _record_task_failure(
|
|||
conn, task_id,
|
||||
outcome=outcome, status=outcome,
|
||||
error=error[:500],
|
||||
metadata={"failures": failures},
|
||||
metadata={
|
||||
"failures": failures,
|
||||
"retry_status": retry_status,
|
||||
},
|
||||
)
|
||||
_append_event(
|
||||
conn, task_id, outcome,
|
||||
{"error": error[:500], "failures": failures},
|
||||
{
|
||||
"error": error[:500],
|
||||
"failures": failures,
|
||||
"retry_status": retry_status,
|
||||
},
|
||||
run_id=run_id,
|
||||
)
|
||||
# Timeout/crash path's caller already emitted its own event.
|
||||
|
|
@ -8456,9 +8580,9 @@ _clear_spawn_failures = _clear_failure_counter
|
|||
def check_respawn_guard(conn: sqlite3.Connection, task_id: str) -> Optional[str]:
|
||||
"""Return a guard reason if ``task_id`` should NOT be re-spawned, else None.
|
||||
|
||||
Called per ready task in ``dispatch_once`` before any claim attempt.
|
||||
Returning a reason defers the spawn this tick; the task stays in
|
||||
``ready`` and gets another chance on the next dispatcher tick.
|
||||
Called per ready/review task in ``dispatch_once`` before any claim attempt.
|
||||
Returning a reason defers the spawn this tick; the task stays in its
|
||||
source phase and gets another chance on the next dispatcher tick.
|
||||
|
||||
Checks in priority order:
|
||||
|
||||
|
|
@ -8811,6 +8935,14 @@ def _dispatch_once_locked(
|
|||
result.timed_out = enforce_max_runtime(conn)
|
||||
result.promoted = recompute_ready(conn, failure_limit=failure_limit)
|
||||
|
||||
# Both knobs are total in-flight caps. Collapse them before either lane
|
||||
# dispatches so ready and review workers consume the same budget without
|
||||
# subtracting the already-running count twice.
|
||||
if max_in_progress is not None and (
|
||||
max_spawn is None or max_in_progress < max_spawn
|
||||
):
|
||||
max_spawn = max_in_progress
|
||||
|
||||
# Count tasks already running so max_spawn enforces concurrency rather
|
||||
# than a per-tick spawn budget. See the docstring above for the full
|
||||
# rationale; the short version is that a 60-second tick interval with a
|
||||
|
|
@ -8831,20 +8963,6 @@ def _dispatch_once_locked(
|
|||
"WHERE status = 'ready' AND claim_lock IS NULL "
|
||||
"ORDER BY priority DESC, created_at ASC"
|
||||
).fetchall()
|
||||
# Honour kanban.max_in_progress: if the board already has enough running
|
||||
# tasks, skip spawning this tick so slow workers (local LLMs,
|
||||
# resource-constrained hosts) can finish what they have before more tasks
|
||||
# pile up and time out.
|
||||
if max_in_progress is not None and ready_rows:
|
||||
in_progress = conn.execute(
|
||||
"SELECT COUNT(*) FROM tasks WHERE status = 'running'"
|
||||
).fetchone()[0]
|
||||
if in_progress >= max_in_progress:
|
||||
return result
|
||||
# Only spawn enough to reach the cap, respecting max_spawn too.
|
||||
remaining = max_in_progress - in_progress
|
||||
if max_spawn is None or max_spawn > remaining:
|
||||
max_spawn = remaining
|
||||
spawned = 0
|
||||
# Per-profile concurrency cap (#21582): when set, track how many
|
||||
# workers each assignee already has in flight, and refuse to spawn
|
||||
|
|
@ -9062,8 +9180,8 @@ def _dispatch_once_locked(
|
|||
# ---- review column dispatch ----
|
||||
# Review tasks are tasks that a worker moved to 'review' after
|
||||
# creating a PR. The dispatcher spawns a review agent (loading
|
||||
# sdlc-review skill) that verifies the PR and either merges (→ done)
|
||||
# or rejects (→ back to running for the worker to fix).
|
||||
# sdlc-review skill) that verifies the candidate and either approves
|
||||
# (→ done) or requests changes (→ ready/todo for the implementer).
|
||||
#
|
||||
# Same concurrency model as ready dispatch: review spawns count
|
||||
# against max_spawn alongside ready tasks, so the total number of
|
||||
|
|
@ -9092,8 +9210,29 @@ def _dispatch_once_locked(
|
|||
if profile_exists is not None and not profile_exists(row["assignee"]):
|
||||
result.skipped_nonspawnable.append(row["id"])
|
||||
continue
|
||||
if _per_profile_cap is not None:
|
||||
current = _per_profile_running.get(row["assignee"], 0)
|
||||
if current >= _per_profile_cap:
|
||||
result.skipped_per_profile_capped.append(
|
||||
(row["id"], row["assignee"], current)
|
||||
)
|
||||
continue
|
||||
guard_reason = check_respawn_guard(conn, row["id"])
|
||||
if guard_reason is not None:
|
||||
result.respawn_guarded.append((row["id"], guard_reason))
|
||||
if not dry_run:
|
||||
with write_txn(conn):
|
||||
_append_event(
|
||||
conn, row["id"], "respawn_guarded",
|
||||
{"reason": guard_reason},
|
||||
)
|
||||
continue
|
||||
if dry_run:
|
||||
result.spawned.append((row["id"], row["assignee"], ""))
|
||||
if _per_profile_cap is not None:
|
||||
_per_profile_running[row["assignee"]] = (
|
||||
_per_profile_running.get(row["assignee"], 0) + 1
|
||||
)
|
||||
continue
|
||||
claimed = claim_review_task(conn, row["id"], ttl_seconds=ttl_seconds)
|
||||
if claimed is None:
|
||||
|
|
@ -9140,6 +9279,10 @@ def _dispatch_once_locked(
|
|||
_set_worker_pid(conn, claimed.id, int(pid))
|
||||
result.spawned.append((claimed.id, claimed.assignee or "", str(workspace)))
|
||||
spawned += 1
|
||||
if _per_profile_cap is not None and claimed.assignee:
|
||||
_per_profile_running[claimed.assignee] = (
|
||||
_per_profile_running.get(claimed.assignee, 0) + 1
|
||||
)
|
||||
except Exception as exc:
|
||||
auto = _record_spawn_failure(
|
||||
conn, claimed.id, str(exc),
|
||||
|
|
|
|||
|
|
@ -365,12 +365,79 @@ def test_review_dispatch_preserves_task_skills_and_adds_reviewer_skill(
|
|||
summary="ready",
|
||||
expected_run_id=implementation.current_run_id,
|
||||
)
|
||||
monkeypatch.setattr(
|
||||
kb,
|
||||
"check_respawn_guard",
|
||||
lambda _conn, _task_id: "rate_limit_cooldown",
|
||||
)
|
||||
guarded = kb.dispatch_once(conn, spawn_fn=spawn)
|
||||
assert guarded.respawn_guarded == [(task_id, "rate_limit_cooldown")]
|
||||
assert not guarded.spawned
|
||||
guarded_task = kb.get_task(conn, task_id)
|
||||
assert guarded_task is not None
|
||||
assert guarded_task.status == "review"
|
||||
|
||||
monkeypatch.setattr(kb, "check_respawn_guard", lambda _conn, _task_id: None)
|
||||
result = kb.dispatch_once(conn, spawn_fn=spawn)
|
||||
|
||||
assert task_id in [task[0] for task in result.spawned]
|
||||
assert captured == [["domain-specific-review", "sdlc-review"]]
|
||||
|
||||
|
||||
def test_review_dispatch_honors_global_and_per_profile_caps(
|
||||
kanban_home: Path,
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
import hermes_cli.config as cfgmod
|
||||
import hermes_cli.profiles as profmod
|
||||
|
||||
monkeypatch.setattr(profmod, "profile_exists", lambda _name: True)
|
||||
monkeypatch.setattr(
|
||||
cfgmod,
|
||||
"load_config",
|
||||
lambda *args, **kwargs: {"kanban": {"review_dispatch": True}},
|
||||
)
|
||||
|
||||
with kb.connect() as conn:
|
||||
running_id = kb.create_task(conn, title="already running", assignee="builder")
|
||||
assert kb.claim_task(conn, running_id) is not None
|
||||
|
||||
review_ids: list[str] = []
|
||||
for title in ("review one", "review two"):
|
||||
task_id = kb.create_task(conn, title=title, assignee="reviewer")
|
||||
implementation = kb.claim_task(conn, task_id)
|
||||
assert implementation is not None
|
||||
assert kb.request_review(
|
||||
conn,
|
||||
task_id,
|
||||
summary="ready",
|
||||
expected_run_id=implementation.current_run_id,
|
||||
)
|
||||
review_ids.append(task_id)
|
||||
|
||||
globally_capped = kb.dispatch_once(
|
||||
conn,
|
||||
dry_run=True,
|
||||
max_in_progress=1,
|
||||
)
|
||||
assert not [
|
||||
task for task in globally_capped.spawned if task[0] in review_ids
|
||||
]
|
||||
|
||||
per_profile_capped = kb.dispatch_once(
|
||||
conn,
|
||||
dry_run=True,
|
||||
max_in_progress=10,
|
||||
max_in_progress_per_profile=1,
|
||||
)
|
||||
spawned_reviews = [
|
||||
task for task in per_profile_capped.spawned if task[0] in review_ids
|
||||
]
|
||||
assert len(spawned_reviews) == 1
|
||||
assert len(per_profile_capped.skipped_per_profile_capped) == 1
|
||||
assert per_profile_capped.skipped_per_profile_capped[0][0] in review_ids
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# reopen: a follow-up sends a review task back out for another pass
|
||||
# ---------------------------------------------------------------------------
|
||||
|
|
|
|||
|
|
@ -11,6 +11,7 @@ These tests cover the two review models that must coexist:
|
|||
|
||||
from __future__ import annotations
|
||||
|
||||
import time
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
|
@ -36,6 +37,37 @@ def _run(runs, outcome: str):
|
|||
return [run for run in runs if run.outcome == outcome][-1]
|
||||
|
||||
|
||||
def _claimed_review(
|
||||
conn,
|
||||
title: str,
|
||||
*,
|
||||
ttl_seconds: int | None = None,
|
||||
max_runtime_seconds: int | None = None,
|
||||
):
|
||||
task_id = kb.create_task(
|
||||
conn,
|
||||
title=title,
|
||||
assignee="builder",
|
||||
max_runtime_seconds=max_runtime_seconds,
|
||||
)
|
||||
implementation = kb.claim_task(conn, task_id, claimer="builder:test")
|
||||
assert implementation is not None
|
||||
assert kb.request_review(
|
||||
conn,
|
||||
task_id,
|
||||
summary="ready for independent review",
|
||||
reviewer="reviewer",
|
||||
expected_run_id=implementation.current_run_id,
|
||||
)
|
||||
review = kb.claim_review_task(
|
||||
conn,
|
||||
task_id,
|
||||
ttl_seconds=ttl_seconds,
|
||||
)
|
||||
assert review is not None
|
||||
return task_id, review
|
||||
|
||||
|
||||
def test_same_card_review_supports_changes_and_approval_without_block_loop(conn):
|
||||
task_id = kb.create_task(conn, title="Implement guarded export", assignee="builder")
|
||||
implementation = kb.claim_task(conn, task_id, claimer="builder:1")
|
||||
|
|
@ -180,6 +212,178 @@ def test_request_changes_fails_closed_on_malformed_review_provenance(conn):
|
|||
assert task.current_run_id == review.current_run_id
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
"reclaim_kind",
|
||||
["spawn_failure", "expired_claim", "manual_reclaim", "stale_heartbeat"],
|
||||
)
|
||||
def test_interrupted_review_runs_retry_in_review_phase(
|
||||
conn,
|
||||
reclaim_kind: str,
|
||||
) -> None:
|
||||
task_id, review = _claimed_review(
|
||||
conn,
|
||||
f"Retry review after {reclaim_kind}",
|
||||
ttl_seconds=-1 if reclaim_kind == "expired_claim" else None,
|
||||
)
|
||||
|
||||
if reclaim_kind == "spawn_failure":
|
||||
assert not kb._record_spawn_failure(
|
||||
conn,
|
||||
task_id,
|
||||
"reviewer process failed to spawn",
|
||||
failure_limit=3,
|
||||
)
|
||||
elif reclaim_kind == "expired_claim":
|
||||
with kb.write_txn(conn):
|
||||
conn.execute(
|
||||
"UPDATE tasks SET claim_expires = ? WHERE id = ?",
|
||||
(int(time.time()) - 1, task_id),
|
||||
)
|
||||
assert kb.release_stale_claims(conn) == 1
|
||||
elif reclaim_kind == "manual_reclaim":
|
||||
assert kb.reclaim_task(conn, task_id, reason="operator retry")
|
||||
else:
|
||||
old = int(time.time()) - 1_000
|
||||
with kb.write_txn(conn):
|
||||
conn.execute(
|
||||
"UPDATE tasks SET started_at = ?, last_heartbeat_at = NULL "
|
||||
"WHERE id = ?",
|
||||
(old, task_id),
|
||||
)
|
||||
conn.execute(
|
||||
"UPDATE task_runs SET started_at = ? WHERE id = ?",
|
||||
(old, review.current_run_id),
|
||||
)
|
||||
assert kb.detect_stale_running(conn, stale_timeout_seconds=1) == [task_id]
|
||||
|
||||
retried = kb.get_task(conn, task_id)
|
||||
assert retried is not None
|
||||
assert retried.status == "review"
|
||||
assert retried.current_run_id is None
|
||||
event = kb.list_events(conn, task_id=task_id)[-1]
|
||||
assert event.payload is not None
|
||||
assert event.payload.get("retry_status") == "review"
|
||||
|
||||
|
||||
def test_review_retry_still_trips_the_failure_breaker(conn) -> None:
|
||||
task_id, _review = _claimed_review(conn, "Reviewer repeatedly fails")
|
||||
assert kb._record_spawn_failure(
|
||||
conn,
|
||||
task_id,
|
||||
"reviewer cannot start",
|
||||
failure_limit=1,
|
||||
)
|
||||
blocked = kb.get_task(conn, task_id)
|
||||
assert blocked is not None
|
||||
assert blocked.status == "blocked"
|
||||
gave_up = _event(kb.list_events(conn, task_id), "gave_up")
|
||||
assert gave_up.payload is not None
|
||||
assert gave_up.payload["retry_status"] == "review"
|
||||
assert kb.unblock_task(conn, task_id)
|
||||
unblocked = kb.get_task(conn, task_id)
|
||||
assert unblocked is not None
|
||||
assert unblocked.status == "review"
|
||||
|
||||
|
||||
def test_review_escalation_unblocks_back_to_review(conn) -> None:
|
||||
task_id, review = _claimed_review(conn, "External review escalation")
|
||||
assert kb.block_task(
|
||||
conn,
|
||||
task_id,
|
||||
reason="needs_input: maintainer decision required",
|
||||
kind="needs_input",
|
||||
expected_run_id=review.current_run_id,
|
||||
)
|
||||
blocked_event = _event(kb.list_events(conn, task_id), "blocked")
|
||||
assert blocked_event.payload is not None
|
||||
assert blocked_event.payload["source_status"] == "review"
|
||||
assert kb.unblock_task(conn, task_id)
|
||||
resumed = kb.get_task(conn, task_id)
|
||||
assert resumed is not None
|
||||
assert resumed.status == "review"
|
||||
|
||||
|
||||
def test_review_dependency_wait_reenters_review_after_parent_finishes(conn) -> None:
|
||||
parent_id = kb.create_task(conn, title="Parent", assignee="planner")
|
||||
assert kb.complete_task(conn, parent_id)
|
||||
task_id = kb.create_task(
|
||||
conn,
|
||||
title="Review after dependency refresh",
|
||||
assignee="builder",
|
||||
parents=[parent_id],
|
||||
)
|
||||
implementation = kb.claim_task(conn, task_id)
|
||||
assert implementation is not None
|
||||
assert kb.request_review(
|
||||
conn,
|
||||
task_id,
|
||||
summary="ready",
|
||||
reviewer="reviewer",
|
||||
expected_run_id=implementation.current_run_id,
|
||||
)
|
||||
review = kb.claim_review_task(conn, task_id)
|
||||
assert review is not None
|
||||
with kb.write_txn(conn):
|
||||
conn.execute("UPDATE tasks SET status = 'ready' WHERE id = ?", (parent_id,))
|
||||
assert kb.block_task(
|
||||
conn,
|
||||
task_id,
|
||||
reason="dependency: parent contract is being refreshed",
|
||||
kind="dependency",
|
||||
expected_run_id=review.current_run_id,
|
||||
)
|
||||
waiting = kb.get_task(conn, task_id)
|
||||
assert waiting is not None
|
||||
assert waiting.status == "todo"
|
||||
assert kb.complete_task(conn, parent_id)
|
||||
resumed = kb.get_task(conn, task_id)
|
||||
assert resumed is not None
|
||||
assert resumed.status == "review"
|
||||
|
||||
|
||||
def test_crashed_and_timed_out_review_runs_retry_in_review_phase(
|
||||
conn,
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
monkeypatch.setattr(kb, "_pid_alive", lambda _pid: False)
|
||||
monkeypatch.setattr(kb, "_classify_worker_exit", lambda _pid: ("nonzero_exit", 1))
|
||||
old = int(time.time()) - 1_000
|
||||
|
||||
timed_out_id, timed_out_run = _claimed_review(
|
||||
conn,
|
||||
"Timeout during review",
|
||||
max_runtime_seconds=1,
|
||||
)
|
||||
with kb.write_txn(conn):
|
||||
conn.execute(
|
||||
"UPDATE tasks SET worker_pid = ?, started_at = ? WHERE id = ?",
|
||||
(999_998, old, timed_out_id),
|
||||
)
|
||||
conn.execute(
|
||||
"UPDATE task_runs SET worker_pid = ?, started_at = ? WHERE id = ?",
|
||||
(999_998, old, timed_out_run.current_run_id),
|
||||
)
|
||||
assert timed_out_id in kb.enforce_max_runtime(conn, signal_fn=lambda *_: None)
|
||||
timed_out = kb.get_task(conn, timed_out_id)
|
||||
assert timed_out is not None
|
||||
assert timed_out.status == "review"
|
||||
|
||||
crashed_id, crashed_run = _claimed_review(conn, "Crash during review")
|
||||
with kb.write_txn(conn):
|
||||
conn.execute(
|
||||
"UPDATE tasks SET worker_pid = ?, started_at = ? WHERE id = ?",
|
||||
(999_999, old, crashed_id),
|
||||
)
|
||||
conn.execute(
|
||||
"UPDATE task_runs SET worker_pid = ?, started_at = ? WHERE id = ?",
|
||||
(999_999, old, crashed_run.current_run_id),
|
||||
)
|
||||
assert crashed_id in kb.detect_crashed_workers(conn)
|
||||
crashed = kb.get_task(conn, crashed_id)
|
||||
assert crashed is not None
|
||||
assert crashed.status == "review"
|
||||
|
||||
|
||||
def test_legacy_review_child_deadlock_is_reported_immediately(conn):
|
||||
implementation_id = kb.create_task(
|
||||
conn,
|
||||
|
|
|
|||
Loading…
Reference in New Issue