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
113
114
115
116
117
118
119
120
121
122
123
124
use sync::spsc::ring::{channel, Receiver, SendError, SendErrorKind};
use thread::prelude::*;
#[must_use]
pub struct RoutineStreamRing<I, E> {
rx: Receiver<I, E>,
}
impl<I, E> RoutineStreamRing<I, E> {
pub(crate) fn new<T, G, O>(
thread: &T,
capacity: usize,
mut generator: G,
overflow: O,
) -> Self
where
T: Thread,
G: Generator<Yield = Option<I>, Return = Result<Option<I>, E>>,
O: Fn(I) -> Result<(), E>,
G: Send + 'static,
I: Send + 'static,
E: Send + 'static,
O: Send + 'static,
{
let (mut tx, rx) = channel(capacity);
thread.routines().push(move || loop {
if tx.is_canceled() {
break;
}
match generator.resume() {
Yielded(None) => {}
Yielded(Some(value)) => match tx.send(value) {
Ok(()) => {}
Err(SendError { value, kind }) => match kind {
SendErrorKind::Canceled => {
break;
}
SendErrorKind::Overflow => match overflow(value) {
Ok(()) => {}
Err(err) => {
tx.send_err(err).ok();
break;
}
},
},
},
Complete(Ok(None)) => {
break;
}
Complete(Ok(Some(value))) => {
tx.send(value).ok();
break;
}
Complete(Err(err)) => {
tx.send_err(err).ok();
break;
}
}
yield;
});
Self { rx }
}
pub(crate) fn new_overwrite<T, G>(
thread: &T,
capacity: usize,
mut generator: G,
) -> Self
where
T: Thread,
G: Generator<Yield = Option<I>, Return = Result<Option<I>, E>>,
G: Send + 'static,
I: Send + 'static,
E: Send + 'static,
{
let (mut tx, rx) = channel(capacity);
thread.routines().push(move || loop {
if tx.is_canceled() {
break;
}
match generator.resume() {
Yielded(None) => {}
Yielded(Some(value)) => match tx.send_overwrite(value) {
Ok(()) => (),
Err(_) => break,
},
Complete(Ok(None)) => {
break;
}
Complete(Ok(Some(value))) => {
tx.send_overwrite(value).ok();
break;
}
Complete(Err(err)) => {
tx.send_err(err).ok();
break;
}
}
yield;
});
Self { rx }
}
#[inline(always)]
pub fn close(&mut self) {
self.rx.close()
}
}
impl<I, E> Stream for RoutineStreamRing<I, E> {
type Item = I;
type Error = E;
#[inline(always)]
fn poll(&mut self) -> Poll<Option<I>, E> {
self.rx.poll()
}
}