1use crate::completion;
11use crate::driver::Model;
12use crate::error::{EncodeError, ProviderError};
13use crate::http_client::{self, NoBody};
14use crate::operation::{Completion, Turn};
15use crate::providers::openai::responses_api::streaming::{
16 ItemChunk, ResponseChunk, ResponseChunkKind, ResponsesDecoder, StreamingCompletionChunk,
17 classify_responses_frame,
18};
19use crate::providers::openai::responses_api::wire::Responses;
20use crate::streaming::Item;
21use crate::wire::{Flow, Reply, Shared, Wire, WireFrame};
22use crate::ws_client::{
23 BoxedWebSocketConnection, ConnectOptions, Frame, WebSocketClientExt, WebSocketConnection,
24};
25use serde::{Deserialize, Serialize};
26use serde_json::{Map, Value};
27use std::time::Duration;
28
29use crate::providers::openai::responses_api::{CompletionResponse, ResponseStatus};
30
31const WEBSOCKET_PATH: &str = "responses";
33
34const DEFAULT_CONNECT_TIMEOUT: Duration = Duration::from_secs(30);
35
36const REQUEST_ID_HEADER: Option<&'static str> =
38 crate::providers::openai::wire::OPENAI.request_id_header;
39
40#[derive(Debug, Clone, Default, Serialize, Deserialize)]
42pub struct ResponsesWebSocketCreateOptions {
43 #[serde(skip_serializing_if = "Option::is_none")]
47 pub generate: Option<bool>,
48}
49
50impl ResponsesWebSocketCreateOptions {
51 #[must_use]
53 pub fn warmup() -> Self {
54 Self {
55 generate: Some(false),
56 }
57 }
58}
59
60#[derive(Debug, Clone, Serialize)]
61struct ResponsesWebSocketClientEvent {
62 #[serde(rename = "type")]
63 kind: ResponsesWebSocketClientEventKind,
64 #[serde(flatten)]
65 request: crate::providers::openai::responses_api::CompletionRequest,
66 #[serde(skip_serializing_if = "Option::is_none")]
67 generate: Option<bool>,
68}
69
70#[derive(Debug, Clone, Serialize)]
71enum ResponsesWebSocketClientEventKind {
72 #[serde(rename = "response.create")]
73 ResponseCreate,
74}
75
76#[derive(Debug, Clone, Serialize, Deserialize)]
78pub struct ResponsesWebSocketErrorEvent {
79 #[serde(rename = "type")]
81 pub kind: ResponsesWebSocketErrorEventKind,
82 pub error: ResponsesWebSocketErrorPayload,
84}
85
86impl std::fmt::Display for ResponsesWebSocketErrorEvent {
87 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
88 self.error.fmt(f)
89 }
90}
91
92#[derive(Debug, Clone, Serialize, Deserialize)]
94pub enum ResponsesWebSocketErrorEventKind {
95 #[serde(rename = "error")]
96 Error,
97}
98
99#[derive(Debug, Clone, Default, Serialize, Deserialize)]
101pub struct ResponsesWebSocketErrorPayload {
102 #[serde(skip_serializing_if = "Option::is_none")]
104 pub code: Option<String>,
105 #[serde(skip_serializing_if = "Option::is_none")]
107 pub message: Option<String>,
108 #[serde(flatten, default)]
110 pub extra: Map<String, Value>,
111}
112
113impl std::fmt::Display for ResponsesWebSocketErrorPayload {
114 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
115 match (&self.code, &self.message) {
116 (Some(code), Some(message)) => write!(f, "{code}: {message}"),
117 (None, Some(message)) => f.write_str(message),
118 (Some(code), None) => f.write_str(code),
119 (None, None) => f.write_str("OpenAI websocket error"),
120 }
121 }
122}
123
124#[derive(Debug, Clone, Serialize, Deserialize)]
126pub struct ResponsesWebSocketDoneEvent {
127 #[serde(rename = "type")]
129 pub kind: ResponsesWebSocketDoneEventKind,
130 pub response: Value,
132}
133
134impl ResponsesWebSocketDoneEvent {
135 #[must_use]
137 pub fn response_id(&self) -> Option<&str> {
138 self.response.get("id").and_then(Value::as_str)
139 }
140
141 fn status(&self) -> Option<ResponseStatus> {
142 self.response
143 .get("status")
144 .cloned()
145 .and_then(|status| serde_json::from_value(status).ok())
146 }
147
148 fn as_completion_response(&self) -> Option<CompletionResponse> {
149 serde_json::from_value(self.response.clone()).ok()
150 }
151}
152
153#[derive(Debug, Clone, Serialize, Deserialize)]
155pub enum ResponsesWebSocketDoneEventKind {
156 #[serde(rename = "response.done")]
157 ResponseDone,
158}
159
160#[derive(Debug, Clone)]
162pub enum ResponsesWebSocketEvent {
163 Response(ResponseChunk),
165 Item(ItemChunk),
167 Error(ResponsesWebSocketErrorEvent),
169 Done(ResponsesWebSocketDoneEvent),
171 Unknown(crate::streaming::UnknownPayload),
173}
174
175impl ResponsesWebSocketEvent {
176 #[must_use]
178 pub fn response_id(&self) -> Option<&str> {
179 match self {
180 Self::Response(chunk) => Some(&chunk.response.id),
181 Self::Done(done) => done.response_id(),
182 Self::Item(_) | Self::Error(_) | Self::Unknown(_) => None,
183 }
184 }
185
186 #[must_use]
188 pub fn is_terminal(&self) -> bool {
189 match self {
190 Self::Response(chunk) => matches!(
191 chunk.kind,
192 ResponseChunkKind::ResponseCompleted
193 | ResponseChunkKind::ResponseFailed
194 | ResponseChunkKind::ResponseIncomplete
195 ),
196 Self::Error(_) | Self::Done(_) => true,
197 Self::Item(_) | Self::Unknown(_) => false,
198 }
199 }
200}
201
202pub struct ResponsesWebSocketSessionBuilder {
207 wire: Responses,
208 connect_timeout: Option<Duration>,
209 event_timeout: Option<Duration>,
210}
211
212impl ResponsesWebSocketSessionBuilder {
213 pub(crate) fn new(wire: Responses) -> Self {
214 Self {
215 wire,
216 connect_timeout: Some(DEFAULT_CONNECT_TIMEOUT),
217 event_timeout: None,
218 }
219 }
220
221 #[must_use]
223 pub fn connect_timeout(mut self, timeout: Duration) -> Self {
224 self.connect_timeout = Some(timeout);
225 self
226 }
227
228 #[must_use]
230 pub fn without_connect_timeout(mut self) -> Self {
231 self.connect_timeout = None;
232 self
233 }
234
235 #[must_use]
237 pub fn event_timeout(mut self, timeout: Duration) -> Self {
238 self.event_timeout = Some(timeout);
239 self
240 }
241
242 #[must_use]
244 pub fn without_event_timeout(mut self) -> Self {
245 self.event_timeout = None;
246 self
247 }
248}
249
250impl ResponsesWebSocketSessionBuilder {
251 #[cfg(all(feature = "tungstenite", not(target_family = "wasm")))]
254 #[cfg_attr(docsrs, doc(cfg(feature = "tungstenite")))]
255 pub async fn connect(self) -> Result<ResponsesWebSocketSession, ProviderError> {
256 self.connect_with(&rig_tungstenite::TungsteniteClient::new())
257 .await
258 }
259
260 pub async fn connect_with<W>(
263 self,
264 backend: &W,
265 ) -> Result<ResponsesWebSocketSession, ProviderError>
266 where
267 W: WebSocketClientExt,
268 {
269 ResponsesWebSocketSession::connect_with_timeouts(
270 backend,
271 self.wire,
272 self.connect_timeout,
273 self.event_timeout,
274 )
275 .await
276 }
277}
278
279pub struct ResponsesWebSocketSession {
283 wire: Responses,
284 previous_response_id: Option<String>,
285 pending_done_response_id: Option<String>,
286 socket: BoxedWebSocketConnection,
287 in_flight: bool,
288 event_timeout: Option<Duration>,
289 closed: bool,
290 failed: bool,
291}
292
293impl ResponsesWebSocketSession {
294 async fn connect_with_timeouts<W>(
295 backend: &W,
296 wire: Responses,
297 connect_timeout: Option<Duration>,
298 event_timeout: Option<Duration>,
299 ) -> Result<Self, ProviderError>
300 where
301 W: WebSocketClientExt,
302 {
303 let request = websocket_request(&wire)?;
304 let socket = backend
305 .connect(request, ConnectOptions::new().with_timeout(connect_timeout))
306 .await
307 .map_err(websocket_provider_error)?;
308
309 Ok(Self::from_connection(wire, socket, event_timeout))
310 }
311
312 pub fn from_connection(
315 wire: Responses,
316 connection: BoxedWebSocketConnection,
317 event_timeout: Option<Duration>,
318 ) -> Self {
319 Self {
320 wire,
321 previous_response_id: None,
322 pending_done_response_id: None,
323 socket: connection,
324 in_flight: false,
325 event_timeout,
326 closed: false,
327 failed: false,
328 }
329 }
330
331 #[must_use]
333 pub fn previous_response_id(&self) -> Option<&str> {
334 self.previous_response_id.as_deref()
335 }
336
337 pub fn clear_previous_response_id(&mut self) {
339 self.previous_response_id = None;
340 }
341
342 pub async fn send(
344 &mut self,
345 completion_request: crate::completion::CompletionRequest,
346 ) -> Result<(), ProviderError> {
347 self.send_with_options(
348 completion_request,
349 ResponsesWebSocketCreateOptions::default(),
350 )
351 .await
352 }
353
354 pub async fn send_with_options(
356 &mut self,
357 completion_request: crate::completion::CompletionRequest,
358 options: ResponsesWebSocketCreateOptions,
359 ) -> Result<(), ProviderError> {
360 self.ensure_open()?;
361
362 if self.in_flight {
363 return Err(ProviderError::Provider(
364 "An OpenAI websocket response is already in flight on this session".to_string(),
365 ));
366 }
367
368 let payload = ResponsesWebSocketClientEvent {
369 kind: ResponsesWebSocketClientEventKind::ResponseCreate,
370 request: self.prepare_request(completion_request)?,
371 generate: options.generate,
372 };
373
374 crate::providers::internal::trace_json(
375 crate::providers::internal::LogTarget::Completions,
376 "OpenAI websocket request",
377 &payload,
378 );
379
380 let payload = serde_json::to_string(&payload)?;
381
382 if let Err(error) = self.socket.send(Frame::Text(payload)).await {
383 return Err(self.fail_session(websocket_provider_error(error)));
384 }
385 self.in_flight = true;
386
387 Ok(())
388 }
389
390 pub async fn next_event(&mut self) -> Result<ResponsesWebSocketEvent, ProviderError> {
392 self.next_event_with_payload().await.map(|(event, _)| event)
393 }
394
395 async fn next_event_with_payload(
398 &mut self,
399 ) -> Result<(ResponsesWebSocketEvent, String), ProviderError> {
400 self.ensure_open()?;
401
402 if !self.in_flight {
403 return Err(ProviderError::Provider(
404 "No OpenAI websocket response is currently in flight on this session".to_string(),
405 ));
406 }
407
408 loop {
409 let message = match self.read_next_frame().await? {
410 Ok(message) => message,
411 Err(error) => return Err(self.fail_session(websocket_provider_error(error))),
412 };
413
414 let Some(message) = message else {
415 self.mark_closed();
416 return Err(ProviderError::Provider(
417 "The OpenAI websocket connection closed before the turn finished".to_string(),
418 ));
419 };
420
421 let payload = match websocket_frame_to_text(message) {
422 Ok(Some(payload)) => payload,
423 Ok(None) => continue,
424 Err(error) => return Err(self.fail_session(error)),
425 };
426 let event = match parse_server_event(&payload) {
427 Ok(Some(event)) => event,
428 Ok(None) => continue,
429 Err(error) => return Err(self.fail_session(error)),
430 };
431 if let ResponsesWebSocketEvent::Done(done) = &event {
432 if self.pending_done_response_id.as_deref() == done.response_id() {
435 self.pending_done_response_id = None;
436 continue;
437 }
438 }
439 self.update_state_for_event(&event);
440 return Ok((event, payload));
441 }
442 }
443
444 pub async fn warmup(
446 &mut self,
447 completion_request: crate::completion::CompletionRequest,
448 ) -> Result<String, ProviderError> {
449 self.send_with_options(
450 completion_request,
451 ResponsesWebSocketCreateOptions::warmup(),
452 )
453 .await?;
454 let response = self.wait_for_completed_response().await?;
455 Ok(response.id)
456 }
457
458 pub async fn completion(
461 &mut self,
462 completion_request: crate::completion::CompletionRequest,
463 ) -> Result<completion::CompletionResponse, ProviderError> {
464 let provider = self.wire.describe().name.to_owned();
465 self.send(completion_request).await?;
466 let (response, folded) = self.wait_for_terminal_response().await?;
467 if folded.choice.is_empty() {
468 return super::wire::fold_body(&provider, response);
473 }
474 Ok(folded)
475 }
476
477 pub async fn close(&mut self) -> Result<(), ProviderError> {
482 if self.closed {
483 return Ok(());
484 }
485
486 let result = self
487 .socket
488 .close(None)
489 .await
490 .map_err(websocket_provider_error);
491 self.mark_closed();
492 result
493 }
494
495 fn prepare_request(
496 &self,
497 completion_request: crate::completion::CompletionRequest,
498 ) -> Result<crate::providers::openai::responses_api::CompletionRequest, ProviderError> {
499 completion_request.validate_message_content()?;
501 let (completion_request, issuers) = crate::providers::openai::wire::scope_reasoning(
502 &self.wire.provider.dialect,
503 &self.wire.model,
504 completion_request,
505 )?;
506 let mut request = self
507 .wire
508 .responses_request(completion_request, issuers, false)?;
509
510 request.stream = None;
513 request.additional_parameters.background = None;
514
515 if request.additional_parameters.previous_response_id.is_none() {
516 request
517 .additional_parameters
518 .previous_response_id
519 .clone_from(&self.previous_response_id);
520 }
521
522 Ok(request)
523 }
524
525 async fn wait_for_completed_response(&mut self) -> Result<CompletionResponse, ProviderError> {
526 Ok(self.wait_for_terminal_response().await?.0)
527 }
528
529 async fn wait_for_terminal_response(
534 &mut self,
535 ) -> Result<(CompletionResponse, completion::CompletionResponse), ProviderError> {
536 let provider = self.wire.describe().name.to_owned();
537 let wire = self.wire.clone();
538 let reply = std::sync::Mutex::new(Shared::new(Turn::new(provider.clone())));
541 let mut decoder = wire.decoder();
542 loop {
543 let (event, payload) = self.next_event_with_payload().await?;
544 match event {
545 ResponsesWebSocketEvent::Response(chunk) => {
546 let terminal = matches!(
547 chunk.kind,
548 ResponseChunkKind::ResponseCompleted
549 | ResponseChunkKind::ResponseFailed
550 | ResponseChunkKind::ResponseIncomplete
551 );
552 if !terminal {
553 feed(&mut decoder, &reply, payload)?;
554 continue;
555 }
556 let response = terminal_response_result(chunk.response)?;
560 let ended = feed(&mut decoder, &reply, payload)?;
561 let folded = fold_reply(reply, ended, &provider, &response)?;
562 return Ok((response, folded));
563 }
564 ResponsesWebSocketEvent::Done(done) => {
565 if let Some(response) = done.as_completion_response() {
566 let response = terminal_response_result(response)?;
569 let body = serde_json::to_string(&done.response)?;
573 let ended = feed(&mut decoder, &reply, body)?;
574 let folded = fold_reply(reply, ended, &provider, &response)?;
575 return Ok((response, folded));
576 }
577
578 let message = if let Some(response_id) = done.response_id() {
579 format!(
580 "OpenAI websocket turn ended with response.done before a terminal response body was available (response_id={response_id})"
581 )
582 } else {
583 "OpenAI websocket turn ended with response.done before a terminal response body was available"
584 .to_string()
585 };
586
587 return Err(ProviderError::Provider(message));
588 }
589 ResponsesWebSocketEvent::Error(error) => {
590 return Err(provider_error_from_event(&error));
595 }
596 ResponsesWebSocketEvent::Item(_) | ResponsesWebSocketEvent::Unknown(_) => {
598 feed(&mut decoder, &reply, payload)?;
599 }
600 }
601 }
602 }
603
604 fn update_state_for_event(&mut self, event: &ResponsesWebSocketEvent) {
605 match event {
606 ResponsesWebSocketEvent::Response(chunk) => match chunk.kind {
607 ResponseChunkKind::ResponseCompleted | ResponseChunkKind::ResponseIncomplete => {
611 let response_id = chunk.response.id.clone();
612 self.previous_response_id = Some(response_id.clone());
613 self.pending_done_response_id = Some(response_id);
614 self.in_flight = false;
615 }
616 ResponseChunkKind::ResponseFailed => {
617 self.pending_done_response_id = Some(chunk.response.id.clone());
618 self.previous_response_id = None;
619 self.in_flight = false;
620 }
621 ResponseChunkKind::ResponseCreated | ResponseChunkKind::ResponseInProgress => {}
622 },
623 ResponsesWebSocketEvent::Done(done) => {
624 match done.status() {
625 Some(ResponseStatus::Completed) | Some(ResponseStatus::Incomplete) => {
626 if let Some(response_id) = done.response_id() {
627 self.previous_response_id = Some(response_id.to_string());
628 }
629 }
630 Some(ResponseStatus::Failed)
631 | Some(ResponseStatus::Cancelled)
632 | Some(ResponseStatus::Other(_)) => {
633 self.previous_response_id = None;
634 }
635 Some(ResponseStatus::InProgress | ResponseStatus::Queued) | None => {}
636 }
637 self.pending_done_response_id = None;
638 self.in_flight = false;
639 }
640 ResponsesWebSocketEvent::Error(_) => {
641 self.previous_response_id = None;
642 self.pending_done_response_id = None;
643 self.in_flight = false;
644 }
645 ResponsesWebSocketEvent::Item(_) | ResponsesWebSocketEvent::Unknown(_) => {}
647 }
648 }
649
650 fn abort_turn(&mut self) {
651 self.previous_response_id = None;
652 self.pending_done_response_id = None;
653 self.in_flight = false;
654 }
655
656 fn mark_closed(&mut self) {
657 self.abort_turn();
658 self.closed = true;
659 self.failed = false;
660 }
661
662 fn mark_failed(&mut self) {
663 self.abort_turn();
664 self.failed = true;
665 }
666
667 fn ensure_open(&self) -> Result<(), ProviderError> {
668 if self.closed || self.failed {
669 return Err(ProviderError::Provider(
670 "The OpenAI websocket session is closed".to_string(),
671 ));
672 }
673
674 Ok(())
675 }
676
677 fn fail_session(&mut self, error: ProviderError) -> ProviderError {
678 self.mark_failed();
679 error
680 }
681
682 async fn read_next_frame(
685 &mut self,
686 ) -> Result<http_client::Result<Option<Frame>>, ProviderError> {
687 let Some(timeout_duration) = self.event_timeout else {
688 return Ok(self.socket.recv().await);
689 };
690
691 match crate::wasm_compat::timeout(timeout_duration, self.socket.recv()).await {
692 Ok(message) => Ok(message),
693 Err(_) => Err(self.fail_session(event_timeout_error(timeout_duration))),
694 }
695 }
696}
697
698impl Drop for ResponsesWebSocketSession {
699 fn drop(&mut self) {
700 if !self.closed {
701 tracing::warn!(
702 target: "rig::completions",
703 in_flight = self.in_flight,
704 "Dropping an OpenAI websocket session without calling close(); the connection will end without a close handshake"
705 );
706 }
707 }
708}
709
710fn feed<'id>(
713 decoder: &mut ResponsesDecoder<'id>,
714 reply: &'id std::sync::Mutex<Shared<Completion>>,
715 payload: String,
716) -> Result<bool, ProviderError> {
717 crate::driver::step(decoder, reply, WireFrame::Text(payload))
718 .map(|step| matches!(step, Flow::Ended(_)))
719}
720
721fn fold_reply(
724 reply: std::sync::Mutex<Shared<Completion>>,
725 ended: bool,
726 provider: &str,
727 response: &CompletionResponse,
728) -> Result<completion::CompletionResponse, ProviderError> {
729 let fed = if ended {
730 Ok(())
731 } else {
732 Err(ProviderError::Truncated)
733 };
734 let reply_of = Reply {
735 provider: provider.to_owned(),
736 raw: serde_json::to_value(response)?,
737 provider_request_id: None,
739 };
740 crate::driver::settle(reply, fed, reply_of).outcome
741}
742
743fn terminal_response_result(
744 response: CompletionResponse,
745) -> Result<CompletionResponse, ProviderError> {
746 match response.status {
747 ResponseStatus::Completed => Ok(response),
748 ResponseStatus::Failed => match response.error.as_ref() {
751 Some(error) => Err(ProviderError::from_provider_body(
752 serde_json::to_string(&response).unwrap_or_else(|_| error.message.clone()),
753 )),
754 None => Err(ProviderError::Provider(response_error_message(
755 "failed response",
756 ))),
757 },
758 ResponseStatus::Incomplete => Ok(response),
763 other => Err(ProviderError::Provider(format!(
764 "OpenAI websocket response ended in state {other:?}"
765 ))),
766 }
767}
768
769fn response_error_message(fallback: &str) -> String {
770 format!("OpenAI websocket returned a {fallback}")
771}
772
773fn provider_error_from_event(error: &ResponsesWebSocketErrorEvent) -> ProviderError {
776 ProviderError::from_provider_body(
777 serde_json::to_string(&error).unwrap_or_else(|_| error.to_string()),
778 )
779}
780
781fn parse_server_event(payload: &str) -> Result<Option<ResponsesWebSocketEvent>, ProviderError> {
784 #[derive(Deserialize)]
785 struct EventType {
786 #[serde(rename = "type")]
787 kind: String,
788 }
789
790 let event_type = serde_json::from_str::<EventType>(payload)?;
791 match event_type.kind.as_str() {
792 "error" => serde_json::from_str(payload)
793 .map(|e| Some(ResponsesWebSocketEvent::Error(e)))
794 .map_err(ProviderError::from),
795 "response.done" => serde_json::from_str(payload)
796 .map(|d| Some(ResponsesWebSocketEvent::Done(d)))
797 .map_err(ProviderError::from),
798 _ => Ok(Some(
799 match crate::driver::triage(classify_responses_frame(payload))? {
800 Item::Event(StreamingCompletionChunk::Response(response)) => {
801 ResponsesWebSocketEvent::Response(response)
802 }
803 Item::Event(StreamingCompletionChunk::Delta(item)) => {
804 ResponsesWebSocketEvent::Item(item)
805 }
806 Item::Unknown(value) => ResponsesWebSocketEvent::Unknown(value),
807 },
808 )),
809 }
810}
811
812fn websocket_frame_to_text(frame: Frame) -> Result<Option<String>, ProviderError> {
817 match frame {
818 Frame::Text(text) => Ok(Some(text)),
819 Frame::Binary(bytes) => String::from_utf8(bytes.to_vec())
820 .map(Some)
821 .map_err(|error| ProviderError::Response(error.to_string())),
822 Frame::Ping(_) | Frame::Pong(_) => Ok(None),
823 Frame::Close(frame) => {
824 let reason = frame
825 .map(|frame| frame.reason)
826 .filter(|reason| !reason.is_empty())
827 .unwrap_or_else(|| "without a close reason".to_string());
828 Err(ProviderError::Provider(format!(
829 "The OpenAI websocket connection closed {reason}"
830 )))
831 }
832 }
833}
834
835fn websocket_request(wire: &Responses) -> Result<http_client::Request<NoBody>, EncodeError> {
841 let url = crate::ws_client::websocket_url(&wire.provider.base_url, WEBSOCKET_PATH)
842 .map_err(EncodeError::request)?;
843
844 let request = wire.provider.headers(
845 http_client::Request::builder()
846 .method(http::Method::GET)
847 .uri(url),
848 );
849
850 request.body(NoBody).map_err(|error| {
851 EncodeError::request(format!("Failed to build OpenAI websocket request: {error}"))
852 })
853}
854
855fn event_timeout_error(timeout: Duration) -> ProviderError {
856 ProviderError::Provider(format!(
857 "Timed out waiting for the next OpenAI websocket event after {timeout:?}"
858 ))
859}
860
861fn websocket_provider_error(error: http_client::Error) -> ProviderError {
864 let provider_request_id = error.non_success_headers().and_then(|headers| {
865 crate::providers::internal::request_id_from_headers(headers, REQUEST_ID_HEADER)
866 });
867 ProviderError::from_transport_error(error).with_provider_request_id(provider_request_id)
868}
869
870impl<T> Model<Responses, T> {
871 pub fn responses_websocket(&self) -> ResponsesWebSocketSessionBuilder {
876 ResponsesWebSocketSessionBuilder::new(self.wire.clone())
877 }
878}
879
880#[cfg(not(target_family = "wasm"))]
882const _: fn() = || {
883 fn assert_send_sync<T: Send + Sync>() {}
884 assert_send_sync::<ResponsesWebSocketSession>();
885 assert_send_sync::<ResponsesWebSocketSessionBuilder>();
886};
887
888#[cfg(test)]
889#[allow(clippy::expect_used, clippy::panic)]
890mod tests;