mirror of
https://github.com/openswarm-ai/openswarm.git
synced 2026-09-06 01:37:43 +02:00
[eric] service: spool 429/timeout/5xx instead of dropping, warn on 4xx rejects + spool evictions
This commit is contained in:
@@ -72,17 +72,21 @@ def enqueue(spool_path: str, kind: str, payload: dict, *, now: float) -> None:
|
||||
target = int(_MAX_BYTES * _TRIM_TARGET_FRACTION)
|
||||
# Delete oldest rows until we're back under target. Use a
|
||||
# reasonable batch size so we don't block forever.
|
||||
dropped = 0
|
||||
for _ in range(64):
|
||||
row = c.execute("SELECT id FROM spool ORDER BY id ASC LIMIT 1").fetchone()
|
||||
if not row:
|
||||
break
|
||||
c.execute("DELETE FROM spool WHERE id = ?", (row[0],))
|
||||
dropped += 1
|
||||
try:
|
||||
new_size = os.path.getsize(spool_path)
|
||||
except OSError:
|
||||
new_size = 0
|
||||
if new_size <= target:
|
||||
break
|
||||
if dropped:
|
||||
logger.warning("Spool over %d MB cap; dropped %d oldest entries", _MAX_BYTES // (1024 * 1024), dropped)
|
||||
# VACUUM is expensive; only run if we still appear oversized after
|
||||
# trimming, otherwise free pages get reused on next insert.
|
||||
try:
|
||||
|
||||
@@ -192,15 +192,24 @@ def _base_url() -> str:
|
||||
return _DEFAULT_BASE
|
||||
|
||||
|
||||
async def _post(path: str, body: dict) -> bool:
|
||||
async def _post(path: str, body: dict) -> int | None:
|
||||
url = f"{_base_url()}{path}"
|
||||
try:
|
||||
async with httpx.AsyncClient(timeout=_TIMEOUT_SECONDS) as c:
|
||||
r = await c.post(url, json=body)
|
||||
return 200 <= r.status_code < 500
|
||||
return r.status_code
|
||||
except Exception as e:
|
||||
logger.debug("service POST %s failed: %s", path, e)
|
||||
return False
|
||||
return None
|
||||
|
||||
|
||||
def _delivered(status: int | None) -> bool:
|
||||
return status is not None and 200 <= status < 300
|
||||
|
||||
|
||||
# 429/timeouts/5xx/network are worth retrying; other 4xx means the payload itself is rejected and retrying forever would just poison the spool.
|
||||
def _retryable(status: int | None) -> bool:
|
||||
return status is None or status >= 500 or status in (408, 429)
|
||||
|
||||
|
||||
async def _post_or_spool(path: str, body: dict, kind: str) -> None:
|
||||
@@ -217,9 +226,11 @@ async def _post_or_spool(path: str, body: dict, kind: str) -> None:
|
||||
return
|
||||
_inflight += 1
|
||||
try:
|
||||
ok = await _post(path, body)
|
||||
if not ok:
|
||||
status = await _post(path, body)
|
||||
if _retryable(status):
|
||||
buffer.enqueue(_spool_path(), f"{kind}:{path}", body, now=time.time())
|
||||
elif not _delivered(status):
|
||||
logger.warning("service POST %s rejected with HTTP %s; payload dropped", path, status)
|
||||
finally:
|
||||
async with _inflight_lock:
|
||||
_inflight = max(0, _inflight - 1)
|
||||
@@ -236,11 +247,14 @@ async def drain_spool(batch_size: int = 50) -> int:
|
||||
if not path:
|
||||
succeeded.append(rid)
|
||||
continue
|
||||
ok = await _post(path, body)
|
||||
if ok:
|
||||
status = await _post(path, body)
|
||||
if _delivered(status):
|
||||
succeeded.append(rid)
|
||||
else:
|
||||
elif _retryable(status):
|
||||
break
|
||||
else:
|
||||
logger.warning("service replay %s rejected with HTTP %s; dropping spooled row", path, status)
|
||||
succeeded.append(rid)
|
||||
if succeeded:
|
||||
buffer.acknowledge(_spool_path(), succeeded)
|
||||
return len(succeeded)
|
||||
|
||||
Reference in New Issue
Block a user