pub use cameleon_device::PixelFormat;
use std::time;
use async_channel::{Receiver, Sender};
use super::{StreamError, StreamResult};
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum PayloadType {
Image,
ImageExtendedChunk,
Chunk,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct ImageInfo {
pub width: usize,
pub height: usize,
pub x_offset: usize,
pub y_offset: usize,
pub pixel_format: PixelFormat,
pub image_size: usize,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct Payload {
pub(crate) id: u64,
pub(crate) payload_type: PayloadType,
pub(crate) image_info: Option<ImageInfo>,
pub(crate) payload: Vec<u8>,
pub(crate) valid_payload_size: usize,
pub(crate) timestamp: time::Duration,
}
impl Payload {
pub fn payload_type(&self) -> PayloadType {
self.payload_type
}
pub fn image_info(&self) -> Option<&ImageInfo> {
self.image_info.as_ref()
}
pub fn image(&self) -> Option<&[u8]> {
let image_info = self.image_info()?;
Some(&self.payload[..image_info.image_size])
}
pub fn payload(&self) -> &[u8] {
&self.payload[..self.valid_payload_size]
}
pub fn id(&self) -> u64 {
self.id
}
pub fn timestamp(&self) -> time::Duration {
self.timestamp
}
pub fn into_vec(mut self) -> Vec<u8> {
self.payload.resize(self.valid_payload_size, 0);
self.payload
}
}
#[derive(Debug, Clone)]
pub struct PayloadReceiver {
tx: Sender<Payload>,
rx: Receiver<StreamResult<Payload>>,
}
impl PayloadReceiver {
pub async fn recv(&self) -> StreamResult<Payload> {
self.rx.recv().await?
}
pub fn try_recv(&self) -> StreamResult<Payload> {
self.rx.try_recv()?
}
pub fn recv_blocking(&self) -> StreamResult<Payload> {
self.rx.recv_blocking()?
}
pub fn send_back(&self, payload: Payload) {
self.tx.try_send(payload).ok();
}
}
#[derive(Debug, Clone)]
pub struct PayloadSender {
tx: Sender<StreamResult<Payload>>,
rx: Receiver<Payload>,
}
impl PayloadSender {
pub async fn send(&self, payload: StreamResult<Payload>) -> StreamResult<()> {
Ok(self.tx.send(payload).await?)
}
pub fn try_send(&self, payload: StreamResult<Payload>) -> StreamResult<()> {
Ok(self.tx.try_send(payload)?)
}
pub fn try_recv(&self) -> StreamResult<Payload> {
Ok(self.rx.try_recv()?)
}
}
pub fn channel(payload_cap: usize, buffer_cap: usize) -> (PayloadSender, PayloadReceiver) {
let (device_tx, host_rx) = async_channel::bounded(payload_cap);
let (host_tx, device_rx) = async_channel::bounded(buffer_cap);
(
PayloadSender {
tx: device_tx,
rx: device_rx,
},
PayloadReceiver {
tx: host_tx,
rx: host_rx,
},
)
}
impl From<async_channel::RecvError> for StreamError {
fn from(err: async_channel::RecvError) -> Self {
StreamError::ReceiveError(err.to_string().into())
}
}
impl From<async_channel::TryRecvError> for StreamError {
fn from(err: async_channel::TryRecvError) -> Self {
StreamError::ReceiveError(err.to_string().into())
}
}
impl<T> From<async_channel::SendError<T>> for StreamError {
fn from(err: async_channel::SendError<T>) -> Self {
StreamError::ReceiveError(err.to_string().into())
}
}
impl<T> From<async_channel::TrySendError<T>> for StreamError {
fn from(err: async_channel::TrySendError<T>) -> Self {
StreamError::ReceiveError(err.to_string().into())
}
}