use bytes::{Bytes, BytesMut};
use crate::codec::Decode;
use crate::error::RecvError;
use crate::ids::DeliveryId;
use crate::types::messaging::{DeliveryState, Message};
#[derive(Debug, Clone)]
pub struct Delivery {
pub delivery_id: DeliveryId,
pub delivery_tag: Bytes,
pub settled: bool,
state: Option<DeliveryState>,
body: Bytes,
}
impl Delivery {
pub fn new(delivery_id: DeliveryId, delivery_tag: Bytes, settled: bool, body: Bytes) -> Self {
Delivery {
delivery_id,
delivery_tag,
settled,
state: None,
body,
}
}
pub fn with_state(mut self, state: Option<DeliveryState>) -> Self {
self.state = state;
self
}
pub fn state(&self) -> Option<&DeliveryState> {
self.state.as_ref()
}
pub fn raw(&self) -> &Bytes {
&self.body
}
pub fn into_raw(self) -> Bytes {
self.body
}
pub fn message(&self) -> Result<Message, RecvError> {
let mut buf = self.body.clone();
Message::decode(&mut buf).map_err(RecvError::from)
}
pub fn decode<T: Decode>(&self) -> Result<T, RecvError> {
let mut buf = self.body.clone();
T::decode(&mut buf).map_err(RecvError::from)
}
}
#[derive(Debug)]
pub struct PartialDelivery {
delivery_id: DeliveryId,
delivery_tag: Bytes,
settled: bool,
state: Option<DeliveryState>,
buf: BytesMut,
}
impl PartialDelivery {
pub fn new(
delivery_id: DeliveryId,
delivery_tag: Bytes,
settled: bool,
state: Option<DeliveryState>,
first: &[u8],
) -> Self {
let mut buf = BytesMut::with_capacity(first.len());
buf.extend_from_slice(first);
PartialDelivery {
delivery_id,
delivery_tag,
settled,
state,
buf,
}
}
pub fn delivery_id(&self) -> DeliveryId {
self.delivery_id
}
pub fn len(&self) -> usize {
self.buf.len()
}
pub fn is_empty(&self) -> bool {
self.buf.is_empty()
}
pub fn append(&mut self, payload: &[u8]) {
self.buf.extend_from_slice(payload);
}
pub fn complete(self) -> Delivery {
Delivery::new(
self.delivery_id,
self.delivery_tag,
self.settled,
self.buf.freeze(),
)
.with_state(self.state)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::codec::to_vec;
use crate::types::messaging::Message;
#[test]
fn lazy_body_and_typed_decode() {
let msg = Message::text("hello");
let bytes = Bytes::from(to_vec(&msg));
let d = Delivery::new(
DeliveryId(1),
Bytes::from_static(b"tag"),
false,
bytes.clone(),
);
assert_eq!(d.raw(), &bytes);
assert_eq!(d.message().unwrap(), msg);
}
#[test]
fn multi_frame_assembly() {
let msg = Message::data(Bytes::from_static(b"abcdefgh"));
let full = to_vec(&msg);
let (a, b) = full.split_at(full.len() / 2);
let mut partial =
PartialDelivery::new(DeliveryId(2), Bytes::from_static(b"t"), false, None, a);
partial.append(b);
let delivery = partial.complete();
assert_eq!(delivery.message().unwrap(), msg);
}
}