Skip to main content

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}