std_mel/flow/vec.rs
1use melodium_core::*;
2use melodium_macro::{check, mel_treatment};
3
4/// Flatten a stream of vector.
5///
6/// All the input vectors are turned into continuous stream of scalar values, keeping order.
7/// ```mermaid
8/// graph LR
9/// T("flatten()")
10/// B["[🟦 🟦][🟦][🟦 🟦 🟦]"] -->|vector| T
11///
12/// T -->|value| O["🟦 🟦 🟦 🟦 🟦 🟦"]
13///
14/// style B fill:#ffffff,stroke:#ffffff
15/// style O fill:#ffffff,stroke:#ffffff
16/// ```
17#[mel_treatment(
18 generic T ()
19 input vector Stream<Vec<T>>
20 output value Stream<T>
21)]
22pub async fn flatten() {
23 'main: while let Ok(mut vectors) = vector
24 .recv_many()
25 .await
26 .map(|values| Into::<VecDeque<Value>>::into(values))
27 {
28 while let Some(vector) = vectors.pop_front().map(|val| match val {
29 Value::Vec(vec) => vec,
30 Value::Packed(arr) => arr.into_values(),
31 _ => panic!("Vec expected"),
32 }) {
33 for val in vector {
34 check!('main, value.send_one(val).await)
35 }
36 }
37 }
38}
39
40/// Gives pattern of a stream of vectors.
41///
42/// ```mermaid
43/// graph LR
44/// T("pattern()")
45/// A["…[🟨 🟨][🟨][🟨 🟨 🟨]"] -->|stream| T
46///
47/// T -->|pattern| O["… [🟦 🟦][🟦][🟦 🟦 🟦]"]
48///
49/// style A fill:#ffffff,stroke:#ffffff
50/// style O fill:#ffffff,stroke:#ffffff
51/// ```
52#[mel_treatment(
53 generic T ()
54 input stream Stream<Vec<T>>
55 output pattern Stream<Vec<void>>
56)]
57pub async fn pattern() {
58 'main: while let Ok(vectors) = stream
59 .recv_many()
60 .await
61 .map(|values| Into::<VecDeque<Value>>::into(values))
62 {
63 for val in vectors {
64 let len = match val {
65 Value::Vec(vec) => vec.len(),
66 Value::Packed(arr) => arr.len(),
67 _ => panic!("Vec expected"),
68 };
69 check!('main, pattern.send_one_as(vec![(); len]).await)
70 }
71 }
72}
73
74/// Fit a stream of raw values into stream of vectors using a pattern.
75///
76/// ℹ️ If some remaining values doesn't fit into the pattern, they are trashed.
77/// If there are not enough values to fit the pattern, uncomplete vector is trashed.
78///
79/// ```mermaid
80/// graph LR
81/// T("fit()")
82/// A["… 🟨 🟨 🟨 🟨 🟨 🟨"] -->|value| T
83/// B["[🟦 🟦][🟦][🟦 🟦 🟦]"] -->|pattern| T
84///
85/// T -->|fitted| O["[🟨 🟨][🟨][🟨 🟨 🟨]"]
86///
87/// style A fill:#ffffff,stroke:#ffffff
88/// style B fill:#ffffff,stroke:#ffffff
89/// style O fill:#ffffff,stroke:#ffffff
90/// ```
91#[mel_treatment(
92 generic T ()
93 input value Stream<T>
94 input pattern Stream<Vec<void>>
95 output fitted Stream<Vec<T>>
96)]
97pub async fn fit() {
98 'main: while let Ok(patterns) = pattern
99 .recv_many()
100 .await
101 .map(|values| Into::<VecDeque<Value>>::into(values))
102 {
103 for pattern in patterns {
104 let pattern_len = match pattern {
105 Value::Vec(pattern) => pattern.len(),
106 Value::Packed(arr) => arr.len(),
107 _ => panic!("Vec expected"),
108 };
109 let mut vector = Vec::with_capacity(pattern_len);
110 for _ in 0..pattern_len {
111 if let Ok(val) = value.recv_one().await {
112 vector.push(val);
113 } else {
114 // Uncomplete, we 'trash' vector
115 break 'main;
116 }
117 }
118 check!('main, fitted.send_one_as(vector).await)
119 }
120 }
121}
122
123/// Fill a pattern stream with a `i64` value.
124///
125/// ```mermaid
126/// graph LR
127/// T("fill(value=🟧)")
128/// B["…[🟦 🟦][🟦][🟦 🟦 🟦]…"] -->|pattern| T
129///
130/// T -->|filled| O["…[🟧 🟧][🟧][🟧 🟧 🟧]…"]
131///
132/// style B fill:#ffffff,stroke:#ffffff
133/// style O fill:#ffffff,stroke:#ffffff
134/// ```
135#[mel_treatment(
136 generic T ()
137 input pattern Stream<Vec<void>>
138 output filled Stream<Vec<T>>
139)]
140pub async fn fill(value: T) {
141 'main: while let Ok(patterns) = pattern
142 .recv_many()
143 .await
144 .map(|values| Into::<VecDeque<Value>>::into(values))
145 {
146 for pattern in patterns {
147 let len = match pattern {
148 Value::Vec(pattern) => pattern.len(),
149 Value::Packed(arr) => arr.len(),
150 _ => panic!("Vec expected"),
151 };
152 check!('main, filled.send_one_as(vec![value.clone(); len]).await)
153 }
154 }
155}
156
157/// Gives size of vectors passing through stream.
158///
159/// For each vector one `size` value is sent, giving the number of elements contained within matching vector.
160///
161/// ```mermaid
162/// graph LR
163/// T("size()")
164/// V["[🟦 🟦][🟦][][🟦 🟦 🟦]…"] -->|vector| T
165///
166/// T -->|size| P["2️⃣ 1️⃣ 0️⃣ 3️⃣ …"]
167///
168/// style V fill:#ffffff,stroke:#ffffff
169/// style P fill:#ffffff,stroke:#ffffff
170/// ```
171#[mel_treatment(
172 generic T ()
173 input vector Stream<Vec<T>>
174 output size Stream<u64>
175)]
176pub async fn size() {
177 while let Ok(iter) = vector
178 .recv_many()
179 .await
180 .map(|values| Into::<VecDeque<Value>>::into(values))
181 {
182 check!(
183 size.send_many_as(
184 iter.into_iter()
185 .map(|v| match v {
186 Value::Vec(v) => v.len() as u64,
187 Value::Packed(arr) => arr.len() as u64,
188 _ => panic!("Vec expected"),
189 })
190 .collect::<Vec<_>>()
191 )
192 .await
193 );
194 }
195}
196
197/// Resize vectors according to given streamed size.
198///
199/// If a vector is smaller than expected size, it is extended using the `default` value.
200///
201/// ```mermaid
202/// graph LR
203/// T("resize(default=🟨)")
204/// V["[🟦 🟦][🟦][][🟦 🟦 🟦]…"] -->|vector| T
205/// S["3️⃣ 2️⃣ 3️⃣ 2️⃣ …"] -->|size| T
206///
207/// T -->|resized| P["[🟦 🟦 🟨][🟦 🟨][🟨 🟨 🟨][🟦 🟦]…"]
208///
209/// style V fill:#ffffff,stroke:#ffffff
210/// style S fill:#ffffff,stroke:#ffffff
211/// style P fill:#ffffff,stroke:#ffffff
212/// ```
213#[mel_treatment(
214 generic T ()
215 input vector Stream<Vec<T>>
216 input size Stream<u64>
217 output resized Stream<Vec<T>>
218)]
219pub async fn resize(default: T) {
220 while let Ok(size) = size.recv_one_as::<u64>().await {
221 if let Ok(vec) = vector.recv_one().await {
222 let mut vec = match vec {
223 Value::Vec(vec) => vec,
224 Value::Packed(arr) => arr.into_values(),
225 _ => panic!("Vec expected"),
226 };
227 vec.resize(size as usize, default.clone());
228 check!(resized.send_one_as(vec).await);
229 } else {
230 break;
231 }
232 }
233}