Skip to content
Open
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -739,9 +739,11 @@ def __next__(self):
self._process_complete_response()
else:
# Handle cleanup for other exceptions during stream iteration
self._ensure_cleanup()
if self._span and self._span.is_recording():
self._span.set_attribute(ERROR_TYPE, e.__class__.__name__)
self._span.record_exception(e)
self._span.set_status(Status(StatusCode.ERROR, str(e)))
self._ensure_cleanup(error=True)
raise
else:
self._process_item(chunk)
Expand All @@ -755,9 +757,11 @@ async def __anext__(self):
self._process_complete_response()
else:
# Handle cleanup for other exceptions during stream iteration
self._ensure_cleanup()
if self._span and self._span.is_recording():
self._span.set_attribute(ERROR_TYPE, e.__class__.__name__)
self._span.record_exception(e)
self._span.set_status(Status(StatusCode.ERROR, str(e)))
self._ensure_cleanup(error=True)
raise
else:
self._process_item(chunk)
Expand Down Expand Up @@ -834,7 +838,7 @@ def _process_complete_response(self):
self._cleanup_completed = True

@dont_throw
def _ensure_cleanup(self):
def _ensure_cleanup(self, error=False):
"""Thread-safe cleanup method that handles different cleanup scenarios"""
with self._cleanup_lock:
if self._cleanup_completed:
Expand All @@ -849,7 +853,8 @@ def _ensure_cleanup(self):

# Set span status and close it
if self._span and self._span.is_recording():
self._span.set_status(Status(StatusCode.OK))
if not error:
self._span.set_status(Status(StatusCode.OK))
self._span.end()
logger.debug("ChatStream span closed successfully")

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -208,43 +208,57 @@ def _set_output_messages(span, choices):
@dont_throw
def _build_from_streaming_response(span, request_kwargs, response):
complete_response = {"choices": [], "model": "", "id": ""}
for item in response:
yield item
_accumulate_streaming_response(complete_response, item)

_set_response_attributes(span, complete_response)

_set_token_usage(span, request_kwargs, complete_response)
try:
for item in response:
yield item
_accumulate_streaming_response(complete_response, item)

span.set_status(Status(StatusCode.OK))
except Exception as e:
if span.is_recording():
span.set_attribute(ERROR_TYPE, e.__class__.__name__)
span.record_exception(e)
span.set_status(Status(StatusCode.ERROR, str(e)))
raise
finally:
_set_response_attributes(span, complete_response)
_set_token_usage(span, request_kwargs, complete_response)

if should_emit_events():
_emit_streaming_response_events(complete_response)
else:
if should_send_prompts():
_set_completions(span, complete_response.get("choices"))
if should_emit_events():
_emit_streaming_response_events(complete_response)
else:
if should_send_prompts():
_set_completions(span, complete_response.get("choices"))

span.set_status(Status(StatusCode.OK))
span.end()
span.end()
Comment thread
coderabbitai[bot] marked this conversation as resolved.
Outdated


@dont_throw
async def _abuild_from_streaming_response(span, request_kwargs, response):
complete_response = {"choices": [], "model": "", "id": ""}
async for item in response:
yield item
_accumulate_streaming_response(complete_response, item)

_set_response_attributes(span, complete_response)

_set_token_usage(span, request_kwargs, complete_response)
try:
async for item in response:
yield item
_accumulate_streaming_response(complete_response, item)

span.set_status(Status(StatusCode.OK))
except Exception as e:
if span.is_recording():
span.set_attribute(ERROR_TYPE, e.__class__.__name__)
span.record_exception(e)
span.set_status(Status(StatusCode.ERROR, str(e)))
raise
finally:
_set_response_attributes(span, complete_response)
_set_token_usage(span, request_kwargs, complete_response)

if should_emit_events():
_emit_streaming_response_events(complete_response)
else:
if should_send_prompts():
_set_completions(span, complete_response.get("choices"))
if should_emit_events():
_emit_streaming_response_events(complete_response)
else:
if should_send_prompts():
_set_completions(span, complete_response.get("choices"))

span.set_status(Status(StatusCode.OK))
span.end()
span.end()


def _emit_streaming_response_events(complete_response):
Expand Down