Revert "Merge PR #5886: fix(robomp): defer rate-limited submissions (@roboomp)"
This reverts commit81102634c8, reversing changes made to5c9b5f7b64.
This commit is contained in:
@@ -129,11 +129,6 @@ 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`.
|
||||
|
||||
+9
-184
@@ -14,7 +14,7 @@ from typing import Any, Literal
|
||||
|
||||
from robomp.github_client import IssueIndexEntry
|
||||
|
||||
EventState = Literal["queued", "deferred", "running", "done", "failed", "skipped"]
|
||||
EventState = Literal["queued", "running", "done", "failed", "skipped"]
|
||||
INACTIVE_EVENT_STATES: tuple[EventState, ...] = ("done", "failed", "skipped")
|
||||
|
||||
IssueState = Literal[
|
||||
@@ -42,7 +42,7 @@ CREATE TABLE IF NOT EXISTS events (
|
||||
payload_json TEXT NOT NULL,
|
||||
received_at TEXT NOT NULL,
|
||||
state TEXT NOT NULL
|
||||
CHECK (state IN ('queued','deferred','running','done','failed','skipped')),
|
||||
CHECK (state IN ('queued','running','done','failed','skipped')),
|
||||
attempts INTEGER NOT NULL DEFAULT 0,
|
||||
last_error TEXT,
|
||||
started_at TEXT,
|
||||
@@ -101,16 +101,6 @@ CREATE TABLE IF NOT EXISTS submissions (
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS submissions_login_ts ON submissions(login, ts);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS deferred_submissions (
|
||||
delivery_id TEXT PRIMARY KEY,
|
||||
login TEXT NOT NULL,
|
||||
repo TEXT,
|
||||
cap INTEGER NOT NULL CHECK (cap > 0),
|
||||
created_at TEXT NOT NULL
|
||||
);
|
||||
CREATE INDEX IF NOT EXISTS deferred_submissions_login_created
|
||||
ON deferred_submissions(login, created_at);
|
||||
|
||||
CREATE TABLE IF NOT EXISTS pending_closures (
|
||||
issue_key TEXT PRIMARY KEY,
|
||||
repo TEXT NOT NULL,
|
||||
@@ -309,47 +299,6 @@ class Database:
|
||||
if "available_at" not in event_cols:
|
||||
self._conn.execute("ALTER TABLE events ADD COLUMN available_at TEXT")
|
||||
|
||||
event_table = self._conn.execute(
|
||||
"SELECT sql FROM sqlite_master WHERE type='table' AND name='events'"
|
||||
).fetchone()
|
||||
event_sql = str(event_table["sql"] or "") if event_table is not None else ""
|
||||
if "'deferred'" not in event_sql:
|
||||
with self._txn() as conn:
|
||||
conn.execute(
|
||||
"""
|
||||
CREATE TABLE events_v2 (
|
||||
delivery_id TEXT PRIMARY KEY,
|
||||
event_type TEXT NOT NULL,
|
||||
repo TEXT,
|
||||
issue_key TEXT,
|
||||
payload_json TEXT NOT NULL,
|
||||
received_at TEXT NOT NULL,
|
||||
state TEXT NOT NULL
|
||||
CHECK (state IN ('queued','deferred','running','done','failed','skipped')),
|
||||
attempts INTEGER NOT NULL DEFAULT 0,
|
||||
last_error TEXT,
|
||||
started_at TEXT,
|
||||
finished_at TEXT,
|
||||
model TEXT,
|
||||
available_at TEXT
|
||||
)
|
||||
"""
|
||||
)
|
||||
conn.execute(
|
||||
"""
|
||||
INSERT INTO events_v2
|
||||
(delivery_id, event_type, repo, issue_key, payload_json, received_at,
|
||||
state, attempts, last_error, started_at, finished_at, model, available_at)
|
||||
SELECT delivery_id, event_type, repo, issue_key, payload_json, received_at,
|
||||
state, attempts, last_error, started_at, finished_at, model, available_at
|
||||
FROM events
|
||||
"""
|
||||
)
|
||||
conn.execute("DROP TABLE events")
|
||||
conn.execute("ALTER TABLE events_v2 RENAME TO events")
|
||||
conn.execute("CREATE INDEX events_state_received ON events(state, received_at)")
|
||||
conn.execute("CREATE INDEX events_issue_state ON events(issue_key, state)")
|
||||
|
||||
def close(self) -> None:
|
||||
with self._lock:
|
||||
self._conn.close()
|
||||
@@ -503,9 +452,8 @@ class Database:
|
||||
|
||||
def remove_event(self, delivery_id: str) -> None:
|
||||
"""Hard-delete an event row. Used to clear stale state before a manual re-trigger."""
|
||||
with self._txn() as conn:
|
||||
conn.execute("DELETE FROM deferred_submissions WHERE delivery_id=?", (delivery_id,))
|
||||
conn.execute("DELETE FROM events WHERE delivery_id=?", (delivery_id,))
|
||||
with self._lock:
|
||||
self._conn.execute("DELETE FROM events WHERE delivery_id=?", (delivery_id,))
|
||||
|
||||
def replace_event_if_state_in(
|
||||
self,
|
||||
@@ -528,7 +476,6 @@ class Database:
|
||||
if row is not None:
|
||||
if row["state"] not in allowed_existing_states:
|
||||
return False
|
||||
conn.execute("DELETE FROM deferred_submissions WHERE delivery_id = ?", (delivery_id,))
|
||||
conn.execute("DELETE FROM events WHERE delivery_id = ?", (delivery_id,))
|
||||
conn.execute(
|
||||
"""
|
||||
@@ -610,7 +557,7 @@ class Database:
|
||||
"""Return current row counts per event state, including states with zero rows."""
|
||||
with self._lock:
|
||||
rows = self._conn.execute("SELECT state, COUNT(*) AS n FROM events GROUP BY state").fetchall()
|
||||
counts: dict[str, int] = dict.fromkeys(("queued", "deferred", "running", "done", "failed", "skipped"), 0)
|
||||
counts: dict[str, int] = dict.fromkeys(("queued", "running", "done", "failed", "skipped"), 0)
|
||||
for row in rows:
|
||||
counts[row["state"]] = int(row["n"])
|
||||
return counts
|
||||
@@ -622,7 +569,7 @@ class Database:
|
||||
run clears an older failure for that issue, and ignored webhook noise
|
||||
does not make a failed issue look skipped.
|
||||
"""
|
||||
counts: dict[str, int] = dict.fromkeys(("queued", "deferred", "running", "done", "failed", "skipped"), 0)
|
||||
counts: dict[str, int] = dict.fromkeys(("queued", "running", "done", "failed", "skipped"), 0)
|
||||
seen: set[str] = set()
|
||||
with self._lock:
|
||||
rows = self._conn.execute(
|
||||
@@ -742,9 +689,9 @@ class Database:
|
||||
keep public retries from mutating queued/running rows while preserving
|
||||
internal recovery of a just-claimed running event.
|
||||
"""
|
||||
with self._txn() as conn:
|
||||
with self._lock:
|
||||
if from_states is None:
|
||||
cur = conn.execute(
|
||||
cur = self._conn.execute(
|
||||
"UPDATE events SET state='queued', available_at=NULL WHERE delivery_id=?",
|
||||
(delivery_id,),
|
||||
)
|
||||
@@ -752,12 +699,10 @@ class Database:
|
||||
return False
|
||||
else:
|
||||
placeholders = ",".join("?" for _ in from_states)
|
||||
cur = conn.execute(
|
||||
cur = self._conn.execute(
|
||||
f"UPDATE events SET state='queued', available_at=NULL WHERE delivery_id=? AND state IN ({placeholders})",
|
||||
(delivery_id, *from_states),
|
||||
)
|
||||
if cur.rowcount > 0:
|
||||
conn.execute("DELETE FROM deferred_submissions WHERE delivery_id=?", (delivery_id,))
|
||||
return cur.rowcount > 0
|
||||
|
||||
def schedule_retry(self, delivery_id: str, *, delay_seconds: float, error: str | None = None) -> bool:
|
||||
@@ -1092,22 +1037,6 @@ class Database:
|
||||
used=int(row["n"]) if row is not None else 0,
|
||||
)
|
||||
|
||||
if cap is not None:
|
||||
deferred = conn.execute(
|
||||
"SELECT 1 FROM deferred_submissions WHERE login=? LIMIT 1",
|
||||
(normalized_login,),
|
||||
).fetchone()
|
||||
if deferred is not None:
|
||||
row = conn.execute(
|
||||
"SELECT COUNT(*) AS n FROM submissions WHERE login=? AND ts>=?",
|
||||
(normalized_login, since),
|
||||
).fetchone()
|
||||
return SubmissionAdmission(
|
||||
accepted=False,
|
||||
duplicate=False,
|
||||
used=int(row["n"]) if row is not None else 0,
|
||||
)
|
||||
|
||||
row = conn.execute(
|
||||
"SELECT COUNT(*) AS n FROM submissions WHERE login=? AND ts>=?",
|
||||
(normalized_login, since),
|
||||
@@ -1122,110 +1051,6 @@ class Database:
|
||||
)
|
||||
return SubmissionAdmission(accepted=True, duplicate=False, used=used + 1)
|
||||
|
||||
def defer_submission_event(
|
||||
self,
|
||||
*,
|
||||
delivery_id: str,
|
||||
event_type: str,
|
||||
login: str,
|
||||
repo: str | None,
|
||||
issue_key: str | None,
|
||||
payload: Mapping[str, Any],
|
||||
cap: int,
|
||||
reason: str,
|
||||
) -> bool:
|
||||
"""Persist bounded rate-limit overflow for later admission.
|
||||
|
||||
Returns True when the event is deferred, including idempotent redelivery
|
||||
of an existing deferred event. Overflow beyond one cap-sized backlog is
|
||||
retained as an ordinary skipped event for operator diagnosis.
|
||||
"""
|
||||
normalized_login = login.lower()
|
||||
now = _utcnow()
|
||||
with self._txn() as conn:
|
||||
existing = conn.execute(
|
||||
"SELECT state FROM events WHERE delivery_id=?",
|
||||
(delivery_id,),
|
||||
).fetchone()
|
||||
if existing is not None:
|
||||
return existing["state"] == "deferred"
|
||||
|
||||
row = conn.execute(
|
||||
"SELECT COUNT(*) AS n FROM deferred_submissions WHERE login=?",
|
||||
(normalized_login,),
|
||||
).fetchone()
|
||||
backlog = int(row["n"]) if row is not None else 0
|
||||
deferred = cap > 0 and backlog < cap
|
||||
state: EventState = "deferred" if deferred else "skipped"
|
||||
last_error = reason if deferred else f"{reason}; deferred backlog full ({backlog}/{cap})"
|
||||
conn.execute(
|
||||
"""
|
||||
INSERT INTO events
|
||||
(delivery_id, event_type, repo, issue_key, payload_json, received_at, state, last_error)
|
||||
VALUES (?, ?, ?, ?, ?, ?, ?, ?)
|
||||
""",
|
||||
(
|
||||
delivery_id,
|
||||
event_type,
|
||||
repo,
|
||||
issue_key,
|
||||
json.dumps(payload, separators=(",", ":")),
|
||||
now,
|
||||
state,
|
||||
last_error,
|
||||
),
|
||||
)
|
||||
if deferred:
|
||||
conn.execute(
|
||||
"""
|
||||
INSERT INTO deferred_submissions (delivery_id, login, repo, cap, created_at)
|
||||
VALUES (?, ?, ?, ?, ?)
|
||||
""",
|
||||
(delivery_id, normalized_login, repo, cap, now),
|
||||
)
|
||||
return deferred
|
||||
|
||||
def promote_deferred_submissions(self, *, since: str) -> int:
|
||||
"""Admit the oldest deferred events when their submitters have capacity."""
|
||||
promoted = 0
|
||||
now = _utcnow()
|
||||
with self._txn() as conn:
|
||||
rows = conn.execute(
|
||||
"""
|
||||
SELECT deferred.delivery_id, deferred.login, deferred.repo, deferred.cap
|
||||
FROM deferred_submissions AS deferred
|
||||
JOIN events ON events.delivery_id = deferred.delivery_id
|
||||
WHERE events.state = 'deferred'
|
||||
ORDER BY events.received_at, events.rowid
|
||||
"""
|
||||
).fetchall()
|
||||
for row in rows:
|
||||
usage = conn.execute(
|
||||
"SELECT COUNT(*) AS n FROM submissions WHERE login=? AND ts>=?",
|
||||
(row["login"], since),
|
||||
).fetchone()
|
||||
used = int(usage["n"]) if usage is not None else 0
|
||||
if used >= int(row["cap"]):
|
||||
continue
|
||||
conn.execute(
|
||||
"INSERT OR IGNORE INTO submissions (delivery_id, login, repo, ts) VALUES (?, ?, ?, ?)",
|
||||
(row["delivery_id"], row["login"], row["repo"], now),
|
||||
)
|
||||
conn.execute(
|
||||
"""
|
||||
UPDATE events
|
||||
SET state='queued', last_error=NULL, available_at=NULL, finished_at=NULL
|
||||
WHERE delivery_id=? AND state='deferred'
|
||||
""",
|
||||
(row["delivery_id"],),
|
||||
)
|
||||
conn.execute(
|
||||
"DELETE FROM deferred_submissions WHERE delivery_id=?",
|
||||
(row["delivery_id"],),
|
||||
)
|
||||
promoted += 1
|
||||
return promoted
|
||||
|
||||
def record_submission(
|
||||
self,
|
||||
*,
|
||||
|
||||
@@ -12,7 +12,7 @@ from contextlib import suppress
|
||||
from robomp import tasks
|
||||
from robomp.cancellation import clear_current_event, set_current_event
|
||||
from robomp.config import Settings
|
||||
from robomp.db import Database, EventRow, iso_seconds_ago
|
||||
from robomp.db import Database, EventRow
|
||||
from robomp.github_backend import GitHubBackend
|
||||
from robomp.sandbox import GitTransport, SandboxManager, _reap_slot
|
||||
from robomp.slot_pool import SlotPool
|
||||
@@ -100,11 +100,6 @@ 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.
|
||||
@@ -191,40 +186,6 @@ 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:
|
||||
@@ -253,11 +214,7 @@ class WorkerPool:
|
||||
# Naive but fine for v1 (small queue).
|
||||
row = await asyncio.to_thread(self.db.claim_next_event)
|
||||
if row is None:
|
||||
if not await self._promote_deferred():
|
||||
return None
|
||||
row = await asyncio.to_thread(self.db.claim_next_event)
|
||||
if row is None:
|
||||
return None
|
||||
return None
|
||||
key = row.issue_key or row.delivery_id
|
||||
if key in self._inflight:
|
||||
# Put it back; another in-flight task is touching the same issue.
|
||||
|
||||
@@ -434,7 +434,6 @@ def create_app(settings: Settings | None = None) -> FastAPI:
|
||||
cap=cap,
|
||||
)
|
||||
if not admission.accepted:
|
||||
assert cap is not None
|
||||
window = int(cfg.rate_limit_window_seconds)
|
||||
reason = f"rate limit: @{submitter} has used {admission.used}/{cap} submissions in the last {window}s"
|
||||
log.info(
|
||||
@@ -448,21 +447,17 @@ def create_app(settings: Settings | None = None) -> FastAPI:
|
||||
"cap": cap,
|
||||
},
|
||||
)
|
||||
deferred = db.defer_submission_event(
|
||||
db.record_event(
|
||||
delivery_id=x_github_delivery,
|
||||
event_type=x_github_event,
|
||||
login=submitter,
|
||||
repo=decision.repo,
|
||||
issue_key=decision.issue_key,
|
||||
payload=payload,
|
||||
cap=cap,
|
||||
reason=reason,
|
||||
state="skipped",
|
||||
last_error=reason,
|
||||
)
|
||||
event_state = "deferred" if deferred else "skipped"
|
||||
if deferred:
|
||||
bag["pool"].wake()
|
||||
return JSONResponse(
|
||||
{"delivery": x_github_delivery, "state": event_state, "reason": "rate_limited"},
|
||||
{"delivery": x_github_delivery, "state": "skipped", "reason": "rate_limited"},
|
||||
status_code=202,
|
||||
)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user