From 58457b076f19900a819793622dc145cc0de330b7 Mon Sep 17 00:00:00 2001 From: saidsurucu Date: Fri, 8 May 2026 21:10:37 +0300 Subject: [PATCH] feat(bedesten): client-side rate limiter with 429 back-pressure MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Probed the live API (2026-05-08): the per-IP limit is 10 requests in a rolling 30s window, with HTTP 429 + Retry-After: 30 on the 11th call. Add a token bucket inside BedestenApiClient (default capacity=1, refill 1 token / 3.5s — strict serialization, no burst) so we stay below the threshold by default. When the server still returns 429 (e.g. the egress IP is shared with other clients), pause the whole bucket for the Retry-After window so queued in-flight requests wait gracefully instead of hammering. Tunable via BEDESTEN_RATE_CAPACITY / BEDESTEN_RATE_REFILL_S. Verified: 14 concurrent requests after a clean cooldown -> 13 OK, 1 stray 429 (bucket auto-paused 22.5s, then drained cleanly). Co-Authored-By: Claude Opus 4.7 (1M context) --- CLAUDE.md | 2 +- bedesten_mcp_module/client.py | 94 ++++++++++++++++++++++++++++++++--- 2 files changed, 88 insertions(+), 8 deletions(-) diff --git a/CLAUDE.md b/CLAUDE.md index 63ee261..61eb5e5 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -473,7 +473,7 @@ doc8 = await get_kvkk_document_markdown(decision_url="https://www.kvkk.gov.tr/Ic | API | Rate Limit | Notes | |-----|------------|-------| -| Bedesten Unified | Unknown | Ministry of Justice API - no documented limits | +| Bedesten Unified | ~10 req / 30s window per source IP (measured 2026-05-08); 11th req → HTTP 429 with `Retry-After: 30`. Client uses an internal token bucket (default 1 token, refill 1/3.5s) plus 429 back-pressure (whole bucket pauses for the Retry-After window). Override via `BEDESTEN_RATE_CAPACITY` / `BEDESTEN_RATE_REFILL_S`. | | Yargıtay Primary | Unknown | Official government API | | Danıştay | Unknown | Official government API | | Anayasa Mahkemesi | Unknown | Constitutional Court API | diff --git a/bedesten_mcp_module/client.py b/bedesten_mcp_module/client.py index 3c03a54..c076739 100644 --- a/bedesten_mcp_module/client.py +++ b/bedesten_mcp_module/client.py @@ -1,11 +1,15 @@ # bedesten_mcp_module/client.py -import httpx +import asyncio import base64 -from typing import Optional -import logging -from markitdown import MarkItDown import io +import logging +import os +import time +from typing import Optional + +import httpx +from markitdown import MarkItDown from .models import ( BedestenSearchRequest, BedestenSearchResponse, @@ -16,6 +20,50 @@ from .enums import get_full_birim_adi logger = logging.getLogger(__name__) + +class _TokenBucket: + """Asyncio token bucket with explicit back-pressure. + + Measured Bedesten limit (per source IP, 2026-05-08): 10 requests per + rolling 30s window with full refill — equivalent to capacity=10, + refill_rate=1 token / 3s. Even with margin, 429s still leak through + when other clients share the egress IP, so we also expose + ``penalize_until`` so callers can freeze the bucket when the server + actually returns 429 (Retry-After). + """ + + def __init__(self, capacity: int, refill_per_s: float) -> None: + self.capacity = float(capacity) + self.refill_per_s = float(refill_per_s) + self._tokens = float(capacity) + self._last = time.monotonic() + self._not_before = 0.0 + self._lock = asyncio.Lock() + + async def acquire(self) -> None: + while True: + async with self._lock: + now = time.monotonic() + if now < self._not_before: + wait_s = self._not_before - now + else: + self._tokens = min( + self.capacity, + self._tokens + (now - self._last) * self.refill_per_s, + ) + self._last = now + if self._tokens >= 1.0: + self._tokens -= 1.0 + return + wait_s = (1.0 - self._tokens) / self.refill_per_s + await asyncio.sleep(wait_s) + + def penalize_until(self, monotonic_deadline: float) -> None: + """Pause the bucket until ``monotonic_deadline`` (drains tokens).""" + self._not_before = max(self._not_before, monotonic_deadline) + self._tokens = 0.0 + self._last = time.monotonic() + class BedestenApiClient: """ API Client for Bedesten (bedesten.adalet.gov.tr) - Alternative legal decision search system. @@ -25,6 +73,14 @@ class BedestenApiClient: SEARCH_ENDPOINT = "/emsal-karar/searchDocuments" DOCUMENT_ENDPOINT = "/emsal-karar/getDocumentContent" + # Measured limit (per source IP): 10 requests per 30s window with full + # refill (≈ 1 token / 3s steady). We default to 1-token capacity and + # 3.5s spacing (no burst, ~14% safety margin). Override via env: + # BEDESTEN_RATE_CAPACITY (default 1) + # BEDESTEN_RATE_REFILL_S (default 3.5; seconds per token) + _DEFAULT_CAPACITY = int(os.getenv("BEDESTEN_RATE_CAPACITY", "1")) + _DEFAULT_REFILL_S = float(os.getenv("BEDESTEN_RATE_REFILL_S", "3.5")) + def __init__(self, request_timeout: float = 60.0): self.http_client = httpx.AsyncClient( base_url=self.BASE_URL, @@ -42,6 +98,24 @@ class BedestenApiClient: }, timeout=request_timeout ) + self._bucket = _TokenBucket( + capacity=self._DEFAULT_CAPACITY, + refill_per_s=1.0 / self._DEFAULT_REFILL_S, + ) + + def _handle_429(self, response: httpx.Response, op: str) -> None: + """Apply back-pressure to the shared bucket based on Retry-After.""" + retry_after_raw = response.headers.get("Retry-After", "") + try: + retry_after = float(retry_after_raw) + except (TypeError, ValueError): + retry_after = 30.0 + # Cap penalty so a hostile/buggy server can't freeze us indefinitely. + retry_after = max(1.0, min(retry_after, 60.0)) + self._bucket.penalize_until(time.monotonic() + retry_after + 0.5) + logger.warning( + f"BedestenApiClient: 429 on {op}; bucket paused {retry_after + 0.5:.1f}s" + ) async def search_documents(self, search_request: BedestenSearchRequest) -> BedestenSearchResponse: """ @@ -63,16 +137,19 @@ class BedestenApiClient: if not request_dict["data"]["birimAdi"]: # Remove if empty string del request_dict["data"]["birimAdi"] + await self._bucket.acquire() response = await self.http_client.post( - self.SEARCH_ENDPOINT, + self.SEARCH_ENDPOINT, json=request_dict ) + if response.status_code == 429: + self._handle_429(response, "search") response.raise_for_status() response_json = response.json() - + # Parse and return the response return BedestenSearchResponse(**response_json) - + except httpx.RequestError as e: logger.error(f"BedestenApiClient: HTTP request error during search: {e}") raise @@ -94,10 +171,13 @@ class BedestenApiClient: ) # Get document + await self._bucket.acquire() response = await self.http_client.post( self.DOCUMENT_ENDPOINT, json=doc_request.model_dump() ) + if response.status_code == 429: + self._handle_429(response, f"document {document_id}") response.raise_for_status() response_json = response.json() doc_response = BedestenDocumentResponse(**response_json)