feat(python/robomp): added automatic event retry scheduling with backoff delays
- Added event retry settings with parsed delay schedules and jittered delay computation. - Extended event persistence to persist an `available_at` timestamp, honor it during dequeuing, and clear it when re-queuing. - Updated worker failure handling to queue bounded retries with backoff and transition to failed only when the retry budget is exhausted.
This commit is contained in:
@@ -67,6 +67,17 @@ class Settings(BaseSettings):
|
||||
task_timeout_seconds: float = Field(2400.0, alias="ROBOMP_TASK_TIMEOUT_SECONDS")
|
||||
task_timeout_hard_grace_seconds: float = Field(60.0, alias="ROBOMP_TASK_TIMEOUT_HARD_GRACE_SECONDS")
|
||||
request_timeout_seconds: float = Field(120.0, alias="ROBOMP_REQUEST_TIMEOUT_SECONDS")
|
||||
|
||||
# Automatic retry of transiently-failed events. When an event handler
|
||||
# raises (and it isn't an operator cancel or a shutdown interrupt), the
|
||||
# dispatcher re-queues the delivery with escalating backoff instead of
|
||||
# giving up, so ephemeral failures (git fetch timeouts, upstream 5xx/429,
|
||||
# flaky RPC startup) self-heal. After `event_max_retries` retries the row
|
||||
# stays `failed`. `event_retry_delays_seconds` is a comma-separated backoff
|
||||
# schedule: the Nth retry waits the Nth value (last value repeats), jittered.
|
||||
# Set `event_max_retries=0` to restore fail-fast behavior.
|
||||
event_max_retries: int = Field(3, alias="ROBOMP_EVENT_MAX_RETRIES")
|
||||
event_retry_delays_raw: str = Field("30,120,600", alias="ROBOMP_EVENT_RETRY_DELAYS_SECONDS")
|
||||
# Premature-end reminder. When a `triage_issue` turn ends without the
|
||||
# agent having reached a terminal tool (`gh_open_pr`,
|
||||
# `mark_unable_to_reproduce`, `abort_task`) for a `bug`/`documentation`
|
||||
@@ -297,6 +308,41 @@ class Settings(BaseSettings):
|
||||
"""Random selection from the pool (uniform). One-element pools return that one."""
|
||||
return random.choice(self.model_pool)
|
||||
|
||||
@field_validator("event_retry_delays_raw", mode="before")
|
||||
@classmethod
|
||||
def _coerce_retry_delays(cls, v: object) -> str:
|
||||
if v is None:
|
||||
return ""
|
||||
if isinstance(v, (list, tuple)):
|
||||
return ",".join(str(item) for item in v)
|
||||
return str(v)
|
||||
|
||||
@property
|
||||
def event_retry_delays(self) -> tuple[float, ...]:
|
||||
"""Parsed backoff schedule in seconds; always non-empty."""
|
||||
vals: list[float] = []
|
||||
for piece in self.event_retry_delays_raw.split(","):
|
||||
piece = piece.strip()
|
||||
if not piece:
|
||||
continue
|
||||
try:
|
||||
seconds = float(piece)
|
||||
except ValueError:
|
||||
continue
|
||||
if seconds >= 0:
|
||||
vals.append(seconds)
|
||||
return tuple(vals) or (30.0,)
|
||||
|
||||
def retry_delay_seconds(self, retry_index: int) -> float:
|
||||
"""Backoff before the `retry_index`-th retry (1-based), with jitter.
|
||||
|
||||
Clamps to the last configured delay; applies ±20% jitter so a
|
||||
fleet-wide outage doesn't replay every event in lockstep.
|
||||
"""
|
||||
delays = self.event_retry_delays
|
||||
idx = min(max(retry_index, 1), len(delays)) - 1
|
||||
return delays[idx] * (0.8 + random.random() * 0.4)
|
||||
|
||||
@property
|
||||
def resolved_author_name(self) -> str:
|
||||
"""Falls back to bot_login if ROBOMP_GIT_AUTHOR_NAME isn't set."""
|
||||
|
||||
+32
-5
@@ -119,6 +119,11 @@ def _utcnow() -> str:
|
||||
return datetime.now(UTC).strftime("%Y-%m-%dT%H:%M:%S.%fZ")
|
||||
|
||||
|
||||
def _utc_after(seconds: float) -> str:
|
||||
"""UTC timestamp `seconds` in the future, same sortable format as `_utcnow`."""
|
||||
return (datetime.now(UTC) + timedelta(seconds=max(seconds, 0.0))).strftime("%Y-%m-%dT%H:%M:%S.%fZ")
|
||||
|
||||
|
||||
def iso_seconds_ago(seconds: float) -> str:
|
||||
"""ISO-UTC timestamp for `seconds` ago, matching the format `_utcnow` writes."""
|
||||
return (datetime.now(UTC) - timedelta(seconds=seconds)).strftime("%Y-%m-%dT%H:%M:%S.%fZ")
|
||||
@@ -241,6 +246,8 @@ class Database:
|
||||
event_cols = {row[1] for row in self._conn.execute("PRAGMA table_info(events)").fetchall()}
|
||||
if "model" not in event_cols:
|
||||
self._conn.execute("ALTER TABLE events ADD COLUMN model TEXT")
|
||||
if "available_at" not in event_cols:
|
||||
self._conn.execute("ALTER TABLE events ADD COLUMN available_at TEXT")
|
||||
|
||||
def close(self) -> None:
|
||||
with self._lock:
|
||||
@@ -298,6 +305,7 @@ class Database:
|
||||
def claim_next_event(self) -> EventRow | None:
|
||||
"""Atomically dequeue one unblocked queued event into running state."""
|
||||
with self._txn() as conn:
|
||||
now = _utcnow()
|
||||
row = conn.execute(
|
||||
"""
|
||||
SELECT queued.delivery_id, queued.event_type, queued.repo, queued.issue_key,
|
||||
@@ -305,6 +313,7 @@ class Database:
|
||||
queued.last_error
|
||||
FROM events AS queued
|
||||
WHERE queued.state = 'queued'
|
||||
AND (queued.available_at IS NULL OR queued.available_at <= ?)
|
||||
AND (
|
||||
queued.issue_key IS NULL
|
||||
OR NOT EXISTS (
|
||||
@@ -316,11 +325,11 @@ class Database:
|
||||
)
|
||||
ORDER BY queued.received_at
|
||||
LIMIT 1
|
||||
"""
|
||||
""",
|
||||
(now,),
|
||||
).fetchone()
|
||||
if row is None:
|
||||
return None
|
||||
now = _utcnow()
|
||||
conn.execute(
|
||||
"UPDATE events SET state='running', attempts=attempts+1, started_at=? WHERE delivery_id=?",
|
||||
(now, row["delivery_id"]),
|
||||
@@ -360,7 +369,7 @@ class Database:
|
||||
"""Recover events that were running at shutdown."""
|
||||
with self._lock:
|
||||
cur = self._conn.execute(
|
||||
"UPDATE events SET state='queued' WHERE state='running'",
|
||||
"UPDATE events SET state='queued', available_at=NULL WHERE state='running'",
|
||||
)
|
||||
return cur.rowcount
|
||||
|
||||
@@ -613,7 +622,7 @@ class Database:
|
||||
with self._lock:
|
||||
if from_states is None:
|
||||
cur = self._conn.execute(
|
||||
"UPDATE events SET state='queued' WHERE delivery_id=?",
|
||||
"UPDATE events SET state='queued', available_at=NULL WHERE delivery_id=?",
|
||||
(delivery_id,),
|
||||
)
|
||||
elif not from_states:
|
||||
@@ -621,11 +630,29 @@ class Database:
|
||||
else:
|
||||
placeholders = ",".join("?" for _ in from_states)
|
||||
cur = self._conn.execute(
|
||||
f"UPDATE events SET state='queued' WHERE delivery_id=? AND state IN ({placeholders})",
|
||||
f"UPDATE events SET state='queued', available_at=NULL WHERE delivery_id=? AND state IN ({placeholders})",
|
||||
(delivery_id, *from_states),
|
||||
)
|
||||
return cur.rowcount > 0
|
||||
|
||||
def schedule_retry(self, delivery_id: str, *, delay_seconds: float, error: str | None = None) -> bool:
|
||||
"""Re-queue a delivery for a future retry with backoff.
|
||||
|
||||
Flips state back to 'queued' but stamps `available_at` so
|
||||
`claim_next_event` skips the row until the backoff elapses. `attempts`
|
||||
is left untouched (it was already incremented at claim) so the retry
|
||||
budget keeps counting down; `last_error` retains the failure reason for
|
||||
the dashboard. Only transitions a 'running'/'failed' row; returns
|
||||
whether a row changed.
|
||||
"""
|
||||
with self._lock:
|
||||
cur = self._conn.execute(
|
||||
"UPDATE events SET state='queued', last_error=?, available_at=?, finished_at=NULL "
|
||||
"WHERE delivery_id=? AND state IN ('running','failed')",
|
||||
(error, _utc_after(delay_seconds), delivery_id),
|
||||
)
|
||||
return cur.rowcount > 0
|
||||
|
||||
# ---- issues ----
|
||||
def upsert_issue(
|
||||
self,
|
||||
|
||||
@@ -301,8 +301,25 @@ class WorkerPool:
|
||||
self.db.mark_event(row.delivery_id, "failed", error="cancelled by operator")
|
||||
else:
|
||||
tb = traceback.format_exc(limit=20)
|
||||
log.exception("event handler failed", extra={"delivery": row.delivery_id})
|
||||
self.db.mark_event(row.delivery_id, "failed", error=f"{exc}\n{tb}")
|
||||
err = f"{exc}\n{tb}"
|
||||
max_retries = self.settings.event_max_retries
|
||||
delay = self.settings.retry_delay_seconds(row.attempts)
|
||||
if 0 < row.attempts <= max_retries and self.db.schedule_retry(
|
||||
row.delivery_id, delay_seconds=delay, error=err
|
||||
):
|
||||
log.warning(
|
||||
"event retry scheduled",
|
||||
extra={
|
||||
"delivery": row.delivery_id,
|
||||
"key": row.issue_key,
|
||||
"attempt": row.attempts,
|
||||
"max_retries": max_retries,
|
||||
"retry_in_seconds": round(delay, 1),
|
||||
},
|
||||
)
|
||||
else:
|
||||
log.exception("event handler failed", extra={"delivery": row.delivery_id})
|
||||
self.db.mark_event(row.delivery_id, "failed", error=err)
|
||||
finally:
|
||||
self._cancelled.discard(row.delivery_id)
|
||||
self._shutdown_cancelled.discard(row.delivery_id)
|
||||
|
||||
@@ -161,6 +161,7 @@ async def test_non_cancelled_failure_keeps_real_traceback(
|
||||
) -> None:
|
||||
"""A garden-variety dispatch failure still records the traceback path."""
|
||||
pool = _make_pool(settings, db)
|
||||
monkeypatch.setattr(settings, "event_max_retries", 0) # assert terminal failure, not retry
|
||||
db.record_event(
|
||||
delivery_id="d4",
|
||||
event_type="issues",
|
||||
@@ -191,6 +192,7 @@ async def test_run_event_marks_failed_when_not_shutting_down(
|
||||
) -> None:
|
||||
"""When `_shutting_down` is False, a dispatch failure still marks the row failed."""
|
||||
pool = _make_pool(settings, db)
|
||||
monkeypatch.setattr(settings, "event_max_retries", 0) # assert terminal failure, not retry
|
||||
assert pool._shutting_down is False # noqa: SLF001
|
||||
db.record_event(
|
||||
delivery_id="d5",
|
||||
|
||||
@@ -256,6 +256,7 @@ async def test_run_event_marks_failed_for_unrelated_failure_during_drain(
|
||||
"""
|
||||
pool = _make_pool(settings, db)
|
||||
pool._shutting_down = True # noqa: SLF001
|
||||
monkeypatch.setattr(settings, "event_max_retries", 0) # assert terminal failure, not retry
|
||||
# Crucially: this delivery is NOT in `_shutdown_cancelled` — stop()
|
||||
# never targeted it. Its failure is its own.
|
||||
|
||||
|
||||
@@ -0,0 +1,218 @@
|
||||
"""Automatic backoff-retry of transiently-failed events.
|
||||
|
||||
Covers the three layers of the feature:
|
||||
- `Settings`: backoff schedule parsing + per-retry delay (escalation/clamp).
|
||||
- `Database`: `schedule_retry` re-queues with an `available_at` gate that
|
||||
`claim_next_event` honors, and a manual requeue clears that gate.
|
||||
- `WorkerPool._run_event`: a raising handler is retried up to the budget,
|
||||
then marked `failed`.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import pytest
|
||||
|
||||
from robomp import config as config_mod
|
||||
from robomp.config import Settings, reset_settings_cache
|
||||
from robomp.db import Database, issue_key
|
||||
from robomp.queue import WorkerPool
|
||||
from robomp.slot_pool import SlotPool
|
||||
|
||||
|
||||
def _record(db: Database, delivery: str = "d1") -> None:
|
||||
db.record_event(
|
||||
delivery_id=delivery,
|
||||
event_type="issues",
|
||||
repo="octo/widget",
|
||||
issue_key=issue_key("octo/widget", 1),
|
||||
payload={"action": "opened"},
|
||||
)
|
||||
|
||||
|
||||
# ---- Settings: backoff schedule ----------------------------------------------
|
||||
|
||||
|
||||
def test_event_retry_delays_parsing_skips_garbage(env: dict[str, str], monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
monkeypatch.setenv("ROBOMP_EVENT_RETRY_DELAYS_SECONDS", "30, 120 ,600,,abc,-5")
|
||||
reset_settings_cache()
|
||||
cfg = Settings() # type: ignore[call-arg]
|
||||
assert cfg.event_retry_delays == (30.0, 120.0, 600.0)
|
||||
|
||||
|
||||
def test_event_retry_delays_defaults_when_empty(env: dict[str, str], monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
monkeypatch.setenv("ROBOMP_EVENT_RETRY_DELAYS_SECONDS", " ")
|
||||
reset_settings_cache()
|
||||
cfg = Settings() # type: ignore[call-arg]
|
||||
assert cfg.event_retry_delays == (30.0,)
|
||||
|
||||
|
||||
def test_retry_delay_escalates_and_clamps(env: dict[str, str], monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
monkeypatch.setenv("ROBOMP_EVENT_RETRY_DELAYS_SECONDS", "1,2,3")
|
||||
reset_settings_cache()
|
||||
cfg = Settings() # type: ignore[call-arg]
|
||||
# Pin jitter to its midpoint (0.8 + 0.5*0.4 = 1.0) so we assert exact bases.
|
||||
monkeypatch.setattr(config_mod.random, "random", lambda: 0.5)
|
||||
assert cfg.retry_delay_seconds(1) == 1.0
|
||||
assert cfg.retry_delay_seconds(2) == 2.0
|
||||
assert cfg.retry_delay_seconds(3) == 3.0
|
||||
assert cfg.retry_delay_seconds(4) == 3.0 # clamps to the last delay
|
||||
assert cfg.retry_delay_seconds(0) == 1.0 # clamps to the first delay
|
||||
|
||||
|
||||
def test_retry_delay_jitter_stays_in_band(env: dict[str, str], monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
monkeypatch.setenv("ROBOMP_EVENT_RETRY_DELAYS_SECONDS", "100")
|
||||
reset_settings_cache()
|
||||
cfg = Settings() # type: ignore[call-arg]
|
||||
for _ in range(200):
|
||||
assert 80.0 <= cfg.retry_delay_seconds(1) <= 120.0
|
||||
|
||||
|
||||
# ---- Database: schedule_retry + claim gating ---------------------------------
|
||||
|
||||
|
||||
def test_schedule_retry_gates_claim_until_available(db: Database) -> None:
|
||||
_record(db)
|
||||
claimed = db.claim_next_event()
|
||||
assert claimed is not None and claimed.attempts == 1
|
||||
|
||||
assert db.schedule_retry("d1", delay_seconds=3600, error="ephemeral boom")
|
||||
ev = db.get_event("d1")
|
||||
assert ev is not None
|
||||
assert ev.state == "queued"
|
||||
assert ev.attempts == 1 # claim budget preserved, not reset
|
||||
assert ev.last_error == "ephemeral boom"
|
||||
|
||||
# Backed off into the future -> not yet claimable.
|
||||
assert db.claim_next_event() is None
|
||||
|
||||
|
||||
def test_schedule_retry_zero_delay_is_immediately_claimable(db: Database) -> None:
|
||||
_record(db)
|
||||
db.claim_next_event()
|
||||
assert db.schedule_retry("d1", delay_seconds=0, error="boom")
|
||||
again = db.claim_next_event()
|
||||
assert again is not None
|
||||
assert again.attempts == 2 # re-claim advances the attempt counter
|
||||
|
||||
|
||||
def test_schedule_retry_only_transitions_running_or_failed(db: Database) -> None:
|
||||
_record(db)
|
||||
# A still-queued row must not be touched by schedule_retry.
|
||||
assert not db.schedule_retry("d1", delay_seconds=0)
|
||||
assert db.get_event("d1").state == "queued"
|
||||
|
||||
# A terminally-failed row can be revived.
|
||||
db.claim_next_event()
|
||||
db.mark_event("d1", "failed", error="x")
|
||||
assert db.schedule_retry("d1", delay_seconds=3600, error="retry me")
|
||||
assert db.get_event("d1").state == "queued"
|
||||
|
||||
|
||||
def test_manual_requeue_clears_retry_backoff(db: Database) -> None:
|
||||
_record(db)
|
||||
db.claim_next_event()
|
||||
db.schedule_retry("d1", delay_seconds=3600, error="boom")
|
||||
assert db.claim_next_event() is None # still backed off
|
||||
|
||||
assert db.requeue_event("d1") # operator override
|
||||
assert db.claim_next_event() is not None # available_at cleared -> claimable
|
||||
|
||||
|
||||
# ---- WorkerPool: retry-then-exhaust through the real failure path -------------
|
||||
|
||||
|
||||
class _StubGitHub:
|
||||
pass
|
||||
|
||||
|
||||
class _StubSandbox:
|
||||
natives_cache = None
|
||||
|
||||
|
||||
class _StubGitTransport:
|
||||
pass
|
||||
|
||||
|
||||
def _retry_settings(monkeypatch: pytest.MonkeyPatch, *, max_retries: int) -> Settings:
|
||||
monkeypatch.setenv("ROBOMP_EVENT_MAX_RETRIES", str(max_retries))
|
||||
monkeypatch.setenv("ROBOMP_EVENT_RETRY_DELAYS_SECONDS", "0")
|
||||
reset_settings_cache()
|
||||
cfg = Settings() # type: ignore[call-arg]
|
||||
cfg.ensure_paths()
|
||||
return cfg
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_run_event_retries_then_marks_failed(
|
||||
env: dict[str, str], monkeypatch: pytest.MonkeyPatch, db: Database
|
||||
) -> None:
|
||||
cfg = _retry_settings(monkeypatch, max_retries=1)
|
||||
monkeypatch.setattr("robomp.queue._reap_slot", lambda uid: None)
|
||||
pool = WorkerPool(
|
||||
settings=cfg,
|
||||
db=db,
|
||||
github=_StubGitHub(), # type: ignore[arg-type]
|
||||
sandbox=_StubSandbox(), # type: ignore[arg-type]
|
||||
git_transport=_StubGitTransport(), # type: ignore[arg-type]
|
||||
slot_pool=SlotPool([2001]),
|
||||
)
|
||||
|
||||
async def boom(*_args: object, **_kwargs: object) -> None:
|
||||
raise ValueError("ephemeral boom")
|
||||
|
||||
monkeypatch.setattr(pool, "_dispatch", boom)
|
||||
_record(db)
|
||||
|
||||
# Attempt 1 fails -> scheduled for retry (queued), not failed.
|
||||
row1 = db.claim_next_event()
|
||||
assert row1 is not None and row1.attempts == 1
|
||||
await pool._run_event(row1)
|
||||
ev = db.get_event("d1")
|
||||
assert ev is not None and ev.state == "queued"
|
||||
assert "ephemeral boom" in (ev.last_error or "")
|
||||
|
||||
# Attempt 2 fails with the retry budget exhausted -> failed.
|
||||
row2 = db.claim_next_event()
|
||||
assert row2 is not None and row2.attempts == 2
|
||||
await pool._run_event(row2)
|
||||
ev = db.get_event("d1")
|
||||
assert ev is not None and ev.state == "failed"
|
||||
assert "ephemeral boom" in (ev.last_error or "")
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_run_event_success_after_transient_failure(
|
||||
env: dict[str, str], monkeypatch: pytest.MonkeyPatch, db: Database
|
||||
) -> None:
|
||||
"""A handler that fails once then succeeds ends `done`, not `failed`."""
|
||||
cfg = _retry_settings(monkeypatch, max_retries=3)
|
||||
monkeypatch.setattr("robomp.queue._reap_slot", lambda uid: None)
|
||||
pool = WorkerPool(
|
||||
settings=cfg,
|
||||
db=db,
|
||||
github=_StubGitHub(), # type: ignore[arg-type]
|
||||
sandbox=_StubSandbox(), # type: ignore[arg-type]
|
||||
git_transport=_StubGitTransport(), # type: ignore[arg-type]
|
||||
slot_pool=SlotPool([2001]),
|
||||
)
|
||||
|
||||
calls = {"n": 0}
|
||||
|
||||
async def flaky(*_args: object, **_kwargs: object) -> None:
|
||||
calls["n"] += 1
|
||||
if calls["n"] == 1:
|
||||
raise ValueError("ephemeral boom")
|
||||
|
||||
monkeypatch.setattr(pool, "_dispatch", flaky)
|
||||
_record(db)
|
||||
|
||||
row1 = db.claim_next_event()
|
||||
assert row1 is not None
|
||||
await pool._run_event(row1)
|
||||
assert db.get_event("d1").state == "queued" # retry scheduled
|
||||
|
||||
row2 = db.claim_next_event()
|
||||
assert row2 is not None
|
||||
await pool._run_event(row2)
|
||||
assert db.get_event("d1").state == "done"
|
||||
assert calls["n"] == 2
|
||||
Reference in New Issue
Block a user