use bytes::Bytes;
use serde::Serialize;
use super::super::super::types::{
ResponsesCreateResponse, ResponsesFunctionCallOutputItem, ResponsesMessageOutputItem,
ResponsesOutputItem, ResponsesStreamEvent,
};
pub(super) fn created_event(
response_id: String,
model: String,
created_at: u64,
) -> ResponsesStreamEvent<ResponsesCreateResponse> {
ResponsesStreamEvent {
event_type: "response.created".to_string(),
response: Some(ResponsesCreateResponse {
id: response_id,
object: "response".to_string(),
created_at,
model,
status: "in_progress".to_string(),
output: Vec::new(),
usage: None,
}),
..Default::default()
}
}
pub(super) fn output_text_delta_event(
response_id: &str,
item_id: &str,
delta: String,
) -> ResponsesStreamEvent<ResponsesCreateResponse> {
ResponsesStreamEvent {
event_type: "response.output_text.delta".to_string(),
response_id: Some(response_id.to_string()),
item_id: Some(item_id.to_string()),
output_index: Some(0),
content_index: Some(0),
delta: Some(delta),
..Default::default()
}
}
pub(super) fn completed_event(
response: ResponsesCreateResponse,
) -> ResponsesStreamEvent<ResponsesCreateResponse> {
ResponsesStreamEvent {
event_type: "response.completed".to_string(),
response: Some(response),
..Default::default()
}
}
fn item_event(
event_type: &str,
response_id: &str,
item_id: &str,
output_index: u32,
) -> ResponsesStreamEvent<ResponsesCreateResponse> {
ResponsesStreamEvent {
event_type: event_type.to_string(),
response_id: Some(response_id.to_string()),
item_id: Some(item_id.to_string()),
output_index: Some(output_index),
..Default::default()
}
}
pub(super) fn message_item_added_event(
response_id: &str,
message_id: &str,
) -> ResponsesStreamEvent<ResponsesCreateResponse> {
let mut event = item_event("response.output_item.added", response_id, message_id, 0);
event.item = Some(ResponsesOutputItem::Message(ResponsesMessageOutputItem {
id: message_id.to_string(),
item_type: "message".to_string(),
role: "assistant".to_string(),
content: Vec::new(),
status: Some("in_progress".to_string()),
}));
event
}
pub(super) fn message_item_done_event(
response_id: &str,
item: &ResponsesMessageOutputItem,
) -> ResponsesStreamEvent<ResponsesCreateResponse> {
let mut event = item_event("response.output_item.done", response_id, &item.id, 0);
event.item = Some(ResponsesOutputItem::Message(item.clone()));
event
}
pub(super) fn function_call_item_events(
response_id: &str,
item: &ResponsesFunctionCallOutputItem,
output_index: u32,
) -> Vec<ResponsesStreamEvent<ResponsesCreateResponse>> {
let mut added = item_event(
"response.output_item.added",
response_id,
&item.id,
output_index,
);
added.item = Some(ResponsesOutputItem::FunctionCall(
ResponsesFunctionCallOutputItem {
arguments: String::new(),
status: Some("in_progress".to_string()),
..item.clone()
},
));
let mut arguments_delta = item_event(
"response.function_call_arguments.delta",
response_id,
&item.id,
output_index,
);
arguments_delta.delta = Some(item.arguments.clone());
let mut arguments_done = item_event(
"response.function_call_arguments.done",
response_id,
&item.id,
output_index,
);
arguments_done.arguments = Some(item.arguments.clone());
let mut done = item_event(
"response.output_item.done",
response_id,
&item.id,
output_index,
);
done.item = Some(ResponsesOutputItem::FunctionCall(
ResponsesFunctionCallOutputItem {
status: Some("completed".to_string()),
..item.clone()
},
));
vec![added, arguments_delta, arguments_done, done]
}
pub(super) fn event_to_sse_bytes<T: Serialize>(event: &ResponsesStreamEvent<T>) -> Bytes {
let payload = serde_json::to_string(event).unwrap_or_else(|_| "{}".to_string());
Bytes::from(format!(
"event: {}\ndata: {}\n\n",
event.event_type, payload
))
}
pub(super) fn failed_sse_bytes(message: &str) -> Bytes {
let payload = serde_json::json!({
"type": "response.failed",
"response": {
"status": "failed",
"error": { "message": message, "type": "api_error" },
},
});
Bytes::from(format!("event: response.failed\ndata: {payload}\n\n"))
}
pub(super) fn done_sse_bytes() -> Bytes {
Bytes::from_static(b"data: [DONE]\n\n")
}