use crate::proto::streaming_api::Datagram;
use std::io::{Error, ErrorKind};
pub fn encode_datagrams(datagrams: &[Datagram]) -> Vec<u8> {
let mut out = Vec::new();
for d in datagrams {
encode_datagram_into(d, &mut out);
}
out
}
fn encode_datagram_into(d: &Datagram, out: &mut Vec<u8>) {
let mut payload = Vec::new();
payload.extend_from_slice(&d.timestamp.to_be_bytes());
payload.extend_from_slice(&(d.sample_count as u16).to_be_bytes());
payload.push(d.clock_synch as u8);
let flags_len = d.flags.len() as u16;
payload.extend_from_slice(&flags_len.to_be_bytes());
payload.extend_from_slice(&d.flags);
let gm_id_len = d.gm_identity.len() as u8;
payload.push(gm_id_len);
payload.extend_from_slice(&d.gm_identity);
for v in &d.values {
payload.extend_from_slice(&v.to_be_bytes());
}
let payload_len = payload.len() as u16;
out.extend_from_slice(&payload_len.to_be_bytes());
out.extend_from_slice(&payload);
}
pub fn decode_datagrams(mut bytes: &[u8]) -> Result<Vec<Datagram>, Error> {
let mut datagrams = Vec::new();
while !bytes.is_empty() {
let datagram;
(datagram, bytes) = decode_one(bytes)?;
datagrams.push(datagram);
}
Ok(datagrams)
}
fn decode_one(bytes: &[u8]) -> Result<(Datagram, &[u8]), Error> {
let (payload_len, bytes) = read_u16(bytes)?;
let payload_len = payload_len as usize;
if bytes.len() < payload_len {
return Err(Error::new(
ErrorKind::InvalidData,
"truncated datagram payload",
));
}
let (payload, rest) = bytes.split_at(payload_len);
let datagram = decode_payload(payload)?;
Ok((datagram, rest))
}
fn decode_payload(bytes: &[u8]) -> Result<Datagram, Error> {
let (timestamp, bytes) = read_i64(bytes)?;
let (sample_count_raw, bytes) = read_u16(bytes)?;
let (clock_synch_raw, bytes) = read_u8(bytes)?;
let (flags_len, bytes) = read_u16(bytes)?;
let (flags, bytes) = read_bytes(bytes, flags_len as usize)?;
let (gm_id_len, bytes) = read_u8(bytes)?;
let (gm_identity, bytes) = read_bytes(bytes, gm_id_len as usize)?;
let values = decode_values(bytes)?;
Ok(Datagram {
timestamp,
sample_count: sample_count_raw as u32,
clock_synch: clock_synch_raw as u32,
flags: flags.to_vec(),
gm_identity: gm_identity.to_vec(),
values,
})
}
fn decode_values(bytes: &[u8]) -> Result<Vec<i64>, Error> {
if !bytes.len().is_multiple_of(8) {
return Err(Error::new(
ErrorKind::InvalidData,
"values payload length is not a multiple of 8",
));
}
let mut values = Vec::with_capacity(bytes.len() / 8);
let mut remaining = bytes;
loop {
let (v, rest) = read_i64(remaining)?;
values.push(v);
remaining = rest;
if remaining.is_empty() {
break;
}
}
Ok(values)
}
fn read_u8(bytes: &[u8]) -> Result<(u8, &[u8]), Error> {
match bytes.split_first() {
Some((&b, rest)) => Ok((b, rest)),
None => Err(Error::new(ErrorKind::InvalidData, "unexpected end of data")),
}
}
fn read_u16(bytes: &[u8]) -> Result<(u16, &[u8]), Error> {
if bytes.len() < 2 {
return Err(Error::new(ErrorKind::InvalidData, "unexpected end of data"));
}
let v = u16::from_be_bytes([bytes[0], bytes[1]]);
Ok((v, &bytes[2..]))
}
fn read_i64(bytes: &[u8]) -> Result<(i64, &[u8]), Error> {
if bytes.len() < 8 {
return Err(Error::new(ErrorKind::InvalidData, "unexpected end of data"));
}
let v = i64::from_be_bytes(bytes[..8].try_into().unwrap());
Ok((v, &bytes[8..]))
}
fn read_bytes(bytes: &[u8], n: usize) -> Result<(&[u8], &[u8]), Error> {
if bytes.len() < n {
return Err(Error::new(ErrorKind::InvalidData, "unexpected end of data"));
}
Ok(bytes.split_at(n))
}