Execute a fully-built streaming chat-completions request (sets `stream`).
(
&self,
mut request: serde_json::Value,
cancel_token: tokio_util::sync::CancellationToken,
)
| 542 | impl OpenAiClient { |
| 543 | /// Execute a fully-built streaming chat-completions request (sets `stream`). |
| 544 | async fn send_streaming( |
| 545 | &self, |
| 546 | mut request: serde_json::Value, |
| 547 | cancel_token: tokio_util::sync::CancellationToken, |
| 548 | ) -> Result<mpsc::Receiver<StreamEvent>> { |
| 549 | { |
| 550 | request["stream"] = serde_json::json!(true); |
| 551 | request["stream_options"] = serde_json::json!({ "include_usage": true }); |
| 552 | let request_started_at = Instant::now(); |
| 553 | let url = format!("{}{}", self.base_url, self.chat_completions_path); |
| 554 | let request_headers = self.request_headers(); |
| 555 | |
| 556 | let streaming_resp = crate::retry::with_retry(&self.retry_config, |_attempt| { |
| 557 | let http = &self.http; |
| 558 | let url = &url; |
| 559 | let request_headers = request_headers.clone(); |
| 560 | let request = &request; |
| 561 | let cancel_token = cancel_token.clone(); |
| 562 | async move { |
| 563 | let headers = request_headers |
| 564 | .iter() |
| 565 | .map(|(key, value)| (key.as_str(), value.as_str())) |
| 566 | .collect::<Vec<_>>(); |
| 567 | // Wrap in tokio::select! so cancellation aborts the HTTP request mid-flight |
| 568 | let resp = tokio::select! { |
| 569 | _ = cancel_token.cancelled() => { |
| 570 | return AttemptOutcome::Fatal(anyhow::anyhow!("HTTP request cancelled")); |
| 571 | } |
| 572 | result = http.post_streaming(url, headers, request, cancel_token.clone()) => { |
| 573 | match result { |
| 574 | Ok(r) => r, |
| 575 | Err(e) => { |
| 576 | // Transient network error (timeout, reset, |
| 577 | // mid-flight drop — common on throttled |
| 578 | // endpoints): retry with backoff like 429/5xx |
| 579 | // instead of failing the turn. GLM and other |
| 580 | // OpenAI-compatible endpoints hit this most. |
| 581 | return if crate::retry::is_transient_error(&e) { |
| 582 | AttemptOutcome::Retryable { |
| 583 | status: reqwest::StatusCode::SERVICE_UNAVAILABLE, |
| 584 | body: format!("network error: {e}"), |
| 585 | retry_after: None, |
| 586 | } |
| 587 | } else { |
| 588 | AttemptOutcome::Fatal(anyhow::anyhow!( |
| 589 | "HTTP request failed: {}", |
| 590 | e |
| 591 | )) |
| 592 | }; |
| 593 | } |
| 594 | } |
| 595 | } |
| 596 | }; |
| 597 | let status = reqwest::StatusCode::from_u16(resp.status) |
| 598 | .unwrap_or(reqwest::StatusCode::INTERNAL_SERVER_ERROR); |
| 599 | if status.is_success() { |
| 600 | AttemptOutcome::Success(resp) |
| 601 | } else { |
no test coverage detected