Skip to main content

tower/util/call_all/
unordered.rs

1//! [`Stream<Item = Request>`][stream] + [`Service<Request>`] => [`Stream<Item = Response>`][stream].
2//!
3//! [`Service<Request>`]: crate::Service
4//! [stream]: https://docs.rs/futures/latest/futures/stream/trait.Stream.html
5
6use super::common;
7use futures_core::Stream;
8use futures_util::stream::FuturesUnordered;
9use pin_project_lite::pin_project;
10use std::{
11    future::Future,
12    pin::Pin,
13    task::{Context, Poll},
14};
15use tower_service::Service;
16
17pin_project! {
18    /// A stream of responses received from the inner service in received order.
19    ///
20    /// Similar to [`CallAll`] except, instead of yielding responses in request order,
21    /// responses are returned as they are available.
22    ///
23    /// [`CallAll`]: crate::util::CallAll
24    #[derive(Debug)]
25    pub struct CallAllUnordered<Svc, S>
26    where
27        Svc: Service<S::Item>,
28        S: Stream,
29    {
30        #[pin]
31        inner: common::CallAll<Svc, S, FuturesUnordered<Svc::Future>>,
32    }
33}
34
35impl<Svc, S> CallAllUnordered<Svc, S>
36where
37    Svc: Service<S::Item>,
38    S: Stream,
39{
40    /// Create new [`CallAllUnordered`] combinator.
41    pub fn new(service: Svc, stream: S) -> CallAllUnordered<Svc, S> {
42        CallAllUnordered {
43            inner: common::CallAll::new(service, stream, FuturesUnordered::new()),
44        }
45    }
46
47    /// Create a new [`CallAllUnordered`] from an existing, partially-evaluated [`CallAll`] instance.
48    ///
49    /// This constructor allows type-safe state migrations across combinators
50    /// without spilling buffered requests from `curr_req`.
51    pub(crate) fn from_inner(
52        inner: common::CallAll<Svc, S, FuturesUnordered<Svc::Future>>,
53    ) -> Self {
54        CallAllUnordered { inner }
55    }
56
57    /// Extract the wrapped [`Service`].
58    ///
59    /// # Panics
60    ///
61    /// Panics if [`take_service`] was already called.
62    ///
63    /// [`take_service`]: crate::util::CallAllUnordered::take_service
64    pub fn into_inner(self) -> Svc {
65        self.inner.into_inner()
66    }
67
68    /// Extract the wrapped `Service`.
69    ///
70    /// This [`CallAllUnordered`] can no longer be used after this function has been called.
71    ///
72    /// # Panics
73    ///
74    /// Panics if [`take_service`] was already called.
75    ///
76    /// [`take_service`]: crate::util::CallAllUnordered::take_service
77    pub fn take_service(self: Pin<&mut Self>) -> Svc {
78        self.project().inner.take_service()
79    }
80}
81
82impl<Svc, S> Stream for CallAllUnordered<Svc, S>
83where
84    Svc: Service<S::Item>,
85    S: Stream,
86{
87    type Item = Result<Svc::Response, Svc::Error>;
88
89    fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Option<Self::Item>> {
90        self.project().inner.poll_next(cx)
91    }
92}
93
94impl<F: Future> common::Drive<F> for FuturesUnordered<F> {
95    fn is_empty(&self) -> bool {
96        FuturesUnordered::is_empty(self)
97    }
98
99    fn push(&mut self, future: F) {
100        FuturesUnordered::push(self, future)
101    }
102
103    fn poll(&mut self, cx: &mut Context<'_>) -> Poll<Option<F::Output>> {
104        Stream::poll_next(Pin::new(self), cx)
105    }
106}