use melodium_core::*;
use melodium_macro::{check, mel_treatment};
#[mel_treatment(
generic T ()
input vector Stream<Vec<T>>
output value Stream<T>
)]
pub async fn flatten() {
'main: while let Ok(mut vectors) = vector
.recv_many()
.await
.map(|values| Into::<VecDeque<Value>>::into(values))
{
while let Some(vector) = vectors.pop_front().map(|val| match val {
Value::Vec(vec) => vec,
Value::Packed(arr) => arr.into_values(),
_ => panic!("Vec expected"),
}) {
for val in vector {
check!('main, value.send_one(val).await)
}
}
}
}
#[mel_treatment(
generic T ()
input stream Stream<Vec<T>>
output pattern Stream<Vec<void>>
)]
pub async fn pattern() {
'main: while let Ok(vectors) = stream
.recv_many()
.await
.map(|values| Into::<VecDeque<Value>>::into(values))
{
for val in vectors {
let len = match val {
Value::Vec(vec) => vec.len(),
Value::Packed(arr) => arr.len(),
_ => panic!("Vec expected"),
};
check!('main, pattern.send_one_as(vec![(); len]).await)
}
}
}
#[mel_treatment(
generic T ()
input value Stream<T>
input pattern Stream<Vec<void>>
output fitted Stream<Vec<T>>
)]
pub async fn fit() {
'main: while let Ok(patterns) = pattern
.recv_many()
.await
.map(|values| Into::<VecDeque<Value>>::into(values))
{
for pattern in patterns {
let pattern_len = match pattern {
Value::Vec(pattern) => pattern.len(),
Value::Packed(arr) => arr.len(),
_ => panic!("Vec expected"),
};
let mut vector = Vec::with_capacity(pattern_len);
for _ in 0..pattern_len {
if let Ok(val) = value.recv_one().await {
vector.push(val);
} else {
break 'main;
}
}
check!('main, fitted.send_one_as(vector).await)
}
}
}
#[mel_treatment(
generic T ()
input pattern Stream<Vec<void>>
output filled Stream<Vec<T>>
)]
pub async fn fill(value: T) {
'main: while let Ok(patterns) = pattern
.recv_many()
.await
.map(|values| Into::<VecDeque<Value>>::into(values))
{
for pattern in patterns {
let len = match pattern {
Value::Vec(pattern) => pattern.len(),
Value::Packed(arr) => arr.len(),
_ => panic!("Vec expected"),
};
check!('main, filled.send_one_as(vec![value.clone(); len]).await)
}
}
}
#[mel_treatment(
generic T ()
input vector Stream<Vec<T>>
output size Stream<u64>
)]
pub async fn size() {
while let Ok(iter) = vector
.recv_many()
.await
.map(|values| Into::<VecDeque<Value>>::into(values))
{
check!(
size.send_many_as(
iter.into_iter()
.map(|v| match v {
Value::Vec(v) => v.len() as u64,
Value::Packed(arr) => arr.len() as u64,
_ => panic!("Vec expected"),
})
.collect::<Vec<_>>()
)
.await
);
}
}
#[mel_treatment(
generic T ()
input vector Stream<Vec<T>>
input size Stream<u64>
output resized Stream<Vec<T>>
)]
pub async fn resize(default: T) {
while let Ok(size) = size.recv_one_as::<u64>().await {
if let Ok(vec) = vector.recv_one().await {
let mut vec = match vec {
Value::Vec(vec) => vec,
Value::Packed(arr) => arr.into_values(),
_ => panic!("Vec expected"),
};
vec.resize(size as usize, default.clone());
check!(resized.send_one_as(vec).await);
} else {
break;
}
}
}