From aa3f106fb96679cffdb998aa5d1ebc1c8ffed863 Mon Sep 17 00:00:00 2001 From: Alex Cheema <41707476+AlexCheema@users.noreply.github.com> Date: Thu, 19 Feb 2026 05:40:24 -0800 Subject: [PATCH] fix: import ResponsesStreamEvent and DRY up SSE formatting (#1499) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit ## 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 --- 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)