#![allow(dead_code)]
use futures::StreamExt;
use open_agent::{AgentOptions, ContentBlock, FinishReason, Message, StreamEvent, query};
use wiremock::matchers::{method, path};
use wiremock::{Mock, MockServer, ResponseTemplate};
pub async fn sse_server(body: impl Into<String>) -> MockServer {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/v1/chat/completions"))
.respond_with(
ResponseTemplate::new(200)
.insert_header("content-type", "text/event-stream")
.set_body_string(body.into()),
)
.mount(&server)
.await;
server
}
pub fn options_for(server: &MockServer) -> AgentOptions {
options_with_reasoning(server, false)
}
pub fn options_with_reasoning(server: &MockServer, include_reasoning: bool) -> AgentOptions {
AgentOptions::builder()
.model("m")
.base_url(format!("{}/v1", server.uri()))
.api_key("k")
.include_reasoning(include_reasoning)
.build()
.expect("minimal options build")
}
pub fn text_chunk(content: &str, finish: Option<&str>) -> String {
sse_frame(serde_json::json!({
"index": 0,
"delta": { "content": content },
"finish_reason": finish,
}))
}
pub fn tool_chunk(id: &str, name: &str, arguments: &str) -> String {
sse_frame(serde_json::json!({
"index": 0,
"delta": {
"tool_calls": [{
"index": 0,
"id": id,
"type": "function",
"function": { "name": name, "arguments": arguments },
}],
},
"finish_reason": serde_json::Value::Null,
}))
}
pub fn sse_frame(choice: serde_json::Value) -> String {
let chunk = serde_json::json!({
"id": "1",
"object": "chat.completion.chunk",
"created": 0,
"model": "m",
"choices": [choice],
});
format!("data: {chunk}\n\n")
}
pub const DONE: &str = "data: [DONE]\n\n";
pub async fn anthropic_sse_server(body: impl Into<String>) -> MockServer {
let server = MockServer::start().await;
Mock::given(method("POST"))
.and(path("/v1/messages"))
.respond_with(
ResponseTemplate::new(200)
.insert_header("content-type", "text/event-stream")
.set_body_string(body.into()),
)
.mount(&server)
.await;
server
}
pub fn anthropic_frame(event: &str, data: serde_json::Value) -> String {
format!("event: {event}\ndata: {data}\n\n")
}
pub fn anthropic_text_response(text: &str, stop_reason: &str) -> String {
[
anthropic_frame(
"message_start",
serde_json::json!({
"type": "message_start",
"message": {
"id": "msg_1", "type": "message", "role": "assistant", "content": [],
"model": "m", "stop_reason": null, "stop_sequence": null,
"usage": { "input_tokens": 1, "output_tokens": 1 },
},
}),
),
anthropic_frame(
"content_block_start",
serde_json::json!({
"type": "content_block_start", "index": 0,
"content_block": { "type": "text", "text": "" },
}),
),
anthropic_frame(
"content_block_delta",
serde_json::json!({
"type": "content_block_delta", "index": 0,
"delta": { "type": "text_delta", "text": text },
}),
),
anthropic_frame(
"content_block_stop",
serde_json::json!({ "type": "content_block_stop", "index": 0 }),
),
anthropic_frame(
"message_delta",
serde_json::json!({
"type": "message_delta",
"delta": { "stop_reason": stop_reason, "stop_sequence": null },
"usage": { "output_tokens": 2 },
}),
),
anthropic_frame(
"message_stop",
serde_json::json!({ "type": "message_stop" }),
),
]
.concat()
}
pub fn text_of(blocks: &[ContentBlock]) -> String {
blocks
.iter()
.filter_map(|block| match block {
ContentBlock::Text(text) => Some(text.text.as_str()),
_ => None,
})
.collect()
}
pub fn message_text(message: &Message) -> String {
text_of(&message.content)
}
pub fn reasoning_chunk(field: &str, value: &str, finish: Option<&str>) -> String {
sse_frame(serde_json::json!({
"index": 0,
"delta": { field: value },
"finish_reason": finish,
}))
}
pub fn mixed_chunk(reasoning_field: &str, reasoning: &str, content: &str) -> String {
sse_frame(serde_json::json!({
"index": 0,
"delta": { reasoning_field: reasoning, "content": content },
"finish_reason": serde_json::Value::Null,
}))
}
pub async fn collect_events(body: String, include_reasoning: bool) -> Vec<StreamEvent> {
let server = sse_server(body).await;
let options = options_with_reasoning(&server, include_reasoning);
let mut stream = query("hi", &options).await.expect("start query");
let mut events = Vec::new();
while let Some(event) = stream.next().await {
events.push(event.expect("stream yields no errors"));
}
events
}
pub fn blocks_of(events: &[StreamEvent]) -> Vec<ContentBlock> {
events
.iter()
.filter_map(|event| event.as_block().cloned())
.collect()
}
pub fn reasoning_of(events: &[StreamEvent]) -> String {
events
.iter()
.filter_map(StreamEvent::as_reasoning)
.collect()
}
pub fn text_of_events(events: &[StreamEvent]) -> String {
events.iter().filter_map(StreamEvent::as_text).collect()
}
pub fn sole_finish_reason(events: &[StreamEvent]) -> FinishReason {
let reasons: Vec<&FinishReason> = events
.iter()
.filter_map(StreamEvent::finish_reason)
.collect();
assert_eq!(
reasons.len(),
1,
"expected exactly one Finish event, got {events:?}"
);
assert!(
events
.last()
.is_some_and(|event| event.finish_reason().is_some()),
"Finish must be the last event, got {events:?}"
);
reasons[0].clone()
}