-
Notifications
You must be signed in to change notification settings - Fork 419
Expand file tree
/
Copy pathmerged.rs
More file actions
143 lines (127 loc) · 4.5 KB
/
Copy pathmerged.rs
File metadata and controls
143 lines (127 loc) · 4.5 KB
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
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
//! Merge external stream into an [`EventStream`][super::EventStream].
//!
//! Sometimes nodes need to listen to external events, in addition to Dora events.
//! This module provides support for that by providing the [`MergeExternal`] trait.
use futures::{Stream, StreamExt};
use futures_concurrency::stream::Merge;
/// A Dora event or an event from an external source.
#[derive(Debug)]
pub enum MergedEvent<E> {
/// A Dora event
Dora(super::Event),
/// An external event
///
/// Yielded by the stream that was merged into the Dora [`EventStream`][super::EventStream].
External(E),
}
/// A general enum to represent a value of two possible types.
pub enum Either<A, B> {
/// Value is of the first type, type `A`.
First(A),
/// Value is of the second type, type `B`.
Second(B),
}
impl<A> Either<A, A> {
/// Unwraps an `Either` instance where both types are identical.
pub fn flatten(self) -> A {
match self {
Either::First(a) => a,
Either::Second(a) => a,
}
}
}
/// Allows merging an external event stream into an existing event stream.
// TODO: use impl trait return type once stable
pub trait MergeExternal<'a, E> {
/// The item type yielded from the merged stream.
type Item;
/// Merge the given stream into an existing event stream.
///
/// Returns a new event stream that yields items from both streams.
/// The ordering between the two streams is not guaranteed.
fn merge_external(
self,
external_events: impl Stream<Item = E> + Unpin + 'a,
) -> Box<dyn Stream<Item = Self::Item> + Unpin + 'a>;
}
/// Allows merging a sendable external event stream into an existing (sendable) event stream.
///
/// By implementing [`Send`], the streams can be sent to different threads.
pub trait MergeExternalSend<'a, E> {
/// The item type yielded from the merged stream.
type Item;
/// Merge the given stream into an existing event stream.
///
/// Returns a new event stream that yields items from both streams.
/// The ordering between the two streams is not guaranteed.
fn merge_external_send(
self,
external_events: impl Stream<Item = E> + Unpin + Send + Sync + 'a,
) -> Box<dyn Stream<Item = Self::Item> + Unpin + Send + Sync + 'a>;
}
impl<'a, E> MergeExternal<'a, E> for super::EventStream
where
E: 'static,
{
type Item = MergedEvent<E>;
fn merge_external(
self,
external_events: impl Stream<Item = E> + Unpin + 'a,
) -> Box<dyn Stream<Item = Self::Item> + Unpin + 'a> {
let dora = self.map(MergedEvent::Dora);
let external = external_events.map(MergedEvent::External);
Box::new((dora, external).merge())
}
}
impl<'a, E> MergeExternalSend<'a, E> for super::EventStream
where
E: 'static,
{
type Item = MergedEvent<E>;
fn merge_external_send(
self,
external_events: impl Stream<Item = E> + Unpin + Send + Sync + 'a,
) -> Box<dyn Stream<Item = Self::Item> + Unpin + Send + Sync + 'a> {
let dora = self.map(MergedEvent::Dora);
let external = external_events.map(MergedEvent::External);
Box::new((dora, external).merge())
}
}
impl<'a, E, F, S> MergeExternal<'a, F> for S
where
S: Stream<Item = MergedEvent<E>> + Unpin + 'a,
E: 'a,
F: 'a,
{
type Item = MergedEvent<Either<E, F>>;
fn merge_external(
self,
external_events: impl Stream<Item = F> + Unpin + 'a,
) -> Box<dyn Stream<Item = Self::Item> + Unpin + 'a> {
let first = self.map(|e| match e {
MergedEvent::Dora(d) => MergedEvent::Dora(d),
MergedEvent::External(e) => MergedEvent::External(Either::First(e)),
});
let second = external_events.map(|e| MergedEvent::External(Either::Second(e)));
Box::new((first, second).merge())
}
}
impl<'a, E, F, S> MergeExternalSend<'a, F> for S
where
S: Stream<Item = MergedEvent<E>> + Unpin + Send + Sync + 'a,
E: 'a,
F: 'a,
{
type Item = MergedEvent<Either<E, F>>;
fn merge_external_send(
self,
external_events: impl Stream<Item = F> + Unpin + Send + Sync + 'a,
) -> Box<dyn Stream<Item = Self::Item> + Unpin + Send + Sync + 'a> {
let first = self.map(|e| match e {
MergedEvent::Dora(d) => MergedEvent::Dora(d),
MergedEvent::External(e) => MergedEvent::External(Either::First(e)),
});
let second = external_events.map(|e| MergedEvent::External(Either::Second(e)));
Box::new((first, second).merge())
}
}