use bytes::Bytes;
use futures::{pin_mut, Stream, StreamExt};
use std::mem;
use std::pin::Pin;
use async_stream::stream;
use async_trait::async_trait;
use crate::{dds::ConstrainedVariable, Constraint};
pub fn xdr_length(len: u32) -> [u8; 8] {
let len = len.to_be();
let x: [u32; 2] = [len, len];
unsafe { mem::transmute(x) }
}
#[async_trait]
pub trait Dods: crate::Dap2 + Send + Sync + Clone + 'static {
async fn dods(
&self,
constraint: Constraint,
) -> Result<
(
u64,
Pin<Box<dyn Stream<Item = Result<Bytes, anyhow::Error>> + Send + 'static>>,
),
anyhow::Error,
> {
let dds = self.dds().await.dds(&constraint)?;
let dds_bytes = Bytes::from(dds.to_string());
let content_length = (dds.dods_size() + dds_bytes.len() + 8) as u64;
let slf = self.clone();
Ok((content_length, stream! {
yield Ok::<_, anyhow::Error>(dds_bytes);
yield Ok(Bytes::from_static(b"\n\nData:\n"));
for c in dds.variables {
match c {
ConstrainedVariable::Variable(v) |
ConstrainedVariable::Structure { variable: _, member: v }
=> {
if !v.is_scalar() {
yield Ok(Bytes::from(Vec::from(xdr_length(v.len() as u32))));
}
let reader = slf.variable(&v).await?;
pin_mut!(reader);
while let Some(b) = reader.next().await {
yield b;
}
},
ConstrainedVariable::Grid {
variable,
dimensions,
} => {
for variable in std::iter::once(variable).chain(dimensions) {
if !variable.is_scalar() {
yield Ok(Bytes::from(Vec::from(xdr_length(variable.len() as u32))));
}
let reader = slf.variable(&variable).await?;
pin_mut!(reader);
while let Some(b) = reader.next().await {
yield b;
}
}
}
}
}
}.boxed()))
}
}
impl<T: crate::Dap2 + Send + Sync + Clone + 'static> Dods for T {}
#[cfg(test)]
mod test {
use super::*;
#[test]
fn length() {
let x: u32 = 2;
let b = xdr_length(x);
assert_eq!(b, [0u8, 0, 0, 2, 0, 0, 0, 2]);
}
}