fix(streaming): remove expensive debug logging and optimize usage stripping

- Remove print_verbose calls that format chunk/response Pydantic objects,
  triggering millions of __repr__ calls (8% of CPU in profiling)
- Guard remaining verbose_logger.debug with isEnabledFor(DEBUG) and use
  lazy %s formatting instead of f-strings
- Replace usage stripping round-trip (model_dump + delete + reconstruct)
  with a _usage_stripped flag, deferring exclusion to serialization time
This commit is contained in:
Ishaan Jaffer 2026-02-18 12:33:08 -08:00
parent 2d20d609ad
commit 8e46335e16

View file

@ -2,6 +2,7 @@ import asyncio
import collections.abc
import datetime
import json
import logging
import threading
import time
import traceback
@ -435,7 +436,7 @@ class CustomStreamWrapper:
def handle_openai_chat_completion_chunk(self, chunk):
try:
print_verbose(f"\nRaw OpenAI Chunk\n{chunk}\n")
str_line = chunk
text = ""
is_finished = False
@ -485,7 +486,7 @@ class CustomStreamWrapper:
def handle_azure_text_completion_chunk(self, chunk):
try:
print_verbose(f"\nRaw OpenAI Chunk\n{chunk}\n")
text = ""
is_finished = False
finish_reason = None
@ -506,7 +507,7 @@ class CustomStreamWrapper:
def handle_openai_text_completion_chunk(self, chunk):
try:
print_verbose(f"\nRaw OpenAI Chunk\n{chunk}\n")
text = ""
is_finished = False
finish_reason = None
@ -870,9 +871,6 @@ class CustomStreamWrapper:
preserve_upstream_non_openai_attributes,
)
print_verbose(
f"completion_obj: {completion_obj}, model_response.choices[0]: {model_response.choices[0]}, response_obj: {response_obj}"
)
is_chunk_non_empty = self.is_chunk_non_empty(
completion_obj, model_response, response_obj
)
@ -899,11 +897,9 @@ class CustomStreamWrapper:
choice_json.pop(
"finish_reason", None
) # for mistral etc. which return a value in their last chunk (not-openai compatible).
print_verbose(f"choice_json: {choice_json}")
choices.append(StreamingChoices(**choice_json))
except Exception:
choices.append(StreamingChoices())
print_verbose(f"choices in streaming: {choices}")
setattr(model_response, "choices", choices)
else:
return
@ -921,8 +917,10 @@ class CustomStreamWrapper:
)
model_response = self.strip_role_from_delta(model_response)
verbose_logger.debug(
f"model_response.choices[0].delta inside is_chunk_non_empty: {model_response.choices[0].delta}"
if verbose_logger.isEnabledFor(logging.DEBUG):
verbose_logger.debug(
"model_response.choices[0].delta: %s",
model_response.choices[0].delta,
)
else:
## else
@ -1370,9 +1368,6 @@ class CustomStreamWrapper:
)
model_response.model = self.model
print_verbose(
f"model_response finish reason 3: {self.received_finish_reason}; response_obj={response_obj}"
)
## FUNCTION CALL PARSING
original_chunk = (
response_obj.get("original_chunk") if response_obj is not None else None
@ -1432,7 +1427,6 @@ class CustomStreamWrapper:
):
t.function.arguments = ""
_json_delta = delta.model_dump()
print_verbose(f"_json_delta: {_json_delta}")
if "role" not in _json_delta or _json_delta["role"] is None:
_json_delta[
"role"
@ -1466,11 +1460,7 @@ class CustomStreamWrapper:
if original_chunk.choices[0].delta is None
else dict(original_chunk.choices[0].delta)
)
print_verbose(f"original delta: {delta}")
model_response.choices[0].delta = Delta(**delta)
print_verbose(
f"new delta: {model_response.choices[0].delta}"
)
except Exception:
model_response.choices[0].delta = Delta()
else:
@ -1480,11 +1470,6 @@ class CustomStreamWrapper:
):
return model_response
return
print_verbose(
f"model_response.choices[0].delta: {model_response.choices[0].delta}; completion_obj: {completion_obj}"
)
print_verbose(f"self.sent_first_chunk: {self.sent_first_chunk}")
## CHECK FOR TOOL USE
if "tool_calls" in completion_obj and len(completion_obj["tool_calls"]) > 0:
@ -1915,18 +1900,9 @@ class CustomStreamWrapper:
and len(chunk.parts) == 0
):
continue
# chunk_creator() does logging/stream chunk building. We need to let it know its being called in_async_func, so we don't double add chunks.
# __anext__ also calls async_success_handler, which does logging
verbose_logger.debug(
f"PROCESSED ASYNC CHUNK PRE CHUNK CREATOR: {chunk}"
)
processed_chunk: Optional[ModelResponseStream] = self.chunk_creator(
chunk=chunk
)
verbose_logger.debug(
f"PROCESSED ASYNC CHUNK POST CHUNK CREATOR: {processed_chunk}"
)
if processed_chunk is None:
continue
@ -1955,17 +1931,14 @@ class CustomStreamWrapper:
hasattr(processed_chunk, "usage")
and getattr(processed_chunk, "usage", None) is not None
):
# Set usage to None so model_dump_json(exclude_none=True)
# drops it. The original usage is already preserved in
# self.chunks (appended above) for calculate_total_usage().
processed_chunk.usage = None # type: ignore
# Flag for proxy to exclude usage during serialization
# instead of expensive model_dump() + reconstruct round-trip
processed_chunk._usage_stripped = True # type: ignore
is_empty = is_model_response_stream_empty(
model_response=cast(ModelResponseStream, processed_chunk)
)
if is_empty:
continue
print_verbose(f"final returned processed chunk: {processed_chunk}")
# add usage as hidden param
if self.sent_last_chunk is True and self.stream_options is None:
@ -1980,8 +1953,7 @@ class CustomStreamWrapper:
)
)
# Add MCP metadata to final chunk if present (after hooks)
if processed_chunk is not None:
processed_chunk = self._add_mcp_metadata_to_final_chunk(processed_chunk)
processed_chunk = self._add_mcp_metadata_to_final_chunk(processed_chunk) # type: ignore[reportArgumentType]
return processed_chunk
raise StopAsyncIteration
@ -1995,13 +1967,9 @@ class CustomStreamWrapper:
else:
chunk = next(self.completion_stream)
if chunk is not None and chunk != b"":
print_verbose(f"PROCESSED CHUNK PRE CHUNK CREATOR: {chunk}")
processed_chunk: Optional[
ModelResponseStream
] = self.chunk_creator(chunk=chunk)
print_verbose(
f"PROCESSED CHUNK POST CHUNK CREATOR: {processed_chunk}"
)
if processed_chunk is None:
continue