1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
use std::prelude::v1::*;
use {Async, IntoFuture, Poll, Future};
use stream::{Stream, Fuse};
#[must_use = "streams do nothing unless polled"]
pub struct BufferUnordered<S>
where S: Stream,
S::Item: IntoFuture,
{
stream: Fuse<S>,
futures: Vec<Option<<S::Item as IntoFuture>::Future>>,
cur: usize,
}
pub fn new<S>(s: S, amt: usize) -> BufferUnordered<S>
where S: Stream,
S::Item: IntoFuture<Error=<S as Stream>::Error>,
{
BufferUnordered {
stream: super::fuse::new(s),
futures: (0..amt).map(|_| None).collect(),
cur: 0,
}
}
impl<S> Stream for BufferUnordered<S>
where S: Stream,
S::Item: IntoFuture<Error=<S as Stream>::Error>,
{
type Item = <S::Item as IntoFuture>::Item;
type Error = <S as Stream>::Error;
fn poll(&mut self) -> Poll<Option<Self::Item>, Self::Error> {
for future in &mut self.futures {
if future.is_none() {
match try!(self.stream.poll()) {
Async::Ready(Some(s)) => {
*future = Some(s.into_future());
}
Async::Ready(None) => break,
Async::NotReady => break,
}
}
}
let mut waiting = false;
for i in 0..self.futures.len() {
let mut idx = self.cur + i;
if idx >= self.futures.len() {
idx -= self.futures.len();
}
let future = &mut self.futures[idx];
let result = match *future {
Some(ref mut s) => match s.poll() {
Ok(Async::NotReady) => {
waiting = true;
continue
},
Ok(Async::Ready(e)) => Ok(Async::Ready(Some(e))),
Err(e) => Err(e),
},
None => continue,
};
self.cur = i + 1;
*future = None;
return result;
}
Ok(if waiting || !self.stream.is_done() {
Async::NotReady
} else {
Async::Ready(None)
})
}
}