From 46ea12f1a13c0cd88dde10d2cf57f2fe90e95303 Mon Sep 17 00:00:00 2001 From: can1357 Date: Sun, 14 Jun 2026 09:15:18 +0200 Subject: [PATCH] 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. --- python/robomp/src/config.py | 46 +++++ python/robomp/src/db.py | 37 +++- python/robomp/src/queue.py | 21 +- python/robomp/tests/test_queue_cancel.py | 2 + python/robomp/tests/test_queue_shutdown.py | 1 + python/robomp/tests/test_retry.py | 218 +++++++++++++++++++++ 6 files changed, 318 insertions(+), 7 deletions(-) create mode 100644 python/robomp/tests/test_retry.py diff --git a/python/robomp/src/config.py b/python/robomp/src/config.py index eb7179a4b..2c8f1494e 100644 --- a/python/robomp/src/config.py +++ b/python/robomp/src/config.py @@ -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.""" diff --git a/python/robomp/src/db.py b/python/robomp/src/db.py index 06906b59c..008fef75a 100644 --- a/python/robomp/src/db.py +++ b/python/robomp/src/db.py @@ -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, diff --git a/python/robomp/src/queue.py b/python/robomp/src/queue.py index 96ff5c1c2..49db480e3 100644 --- a/python/robomp/src/queue.py +++ b/python/robomp/src/queue.py @@ -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) diff --git a/python/robomp/tests/test_queue_cancel.py b/python/robomp/tests/test_queue_cancel.py index 98cc77a03..ba2e402bc 100644 --- a/python/robomp/tests/test_queue_cancel.py +++ b/python/robomp/tests/test_queue_cancel.py @@ -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", diff --git a/python/robomp/tests/test_queue_shutdown.py b/python/robomp/tests/test_queue_shutdown.py index 87497e863..c67dd594b 100644 --- a/python/robomp/tests/test_queue_shutdown.py +++ b/python/robomp/tests/test_queue_shutdown.py @@ -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. diff --git a/python/robomp/tests/test_retry.py b/python/robomp/tests/test_retry.py new file mode 100644 index 000000000..a795bb71d --- /dev/null +++ b/python/robomp/tests/test_retry.py @@ -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