tower/util/call_all/
unordered.rs1use 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 #[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 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 pub(crate) fn from_inner(
52 inner: common::CallAll<Svc, S, FuturesUnordered<Svc::Future>>,
53 ) -> Self {
54 CallAllUnordered { inner }
55 }
56
57 pub fn into_inner(self) -> Svc {
65 self.inner.into_inner()
66 }
67
68 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}