fix: address autofix bugs in chatgpt SSE, vertex token cache, rubrik aclose

- chatgpt responses: don't overwrite a meaningful error_message with None
  when a later RESPONSE_FAILED/ERROR event lacks an error object.
- vertex_ai: serve STALE tokens from the lock-free fast path and only
  schedule a deduplicated background refresh, eliminating per-key lock
  contention near token expiry.
- rubrik: aclose() now closes both async_httpx_client and
  tool_blocking_client to avoid leaking connections from the dedicated
  client when the logger shuts down.

Co-authored-by: Yassin Kortam <yassin@berri.ai>
This commit is contained in:
Cursor Agent 2026-05-20 23:33:12 +00:00
parent a0b69aa54a
commit c1916d3551
No known key found for this signature in database
3 changed files with 47 additions and 6 deletions

View file

@ -142,8 +142,9 @@ class RubrikLogger(CustomGuardrail, CustomBatchLogger):
self._flush_task = self._start_periodic_flush_task()
async def aclose(self):
"""Close the dedicated tool blocking HTTP client."""
"""Close the dedicated HTTP clients used by this logger."""
await self.tool_blocking_client.close()
await self.async_httpx_client.close()
# -- Guardrail hook --------------------------------------------------------

View file

@ -203,7 +203,9 @@ class ChatGPTResponsesAPIConfig(OpenAIResponsesAPIConfig):
ResponsesAPIStreamEvents.RESPONSE_FAILED,
ResponsesAPIStreamEvents.ERROR,
):
error_message = self._extract_error_message(parsed_chunk)
extracted_error = self._extract_error_message(parsed_chunk)
if extracted_error is not None:
error_message = extracted_error
return completed_response, error_message

View file

@ -434,6 +434,33 @@ class VertexBase:
return creds.token, resolved_project
return None
def _try_get_usable_cached_token(
self,
credential_cache_key: tuple,
project_id: Optional[str],
) -> Optional[Tuple[str, str, "TokenState", Any, Optional[str]]]:
"""
Look up cached credentials and return usable token info for FRESH or
STALE tokens (both are still valid for outbound requests). STALE
tokens are returned along with their state and the underlying
credentials object so the caller can schedule a background refresh
without holding the per-key async lock.
"""
from google.auth.credentials import TokenState
creds, cached_project_id = self._unpack_cached_credentials(credential_cache_key)
if creds is None:
return None
token_state = self._get_token_state(creds)
if token_state not in (TokenState.FRESH, TokenState.STALE):
return None
if creds.token is None or not isinstance(creds.token, str):
return None
resolved_project = project_id or cached_project_id
if not resolved_project:
return None
return creds.token, resolved_project, token_state, creds, cached_project_id
def _unpack_cached_credentials(
self, credential_cache_key: tuple
) -> Tuple[Any, Optional[str]]:
@ -1010,10 +1037,21 @@ class VertexBase:
credential_cache_key = (cache_credentials, project_id)
# === FAST PATH (no lock) ===
# If credentials are FRESH (valid, not near expiry), return immediately.
cached = self._try_get_cached_token(credential_cache_key, project_id)
if cached is not None:
return cached
# If credentials are FRESH or STALE, return immediately without
# touching the per-key async lock. STALE tokens are still usable;
# we kick off a deduplicated background refresh so subsequent
# requests get a fresh token, but we must not serialize concurrent
# callers on the lock just to schedule that refresh.
usable = self._try_get_usable_cached_token(credential_cache_key, project_id)
if usable is not None:
cached_token, resolved_project, token_state, creds, cached_project_id = (
usable
)
if token_state == TokenState.STALE:
self._schedule_background_refresh(
creds, credential_cache_key, cached_project_id
)
return cached_token, resolved_project
# === SLOW PATH (per-key lock) ===
lock = self._acquire_async_refresh_lock(credential_cache_key)