sinusoidal 0.9.0

The official SDK to write rust apps for the Sinusoidal Systems Digital Measurement Platform
Documentation
//! Encoding and decoding of the custom binary datagram format used in the
//! `datagrams` bytes field of `StreamPacket`.
//!
//! Wire format for a sequence of datagrams (concatenated, no count prefix):
//!
//! ```text
//! for each datagram:
//!   [payload_len: u16 big-endian]
//!   [timestamp:   i64 big-endian]
//!   [sample_count: u16 big-endian]
//!   [clock_synch:  u8]
//!   [flags_len:    u16 big-endian]
//!   [flags:        flags_len bytes]
//!   [gm_id_len:    u8]
//!   [gm_identity:  gm_id_len bytes]
//!   [values:       N × i64 big-endian]
//! ```

use crate::proto::streaming_api::Datagram;
use std::io::{Error, ErrorKind};

/// Encode a slice of datagrams into the custom binary format.
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);
}

/// Decode a sequence of datagrams from the custom binary format.
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))
}