From df833940f2698cd60e4621eb515a812cf45732ff Mon Sep 17 00:00:00 2001 From: Jugal Bhatt Date: Tue, 12 Aug 2025 17:09:45 -0700 Subject: [PATCH] refactor(http_parsing_utils.py): streamline request body handling and memory management - Removed the weak reference cache for request bodies and replaced it with direct storage on the request object to prevent memory leaks. - Implemented immediate garbage collection every 100 requests to manage memory usage effectively. - Added a cleanup function to explicitly free memory associated with request processing. - Updated the logic for retrieving and storing parsed request bodies to enhance performance and reliability. --- .../proxy/common_utils/http_parsing_utils.py | 85 ++++++++++--------- 1 file changed, 43 insertions(+), 42 deletions(-) diff --git a/litellm/proxy/common_utils/http_parsing_utils.py b/litellm/proxy/common_utils/http_parsing_utils.py index 38b53a52c2e..ae71a39ec88 100644 --- a/litellm/proxy/common_utils/http_parsing_utils.py +++ b/litellm/proxy/common_utils/http_parsing_utils.py @@ -1,5 +1,4 @@ import json -import weakref from typing import Any, Dict, List, Optional import orjson @@ -9,10 +8,6 @@ from litellm._logging import verbose_proxy_logger from litellm.proxy._types import ProxyException from litellm.types.router import Deployment -# Use a regular dictionary with manual cleanup -# WeakValueDictionary doesn't work well with our setup since request objects may be kept alive -_request_body_cache: Dict[int, Dict] = {} - async def _read_request_body(request: Optional[Request]) -> Dict: """ @@ -49,7 +44,19 @@ async def _read_request_body(request: Optional[Request]) -> Dict: parsed_body = {} else: try: + # Parse and immediately clear body reference to prevent memory accumulation parsed_body = orjson.loads(body) + del body # Clear body reference immediately + + # Force garbage collection every 100 requests to prevent accumulation + import gc + if hasattr(_read_request_body, '_gc_counter'): + _read_request_body._gc_counter += 1 + else: + _read_request_body._gc_counter = 1 + + if _read_request_body._gc_counter % 100 == 0: + gc.collect() except orjson.JSONDecodeError as e: # First try the standard json module which is more forgiving # First decode bytes to string if needed @@ -95,36 +102,15 @@ async def _read_request_body(request: Optional[Request]) -> Dict: return {} -def clear_request_cache(request: Optional[Request] = None) -> None: - """ - Clear cached request bodies to free memory. - - Parameters: - - request: If provided, only clear cache for this specific request. - If None, clear entire cache. - """ - global _request_body_cache - - if request is not None: - request_id = id(request) - if request_id in _request_body_cache: - try: - del _request_body_cache[request_id] - except KeyError: - pass - else: - # Clear entire cache - _request_body_cache.clear() def _safe_get_request_parsed_body(request: Optional[Request]) -> Optional[dict]: if request is None: return None - # Try to get from cache first - request_id = id(request) - if request_id in _request_body_cache: - return _request_body_cache[request_id] + # Check if we already have a parsed body stored directly on the request object + if hasattr(request, "_litellm_parsed_body"): + return getattr(request, "_litellm_parsed_body") # Fallback to checking request.scope for backward compatibility if ( @@ -134,13 +120,34 @@ def _safe_get_request_parsed_body(request: Optional[Request]) -> Optional[dict]: ): accepted_keys, parsed_body = request.scope["parsed_body"] result = {key: parsed_body[key] for key in accepted_keys} - # Clean up the scope to free memory + # Clean up the scope to free memory and store on request object del request.scope["parsed_body"] - # Store in our cache for consistency - _request_body_cache[request_id] = result + setattr(request, "_litellm_parsed_body", result) return result return None + +def cleanup_request_memory(request: Optional[Request]) -> None: + """ + Explicitly cleanup request memory to prevent leaks. + Call this after request processing is complete. + """ + if request is None: + return + + try: + # Remove parsed body from request object + if hasattr(request, "_litellm_parsed_body"): + delattr(request, "_litellm_parsed_body") + + # Clean up any remaining scope data + if hasattr(request, "scope") and "parsed_body" in request.scope: + del request.scope["parsed_body"] + + except Exception: + pass # Silent cleanup - don't break request processing + + def _safe_get_request_query_params(request: Optional[Request]) -> Dict: if request is None: return {} @@ -162,16 +169,10 @@ def _safe_set_request_parsed_body( if request is None: return - # Store in cache with size limit to prevent unbounded growth - request_id = id(request) - _request_body_cache[request_id] = parsed_body - - # If cache gets too large, remove oldest entries - if len(_request_body_cache) > 1000: # Maximum 1000 cached requests - # Remove the oldest 100 entries - keys_to_remove = list(_request_body_cache.keys())[:100] - for key in keys_to_remove: - del _request_body_cache[key] + # Store the parsed body directly on the request object + # This prevents memory leaks and data cross-contamination since + # the data is tied to the specific request object lifecycle + setattr(request, "_litellm_parsed_body", parsed_body) except Exception as e: verbose_proxy_logger.debug(