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
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
use futures::{Stream, Poll, Async};
use std::rc::{Rc, Weak};
use std::collections::VecDeque;
use std::cell::RefCell;
pub fn unsync_cloneable<S: Stream>(stream: S) -> UnsyncCloneable<S> {
let queue = Rc::new(RefCell::new(VecDeque::new()));
let receivers = vec![Rc::downgrade(&queue)];
let shared = Shared {
stream: stream,
receivers: receivers,
};
UnsyncCloneable {
queue: queue,
shared: Rc::new(RefCell::new(shared)),
}
}
struct Shared<S: Stream> {
stream: S,
receivers: Vec<Weak<Queue<S>>>,
}
type Queue<S: Stream> = RefCell<VecDeque<Result<Option<S::Item>, S::Error>>>;
pub struct UnsyncCloneable<S: Stream> {
queue: Rc<Queue<S>>,
shared: Rc<RefCell<Shared<S>>>,
}
impl<S> Stream for UnsyncCloneable<S>
where
S: Stream,
S::Item: Clone,
S::Error: Clone,
{
type Item = S::Item;
type Error = S::Error;
fn poll(&mut self) -> Poll<Option<S::Item>, S::Error> {
match self.queue.borrow_mut().pop_front() {
Some(Ok(Some(msg))) => return Ok(Async::Ready(Some(msg))),
Some(Ok(None)) => return Ok(Async::Ready(None)),
Some(Err(e)) => return Err(e),
None => (),
}
{
let mut shared = self.shared.borrow_mut();
let poll = shared.stream.poll();
let msg = match poll {
Err(e) => Err(e.clone()),
Ok(Async::Ready(Some(msg))) => Ok(Some(msg)),
Ok(Async::Ready(None)) => Ok(None),
Ok(Async::NotReady) => {
return Ok(Async::NotReady);
}
};
for rx in shared.receivers.iter().filter_map(Weak::upgrade) {
rx.borrow_mut().push_back(msg.clone());
}
}
self.poll()
}
}
impl<S: Stream> Clone for UnsyncCloneable<S> {
fn clone(&self) -> Self {
let queue = Rc::new(RefCell::new(VecDeque::new()));
let mut shared = self.shared.borrow_mut();
shared.receivers.retain(|weak| weak.upgrade().is_some());
shared.receivers.push(Rc::downgrade(&queue));
drop(shared);
UnsyncCloneable {
queue: queue,
shared: self.shared.clone(),
}
}
}
impl<S: Stream> ::std::fmt::Debug for UnsyncCloneable<S> {
fn fmt(&self, f: &mut ::std::fmt::Formatter) -> Result<(), ::std::fmt::Error> {
write!(f, "UnsyncCloneable(..)")
}
}