#![cfg(not(target_family = "wasm"))]
#![allow(clippy::expect_used, clippy::panic)]
use rig_core::completion::CompletionRequest;
use rig_core::providers::openai::OpenAIConfig;
use rig_core::test_utils::RecordingHttpClient;
use std::sync::mpsc;
use std::time::Duration;
fn block_on<F: std::future::Future>(future: F) -> F::Output {
futures::executor::block_on(rig_core::wasm_compat::timeout(
std::time::Duration::from_secs(10),
future,
))
.expect("client operation deadline")
}
fn serve_one_turn(events: Vec<String>) -> String {
serve_one_turn_after(None, events)
}
fn serve_one_turn_after(
release: Option<futures::channel::oneshot::Receiver<()>>,
events: Vec<String>,
) -> String {
use futures::{SinkExt, StreamExt};
let (address_tx, address_rx) = mpsc::channel();
std::thread::spawn(move || {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("server runtime should build");
runtime.block_on(async move {
let _ = tokio::time::timeout(Duration::from_secs(30), async move {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind");
address_tx
.send(listener.local_addr().expect("address"))
.expect("address should send");
let (stream, _) = listener.accept().await.expect("accept");
let mut socket = rig_tungstenite::tokio_tungstenite::accept_async(stream)
.await
.expect("upgrade");
let request = socket
.next()
.await
.expect("request should arrive")
.expect("request should be valid");
assert!(
request
.into_text()
.expect("request should be text")
.contains("\"type\":\"response.create\""),
"the session should open the turn with response.create"
);
if let Some(release) = release {
release.await.expect("release events");
}
for event in events {
socket
.send(rig_tungstenite::tokio_tungstenite::tungstenite::Message::text(event))
.await
.expect("event should send");
}
while let Some(Ok(message)) = socket.next().await {
if message.is_close() {
break;
}
}
})
.await;
});
});
let address = address_rx
.recv_timeout(Duration::from_secs(10))
.expect("server should report its address");
format!("http://{address}/v1")
}
#[test]
fn a_whole_session_runs_without_a_tokio_runtime() {
let completed = serde_json::json!({
"type": "response.completed",
"sequence_number": 2,
"response": {
"id": "resp_off_runtime",
"object": "response",
"created_at": 0,
"status": "completed",
"error": null,
"incomplete_details": null,
"instructions": null,
"max_output_tokens": null,
"model": "gpt-5.4",
"usage": null,
"output": [],
"tools": []
}
})
.to_string();
let delta = serde_json::json!({
"type": "response.output_text.delta",
"content_index": 0,
"delta": "off runtime",
"item_id": "msg_1",
"logprobs": [],
"output_index": 0,
"sequence_number": 1
})
.to_string();
let base_url = serve_one_turn(vec![delta, completed]);
assert!(
tokio::runtime::Handle::try_current().is_err(),
"this test is meaningless inside a tokio runtime"
);
block_on(async move {
let bound = OpenAIConfig::new("test-key")
.with_base_url(&base_url)
.connect(RecordingHttpClient::new("{}"))
.responses("gpt-5.4");
let mut session = match bound.responses_websocket().connect().await {
Ok(session) => session,
Err(error) => panic!("session should connect off-runtime: {error}"),
};
let response = session
.completion(CompletionRequest::new("hello"))
.await
.expect("the turn should complete off-runtime");
assert!(
matches!(
response.choice.first(),
Some(rig_core::completion::AssistantContent::Text(text))
if text.text == "off runtime"
),
"the streamed delta should arrive off-runtime, got {:?}",
response.choice
);
session.close().await.expect("close should succeed");
});
}
#[test]
fn an_event_timeout_still_allows_close_without_a_tokio_runtime() {
let base_url = serve_one_turn(Vec::new());
assert!(
tokio::runtime::Handle::try_current().is_err(),
"this test is meaningless inside a tokio runtime"
);
block_on(async move {
let bound = OpenAIConfig::new("test-key")
.with_base_url(&base_url)
.connect(RecordingHttpClient::new("{}"))
.responses("gpt-5.4");
let mut session = match bound
.responses_websocket()
.event_timeout(Duration::from_millis(50))
.connect()
.await
{
Ok(session) => session,
Err(error) => panic!("session should connect off-runtime: {error}"),
};
session
.send(CompletionRequest::new("hello"))
.await
.expect("request should send");
let error = session
.next_event()
.await
.expect_err("a silent server should trip the event timeout");
assert!(
error
.to_string()
.contains("Timed out waiting for the next OpenAI websocket event"),
"expected the event timeout, got {error}"
);
rig_core::wasm_compat::timeout(Duration::from_secs(5), session.close())
.await
.expect("close() must not hang after an event timeout")
.expect("close should succeed");
});
}