diff options
| author | Emīls <emils@mullvad.net> | 2020-07-17 15:36:20 +0100 |
|---|---|---|
| committer | Emīls <emils@mullvad.net> | 2020-07-21 10:50:25 +0100 |
| commit | b19967d5d1a30367aeb44ec0d834befd60538fe8 (patch) | |
| tree | afd859cbb23caf31b37a577a6a1c712aea12f512 /mullvad-rpc/src | |
| parent | 4c40147f34fed91410b887de97794fbb1a3bc90f (diff) | |
| download | mullvadvpn-b19967d5d1a30367aeb44ec0d834befd60538fe8.tar.xz mullvadvpn-b19967d5d1a30367aeb44ec0d834befd60538fe8.zip | |
Move Cancellable et al to talpid-core
Diffstat (limited to 'mullvad-rpc/src')
| -rw-r--r-- | mullvad-rpc/src/rest.rs | 85 |
1 files changed, 8 insertions, 77 deletions
diff --git a/mullvad-rpc/src/rest.rs b/mullvad-rpc/src/rest.rs index c2677c7212..4a8220749b 100644 --- a/mullvad-rpc/src/rest.rs +++ b/mullvad-rpc/src/rest.rs @@ -2,7 +2,7 @@ use futures::{ channel::{mpsc, oneshot}, sink::SinkExt, stream::StreamExt, - FutureExt, TryFutureExt, + TryFutureExt, }; use futures01::Future as OldFuture; use hyper::{ @@ -10,16 +10,8 @@ use hyper::{ header::{self, HeaderValue}, Method, Uri, }; -use std::{ - collections::BTreeMap, - future::Future, - mem, - net::IpAddr, - pin::Pin, - str::FromStr, - task::{Context, Poll}, - time::Duration, -}; +use std::{collections::BTreeMap, future::Future, mem, net::IpAddr, str::FromStr, time::Duration}; +use talpid_core::future_cancel::{CancelErr, CancelHandle, Cancellable}; use tokio::runtime::Handle; pub use hyper::StatusCode; @@ -122,12 +114,10 @@ impl<C: Connect + Clone + Send + Sync + 'static> RequestService<C> { ); let future = async move { - let response = tokio::time::timeout( - timeout, - request_future.into_future().map_err(Error::Cancelled), - ) - .await - .map_err(Error::TimeoutError); + let response = + tokio::time::timeout(timeout, request_future.map_err(Error::Cancelled)) + .await + .map_err(Error::TimeoutError); let response = flatten_result(flatten_result(response)); @@ -219,7 +209,7 @@ impl RequestServiceHandle { }); - rx.map_err(|_| Error::Cancelled(CancelErr(()))).flatten() + rx.map_err(|_| Error::ReceiveError).flatten() } /// Spawns a future on the RPC runtime. @@ -425,65 +415,6 @@ impl RequestFactory { } -#[derive(Debug)] -pub struct CancelErr(()); - -pub struct Cancellable<F: Future> { - rx: oneshot::Receiver<()>, - f: F, -} - -pub struct CancelHandle { - tx: oneshot::Sender<()>, -} - -impl CancelHandle { - pub fn cancel(self) { - let _ = self.tx.send(()); - } -} - - -impl<F> Cancellable<F> -where - F: Future, -{ - pub fn new(f: F) -> (Self, CancelHandle) { - let (tx, rx) = oneshot::channel(); - (Self { f, rx }, CancelHandle { tx }) - } - - async fn into_future(self) -> std::result::Result<F::Output, CancelErr> { - futures::select! { - _cancelled = self.rx.fuse() => { - Err(CancelErr(())) - }, - value = self.f.fuse() => { - Ok(value) - } - } - } -} - - -impl<F: Future<Output = T> + Unpin, T: Unpin> Future for Cancellable<F> { - type Output = std::result::Result<T, CancelErr>; - - fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> { - let inner = self.get_mut(); - - if let Poll::Ready(ready) = inner.f.poll_unpin(cx) { - return Poll::Ready(Ok(ready)); - } - - if let Poll::Ready(_) = inner.rx.poll_unpin(cx) { - return Poll::Ready(Err(CancelErr(()))); - } - - Poll::Pending - } -} - pub fn get_request<T: serde::de::DeserializeOwned>( factory: &RequestFactory, service: RequestServiceHandle, |
