[object Object]

← back to Exo

fix: import ResponsesStreamEvent and DRY up SSE formatting (#1499)

aa3f106fb96679cffdb998aa5d1ebc1c8ffed863 · 2026-02-19 05:40:24 -0800 · Alex Cheema

## Summary
- `ResponsesStreamEvent` was defined in `openai_responses.py` as a union
of all 11 streaming event types but never imported or used anywhere in
the codebase
- Import it in the responses adapter and add a `_format_sse(event:
ResponsesStreamEvent) -> str` helper
- Replace 13 hardcoded `f"event: {type}\ndata:
{event.model_dump_json()}\n\n"` strings with `_format_sse()` calls

## Test plan
- [x] `uv run basedpyright` — 0 errors
- [x] `uv run ruff check` — all checks passed
- [x] `nix fmt` — 0 files changed
- [x] `uv run pytest` — 188 passed, 1 skipped

🤖 Generated with [Claude Code](https://claude.com/claude-code)

Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>

Files touched

Diff

commit aa3f106fb96679cffdb998aa5d1ebc1c8ffed863
Author: Alex Cheema <41707476+AlexCheema@users.noreply.github.com>
Date:   Thu Feb 19 05:40:24 2026 -0800

    fix: import ResponsesStreamEvent and DRY up SSE formatting (#1499)
    
    ## Summary
    - `ResponsesStreamEvent` was defined in `openai_responses.py` as a union
    of all 11 streaming event types but never imported or used anywhere in
    the codebase
    - Import it in the responses adapter and add a `_format_sse(event:
    ResponsesStreamEvent) -> str` helper
    - Replace 13 hardcoded `f"event: {type}\ndata:
    {event.model_dump_json()}\n\n"` strings with `_format_sse()` calls
    
    ## Test plan
    - [x] `uv run basedpyright` — 0 errors
    - [x] `uv run ruff check` — all checks passed
    - [x] `nix fmt` — 0 files changed
    - [x] `uv run pytest` — 188 passed, 1 skipped
    
    🤖 Generated with [Claude Code](https://claude.com/claude-code)
    
    Co-authored-by: Claude Opus 4.6 <noreply@anthropic.com>
---
 src/exo/master/adapters/responses.py | 32 +++++++++++++++++++-------------
 1 file changed, 19 insertions(+), 13 deletions(-)

diff --git a/src/exo/master/adapters/responses.py b/src/exo/master/adapters/responses.py
index e55160b6..90fa7732 100644
--- a/src/exo/master/adapters/responses.py
+++ b/src/exo/master/adapters/responses.py
@@ -31,6 +31,7 @@ from exo.shared.types.openai_responses import (
     ResponseOutputText,
     ResponsesRequest,
     ResponsesResponse,
+    ResponsesStreamEvent,
     ResponseTextDeltaEvent,
     ResponseTextDoneEvent,
     ResponseUsage,
@@ -38,6 +39,11 @@ from exo.shared.types.openai_responses import (
 from exo.shared.types.text_generation import InputMessage, TextGenerationTaskParams
 
 
+def _format_sse(event: ResponsesStreamEvent) -> str:
+    """Format a streaming event as an SSE message."""
+    return f"event: {event.type}\ndata: {event.model_dump_json()}\n\n"
+
+
 def _extract_content(content: str | list[ResponseContentPart]) -> str:
     """Extract plain text from a content field that may be a string or list of parts."""
     if isinstance(content, str):
@@ -219,13 +225,13 @@ async def generate_responses_stream(
     created_event = ResponseCreatedEvent(
         sequence_number=next(seq), response=initial_response
     )
-    yield f"event: response.created\ndata: {created_event.model_dump_json()}\n\n"
+    yield _format_sse(created_event)
 
     # response.in_progress
     in_progress_event = ResponseInProgressEvent(
         sequence_number=next(seq), response=initial_response
     )
-    yield f"event: response.in_progress\ndata: {in_progress_event.model_dump_json()}\n\n"
+    yield _format_sse(in_progress_event)
 
     # response.output_item.added
     initial_item = ResponseMessageItem(
@@ -236,7 +242,7 @@ async def generate_responses_stream(
     item_added = ResponseOutputItemAddedEvent(
         sequence_number=next(seq), output_index=0, item=initial_item
     )
-    yield f"event: response.output_item.added\ndata: {item_added.model_dump_json()}\n\n"
+    yield _format_sse(item_added)
 
     # response.content_part.added
     initial_part = ResponseOutputText(text="")
@@ -247,7 +253,7 @@ async def generate_responses_stream(
         content_index=0,
         part=initial_part,
     )
-    yield f"event: response.content_part.added\ndata: {part_added.model_dump_json()}\n\n"
+    yield _format_sse(part_added)
 
     accumulated_text = ""
     function_call_items: list[ResponseFunctionCallItem] = []
@@ -281,7 +287,7 @@ async def generate_responses_stream(
                     output_index=next_output_index,
                     item=fc_item,
                 )
-                yield f"event: response.output_item.added\ndata: {fc_added.model_dump_json()}\n\n"
+                yield _format_sse(fc_added)
 
                 # response.function_call_arguments.delta
                 args_delta = ResponseFunctionCallArgumentsDeltaEvent(
@@ -290,7 +296,7 @@ async def generate_responses_stream(
                     output_index=next_output_index,
                     delta=tool.arguments,
                 )
-                yield f"event: response.function_call_arguments.delta\ndata: {args_delta.model_dump_json()}\n\n"
+                yield _format_sse(args_delta)
 
                 # response.function_call_arguments.done
                 args_done = ResponseFunctionCallArgumentsDoneEvent(
@@ -300,7 +306,7 @@ async def generate_responses_stream(
                     name=tool.name,
                     arguments=tool.arguments,
                 )
-                yield f"event: response.function_call_arguments.done\ndata: {args_done.model_dump_json()}\n\n"
+                yield _format_sse(args_done)
 
                 # response.output_item.done
                 fc_done_item = ResponseFunctionCallItem(
@@ -315,7 +321,7 @@ async def generate_responses_stream(
                     output_index=next_output_index,
                     item=fc_done_item,
                 )
-                yield f"event: response.output_item.done\ndata: {fc_item_done.model_dump_json()}\n\n"
+                yield _format_sse(fc_item_done)
 
                 function_call_items.append(fc_done_item)
                 next_output_index += 1
@@ -331,7 +337,7 @@ async def generate_responses_stream(
             content_index=0,
             delta=chunk.text,
         )
-        yield f"event: response.output_text.delta\ndata: {delta_event.model_dump_json()}\n\n"
+        yield _format_sse(delta_event)
 
     # response.output_text.done
     text_done = ResponseTextDoneEvent(
@@ -341,7 +347,7 @@ async def generate_responses_stream(
         content_index=0,
         text=accumulated_text,
     )
-    yield f"event: response.output_text.done\ndata: {text_done.model_dump_json()}\n\n"
+    yield _format_sse(text_done)
 
     # response.content_part.done
     final_part = ResponseOutputText(text=accumulated_text)
@@ -352,7 +358,7 @@ async def generate_responses_stream(
         content_index=0,
         part=final_part,
     )
-    yield f"event: response.content_part.done\ndata: {part_done.model_dump_json()}\n\n"
+    yield _format_sse(part_done)
 
     # response.output_item.done
     final_message_item = ResponseMessageItem(
@@ -363,7 +369,7 @@ async def generate_responses_stream(
     item_done = ResponseOutputItemDoneEvent(
         sequence_number=next(seq), output_index=0, item=final_message_item
     )
-    yield f"event: response.output_item.done\ndata: {item_done.model_dump_json()}\n\n"
+    yield _format_sse(item_done)
 
     # Create usage from usage data if available
     usage = None
@@ -388,4 +394,4 @@ async def generate_responses_stream(
     completed_event = ResponseCompletedEvent(
         sequence_number=next(seq), response=final_response
     )
-    yield f"event: response.completed\ndata: {completed_event.model_dump_json()}\n\n"
+    yield _format_sse(completed_event)

← 2e296051 fix: finalize cancel tasks (#1498)  ·  back to Exo  ·  bench: restore --danger-delete-downloads planning phase (#15 42e1e732 →