fix(robomp): deferred rate-limited submissions

Stored a bounded per-login overflow backlog and promoted deferred events oldest-first as rolling-window capacity became available. Surfaced the deferred state through the dashboard contract and documented admission behavior.

Fixes #5882
This commit is contained in:
roboomp
2026-07-17 17:26:39 +00:00
parent 0f9fceeea4
commit aecef46436
12 changed files with 297 additions and 32 deletions
+184 -9
View File
@@ -14,7 +14,7 @@ from typing import Any, Literal
from robomp.github_client import IssueIndexEntry
EventState = Literal["queued", "running", "done", "failed", "skipped"]
EventState = Literal["queued", "deferred", "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','running','done','failed','skipped')),
CHECK (state IN ('queued','deferred','running','done','failed','skipped')),
attempts INTEGER NOT NULL DEFAULT 0,
last_error TEXT,
started_at TEXT,
@@ -101,6 +101,16 @@ 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,
@@ -299,6 +309,47 @@ 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()
@@ -452,8 +503,9 @@ 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._lock:
self._conn.execute("DELETE FROM events WHERE delivery_id=?", (delivery_id,))
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,))
def replace_event_if_state_in(
self,
@@ -476,6 +528,7 @@ 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(
"""
@@ -557,7 +610,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", "running", "done", "failed", "skipped"), 0)
counts: dict[str, int] = dict.fromkeys(("queued", "deferred", "running", "done", "failed", "skipped"), 0)
for row in rows:
counts[row["state"]] = int(row["n"])
return counts
@@ -569,7 +622,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", "running", "done", "failed", "skipped"), 0)
counts: dict[str, int] = dict.fromkeys(("queued", "deferred", "running", "done", "failed", "skipped"), 0)
seen: set[str] = set()
with self._lock:
rows = self._conn.execute(
@@ -689,9 +742,9 @@ class Database:
keep public retries from mutating queued/running rows while preserving
internal recovery of a just-claimed running event.
"""
with self._lock:
with self._txn() as conn:
if from_states is None:
cur = self._conn.execute(
cur = conn.execute(
"UPDATE events SET state='queued', available_at=NULL WHERE delivery_id=?",
(delivery_id,),
)
@@ -699,10 +752,12 @@ class Database:
return False
else:
placeholders = ",".join("?" for _ in from_states)
cur = self._conn.execute(
cur = 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:
@@ -1037,6 +1092,22 @@ 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),
@@ -1051,6 +1122,110 @@ 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,
*,
+9 -2
View File
@@ -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
from robomp.db import Database, EventRow, iso_seconds_ago
from robomp.github_backend import GitHubBackend
from robomp.sandbox import GitTransport, SandboxManager, _reap_slot
from robomp.slot_pool import SlotPool
@@ -214,7 +214,14 @@ class WorkerPool:
# Naive but fine for v1 (small queue).
row = await asyncio.to_thread(self.db.claim_next_event)
if row is None:
return 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:
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
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.
+9 -4
View File
@@ -434,6 +434,7 @@ 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(
@@ -447,17 +448,21 @@ def create_app(settings: Settings | None = None) -> FastAPI:
"cap": cap,
},
)
db.record_event(
deferred = db.defer_submission_event(
delivery_id=x_github_delivery,
event_type=x_github_event,
login=submitter,
repo=decision.repo,
issue_key=decision.issue_key,
payload=payload,
state="skipped",
last_error=reason,
cap=cap,
reason=reason,
)
event_state = "deferred" if deferred else "skipped"
if deferred:
bag["pool"].wake()
return JSONResponse(
{"delivery": x_github_delivery, "state": "skipped", "reason": "rate_limited"},
{"delivery": x_github_delivery, "state": event_state, "reason": "rate_limited"},
status_code=202,
)