From 4d494cc5123d4faefda0a5170018d5805e190def Mon Sep 17 00:00:00 2001 From: roboomp Date: Fri, 17 Jul 2026 17:33:23 +0000 Subject: [PATCH] fix(robomp): promote deferred events on periodic sweep Ran deferred-submission promotion on an independent timer in addition to the empty-queue path, so a sustained ordinary queue can no longer starve a rate-limited submitter after their rolling window frees. Added ROBOMP_DEFERRED_PROMOTION_SCAN_SECONDS to tune or disable the sweep. Fixes #5882 --- python/robomp/.env.example | 3 ++ python/robomp/src/config.py | 5 +++ python/robomp/src/queue.py | 44 ++++++++++++++++++++-- python/robomp/tests/test_queue_dispatch.py | 36 ++++++++++++++++++ 4 files changed, 84 insertions(+), 4 deletions(-) diff --git a/python/robomp/.env.example b/python/robomp/.env.example index 320bec941..a2a85b723 100644 --- a/python/robomp/.env.example +++ b/python/robomp/.env.example @@ -150,6 +150,9 @@ ROBOMP_RATE_LIMIT_WINDOW_SECONDS=3600 ROBOMP_RATE_LIMIT_DEFAULT=3 ROBOMP_RATE_LIMIT_CONTRIBUTOR=10 ROBOMP_RATE_LIMIT_UNLIMITED= +# How often (seconds) deferred events are swept back into the queue as windows +# free up, independent of queue depth. Set to 0 to disable the periodic sweep. +ROBOMP_DEFERRED_PROMOTION_SCAN_SECONDS=60 # ============================================================================= diff --git a/python/robomp/src/config.py b/python/robomp/src/config.py index 046e7c656..4040c4b4f 100644 --- a/python/robomp/src/config.py +++ b/python/robomp/src/config.py @@ -129,6 +129,11 @@ class Settings(BaseSettings): rate_limit_default: int = Field(3, alias="ROBOMP_RATE_LIMIT_DEFAULT") rate_limit_contributor: int = Field(10, alias="ROBOMP_RATE_LIMIT_CONTRIBUTOR") rate_limit_unlimited_raw: str = Field("", alias="ROBOMP_RATE_LIMIT_UNLIMITED") + # How often the dispatcher sweeps deferred rate-limited events back into the + # queue as their submitters' rolling windows free up. The empty-queue path + # promotes immediately; this periodic sweep guarantees progress even while a + # sustained ordinary queue keeps `claim_next_event` returning work. + deferred_promotion_scan_seconds: float = Field(60.0, alias="ROBOMP_DEFERRED_PROMOTION_SCAN_SECONDS") # Logins (comma-separated, `@` prefix optional, case-insensitive) whose `@bot_login` # mentions are treated as authoritative directives. These accounts also # bypass rate limiting regardless of `author_association`. diff --git a/python/robomp/src/queue.py b/python/robomp/src/queue.py index 37e24dcbe..d9d69bfd5 100644 --- a/python/robomp/src/queue.py +++ b/python/robomp/src/queue.py @@ -100,6 +100,11 @@ class WorkerPool: # restarted orchestrator doesn't burn CPU on a cold cache. if self.sandbox.natives_cache is not None and self.settings.natives_cache_gc_interval_seconds > 0: self._workers.append(asyncio.create_task(self._natives_cache_gc_loop(), name="robomp-natives-gc")) + # Periodic promotion of deferred rate-limited events. Runs independently + # of the empty-queue path so a sustained ordinary queue can't starve a + # submitter whose rolling window has since freed. + if self.settings.deferred_promotion_scan_seconds > 0: + self._workers.append(asyncio.create_task(self._deferred_promotion_loop(), name="robomp-deferred-promotion")) async def stop(self, *, drain_timeout: float = 25.0, kill_timeout: float = 5.0) -> None: """Halt the dispatcher, then drain (or kill) in-flight `_run_event` tasks. @@ -186,6 +191,40 @@ class WorkerPool: except asyncio.CancelledError: raise + async def _deferred_promotion_loop(self) -> None: + """Sweep deferred rate-limited events back into the queue on a timer. + + Sleeps the configured interval first, then re-admits any deferred event + whose submitter now has rolling-window capacity and wakes the dispatcher + so promoted rows are claimed promptly. Cancellation is the only exit; + any per-sweep failure is logged and the loop continues. + """ + interval = self.settings.deferred_promotion_scan_seconds + log.info("deferred promotion loop online", extra={"interval": interval}) + try: + while not self._stop.is_set(): + try: + await asyncio.wait_for(self._stop.wait(), timeout=interval) + return # stop was set during the wait + except TimeoutError: + pass + try: + promoted = await self._promote_deferred() + if promoted: + self.wake() + except Exception: + log.exception("deferred promotion sweep raised") + except asyncio.CancelledError: + raise + + async def _promote_deferred(self) -> int: + """Re-admit deferred events whose submitters regained window capacity.""" + since = iso_seconds_ago(self.settings.rate_limit_window_seconds) + promoted = await asyncio.to_thread(self.db.promote_deferred_submissions, since=since) + if promoted: + log.info("deferred submissions promoted", extra={"count": promoted}) + return promoted + async def _dispatch_loop(self) -> None: log.info("dispatch loop online") try: @@ -214,11 +253,8 @@ class WorkerPool: # Naive but fine for v1 (small queue). row = await asyncio.to_thread(self.db.claim_next_event) if row is None: - since = iso_seconds_ago(self.settings.rate_limit_window_seconds) - promoted = await asyncio.to_thread(self.db.promote_deferred_submissions, since=since) - if not promoted: + if not await self._promote_deferred(): return None - log.info("deferred submissions promoted", extra={"count": promoted}) row = await asyncio.to_thread(self.db.claim_next_event) if row is None: return None diff --git a/python/robomp/tests/test_queue_dispatch.py b/python/robomp/tests/test_queue_dispatch.py index 12517b95d..1034d8e62 100644 --- a/python/robomp/tests/test_queue_dispatch.py +++ b/python/robomp/tests/test_queue_dispatch.py @@ -114,3 +114,39 @@ async def test_claim_promotes_deferred_submission_after_window_frees(settings: S assert row is not None assert row.delivery_id == "deferred" assert row.state == "running" + + +@pytest.mark.asyncio +async def test_deferred_promotion_sweep_runs_while_queue_is_busy(settings: Settings, db: Database) -> None: + """A sustained ordinary queue must not starve deferred events forever. + + The empty-queue path never fires when `claim_next_event` keeps returning + work, so the independent sweep is the only thing that re-admits a deferred + submission after its rolling window frees. + """ + # Ordinary queued work for another issue keeps `claim_next_event` busy. + assert db.record_event( + delivery_id="busy", + event_type="issues", + repo="octo/widget", + issue_key="octo/widget#1", + payload={"action": "opened"}, + ) + assert db.record_submission(delivery_id="accepted", login="alice", repo="octo/widget") + assert db.defer_submission_event( + delivery_id="deferred", + event_type="issues", + login="alice", + repo="octo/widget", + issue_key="octo/widget#8", + payload={"action": "opened", "issue": {"number": 8}}, + cap=1, + reason="rate limit", + ) + settings.rate_limit_window_seconds = -1 + + pool = _make_pool(settings, db) + promoted = await pool._promote_deferred() # noqa: SLF001 + + assert promoted == 1 + assert db.get_event("deferred").state == "queued"