Files
saidsurucuandClaude Opus 4.7 96a5a538b2 perf(server): unblock event loop on rate-limit waits and markitdown
Two complementary changes to mitigate intermittent TLS handshake
timeouts and "notifications/cancelled: Bad Request" seen against the
single-worker uvicorn deployment.

1. bedesten rate-limiter back-pressure
   - Add optional ``max_wait`` to ``_TokenBucket.acquire``: if the next
     wait would exceed it, raise ``BedestenRateLimited`` immediately
     instead of sleeping. After a server-side 429 the bucket pauses for
     up to 30s; previously a queued request sat in ``asyncio.sleep``
     for that whole window, holding the worker slot and pushing the
     MCP client past its cancellation timeout.
   - ``search_bedesten_unified`` / ``get_bedesten_document_markdown``
     catch ``BedestenRateLimited`` and reuse the existing structured
     429-style response, so callers get a fast, clean retry signal.
   - Tunable via ``BEDESTEN_RATE_MAX_WAIT_S`` (default 8.0s).

2. Offload sync markitdown conversions to a thread
   - Every ``markitdown.convert*`` call site is now wrapped in
     ``asyncio.to_thread(...)`` across 14 modules (bedesten, yargitay,
     danistay, anayasa norm + bireysel, uyusmazlik, emsal, rekabet,
     gib, kvkk, sayistay, bddk, sigorta_tahkim, kik_v2). PDF / large
     HTML parsing was stalling the event loop for seconds, which on a
     single-worker deployment delayed every other in-flight request
     and queued new TLS handshakes until they timed out.

Verified locally:
- ``ast.parse`` + ``importlib.import_module`` on all 15 modified files
- ``mcp_server_main.create_app()`` constructs successfully
- New ``_TokenBucket.acquire(max_wait=...)`` smoke-tested across 6
  paths: capacity-available, no-arg backward compat, max_wait raise,
  max_wait wait+succeed, ``penalize_until`` + max_wait fast-raise.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
2026-05-11 14:31:23 +03:00

308 lines
13 KiB
Python
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
# bedesten_mcp_module/client.py
import asyncio
import base64
import io
import logging
import os
import time
from typing import Optional
import httpx
from markitdown import MarkItDown
from .models import (
BedestenSearchRequest, BedestenSearchResponse,
BedestenDocumentRequest, BedestenDocumentResponse,
BedestenDocumentMarkdown, BedestenDocumentRequestData
)
from .enums import get_full_birim_adi
logger = logging.getLogger(__name__)
class BedestenRateLimited(Exception):
"""Raised when the local rate-limit bucket would block longer than allowed.
Carries the suggested retry-after (seconds) so callers can surface a
structured 429-style response to the MCP client instead of silently
blocking the event-loop slot for the full bucket-pause window.
"""
def __init__(self, retry_after: float) -> None:
self.retry_after = retry_after
super().__init__(f"local bucket would block {retry_after:.1f}s")
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, max_wait: Optional[float] = None) -> None:
"""Acquire one token. If ``max_wait`` is set and the next wait would
exceed it, raise :class:`BedestenRateLimited` immediately instead of
sleeping — keeps a single rate-limited request from holding the
worker-slot for the full bucket-pause window (up to ~30s on 429)."""
deadline = (time.monotonic() + max_wait) if max_wait is not None else 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
if deadline is not None:
remaining = deadline - time.monotonic()
if wait_s > remaining:
raise BedestenRateLimited(retry_after=wait_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.
Currently used for Yargıtay decisions, but can be extended for other court types.
"""
BASE_URL = "https://bedesten.adalet.gov.tr"
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)
# BEDESTEN_RATE_MAX_WAIT_S (default 8.0; max seconds to wait in the
# local bucket before returning a structured 429 to the caller)
_DEFAULT_CAPACITY = int(os.getenv("BEDESTEN_RATE_CAPACITY", "1"))
_DEFAULT_REFILL_S = float(os.getenv("BEDESTEN_RATE_REFILL_S", "3.5"))
_DEFAULT_MAX_WAIT_S = float(os.getenv("BEDESTEN_RATE_MAX_WAIT_S", "8.0"))
def __init__(self, request_timeout: float = 60.0):
self.http_client = httpx.AsyncClient(
base_url=self.BASE_URL,
headers={
"Accept": "*/*",
"Accept-Language": "tr-TR,tr;q=0.9,en-US;q=0.8,en;q=0.7",
"AdaletApplicationName": "UyapMevzuat",
"Content-Type": "application/json; charset=utf-8",
"Origin": "https://mevzuat.adalet.gov.tr",
"Referer": "https://mevzuat.adalet.gov.tr/",
"Sec-Fetch-Dest": "empty",
"Sec-Fetch-Mode": "cors",
"Sec-Fetch-Site": "same-site",
"User-Agent": "Mozilla/5.0 (Macintosh; Intel Mac OS X 10_15_7) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/137.0.0.0 Safari/537.36"
},
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:
"""
Search for documents using Bedesten API.
Currently supports: YARGITAYKARARI, DANISTAYKARARI, YERELHUKMAHKARARI, etc.
"""
logger.info(f"BedestenApiClient: Searching documents with phrase: {search_request.data.phrase}")
# Map abbreviated birimAdi to full Turkish name before sending to API
original_birim_adi = search_request.data.birimAdi
mapped_birim_adi = get_full_birim_adi(original_birim_adi)
search_request.data.birimAdi = mapped_birim_adi
if original_birim_adi != "ALL":
logger.info(f"BedestenApiClient: Mapped birimAdi '{original_birim_adi}' to '{mapped_birim_adi}'")
try:
# Create request dict and remove birimAdi if empty
request_dict = search_request.model_dump()
if not request_dict["data"]["birimAdi"]: # Remove if empty string
del request_dict["data"]["birimAdi"]
await self._bucket.acquire(max_wait=self._DEFAULT_MAX_WAIT_S)
response = await self.http_client.post(
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
except Exception as e:
logger.error(f"BedestenApiClient: Error processing search response: {e}")
raise
async def get_document_as_markdown(self, document_id: str) -> BedestenDocumentMarkdown:
"""
Get document content and convert to markdown.
Handles both HTML (text/html) and PDF (application/pdf) content types.
"""
logger.info(f"BedestenApiClient: Fetching document for markdown conversion (ID: {document_id})")
try:
# Prepare request
doc_request = BedestenDocumentRequest(
data=BedestenDocumentRequestData(documentId=document_id)
)
# Get document
await self._bucket.acquire(max_wait=self._DEFAULT_MAX_WAIT_S)
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)
# Add null safety checks for document data
if not hasattr(doc_response, 'data') or doc_response.data is None:
raise ValueError("Document response does not contain data")
if not hasattr(doc_response.data, 'content') or doc_response.data.content is None:
raise ValueError("Document data does not contain content")
if not hasattr(doc_response.data, 'mimeType') or doc_response.data.mimeType is None:
raise ValueError("Document data does not contain mimeType")
# Decode base64 content with error handling
try:
content_bytes = base64.b64decode(doc_response.data.content)
except Exception as e:
raise ValueError(f"Failed to decode base64 content: {str(e)}")
mime_type = doc_response.data.mimeType
logger.info(f"BedestenApiClient: Document mime type: {mime_type}")
# Convert to markdown based on mime type. markitdown is sync and
# PDF parsing in particular can block the event-loop for seconds,
# which on a single-worker uvicorn deployment stalls every other
# in-flight MCP request and new TLS handshakes. Offload to a
# thread so the event-loop stays responsive.
if mime_type == "text/html":
html_content = content_bytes.decode('utf-8')
markdown_content = await asyncio.to_thread(
self._convert_html_to_markdown, html_content
)
elif mime_type == "application/pdf":
markdown_content = await asyncio.to_thread(
self._convert_pdf_to_markdown, content_bytes
)
else:
logger.warning(f"Unsupported mime type: {mime_type}")
markdown_content = f"Unsupported content type: {mime_type}. Unable to convert to markdown."
return BedestenDocumentMarkdown(
documentId=document_id,
markdown_content=markdown_content,
source_url=f"https://mevzuat.adalet.gov.tr/ictihat/{document_id}",
mime_type=mime_type
)
except httpx.RequestError as e:
logger.error(f"BedestenApiClient: HTTP error fetching document {document_id}: {e}")
raise
except Exception as e:
logger.error(f"BedestenApiClient: Error processing document {document_id}: {e}")
raise
def _convert_html_to_markdown(self, html_content: str) -> Optional[str]:
"""Convert HTML to Markdown using MarkItDown"""
if not html_content:
return None
try:
# Convert HTML string to bytes and create BytesIO stream
html_bytes = html_content.encode('utf-8')
html_stream = io.BytesIO(html_bytes)
# Pass BytesIO stream to MarkItDown to avoid temp file creation
md_converter = MarkItDown()
result = md_converter.convert(html_stream)
markdown_content = result.text_content
logger.info("Successfully converted HTML to Markdown")
return markdown_content
except Exception as e:
logger.error(f"Error converting HTML to Markdown: {e}")
return f"Error converting HTML content: {str(e)}"
def _convert_pdf_to_markdown(self, pdf_bytes: bytes) -> Optional[str]:
"""Convert PDF to Markdown using MarkItDown"""
if not pdf_bytes:
return None
try:
# Create BytesIO stream from PDF bytes
pdf_stream = io.BytesIO(pdf_bytes)
# Pass BytesIO stream to MarkItDown to avoid temp file creation
md_converter = MarkItDown()
result = md_converter.convert(pdf_stream)
markdown_content = result.text_content
logger.info("Successfully converted PDF to Markdown")
return markdown_content
except Exception as e:
logger.error(f"Error converting PDF to Markdown: {e}")
return f"Error converting PDF content: {str(e)}. The document may be corrupted or in an unsupported format."
async def close_client_session(self):
"""Close HTTP client session"""
await self.http_client.aclose()
logger.info("BedestenApiClient: HTTP client session closed.")