| |
| |
| |
| |
| |
| |
|
|
| use std::sync::Arc; |
|
|
| use futures::FutureExt; |
| use futures::future::BoxFuture; |
| use tokio::sync::mpsc; |
|
|
| use super::HttpResponseBodyStream; |
| use super::response_body_stream::HttpBodyStreamRegistration; |
| use crate::HttpClient; |
| use crate::client::ExecServerClient; |
| use crate::client::ExecServerError; |
| use crate::protocol::HTTP_REQUEST_METHOD; |
| use crate::protocol::HttpRequestParams; |
| use crate::protocol::HttpRequestResponse; |
|
|
| |
| const HTTP_BODY_DELTA_CHANNEL_CAPACITY: usize = 256; |
|
|
| impl ExecServerClient { |
| |
| pub async fn http_request( |
| &self, |
| mut params: HttpRequestParams, |
| ) -> Result<HttpRequestResponse, ExecServerError> { |
| params.stream_response = false; |
| self.call(HTTP_REQUEST_METHOD, ¶ms).await |
| } |
|
|
| |
| |
| |
| |
| |
| pub async fn http_request_stream( |
| &self, |
| mut params: HttpRequestParams, |
| ) -> Result<(HttpRequestResponse, HttpResponseBodyStream), ExecServerError> { |
| let rpc_client = self.rpc_client().await?; |
| params.stream_response = true; |
| let request_id = self.inner.next_http_body_stream_request_id(); |
| params.request_id = request_id.clone(); |
| let (tx, rx) = mpsc::channel(HTTP_BODY_DELTA_CHANNEL_CAPACITY); |
| self.inner |
| .insert_http_body_stream(request_id.clone(), tx) |
| .await?; |
| let mut registration = |
| HttpBodyStreamRegistration::new(Arc::clone(&self.inner), request_id.clone()); |
| let response = match self |
| .call_rpc(&rpc_client, HTTP_REQUEST_METHOD, ¶ms) |
| .await |
| { |
| Ok(response) => response, |
| Err(error) => { |
| self.inner.remove_http_body_stream(&request_id).await; |
| registration.disarm(); |
| return Err(error); |
| } |
| }; |
| registration.disarm(); |
| Ok(( |
| response, |
| HttpResponseBodyStream::remote(Arc::clone(&self.inner), request_id, rx), |
| )) |
| } |
| } |
|
|
| impl HttpClient for ExecServerClient { |
| |
| |
| fn http_request( |
| &self, |
| params: HttpRequestParams, |
| ) -> BoxFuture<'_, Result<HttpRequestResponse, ExecServerError>> { |
| async move { ExecServerClient::http_request(self, params).await }.boxed() |
| } |
|
|
| |
| |
| fn http_request_stream( |
| &self, |
| params: HttpRequestParams, |
| ) -> BoxFuture<'_, Result<(HttpRequestResponse, HttpResponseBodyStream), ExecServerError>> { |
| async move { ExecServerClient::http_request_stream(self, params).await }.boxed() |
| } |
| } |
|
|