use bytes::{Buf, BufMut};
use crate::dispatch::{AnyFetchHeader, AnySubgroupHeader};
use crate::error::CodecError;
use crate::varint::VarInt;
use crate::version::DraftVersion;
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AnySubgroupObject {
pub object_id: u64,
pub extension_headers: Vec<u8>,
pub extension_count: Option<u64>,
pub status: Option<u64>,
pub payload: Vec<u8>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct AnySubgroupObjectMeta {
pub object_id: u64,
pub payload_length: u64,
pub status: Option<u64>,
pub extension_headers_len: u64,
pub wire_len: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AnyFetchEndOfRange {
NonExistent,
Unknown,
TimedOut,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AnyFetchGroupOrder {
Ascending,
Descending,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AnyFetchObject {
pub group_id: u64,
pub subgroup_id: u64,
pub has_subgroup_id: bool,
pub object_id: u64,
pub publisher_priority: u8,
pub extension_headers: Vec<u8>,
pub extension_count: Option<u64>,
pub status: Option<u64>,
pub end_of_range: Option<AnyFetchEndOfRange>,
pub payload: Vec<u8>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct AnyFetchObjectMeta {
pub group_id: u64,
pub subgroup_id: u64,
pub has_subgroup_id: bool,
pub object_id: u64,
pub publisher_priority: u8,
pub payload_length: u64,
pub status: Option<u64>,
pub end_of_range: Option<AnyFetchEndOfRange>,
pub extension_headers_len: u64,
pub wire_len: u64,
}
#[cfg(any(
feature = "draft16",
feature = "draft17",
feature = "draft18",
feature = "draft19",
feature = "draft20"
))]
const DEFAULT_PUBLISHER_PRIORITY: u8 = 128;
#[allow(dead_code)]
mod conv {
use super::{AnySubgroupObject, Buf, CodecError};
use crate::varint::VarInt;
pub fn skip(buf: &mut impl Buf, len: u64) -> Result<(), CodecError> {
let len = usize::try_from(len).map_err(|_| CodecError::UnexpectedEnd)?;
if buf.remaining() < len {
return Err(CodecError::UnexpectedEnd);
}
buf.advance(len);
Ok(())
}
pub fn take(buf: &mut impl Buf, len: u64) -> Result<Vec<u8>, CodecError> {
let len = usize::try_from(len).map_err(|_| CodecError::UnexpectedEnd)?;
crate::types::read_bytes(buf, len)
}
pub fn varint(v: u64) -> Result<VarInt, CodecError> {
VarInt::from_u64(v).map_err(|_| CodecError::InvalidField)
}
pub fn status_to_write(object: &AnySubgroupObject) -> Result<Option<u64>, CodecError> {
match (object.status, object.payload.is_empty()) {
(Some(0), false) => Ok(None),
(Some(_), false) => Err(CodecError::InvalidField),
(Some(code), true) => Ok(Some(code)),
(None, true) => Ok(Some(0)),
(None, false) => Ok(None),
}
}
}
macro_rules! legacy_subgroup_glue {
(no_extensions $name:ident, $feat:literal, $draft:ident) => {
#[cfg(feature = $feat)]
mod $name {
use super::conv;
use super::{AnySubgroupObject, AnySubgroupObjectMeta};
use crate::error::CodecError;
use crate::$draft::data_stream::ObjectHeader;
use crate::$draft::types::ObjectStatus;
use bytes::{Buf, BufMut};
pub fn read_object(buf: &mut impl Buf) -> Result<AnySubgroupObject, CodecError> {
let header = ObjectHeader::decode(buf)?;
let payload_length = header.payload_length.into_inner();
let (status, payload) = if payload_length == 0 {
(Some(header.object_status as u64), Vec::new())
} else {
(None, conv::take(buf, payload_length)?)
};
Ok(AnySubgroupObject {
object_id: header.object_id.into_inner(),
extension_headers: Vec::new(),
extension_count: None,
status,
payload,
})
}
pub fn read_object_meta(
buf: &mut impl Buf,
) -> Result<AnySubgroupObjectMeta, CodecError> {
let start = buf.remaining();
let header = ObjectHeader::decode(buf)?;
let payload_length = header.payload_length.into_inner();
let status = if payload_length == 0 {
Some(header.object_status as u64)
} else {
conv::skip(buf, payload_length)?;
None
};
Ok(AnySubgroupObjectMeta {
object_id: header.object_id.into_inner(),
payload_length,
status,
extension_headers_len: 0,
wire_len: (start - buf.remaining()) as u64,
})
}
pub fn write_object(
object: &AnySubgroupObject,
buf: &mut impl BufMut,
) -> Result<(), CodecError> {
if !object.extension_headers.is_empty() {
return Err(CodecError::InvalidField);
}
let object_status = match conv::status_to_write(object)? {
Some(code) => ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?,
None => ObjectStatus::Normal,
};
ObjectHeader {
object_id: conv::varint(object.object_id)?,
payload_length: conv::varint(object.payload.len() as u64)?,
object_status,
}
.encode(buf);
buf.put_slice(&object.payload);
Ok(())
}
}
};
(count_extensions $name:ident, $feat:literal, $draft:ident) => {
#[cfg(feature = $feat)]
mod $name {
use super::conv;
use super::{AnySubgroupObject, AnySubgroupObjectMeta};
use crate::error::CodecError;
use crate::$draft::data_stream::ObjectHeader;
use crate::$draft::types::ObjectStatus;
use bytes::{Buf, BufMut};
pub fn read_object(buf: &mut impl Buf) -> Result<AnySubgroupObject, CodecError> {
let header = ObjectHeader::decode(buf)?;
let payload_length = header.payload_length.into_inner();
let (status, payload) = if payload_length == 0 {
(Some(header.object_status as u64), Vec::new())
} else {
(None, conv::take(buf, payload_length)?)
};
Ok(AnySubgroupObject {
object_id: header.object_id.into_inner(),
extension_headers: header.extensions,
extension_count: Some(header.extension_count.into_inner()),
status,
payload,
})
}
pub fn read_object_meta(
buf: &mut impl Buf,
) -> Result<AnySubgroupObjectMeta, CodecError> {
let start = buf.remaining();
let header = ObjectHeader::decode(buf)?;
let payload_length = header.payload_length.into_inner();
let status = if payload_length == 0 {
Some(header.object_status as u64)
} else {
conv::skip(buf, payload_length)?;
None
};
Ok(AnySubgroupObjectMeta {
object_id: header.object_id.into_inner(),
payload_length,
status,
extension_headers_len: header.extensions.len() as u64,
wire_len: (start - buf.remaining()) as u64,
})
}
pub fn write_object(
object: &AnySubgroupObject,
buf: &mut impl BufMut,
) -> Result<(), CodecError> {
let extension_count = match object.extension_count {
Some(count) => count,
None if object.extension_headers.is_empty() => 0,
None => return Err(CodecError::InvalidField),
};
let object_status = match conv::status_to_write(object)? {
Some(code) => ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?,
None => ObjectStatus::Normal,
};
ObjectHeader {
object_id: conv::varint(object.object_id)?,
extension_count: conv::varint(extension_count)?,
extensions: object.extension_headers.clone(),
payload_length: conv::varint(object.payload.len() as u64)?,
object_status,
}
.encode(buf);
buf.put_slice(&object.payload);
Ok(())
}
}
};
(length_extensions $name:ident, $feat:literal, $draft:ident) => {
#[cfg(feature = $feat)]
mod $name {
use super::conv;
use super::{AnySubgroupObject, AnySubgroupObjectMeta};
use crate::error::CodecError;
use crate::$draft::data_stream::ObjectHeader;
use crate::$draft::types::ObjectStatus;
use bytes::{Buf, BufMut};
pub fn read_object(buf: &mut impl Buf) -> Result<AnySubgroupObject, CodecError> {
let header = ObjectHeader::decode(buf)?;
let payload_length = header.payload_length.into_inner();
let (status, payload) = if payload_length == 0 {
(Some(header.object_status as u64), Vec::new())
} else {
(None, conv::take(buf, payload_length)?)
};
Ok(AnySubgroupObject {
object_id: header.object_id.into_inner(),
extension_headers: header.extensions,
extension_count: None,
status,
payload,
})
}
pub fn read_object_meta(
buf: &mut impl Buf,
) -> Result<AnySubgroupObjectMeta, CodecError> {
let start = buf.remaining();
let header = ObjectHeader::decode(buf)?;
let payload_length = header.payload_length.into_inner();
let status = if payload_length == 0 {
Some(header.object_status as u64)
} else {
conv::skip(buf, payload_length)?;
None
};
Ok(AnySubgroupObjectMeta {
object_id: header.object_id.into_inner(),
payload_length,
status,
extension_headers_len: header.extension_headers_length.into_inner(),
wire_len: (start - buf.remaining()) as u64,
})
}
pub fn write_object(
object: &AnySubgroupObject,
buf: &mut impl BufMut,
) -> Result<(), CodecError> {
let object_status = match conv::status_to_write(object)? {
Some(code) => ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?,
None => ObjectStatus::Normal,
};
ObjectHeader {
object_id: conv::varint(object.object_id)?,
extension_headers_length: conv::varint(object.extension_headers.len() as u64)?,
extensions: object.extension_headers.clone(),
payload_length: conv::varint(object.payload.len() as u64)?,
object_status,
}
.encode(buf);
buf.put_slice(&object.payload);
Ok(())
}
}
};
(gated_extensions $name:ident, $feat:literal, $draft:ident) => {
#[cfg(feature = $feat)]
mod $name {
use super::conv;
use super::{AnySubgroupObject, AnySubgroupObjectMeta};
use crate::error::CodecError;
use crate::$draft::data_stream::ObjectHeader;
use crate::$draft::types::ObjectStatus;
use bytes::{Buf, BufMut};
pub fn read_object(
extensions: bool,
buf: &mut impl Buf,
) -> Result<AnySubgroupObject, CodecError> {
let header = ObjectHeader::decode_with_extensions(extensions, buf)?;
let payload_length = header.payload_length.into_inner();
let (status, payload) = if payload_length == 0 {
(Some(header.object_status as u64), Vec::new())
} else {
(None, conv::take(buf, payload_length)?)
};
Ok(AnySubgroupObject {
object_id: header.object_id.into_inner(),
extension_headers: header.extensions,
extension_count: None,
status,
payload,
})
}
pub fn read_object_meta(
extensions: bool,
buf: &mut impl Buf,
) -> Result<AnySubgroupObjectMeta, CodecError> {
let start = buf.remaining();
let header = ObjectHeader::decode_with_extensions(extensions, buf)?;
let payload_length = header.payload_length.into_inner();
let status = if payload_length == 0 {
Some(header.object_status as u64)
} else {
conv::skip(buf, payload_length)?;
None
};
Ok(AnySubgroupObjectMeta {
object_id: header.object_id.into_inner(),
payload_length,
status,
extension_headers_len: header.extension_headers_length.into_inner(),
wire_len: (start - buf.remaining()) as u64,
})
}
pub fn write_object(
extensions: bool,
object: &AnySubgroupObject,
buf: &mut impl BufMut,
) -> Result<(), CodecError> {
if !extensions && !object.extension_headers.is_empty() {
return Err(CodecError::InvalidField);
}
let object_status = match conv::status_to_write(object)? {
Some(code) => ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?,
None => ObjectStatus::Normal,
};
ObjectHeader {
object_id: conv::varint(object.object_id)?,
extension_headers_length: conv::varint(object.extension_headers.len() as u64)?,
extensions: object.extension_headers.clone(),
payload_length: conv::varint(object.payload.len() as u64)?,
object_status,
}
.encode_with_extensions(extensions, buf);
buf.put_slice(&object.payload);
Ok(())
}
}
};
}
legacy_subgroup_glue!(no_extensions sg07, "draft07", draft07);
legacy_subgroup_glue!(count_extensions sg08, "draft08", draft08);
legacy_subgroup_glue!(length_extensions sg09, "draft09", draft09);
legacy_subgroup_glue!(length_extensions sg10, "draft10", draft10);
legacy_subgroup_glue!(gated_extensions sg11, "draft11", draft11);
legacy_subgroup_glue!(gated_extensions sg12, "draft12", draft12);
legacy_subgroup_glue!(gated_extensions sg13, "draft13", draft13);
macro_rules! modern_subgroup_glue {
(derived_length $name:ident, $feat:literal, $draft:ident) => {
#[cfg(feature = $feat)]
mod $name {
use super::conv;
use super::{AnySubgroupObject, AnySubgroupObjectMeta};
use crate::error::CodecError;
use crate::$draft::data_stream::{SubgroupObject, SubgroupObjectReader};
use crate::$draft::types::ObjectStatus;
use bytes::{Buf, BufMut};
pub fn read_object(
reader: &mut SubgroupObjectReader,
buf: &mut impl Buf,
) -> Result<AnySubgroupObject, CodecError> {
let object = reader.read_object(buf)?;
Ok(AnySubgroupObject {
object_id: object.object_id.into_inner(),
extension_headers: object.extension_headers,
extension_count: None,
status: object.status.map(ObjectStatus::as_u64),
payload: object.payload,
})
}
pub fn read_object_meta(
reader: &mut SubgroupObjectReader,
buf: &mut impl Buf,
) -> Result<AnySubgroupObjectMeta, CodecError> {
let meta = reader.read_object_meta(buf)?;
Ok(AnySubgroupObjectMeta {
object_id: meta.object_id,
payload_length: meta.payload_length,
status: meta.status,
extension_headers_len: meta.extension_headers_len,
wire_len: meta.wire_len,
})
}
pub fn write_object(
writer: &mut SubgroupObjectReader,
object: &AnySubgroupObject,
buf: &mut impl BufMut,
) -> Result<(), CodecError> {
let status = match conv::status_to_write(object)? {
Some(code) => {
Some(ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?)
}
None => None,
};
writer.write_object(
&SubgroupObject {
object_id: conv::varint(object.object_id)?,
extension_headers: object.extension_headers.clone(),
status,
payload: object.payload.clone(),
},
buf,
)
}
}
};
(explicit_length $name:ident, $feat:literal, $draft:ident) => {
#[cfg(feature = $feat)]
mod $name {
use super::conv;
use super::{AnySubgroupObject, AnySubgroupObjectMeta};
use crate::error::CodecError;
use crate::$draft::data_stream::{SubgroupObject, SubgroupObjectReader};
use crate::$draft::types::ObjectStatus;
use bytes::{Buf, BufMut};
pub fn read_object(
reader: &mut SubgroupObjectReader,
buf: &mut impl Buf,
) -> Result<AnySubgroupObject, CodecError> {
let object = reader.read_object(buf)?;
Ok(AnySubgroupObject {
object_id: object.object_id.into_inner(),
extension_headers: object.extension_headers,
extension_count: None,
status: object.object_status.map(ObjectStatus::as_u64),
payload: object.payload,
})
}
pub fn read_object_meta(
reader: &mut SubgroupObjectReader,
buf: &mut impl Buf,
) -> Result<AnySubgroupObjectMeta, CodecError> {
let meta = reader.read_object_meta(buf)?;
Ok(AnySubgroupObjectMeta {
object_id: meta.object_id,
payload_length: meta.payload_length,
status: meta.status,
extension_headers_len: meta.extension_headers_len,
wire_len: meta.wire_len,
})
}
pub fn write_object(
writer: &mut SubgroupObjectReader,
object: &AnySubgroupObject,
buf: &mut impl BufMut,
) -> Result<(), CodecError> {
let object_status = match conv::status_to_write(object)? {
Some(code) => {
Some(ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?)
}
None => None,
};
writer.write_object(
&SubgroupObject {
object_id: conv::varint(object.object_id)?,
extension_headers: object.extension_headers.clone(),
payload_length: conv::varint(object.payload.len() as u64)?,
object_status,
payload: object.payload.clone(),
},
buf,
)
}
}
};
}
modern_subgroup_glue!(derived_length sg14, "draft14", draft14);
modern_subgroup_glue!(explicit_length sg15, "draft15", draft15);
modern_subgroup_glue!(explicit_length sg16, "draft16", draft16);
modern_subgroup_glue!(explicit_length sg17, "draft17", draft17);
modern_subgroup_glue!(explicit_length sg18, "draft18", draft18);
modern_subgroup_glue!(explicit_length sg19, "draft19", draft19);
modern_subgroup_glue!(explicit_length sg20, "draft20", draft20);
macro_rules! fetch_glue {
(no_extensions $name:ident, $feat:literal, $draft:ident) => {
#[cfg(feature = $feat)]
mod $name {
use super::conv;
use super::{AnyFetchObject, AnyFetchObjectMeta};
use crate::error::CodecError;
use crate::$draft::data_stream::FetchObjectHeader;
use bytes::Buf;
pub fn read_object(buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
let header = FetchObjectHeader::decode(buf)?;
let payload_length = header.payload_length.into_inner();
let (status, payload) = if payload_length == 0 {
(Some(header.object_status as u64), Vec::new())
} else {
(None, conv::take(buf, payload_length)?)
};
Ok(AnyFetchObject {
group_id: header.group_id.into_inner(),
subgroup_id: header.subgroup_id.into_inner(),
has_subgroup_id: true,
object_id: header.object_id.into_inner(),
publisher_priority: header.publisher_priority,
extension_headers: Vec::new(),
extension_count: None,
status,
end_of_range: None,
payload,
})
}
pub fn read_object_meta(buf: &mut impl Buf) -> Result<AnyFetchObjectMeta, CodecError> {
let start = buf.remaining();
let header = FetchObjectHeader::decode(buf)?;
let payload_length = header.payload_length.into_inner();
let status = if payload_length == 0 {
Some(header.object_status as u64)
} else {
conv::skip(buf, payload_length)?;
None
};
Ok(AnyFetchObjectMeta {
group_id: header.group_id.into_inner(),
subgroup_id: header.subgroup_id.into_inner(),
has_subgroup_id: true,
object_id: header.object_id.into_inner(),
publisher_priority: header.publisher_priority,
payload_length,
status,
end_of_range: None,
extension_headers_len: 0,
wire_len: (start - buf.remaining()) as u64,
})
}
}
};
(count_extensions $name:ident, $feat:literal, $draft:ident) => {
#[cfg(feature = $feat)]
mod $name {
use super::conv;
use super::{AnyFetchObject, AnyFetchObjectMeta};
use crate::error::CodecError;
use crate::$draft::data_stream::FetchObjectHeader;
use bytes::Buf;
pub fn read_object(buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
let header = FetchObjectHeader::decode(buf)?;
let payload_length = header.payload_length.into_inner();
let (status, payload) = if payload_length == 0 {
(Some(header.object_status as u64), Vec::new())
} else {
(None, conv::take(buf, payload_length)?)
};
Ok(AnyFetchObject {
group_id: header.group_id.into_inner(),
subgroup_id: header.subgroup_id.into_inner(),
has_subgroup_id: true,
object_id: header.object_id.into_inner(),
publisher_priority: header.publisher_priority,
extension_headers: header.extensions,
extension_count: Some(header.extension_count.into_inner()),
status,
end_of_range: None,
payload,
})
}
pub fn read_object_meta(buf: &mut impl Buf) -> Result<AnyFetchObjectMeta, CodecError> {
let start = buf.remaining();
let header = FetchObjectHeader::decode(buf)?;
let payload_length = header.payload_length.into_inner();
let status = if payload_length == 0 {
Some(header.object_status as u64)
} else {
conv::skip(buf, payload_length)?;
None
};
Ok(AnyFetchObjectMeta {
group_id: header.group_id.into_inner(),
subgroup_id: header.subgroup_id.into_inner(),
has_subgroup_id: true,
object_id: header.object_id.into_inner(),
publisher_priority: header.publisher_priority,
payload_length,
status,
end_of_range: None,
extension_headers_len: header.extensions.len() as u64,
wire_len: (start - buf.remaining()) as u64,
})
}
}
};
(length_extensions $name:ident, $feat:literal, $draft:ident) => {
#[cfg(feature = $feat)]
mod $name {
use super::conv;
use super::{AnyFetchObject, AnyFetchObjectMeta};
use crate::error::CodecError;
use crate::$draft::data_stream::FetchObjectHeader;
use bytes::Buf;
pub fn read_object(buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
let header = FetchObjectHeader::decode(buf)?;
let payload_length = header.payload_length.into_inner();
let (status, payload) = if payload_length == 0 {
(Some(header.object_status as u64), Vec::new())
} else {
(None, conv::take(buf, payload_length)?)
};
Ok(AnyFetchObject {
group_id: header.group_id.into_inner(),
subgroup_id: header.subgroup_id.into_inner(),
has_subgroup_id: true,
object_id: header.object_id.into_inner(),
publisher_priority: header.publisher_priority,
extension_headers: header.extensions,
extension_count: None,
status,
end_of_range: None,
payload,
})
}
pub fn read_object_meta(buf: &mut impl Buf) -> Result<AnyFetchObjectMeta, CodecError> {
let start = buf.remaining();
let header = FetchObjectHeader::decode(buf)?;
let payload_length = header.payload_length.into_inner();
let status = if payload_length == 0 {
Some(header.object_status as u64)
} else {
conv::skip(buf, payload_length)?;
None
};
Ok(AnyFetchObjectMeta {
group_id: header.group_id.into_inner(),
subgroup_id: header.subgroup_id.into_inner(),
has_subgroup_id: true,
object_id: header.object_id.into_inner(),
publisher_priority: header.publisher_priority,
payload_length,
status,
end_of_range: None,
extension_headers_len: header.extension_headers_length.into_inner(),
wire_len: (start - buf.remaining()) as u64,
})
}
}
};
}
fetch_glue!(no_extensions fo07, "draft07", draft07);
fetch_glue!(count_extensions fo08, "draft08", draft08);
fetch_glue!(length_extensions fo09, "draft09", draft09);
fetch_glue!(length_extensions fo10, "draft10", draft10);
fetch_glue!(length_extensions fo11, "draft11", draft11);
fetch_glue!(length_extensions fo12, "draft12", draft12);
fetch_glue!(length_extensions fo13, "draft13", draft13);
#[cfg(feature = "draft14")]
mod fo14 {
use super::{AnyFetchObject, AnyFetchObjectMeta};
use crate::draft14::data_stream::FetchObject;
use crate::draft14::types::ObjectStatus;
use crate::error::CodecError;
use bytes::Buf;
pub fn read_object(buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
let object = FetchObject::decode(buf)?;
Ok(AnyFetchObject {
group_id: object.group_id.into_inner(),
subgroup_id: object.subgroup_id.into_inner(),
has_subgroup_id: true,
object_id: object.object_id.into_inner(),
publisher_priority: object.publisher_priority,
extension_headers: object.extension_headers,
extension_count: None,
status: object.status.map(ObjectStatus::as_u64),
end_of_range: None,
payload: object.payload,
})
}
pub fn read_object_meta(buf: &mut impl Buf) -> Result<AnyFetchObjectMeta, CodecError> {
let meta = FetchObject::decode_meta(buf)?;
Ok(AnyFetchObjectMeta {
group_id: meta.group_id,
subgroup_id: meta.subgroup_id,
has_subgroup_id: true,
object_id: meta.object_id,
publisher_priority: meta.publisher_priority,
payload_length: meta.payload_length,
status: meta.status,
end_of_range: None,
extension_headers_len: meta.extension_headers_len,
wire_len: meta.wire_len,
})
}
}
#[cfg(feature = "draft15")]
mod fo15 {
use super::{conv, AnyFetchObject, AnyFetchObjectMeta};
use crate::draft15::data_stream::FetchObjectReader;
use crate::draft15::types::ObjectStatus;
use crate::error::CodecError;
use bytes::Buf;
pub fn read_object(
reader: &mut FetchObjectReader,
buf: &mut impl Buf,
) -> Result<AnyFetchObject, CodecError> {
let header = reader.read_object_header(buf)?;
let payload = conv::take(buf, header.payload_length.into_inner())?;
Ok(AnyFetchObject {
group_id: header.group_id.into_inner(),
subgroup_id: header.subgroup_id.into_inner(),
has_subgroup_id: true,
object_id: header.object_id.into_inner(),
publisher_priority: header.publisher_priority,
extension_headers: header.extension_headers,
extension_count: None,
status: header.object_status.map(ObjectStatus::as_u64),
end_of_range: None,
payload,
})
}
pub fn read_object_frame(
reader: &mut FetchObjectReader,
buf: &mut impl Buf,
) -> Result<super::AnyFetchFrame, CodecError> {
let start = buf.remaining();
let header = reader.read_object_header(buf)?;
let payload_length = header.payload_length.into_inner();
conv::skip(buf, payload_length)?;
let meta = AnyFetchObjectMeta {
group_id: header.group_id.into_inner(),
subgroup_id: header.subgroup_id.into_inner(),
has_subgroup_id: true,
object_id: header.object_id.into_inner(),
publisher_priority: header.publisher_priority,
payload_length,
status: header.object_status.map(ObjectStatus::as_u64),
end_of_range: None,
extension_headers_len: header.extension_headers.len() as u64,
wire_len: (start - buf.remaining()) as u64,
};
Ok(super::AnyFetchFrame {
meta,
draft: crate::version::DraftVersion::Draft15,
shape: super::FetchFrameShape::Draft15(header),
})
}
pub fn read_object_meta(
reader: &mut FetchObjectReader,
buf: &mut impl Buf,
) -> Result<AnyFetchObjectMeta, CodecError> {
read_object_frame(reader, buf).map(|frame| frame.meta)
}
}
#[cfg(feature = "draft16")]
mod fo16 {
use super::DEFAULT_PUBLISHER_PRIORITY;
use super::{conv, AnyFetchEndOfRange, AnyFetchObject, AnyFetchObjectMeta};
use crate::draft16::data_stream::{
FetchEndOfRange, FetchObjectHeader, FetchObjectLocation, FetchObjectReader,
};
use crate::error::CodecError;
use bytes::Buf;
fn parts(
reader: &mut FetchObjectReader,
buf: &mut impl Buf,
) -> Result<(FetchObjectHeader, FetchObjectLocation), CodecError> {
let header = FetchObjectHeader::decode(buf)?;
let location = reader.resolve(&header)?;
Ok((header, location))
}
fn resolved(location: &FetchObjectLocation) -> super::Resolved {
let end_of_range = location.end_of_range;
super::Resolved {
group_id: location.group_id,
subgroup_id: location.subgroup_id.filter(|_| end_of_range.is_none()),
object_id: location.object_id,
publisher_priority: location.publisher_priority.unwrap_or(DEFAULT_PUBLISHER_PRIORITY),
end_of_range: end_of_range.map(|r| match r {
FetchEndOfRange::NonExistent => AnyFetchEndOfRange::NonExistent,
FetchEndOfRange::Unknown => AnyFetchEndOfRange::Unknown,
}),
}
}
pub fn read_object(
reader: &mut FetchObjectReader,
buf: &mut impl Buf,
) -> Result<AnyFetchObject, CodecError> {
let (header, location) = parts(reader, buf)?;
let payload = conv::take(buf, header.payload_length.into_inner())?;
Ok(resolved(&location).into_object(header.extensions.unwrap_or_default(), payload))
}
pub fn read_object_frame(
reader: &mut FetchObjectReader,
buf: &mut impl Buf,
) -> Result<super::AnyFetchFrame, CodecError> {
let start = buf.remaining();
let (header, location) = parts(reader, buf)?;
let payload_length = header.payload_length.into_inner();
conv::skip(buf, payload_length)?;
let meta = resolved(&location).into_meta(
header.extensions.as_ref().map_or(0, |e| e.len() as u64),
payload_length,
(start - buf.remaining()) as u64,
);
Ok(super::AnyFetchFrame {
meta,
draft: crate::version::DraftVersion::Draft16,
shape: super::FetchFrameShape::Draft16(header, location),
})
}
pub fn read_object_meta(
reader: &mut FetchObjectReader,
buf: &mut impl Buf,
) -> Result<AnyFetchObjectMeta, CodecError> {
read_object_frame(reader, buf).map(|frame| frame.meta)
}
}
#[cfg(feature = "draft17")]
mod fo17 {
use super::{conv, AnyFetchEndOfRange, AnyFetchObject, AnyFetchObjectMeta};
use crate::draft17::data_stream::{EndOfRange, FetchObject, FetchObjectReader};
use crate::error::CodecError;
use bytes::Buf;
fn resolved(object: &FetchObject) -> super::Resolved {
super::Resolved {
group_id: object.group_id,
subgroup_id: object.subgroup_id,
object_id: object.object_id,
publisher_priority: object
.publisher_priority
.unwrap_or(super::DEFAULT_PUBLISHER_PRIORITY),
end_of_range: object.header.end_of_range().map(|r| match r {
EndOfRange::NonExistent => AnyFetchEndOfRange::NonExistent,
EndOfRange::Unknown => AnyFetchEndOfRange::Unknown,
}),
}
}
pub fn read_object(
reader: &mut FetchObjectReader,
buf: &mut impl Buf,
) -> Result<AnyFetchObject, CodecError> {
let object = reader.read_object_header(buf)?;
let resolved = resolved(&object);
let payload = conv::take(buf, object.header.payload_length.into_inner())?;
Ok(resolved.into_object(object.header.properties, payload))
}
pub fn read_object_frame(
reader: &mut FetchObjectReader,
buf: &mut impl Buf,
) -> Result<super::AnyFetchFrame, CodecError> {
let start = buf.remaining();
let object = reader.read_object_header(buf)?;
let resolved = resolved(&object);
let payload_length = object.header.payload_length.into_inner();
conv::skip(buf, payload_length)?;
let meta = resolved.into_meta(
object.header.properties.len() as u64,
payload_length,
(start - buf.remaining()) as u64,
);
Ok(super::AnyFetchFrame {
meta,
draft: crate::version::DraftVersion::Draft17,
shape: super::FetchFrameShape::Draft17(object),
})
}
pub fn read_object_meta(
reader: &mut FetchObjectReader,
buf: &mut impl Buf,
) -> Result<AnyFetchObjectMeta, CodecError> {
read_object_frame(reader, buf).map(|frame| frame.meta)
}
}
#[cfg(feature = "draft18")]
mod fo18 {
use super::{conv, AnyFetchEndOfRange, AnyFetchObject, AnyFetchObjectMeta};
use crate::draft18::data_stream::{EndOfRange, FetchObject, FetchObjectReader};
use crate::error::CodecError;
use bytes::Buf;
fn resolved(object: &FetchObject) -> super::Resolved {
super::Resolved {
group_id: object.group_id,
subgroup_id: object.subgroup_id,
object_id: object.object_id,
publisher_priority: object
.publisher_priority
.unwrap_or(super::DEFAULT_PUBLISHER_PRIORITY),
end_of_range: object.header.end_of_range().map(|r| match r {
EndOfRange::NonExistent => AnyFetchEndOfRange::NonExistent,
EndOfRange::Unknown => AnyFetchEndOfRange::Unknown,
}),
}
}
pub fn read_object(
reader: &mut FetchObjectReader,
buf: &mut impl Buf,
) -> Result<AnyFetchObject, CodecError> {
let object = reader.read_object_header(buf)?;
let resolved = resolved(&object);
let payload = conv::take(buf, object.header.payload_length.into_inner())?;
Ok(resolved.into_object(object.header.properties, payload))
}
pub fn read_object_frame(
reader: &mut FetchObjectReader,
buf: &mut impl Buf,
) -> Result<super::AnyFetchFrame, CodecError> {
let start = buf.remaining();
let object = reader.read_object_header(buf)?;
let resolved = resolved(&object);
let payload_length = object.header.payload_length.into_inner();
conv::skip(buf, payload_length)?;
let meta = resolved.into_meta(
object.header.properties.len() as u64,
payload_length,
(start - buf.remaining()) as u64,
);
Ok(super::AnyFetchFrame {
meta,
draft: crate::version::DraftVersion::Draft18,
shape: super::FetchFrameShape::Draft18(object),
})
}
pub fn read_object_meta(
reader: &mut FetchObjectReader,
buf: &mut impl Buf,
) -> Result<AnyFetchObjectMeta, CodecError> {
read_object_frame(reader, buf).map(|frame| frame.meta)
}
}
#[cfg(feature = "draft19")]
mod fo19 {
use super::{conv, AnyFetchEndOfRange, AnyFetchObject, AnyFetchObjectMeta};
use crate::draft19::data_stream::{FetchEndOfRange, FetchObject, FetchObjectReader};
use crate::error::CodecError;
use bytes::Buf;
fn resolved(object: &FetchObject) -> super::Resolved {
super::Resolved {
group_id: object.group_id,
subgroup_id: object.subgroup_id,
object_id: object.object_id,
publisher_priority: object
.publisher_priority
.unwrap_or(super::DEFAULT_PUBLISHER_PRIORITY),
end_of_range: object.header.end_of_range().map(|r| match r {
FetchEndOfRange::NonExistent => AnyFetchEndOfRange::NonExistent,
FetchEndOfRange::Unknown => AnyFetchEndOfRange::Unknown,
}),
}
}
pub fn read_object(
reader: &mut FetchObjectReader,
buf: &mut impl Buf,
) -> Result<AnyFetchObject, CodecError> {
let object = reader.read_object_header(buf)?;
let resolved = resolved(&object);
let payload = conv::take(buf, object.header.payload_length.into_inner())?;
Ok(resolved.into_object(object.header.properties.unwrap_or_default(), payload))
}
pub fn read_object_frame(
reader: &mut FetchObjectReader,
buf: &mut impl Buf,
) -> Result<super::AnyFetchFrame, CodecError> {
let start = buf.remaining();
let object = reader.read_object_header(buf)?;
let resolved = resolved(&object);
let payload_length = object.header.payload_length.into_inner();
conv::skip(buf, payload_length)?;
let meta = resolved.into_meta(
object.header.properties.as_ref().map_or(0, |p| p.len() as u64),
payload_length,
(start - buf.remaining()) as u64,
);
Ok(super::AnyFetchFrame {
meta,
draft: crate::version::DraftVersion::Draft19,
shape: super::FetchFrameShape::Draft19(object),
})
}
pub fn read_object_meta(
reader: &mut FetchObjectReader,
buf: &mut impl Buf,
) -> Result<AnyFetchObjectMeta, CodecError> {
read_object_frame(reader, buf).map(|frame| frame.meta)
}
}
#[cfg(feature = "draft20")]
mod fo20 {
use super::{conv, AnyFetchEndOfRange, AnyFetchObject, AnyFetchObjectMeta};
use crate::draft20::data_stream::{FetchEndOfRange, FetchObject, FetchObjectReader};
use crate::error::CodecError;
use bytes::Buf;
fn resolved(object: &FetchObject) -> super::Resolved {
super::Resolved {
group_id: object.group_id,
subgroup_id: object.subgroup_id,
object_id: object.object_id,
publisher_priority: object
.publisher_priority
.unwrap_or(super::DEFAULT_PUBLISHER_PRIORITY),
end_of_range: object.header.end_of_range().map(|r| match r {
FetchEndOfRange::NonExistent => AnyFetchEndOfRange::NonExistent,
FetchEndOfRange::Unknown => AnyFetchEndOfRange::Unknown,
FetchEndOfRange::TimedOut => AnyFetchEndOfRange::TimedOut,
}),
}
}
pub fn read_object(
reader: &mut FetchObjectReader,
buf: &mut impl Buf,
) -> Result<AnyFetchObject, CodecError> {
let object = reader.read_object_header(buf)?;
let resolved = resolved(&object);
let payload = conv::take(buf, object.header.payload_length.into_inner())?;
Ok(resolved.into_object(object.header.properties.unwrap_or_default(), payload))
}
pub fn read_object_frame(
reader: &mut FetchObjectReader,
buf: &mut impl Buf,
) -> Result<super::AnyFetchFrame, CodecError> {
let start = buf.remaining();
let object = reader.read_object_header(buf)?;
let resolved = resolved(&object);
let payload_length = object.header.payload_length.into_inner();
conv::skip(buf, payload_length)?;
let meta = resolved.into_meta(
object.header.properties.as_ref().map_or(0, |p| p.len() as u64),
payload_length,
(start - buf.remaining()) as u64,
);
Ok(super::AnyFetchFrame {
meta,
draft: crate::version::DraftVersion::Draft20,
shape: super::FetchFrameShape::Draft20(object),
})
}
pub fn read_object_meta(
reader: &mut FetchObjectReader,
buf: &mut impl Buf,
) -> Result<AnyFetchObjectMeta, CodecError> {
read_object_frame(reader, buf).map(|frame| frame.meta)
}
}
#[cfg(any(
feature = "draft16",
feature = "draft17",
feature = "draft18",
feature = "draft19",
feature = "draft20"
))]
struct Resolved {
group_id: u64,
subgroup_id: Option<u64>,
object_id: u64,
publisher_priority: u8,
end_of_range: Option<AnyFetchEndOfRange>,
}
#[cfg(any(
feature = "draft16",
feature = "draft17",
feature = "draft18",
feature = "draft19",
feature = "draft20"
))]
impl Resolved {
fn into_object(self, extension_headers: Vec<u8>, payload: Vec<u8>) -> AnyFetchObject {
AnyFetchObject {
group_id: self.group_id,
subgroup_id: self.subgroup_id.unwrap_or(0),
has_subgroup_id: self.subgroup_id.is_some(),
object_id: self.object_id,
publisher_priority: self.publisher_priority,
extension_headers,
extension_count: None,
status: None,
end_of_range: self.end_of_range,
payload,
}
}
fn into_meta(
self,
extension_headers_len: u64,
payload_length: u64,
wire_len: u64,
) -> AnyFetchObjectMeta {
AnyFetchObjectMeta {
group_id: self.group_id,
subgroup_id: self.subgroup_id.unwrap_or(0),
has_subgroup_id: self.subgroup_id.is_some(),
object_id: self.object_id,
publisher_priority: self.publisher_priority,
payload_length,
status: None,
end_of_range: self.end_of_range,
extension_headers_len,
wire_len,
}
}
}
#[derive(Debug, Clone)]
enum SubgroupReaderState {
#[cfg(feature = "draft07")]
Draft07,
#[cfg(feature = "draft08")]
Draft08,
#[cfg(feature = "draft09")]
Draft09,
#[cfg(feature = "draft10")]
Draft10,
#[cfg(feature = "draft11")]
Draft11 { extensions: bool },
#[cfg(feature = "draft12")]
Draft12 { extensions: bool },
#[cfg(feature = "draft13")]
Draft13 { extensions: bool },
#[cfg(feature = "draft14")]
Draft14(crate::draft14::data_stream::SubgroupObjectReader),
#[cfg(feature = "draft15")]
Draft15(crate::draft15::data_stream::SubgroupObjectReader),
#[cfg(feature = "draft16")]
Draft16(crate::draft16::data_stream::SubgroupObjectReader),
#[cfg(feature = "draft17")]
Draft17(crate::draft17::data_stream::SubgroupObjectReader),
#[cfg(feature = "draft18")]
Draft18(crate::draft18::data_stream::SubgroupObjectReader),
#[cfg(feature = "draft19")]
Draft19(crate::draft19::data_stream::SubgroupObjectReader),
#[cfg(feature = "draft20")]
Draft20(crate::draft20::data_stream::SubgroupObjectReader),
}
#[derive(Debug, Clone)]
pub struct AnySubgroupObjectReader {
state: SubgroupReaderState,
}
impl AnySubgroupObjectReader {
#[allow(unused_variables, unreachable_code)]
pub fn new(header: &AnySubgroupHeader) -> Result<Self, CodecError> {
let state = match header {
#[cfg(feature = "draft07")]
AnySubgroupHeader::Draft07(_) => SubgroupReaderState::Draft07,
#[cfg(feature = "draft08")]
AnySubgroupHeader::Draft08(_) => SubgroupReaderState::Draft08,
#[cfg(feature = "draft09")]
AnySubgroupHeader::Draft09(_) => SubgroupReaderState::Draft09,
#[cfg(feature = "draft10")]
AnySubgroupHeader::Draft10(_) => SubgroupReaderState::Draft10,
#[cfg(feature = "draft11")]
AnySubgroupHeader::Draft11(h) => {
SubgroupReaderState::Draft11 { extensions: subgroup_extensions_11(h)? }
}
#[cfg(feature = "draft12")]
AnySubgroupHeader::Draft12(h) => {
SubgroupReaderState::Draft12 { extensions: subgroup_extensions_12(h)? }
}
#[cfg(feature = "draft13")]
AnySubgroupHeader::Draft13(h) => {
SubgroupReaderState::Draft13 { extensions: subgroup_extensions_13(h)? }
}
#[cfg(feature = "draft14")]
AnySubgroupHeader::Draft14(h) => SubgroupReaderState::Draft14(
crate::draft14::data_stream::SubgroupObjectReader::new(h),
),
#[cfg(feature = "draft15")]
AnySubgroupHeader::Draft15(h) => SubgroupReaderState::Draft15(
crate::draft15::data_stream::SubgroupObjectReader::new(h),
),
#[cfg(feature = "draft16")]
AnySubgroupHeader::Draft16(h) => SubgroupReaderState::Draft16(
crate::draft16::data_stream::SubgroupObjectReader::new(h),
),
#[cfg(feature = "draft17")]
AnySubgroupHeader::Draft17(h) => SubgroupReaderState::Draft17(
crate::draft17::data_stream::SubgroupObjectReader::new(h),
),
#[cfg(feature = "draft18")]
AnySubgroupHeader::Draft18(h) => SubgroupReaderState::Draft18(
crate::draft18::data_stream::SubgroupObjectReader::new(h),
),
#[cfg(feature = "draft19")]
AnySubgroupHeader::Draft19(h) => SubgroupReaderState::Draft19(
crate::draft19::data_stream::SubgroupObjectReader::new(h),
),
#[cfg(feature = "draft20")]
AnySubgroupHeader::Draft20(h) => SubgroupReaderState::Draft20(
crate::draft20::data_stream::SubgroupObjectReader::new(h),
),
#[allow(unreachable_patterns)]
_ => {
return Err(CodecError::UnsupportedDraft(format!(
"draft {:?} not enabled via feature flag",
header.draft()
)));
}
};
Ok(Self { state })
}
#[allow(unreachable_code)]
pub fn draft(&self) -> DraftVersion {
match &self.state {
#[cfg(feature = "draft07")]
SubgroupReaderState::Draft07 => DraftVersion::Draft07,
#[cfg(feature = "draft08")]
SubgroupReaderState::Draft08 => DraftVersion::Draft08,
#[cfg(feature = "draft09")]
SubgroupReaderState::Draft09 => DraftVersion::Draft09,
#[cfg(feature = "draft10")]
SubgroupReaderState::Draft10 => DraftVersion::Draft10,
#[cfg(feature = "draft11")]
SubgroupReaderState::Draft11 { .. } => DraftVersion::Draft11,
#[cfg(feature = "draft12")]
SubgroupReaderState::Draft12 { .. } => DraftVersion::Draft12,
#[cfg(feature = "draft13")]
SubgroupReaderState::Draft13 { .. } => DraftVersion::Draft13,
#[cfg(feature = "draft14")]
SubgroupReaderState::Draft14(_) => DraftVersion::Draft14,
#[cfg(feature = "draft15")]
SubgroupReaderState::Draft15(_) => DraftVersion::Draft15,
#[cfg(feature = "draft16")]
SubgroupReaderState::Draft16(_) => DraftVersion::Draft16,
#[cfg(feature = "draft17")]
SubgroupReaderState::Draft17(_) => DraftVersion::Draft17,
#[cfg(feature = "draft18")]
SubgroupReaderState::Draft18(_) => DraftVersion::Draft18,
#[cfg(feature = "draft19")]
SubgroupReaderState::Draft19(_) => DraftVersion::Draft19,
#[cfg(feature = "draft20")]
SubgroupReaderState::Draft20(_) => DraftVersion::Draft20,
#[allow(unreachable_patterns)]
_ => unreachable!("AnySubgroupObjectReader has no enabled variants"),
}
}
#[allow(unused_variables, unreachable_code)]
pub fn read_object(&mut self, buf: &mut impl Buf) -> Result<AnySubgroupObject, CodecError> {
match &mut self.state {
#[cfg(feature = "draft07")]
SubgroupReaderState::Draft07 => sg07::read_object(buf),
#[cfg(feature = "draft08")]
SubgroupReaderState::Draft08 => sg08::read_object(buf),
#[cfg(feature = "draft09")]
SubgroupReaderState::Draft09 => sg09::read_object(buf),
#[cfg(feature = "draft10")]
SubgroupReaderState::Draft10 => sg10::read_object(buf),
#[cfg(feature = "draft11")]
SubgroupReaderState::Draft11 { extensions } => sg11::read_object(*extensions, buf),
#[cfg(feature = "draft12")]
SubgroupReaderState::Draft12 { extensions } => sg12::read_object(*extensions, buf),
#[cfg(feature = "draft13")]
SubgroupReaderState::Draft13 { extensions } => sg13::read_object(*extensions, buf),
#[cfg(feature = "draft14")]
SubgroupReaderState::Draft14(inner) => sg14::read_object(inner, buf),
#[cfg(feature = "draft15")]
SubgroupReaderState::Draft15(inner) => sg15::read_object(inner, buf),
#[cfg(feature = "draft16")]
SubgroupReaderState::Draft16(inner) => sg16::read_object(inner, buf),
#[cfg(feature = "draft17")]
SubgroupReaderState::Draft17(inner) => sg17::read_object(inner, buf),
#[cfg(feature = "draft18")]
SubgroupReaderState::Draft18(inner) => sg18::read_object(inner, buf),
#[cfg(feature = "draft19")]
SubgroupReaderState::Draft19(inner) => sg19::read_object(inner, buf),
#[cfg(feature = "draft20")]
SubgroupReaderState::Draft20(inner) => sg20::read_object(inner, buf),
#[allow(unreachable_patterns)]
_ => unreachable!("AnySubgroupObjectReader has no enabled variants"),
}
}
#[allow(unused_variables, unreachable_code)]
pub fn read_object_meta(
&mut self,
buf: &mut impl Buf,
) -> Result<AnySubgroupObjectMeta, CodecError> {
match &mut self.state {
#[cfg(feature = "draft07")]
SubgroupReaderState::Draft07 => sg07::read_object_meta(buf),
#[cfg(feature = "draft08")]
SubgroupReaderState::Draft08 => sg08::read_object_meta(buf),
#[cfg(feature = "draft09")]
SubgroupReaderState::Draft09 => sg09::read_object_meta(buf),
#[cfg(feature = "draft10")]
SubgroupReaderState::Draft10 => sg10::read_object_meta(buf),
#[cfg(feature = "draft11")]
SubgroupReaderState::Draft11 { extensions } => sg11::read_object_meta(*extensions, buf),
#[cfg(feature = "draft12")]
SubgroupReaderState::Draft12 { extensions } => sg12::read_object_meta(*extensions, buf),
#[cfg(feature = "draft13")]
SubgroupReaderState::Draft13 { extensions } => sg13::read_object_meta(*extensions, buf),
#[cfg(feature = "draft14")]
SubgroupReaderState::Draft14(inner) => sg14::read_object_meta(inner, buf),
#[cfg(feature = "draft15")]
SubgroupReaderState::Draft15(inner) => sg15::read_object_meta(inner, buf),
#[cfg(feature = "draft16")]
SubgroupReaderState::Draft16(inner) => sg16::read_object_meta(inner, buf),
#[cfg(feature = "draft17")]
SubgroupReaderState::Draft17(inner) => sg17::read_object_meta(inner, buf),
#[cfg(feature = "draft18")]
SubgroupReaderState::Draft18(inner) => sg18::read_object_meta(inner, buf),
#[cfg(feature = "draft19")]
SubgroupReaderState::Draft19(inner) => sg19::read_object_meta(inner, buf),
#[cfg(feature = "draft20")]
SubgroupReaderState::Draft20(inner) => sg20::read_object_meta(inner, buf),
#[allow(unreachable_patterns)]
_ => unreachable!("AnySubgroupObjectReader has no enabled variants"),
}
}
}
#[derive(Debug, Clone)]
enum SubgroupWriterState {
#[cfg(feature = "draft07")]
Draft07 { prev_object_id: Option<u64> },
#[cfg(feature = "draft08")]
Draft08 { prev_object_id: Option<u64> },
#[cfg(feature = "draft09")]
Draft09 { prev_object_id: Option<u64> },
#[cfg(feature = "draft10")]
Draft10 { prev_object_id: Option<u64> },
#[cfg(feature = "draft11")]
Draft11 { extensions: bool, prev_object_id: Option<u64> },
#[cfg(feature = "draft12")]
Draft12 { extensions: bool, prev_object_id: Option<u64> },
#[cfg(feature = "draft13")]
Draft13 { extensions: bool, prev_object_id: Option<u64> },
#[cfg(feature = "draft14")]
Draft14 { inner: crate::draft14::data_stream::SubgroupObjectReader, extensions: bool },
#[cfg(feature = "draft15")]
Draft15 { inner: crate::draft15::data_stream::SubgroupObjectReader, extensions: bool },
#[cfg(feature = "draft16")]
Draft16 { inner: crate::draft16::data_stream::SubgroupObjectReader, extensions: bool },
#[cfg(feature = "draft17")]
Draft17 { inner: crate::draft17::data_stream::SubgroupObjectReader, extensions: bool },
#[cfg(feature = "draft18")]
Draft18 { inner: crate::draft18::data_stream::SubgroupObjectReader, extensions: bool },
#[cfg(feature = "draft19")]
Draft19 { inner: crate::draft19::data_stream::SubgroupObjectReader, extensions: bool },
#[cfg(feature = "draft20")]
Draft20 { inner: crate::draft20::data_stream::SubgroupObjectReader, extensions: bool },
}
#[derive(Debug, Clone)]
pub struct AnySubgroupObjectWriter {
state: SubgroupWriterState,
}
impl AnySubgroupObjectWriter {
#[allow(unused_variables, unreachable_code)]
pub fn new(header: &AnySubgroupHeader) -> Result<Self, CodecError> {
let state = match header {
#[cfg(feature = "draft07")]
AnySubgroupHeader::Draft07(_) => SubgroupWriterState::Draft07 { prev_object_id: None },
#[cfg(feature = "draft08")]
AnySubgroupHeader::Draft08(_) => SubgroupWriterState::Draft08 { prev_object_id: None },
#[cfg(feature = "draft09")]
AnySubgroupHeader::Draft09(_) => SubgroupWriterState::Draft09 { prev_object_id: None },
#[cfg(feature = "draft10")]
AnySubgroupHeader::Draft10(_) => SubgroupWriterState::Draft10 { prev_object_id: None },
#[cfg(feature = "draft11")]
AnySubgroupHeader::Draft11(h) => SubgroupWriterState::Draft11 {
extensions: subgroup_extensions_11(h)?,
prev_object_id: None,
},
#[cfg(feature = "draft12")]
AnySubgroupHeader::Draft12(h) => SubgroupWriterState::Draft12 {
extensions: subgroup_extensions_12(h)?,
prev_object_id: None,
},
#[cfg(feature = "draft13")]
AnySubgroupHeader::Draft13(h) => SubgroupWriterState::Draft13 {
extensions: subgroup_extensions_13(h)?,
prev_object_id: None,
},
#[cfg(feature = "draft14")]
AnySubgroupHeader::Draft14(h) => SubgroupWriterState::Draft14 {
inner: crate::draft14::data_stream::SubgroupObjectReader::new(h),
extensions: h.stream_type.extensions_present(),
},
#[cfg(feature = "draft15")]
AnySubgroupHeader::Draft15(h) => SubgroupWriterState::Draft15 {
inner: crate::draft15::data_stream::SubgroupObjectReader::new(h),
extensions: h.has_extensions(),
},
#[cfg(feature = "draft16")]
AnySubgroupHeader::Draft16(h) => SubgroupWriterState::Draft16 {
inner: crate::draft16::data_stream::SubgroupObjectReader::new(h),
extensions: h.has_extensions(),
},
#[cfg(feature = "draft17")]
AnySubgroupHeader::Draft17(h) => SubgroupWriterState::Draft17 {
inner: crate::draft17::data_stream::SubgroupObjectReader::new(h),
extensions: h.has_properties(),
},
#[cfg(feature = "draft18")]
AnySubgroupHeader::Draft18(h) => SubgroupWriterState::Draft18 {
inner: crate::draft18::data_stream::SubgroupObjectReader::new(h),
extensions: h.has_properties(),
},
#[cfg(feature = "draft19")]
AnySubgroupHeader::Draft19(h) => SubgroupWriterState::Draft19 {
inner: crate::draft19::data_stream::SubgroupObjectReader::new(h),
extensions: h.has_properties(),
},
#[cfg(feature = "draft20")]
AnySubgroupHeader::Draft20(h) => SubgroupWriterState::Draft20 {
inner: crate::draft20::data_stream::SubgroupObjectReader::new(h),
extensions: h.has_properties(),
},
#[allow(unreachable_patterns)]
_ => {
return Err(CodecError::UnsupportedDraft(format!(
"draft {:?} not enabled via feature flag",
header.draft()
)));
}
};
Ok(Self { state })
}
#[allow(unreachable_code)]
pub fn draft(&self) -> DraftVersion {
match &self.state {
#[cfg(feature = "draft07")]
SubgroupWriterState::Draft07 { .. } => DraftVersion::Draft07,
#[cfg(feature = "draft08")]
SubgroupWriterState::Draft08 { .. } => DraftVersion::Draft08,
#[cfg(feature = "draft09")]
SubgroupWriterState::Draft09 { .. } => DraftVersion::Draft09,
#[cfg(feature = "draft10")]
SubgroupWriterState::Draft10 { .. } => DraftVersion::Draft10,
#[cfg(feature = "draft11")]
SubgroupWriterState::Draft11 { .. } => DraftVersion::Draft11,
#[cfg(feature = "draft12")]
SubgroupWriterState::Draft12 { .. } => DraftVersion::Draft12,
#[cfg(feature = "draft13")]
SubgroupWriterState::Draft13 { .. } => DraftVersion::Draft13,
#[cfg(feature = "draft14")]
SubgroupWriterState::Draft14 { .. } => DraftVersion::Draft14,
#[cfg(feature = "draft15")]
SubgroupWriterState::Draft15 { .. } => DraftVersion::Draft15,
#[cfg(feature = "draft16")]
SubgroupWriterState::Draft16 { .. } => DraftVersion::Draft16,
#[cfg(feature = "draft17")]
SubgroupWriterState::Draft17 { .. } => DraftVersion::Draft17,
#[cfg(feature = "draft18")]
SubgroupWriterState::Draft18 { .. } => DraftVersion::Draft18,
#[cfg(feature = "draft19")]
SubgroupWriterState::Draft19 { .. } => DraftVersion::Draft19,
#[cfg(feature = "draft20")]
SubgroupWriterState::Draft20 { .. } => DraftVersion::Draft20,
#[allow(unreachable_patterns)]
_ => unreachable!("AnySubgroupObjectWriter has no enabled variants"),
}
}
#[allow(unused_variables, unreachable_code)]
pub fn write_object(
&mut self,
object: &AnySubgroupObject,
buf: &mut impl BufMut,
) -> Result<(), CodecError> {
match &mut self.state {
#[cfg(feature = "draft07")]
SubgroupWriterState::Draft07 { prev_object_id } => {
advance_absolute_id(prev_object_id, object, |o| sg07::write_object(o, buf))
}
#[cfg(feature = "draft08")]
SubgroupWriterState::Draft08 { prev_object_id } => {
advance_absolute_id(prev_object_id, object, |o| sg08::write_object(o, buf))
}
#[cfg(feature = "draft09")]
SubgroupWriterState::Draft09 { prev_object_id } => {
advance_absolute_id(prev_object_id, object, |o| sg09::write_object(o, buf))
}
#[cfg(feature = "draft10")]
SubgroupWriterState::Draft10 { prev_object_id } => {
advance_absolute_id(prev_object_id, object, |o| sg10::write_object(o, buf))
}
#[cfg(feature = "draft11")]
SubgroupWriterState::Draft11 { extensions, prev_object_id } => {
let extensions = *extensions;
advance_absolute_id(prev_object_id, object, |o| {
sg11::write_object(extensions, o, buf)
})
}
#[cfg(feature = "draft12")]
SubgroupWriterState::Draft12 { extensions, prev_object_id } => {
let extensions = *extensions;
advance_absolute_id(prev_object_id, object, |o| {
sg12::write_object(extensions, o, buf)
})
}
#[cfg(feature = "draft13")]
SubgroupWriterState::Draft13 { extensions, prev_object_id } => {
let extensions = *extensions;
advance_absolute_id(prev_object_id, object, |o| {
sg13::write_object(extensions, o, buf)
})
}
#[cfg(feature = "draft14")]
SubgroupWriterState::Draft14 { inner, extensions } => {
reject_unrepresentable_extensions(*extensions, object)?;
sg14::write_object(inner, object, buf)
}
#[cfg(feature = "draft15")]
SubgroupWriterState::Draft15 { inner, extensions } => {
reject_unrepresentable_extensions(*extensions, object)?;
sg15::write_object(inner, object, buf)
}
#[cfg(feature = "draft16")]
SubgroupWriterState::Draft16 { inner, extensions } => {
reject_unrepresentable_extensions(*extensions, object)?;
sg16::write_object(inner, object, buf)
}
#[cfg(feature = "draft17")]
SubgroupWriterState::Draft17 { inner, extensions } => {
reject_unrepresentable_extensions(*extensions, object)?;
sg17::write_object(inner, object, buf)
}
#[cfg(feature = "draft18")]
SubgroupWriterState::Draft18 { inner, extensions } => {
reject_unrepresentable_extensions(*extensions, object)?;
sg18::write_object(inner, object, buf)
}
#[cfg(feature = "draft19")]
SubgroupWriterState::Draft19 { inner, extensions } => {
reject_unrepresentable_extensions(*extensions, object)?;
sg19::write_object(inner, object, buf)
}
#[cfg(feature = "draft20")]
SubgroupWriterState::Draft20 { inner, extensions } => {
reject_unrepresentable_extensions(*extensions, object)?;
sg20::write_object(inner, object, buf)
}
#[allow(unreachable_patterns)]
_ => unreachable!("AnySubgroupObjectWriter has no enabled variants"),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Reemit {
Verbatim,
Reencoded {
id_bytes_before: usize,
id_bytes_after: usize,
},
}
pub fn reemit_subgroup_object(
draft: DraftVersion,
prev_forwarded: Option<u64>,
object_id: u64,
raw: &[u8],
out: &mut impl BufMut,
) -> Result<Reemit, CodecError> {
if matches!(prev_forwarded, Some(prev) if object_id <= prev) {
return Err(CodecError::InvalidField);
}
let mut cursor: &[u8] = raw;
draft.decode_varint(&mut cursor).map_err(|_| CodecError::InvalidField)?;
let id_bytes_before = raw.len() - cursor.len();
if !delta_encodes_object_ids(draft) {
out.put_slice(raw);
return Ok(Reemit::Verbatim);
}
let delta = match prev_forwarded {
None => object_id,
Some(prev) => object_id
.checked_sub(prev)
.and_then(|v| v.checked_sub(1))
.ok_or(CodecError::InvalidField)?,
};
let field = if draft.uses_moqt_varint() {
VarInt::from_u64_moqt(delta)
} else {
VarInt::from_u64(delta).map_err(|_| CodecError::InvalidField)?
};
let mut encoded = [0u8; 9];
let mut slot: &mut [u8] = &mut encoded;
draft.encode_varint(field, &mut slot);
let id_bytes_after = 9 - slot.len();
let encoded = &encoded[..id_bytes_after];
if encoded == &raw[..id_bytes_before] {
out.put_slice(raw);
return Ok(Reemit::Verbatim);
}
out.put_slice(encoded);
out.put_slice(&raw[id_bytes_before..]);
Ok(Reemit::Reencoded { id_bytes_before, id_bytes_after })
}
fn delta_encodes_object_ids(draft: DraftVersion) -> bool {
matches!(
draft,
DraftVersion::Draft14
| DraftVersion::Draft15
| DraftVersion::Draft16
| DraftVersion::Draft17
| DraftVersion::Draft18
| DraftVersion::Draft19
| DraftVersion::Draft20
)
}
#[cfg(any(
feature = "draft07",
feature = "draft08",
feature = "draft09",
feature = "draft10",
feature = "draft11",
feature = "draft12",
feature = "draft13"
))]
fn advance_absolute_id(
prev_object_id: &mut Option<u64>,
object: &AnySubgroupObject,
write: impl FnOnce(&AnySubgroupObject) -> Result<(), CodecError>,
) -> Result<(), CodecError> {
if matches!(*prev_object_id, Some(prev) if object.object_id <= prev) {
return Err(CodecError::InvalidField);
}
write(object)?;
*prev_object_id = Some(object.object_id);
Ok(())
}
#[cfg(any(
feature = "draft14",
feature = "draft15",
feature = "draft16",
feature = "draft17",
feature = "draft18",
feature = "draft19",
feature = "draft20"
))]
fn reject_unrepresentable_extensions(
extensions: bool,
object: &AnySubgroupObject,
) -> Result<(), CodecError> {
if !extensions && !object.extension_headers.is_empty() {
return Err(CodecError::InvalidField);
}
Ok(())
}
#[derive(Debug, Clone)]
enum FetchReaderState {
#[cfg(feature = "draft07")]
Draft07,
#[cfg(feature = "draft08")]
Draft08,
#[cfg(feature = "draft09")]
Draft09,
#[cfg(feature = "draft10")]
Draft10,
#[cfg(feature = "draft11")]
Draft11,
#[cfg(feature = "draft12")]
Draft12,
#[cfg(feature = "draft13")]
Draft13,
#[cfg(feature = "draft14")]
Draft14,
#[cfg(feature = "draft15")]
Draft15(crate::draft15::data_stream::FetchObjectReader),
#[cfg(feature = "draft16")]
Draft16(crate::draft16::data_stream::FetchObjectReader),
#[cfg(feature = "draft17")]
Draft17(crate::draft17::data_stream::FetchObjectReader),
#[cfg(feature = "draft18")]
Draft18(crate::draft18::data_stream::FetchObjectReader),
#[cfg(feature = "draft19")]
Draft19(crate::draft19::data_stream::FetchObjectReader),
#[cfg(feature = "draft20")]
Draft20(crate::draft20::data_stream::FetchObjectReader),
}
#[derive(Debug, Clone)]
pub struct AnyFetchObjectReader {
state: FetchReaderState,
}
impl AnyFetchObjectReader {
#[allow(unused_variables, unreachable_code)]
pub fn new(
header: &AnyFetchHeader,
group_order: AnyFetchGroupOrder,
) -> Result<Self, CodecError> {
let state = match header {
#[cfg(feature = "draft07")]
AnyFetchHeader::Draft07(_) => FetchReaderState::Draft07,
#[cfg(feature = "draft08")]
AnyFetchHeader::Draft08(_) => FetchReaderState::Draft08,
#[cfg(feature = "draft09")]
AnyFetchHeader::Draft09(_) => FetchReaderState::Draft09,
#[cfg(feature = "draft10")]
AnyFetchHeader::Draft10(_) => FetchReaderState::Draft10,
#[cfg(feature = "draft11")]
AnyFetchHeader::Draft11(_) => FetchReaderState::Draft11,
#[cfg(feature = "draft12")]
AnyFetchHeader::Draft12(_) => FetchReaderState::Draft12,
#[cfg(feature = "draft13")]
AnyFetchHeader::Draft13(_) => FetchReaderState::Draft13,
#[cfg(feature = "draft14")]
AnyFetchHeader::Draft14(_) => FetchReaderState::Draft14,
#[cfg(feature = "draft15")]
AnyFetchHeader::Draft15(_) => {
FetchReaderState::Draft15(crate::draft15::data_stream::FetchObjectReader::new())
}
#[cfg(feature = "draft16")]
AnyFetchHeader::Draft16(_) => {
FetchReaderState::Draft16(crate::draft16::data_stream::FetchObjectReader::new())
}
#[cfg(feature = "draft17")]
AnyFetchHeader::Draft17(_) => {
FetchReaderState::Draft17(crate::draft17::data_stream::FetchObjectReader::new())
}
#[cfg(feature = "draft18")]
AnyFetchHeader::Draft18(_) => FetchReaderState::Draft18(
crate::draft18::data_stream::FetchObjectReader::new(match group_order {
AnyFetchGroupOrder::Ascending => {
crate::draft18::data_stream::GroupOrder::Ascending
}
AnyFetchGroupOrder::Descending => {
crate::draft18::data_stream::GroupOrder::Descending
}
}),
),
#[cfg(feature = "draft19")]
AnyFetchHeader::Draft19(_) => FetchReaderState::Draft19(
crate::draft19::data_stream::FetchObjectReader::new(match group_order {
AnyFetchGroupOrder::Ascending => {
crate::draft19::data_stream::GroupOrder::Ascending
}
AnyFetchGroupOrder::Descending => {
crate::draft19::data_stream::GroupOrder::Descending
}
}),
),
#[cfg(feature = "draft20")]
AnyFetchHeader::Draft20(_) => FetchReaderState::Draft20(
crate::draft20::data_stream::FetchObjectReader::new(match group_order {
AnyFetchGroupOrder::Ascending => {
crate::draft20::data_stream::GroupOrder::Ascending
}
AnyFetchGroupOrder::Descending => {
crate::draft20::data_stream::GroupOrder::Descending
}
}),
),
#[allow(unreachable_patterns)]
_ => {
return Err(CodecError::UnsupportedDraft(format!(
"draft {:?} not enabled via feature flag",
header.draft()
)));
}
};
Ok(Self { state })
}
#[allow(unreachable_code)]
pub fn draft(&self) -> DraftVersion {
match &self.state {
#[cfg(feature = "draft07")]
FetchReaderState::Draft07 => DraftVersion::Draft07,
#[cfg(feature = "draft08")]
FetchReaderState::Draft08 => DraftVersion::Draft08,
#[cfg(feature = "draft09")]
FetchReaderState::Draft09 => DraftVersion::Draft09,
#[cfg(feature = "draft10")]
FetchReaderState::Draft10 => DraftVersion::Draft10,
#[cfg(feature = "draft11")]
FetchReaderState::Draft11 => DraftVersion::Draft11,
#[cfg(feature = "draft12")]
FetchReaderState::Draft12 => DraftVersion::Draft12,
#[cfg(feature = "draft13")]
FetchReaderState::Draft13 => DraftVersion::Draft13,
#[cfg(feature = "draft14")]
FetchReaderState::Draft14 => DraftVersion::Draft14,
#[cfg(feature = "draft15")]
FetchReaderState::Draft15(_) => DraftVersion::Draft15,
#[cfg(feature = "draft16")]
FetchReaderState::Draft16(_) => DraftVersion::Draft16,
#[cfg(feature = "draft17")]
FetchReaderState::Draft17(_) => DraftVersion::Draft17,
#[cfg(feature = "draft18")]
FetchReaderState::Draft18(_) => DraftVersion::Draft18,
#[cfg(feature = "draft19")]
FetchReaderState::Draft19(_) => DraftVersion::Draft19,
#[cfg(feature = "draft20")]
FetchReaderState::Draft20(_) => DraftVersion::Draft20,
#[allow(unreachable_patterns)]
_ => unreachable!("AnyFetchObjectReader has no enabled variants"),
}
}
#[allow(unused_variables, unreachable_code)]
pub fn read_object(&mut self, buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
match &mut self.state {
#[cfg(feature = "draft07")]
FetchReaderState::Draft07 => fo07::read_object(buf),
#[cfg(feature = "draft08")]
FetchReaderState::Draft08 => fo08::read_object(buf),
#[cfg(feature = "draft09")]
FetchReaderState::Draft09 => fo09::read_object(buf),
#[cfg(feature = "draft10")]
FetchReaderState::Draft10 => fo10::read_object(buf),
#[cfg(feature = "draft11")]
FetchReaderState::Draft11 => fo11::read_object(buf),
#[cfg(feature = "draft12")]
FetchReaderState::Draft12 => fo12::read_object(buf),
#[cfg(feature = "draft13")]
FetchReaderState::Draft13 => fo13::read_object(buf),
#[cfg(feature = "draft14")]
FetchReaderState::Draft14 => fo14::read_object(buf),
#[cfg(feature = "draft15")]
FetchReaderState::Draft15(inner) => fo15::read_object(inner, buf),
#[cfg(feature = "draft16")]
FetchReaderState::Draft16(inner) => fo16::read_object(inner, buf),
#[cfg(feature = "draft17")]
FetchReaderState::Draft17(inner) => fo17::read_object(inner, buf),
#[cfg(feature = "draft18")]
FetchReaderState::Draft18(inner) => fo18::read_object(inner, buf),
#[cfg(feature = "draft19")]
FetchReaderState::Draft19(inner) => fo19::read_object(inner, buf),
#[cfg(feature = "draft20")]
FetchReaderState::Draft20(inner) => fo20::read_object(inner, buf),
#[allow(unreachable_patterns)]
_ => unreachable!("AnyFetchObjectReader has no enabled variants"),
}
}
#[allow(unused_variables, unreachable_code)]
pub fn read_object_frame(&mut self, buf: &mut impl Buf) -> Result<AnyFetchFrame, CodecError> {
match &mut self.state {
#[cfg(feature = "draft07")]
FetchReaderState::Draft07 => fo07::read_object_meta(buf)
.map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft07, meta)),
#[cfg(feature = "draft08")]
FetchReaderState::Draft08 => fo08::read_object_meta(buf)
.map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft08, meta)),
#[cfg(feature = "draft09")]
FetchReaderState::Draft09 => fo09::read_object_meta(buf)
.map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft09, meta)),
#[cfg(feature = "draft10")]
FetchReaderState::Draft10 => fo10::read_object_meta(buf)
.map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft10, meta)),
#[cfg(feature = "draft11")]
FetchReaderState::Draft11 => fo11::read_object_meta(buf)
.map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft11, meta)),
#[cfg(feature = "draft12")]
FetchReaderState::Draft12 => fo12::read_object_meta(buf)
.map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft12, meta)),
#[cfg(feature = "draft13")]
FetchReaderState::Draft13 => fo13::read_object_meta(buf)
.map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft13, meta)),
#[cfg(feature = "draft14")]
FetchReaderState::Draft14 => fo14::read_object_meta(buf)
.map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft14, meta)),
#[cfg(feature = "draft15")]
FetchReaderState::Draft15(inner) => fo15::read_object_frame(inner, buf),
#[cfg(feature = "draft16")]
FetchReaderState::Draft16(inner) => fo16::read_object_frame(inner, buf),
#[cfg(feature = "draft17")]
FetchReaderState::Draft17(inner) => fo17::read_object_frame(inner, buf),
#[cfg(feature = "draft18")]
FetchReaderState::Draft18(inner) => fo18::read_object_frame(inner, buf),
#[cfg(feature = "draft19")]
FetchReaderState::Draft19(inner) => fo19::read_object_frame(inner, buf),
#[cfg(feature = "draft20")]
FetchReaderState::Draft20(inner) => fo20::read_object_frame(inner, buf),
#[allow(unreachable_patterns)]
_ => unreachable!("AnyFetchObjectReader has no enabled variants"),
}
}
#[allow(unused_variables, unreachable_code)]
pub fn read_object_meta(
&mut self,
buf: &mut impl Buf,
) -> Result<AnyFetchObjectMeta, CodecError> {
match &mut self.state {
#[cfg(feature = "draft07")]
FetchReaderState::Draft07 => fo07::read_object_meta(buf),
#[cfg(feature = "draft08")]
FetchReaderState::Draft08 => fo08::read_object_meta(buf),
#[cfg(feature = "draft09")]
FetchReaderState::Draft09 => fo09::read_object_meta(buf),
#[cfg(feature = "draft10")]
FetchReaderState::Draft10 => fo10::read_object_meta(buf),
#[cfg(feature = "draft11")]
FetchReaderState::Draft11 => fo11::read_object_meta(buf),
#[cfg(feature = "draft12")]
FetchReaderState::Draft12 => fo12::read_object_meta(buf),
#[cfg(feature = "draft13")]
FetchReaderState::Draft13 => fo13::read_object_meta(buf),
#[cfg(feature = "draft14")]
FetchReaderState::Draft14 => fo14::read_object_meta(buf),
#[cfg(feature = "draft15")]
FetchReaderState::Draft15(inner) => fo15::read_object_meta(inner, buf),
#[cfg(feature = "draft16")]
FetchReaderState::Draft16(inner) => fo16::read_object_meta(inner, buf),
#[cfg(feature = "draft17")]
FetchReaderState::Draft17(inner) => fo17::read_object_meta(inner, buf),
#[cfg(feature = "draft18")]
FetchReaderState::Draft18(inner) => fo18::read_object_meta(inner, buf),
#[cfg(feature = "draft19")]
FetchReaderState::Draft19(inner) => fo19::read_object_meta(inner, buf),
#[cfg(feature = "draft20")]
FetchReaderState::Draft20(inner) => fo20::read_object_meta(inner, buf),
#[allow(unreachable_patterns)]
_ => unreachable!("AnyFetchObjectReader has no enabled variants"),
}
}
}
#[derive(Debug, Clone)]
enum FetchFrameShape {
#[cfg(any(
feature = "draft07",
feature = "draft08",
feature = "draft09",
feature = "draft10",
feature = "draft11",
feature = "draft12",
feature = "draft13",
feature = "draft14"
))]
Absolute,
#[cfg(feature = "draft15")]
Draft15(crate::draft15::data_stream::FetchObjectHeader),
#[cfg(feature = "draft16")]
Draft16(
crate::draft16::data_stream::FetchObjectHeader,
crate::draft16::data_stream::FetchObjectLocation,
),
#[cfg(feature = "draft17")]
Draft17(crate::draft17::data_stream::FetchObject),
#[cfg(feature = "draft18")]
Draft18(crate::draft18::data_stream::FetchObject),
#[cfg(feature = "draft19")]
Draft19(crate::draft19::data_stream::FetchObject),
#[cfg(feature = "draft20")]
Draft20(crate::draft20::data_stream::FetchObject),
}
#[derive(Debug, Clone)]
pub struct AnyFetchFrame {
pub meta: AnyFetchObjectMeta,
draft: DraftVersion,
shape: FetchFrameShape,
}
impl AnyFetchFrame {
#[must_use]
pub fn draft(&self) -> DraftVersion {
self.draft
}
#[cfg(any(
feature = "draft07",
feature = "draft08",
feature = "draft09",
feature = "draft10",
feature = "draft11",
feature = "draft12",
feature = "draft13",
feature = "draft14"
))]
fn absolute(draft: DraftVersion, meta: AnyFetchObjectMeta) -> Self {
Self { meta, draft, shape: FetchFrameShape::Absolute }
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum FetchReemit {
Unchanged,
Reframed {
framing_bytes_before: usize,
framing_bytes_after: usize,
},
}
#[derive(Debug, Clone)]
enum FetchWriterState {
#[cfg(any(
feature = "draft07",
feature = "draft08",
feature = "draft09",
feature = "draft10",
feature = "draft11",
feature = "draft12",
feature = "draft13",
feature = "draft14"
))]
Absolute(DraftVersion),
#[cfg(feature = "draft15")]
Draft15(crate::draft15::data_stream::FetchObjectWriter),
#[cfg(feature = "draft16")]
Draft16(crate::draft16::data_stream::FetchObjectWriter),
#[cfg(feature = "draft17")]
Draft17(crate::draft17::data_stream::FetchObjectWriter),
#[cfg(feature = "draft18")]
Draft18(crate::draft18::data_stream::FetchObjectWriter),
#[cfg(feature = "draft19")]
Draft19(crate::draft19::data_stream::FetchObjectWriter),
#[cfg(feature = "draft20")]
Draft20(crate::draft20::data_stream::FetchObjectWriter),
}
#[derive(Debug, Clone)]
pub struct AnyFetchObjectWriter {
state: FetchWriterState,
}
impl AnyFetchObjectWriter {
#[allow(unused_variables, unreachable_code)]
pub fn new(
header: &AnyFetchHeader,
group_order: AnyFetchGroupOrder,
) -> Result<Self, CodecError> {
let state = match header {
#[cfg(feature = "draft07")]
AnyFetchHeader::Draft07(_) => FetchWriterState::Absolute(DraftVersion::Draft07),
#[cfg(feature = "draft08")]
AnyFetchHeader::Draft08(_) => FetchWriterState::Absolute(DraftVersion::Draft08),
#[cfg(feature = "draft09")]
AnyFetchHeader::Draft09(_) => FetchWriterState::Absolute(DraftVersion::Draft09),
#[cfg(feature = "draft10")]
AnyFetchHeader::Draft10(_) => FetchWriterState::Absolute(DraftVersion::Draft10),
#[cfg(feature = "draft11")]
AnyFetchHeader::Draft11(_) => FetchWriterState::Absolute(DraftVersion::Draft11),
#[cfg(feature = "draft12")]
AnyFetchHeader::Draft12(_) => FetchWriterState::Absolute(DraftVersion::Draft12),
#[cfg(feature = "draft13")]
AnyFetchHeader::Draft13(_) => FetchWriterState::Absolute(DraftVersion::Draft13),
#[cfg(feature = "draft14")]
AnyFetchHeader::Draft14(_) => FetchWriterState::Absolute(DraftVersion::Draft14),
#[cfg(feature = "draft15")]
AnyFetchHeader::Draft15(_) => {
FetchWriterState::Draft15(crate::draft15::data_stream::FetchObjectWriter::new())
}
#[cfg(feature = "draft16")]
AnyFetchHeader::Draft16(_) => {
FetchWriterState::Draft16(crate::draft16::data_stream::FetchObjectWriter::new())
}
#[cfg(feature = "draft17")]
AnyFetchHeader::Draft17(_) => {
FetchWriterState::Draft17(crate::draft17::data_stream::FetchObjectWriter::new())
}
#[cfg(feature = "draft18")]
AnyFetchHeader::Draft18(_) => FetchWriterState::Draft18(
crate::draft18::data_stream::FetchObjectWriter::new(match group_order {
AnyFetchGroupOrder::Ascending => {
crate::draft18::data_stream::GroupOrder::Ascending
}
AnyFetchGroupOrder::Descending => {
crate::draft18::data_stream::GroupOrder::Descending
}
}),
),
#[cfg(feature = "draft19")]
AnyFetchHeader::Draft19(_) => FetchWriterState::Draft19(
crate::draft19::data_stream::FetchObjectWriter::new(match group_order {
AnyFetchGroupOrder::Ascending => {
crate::draft19::data_stream::GroupOrder::Ascending
}
AnyFetchGroupOrder::Descending => {
crate::draft19::data_stream::GroupOrder::Descending
}
}),
),
#[cfg(feature = "draft20")]
AnyFetchHeader::Draft20(_) => FetchWriterState::Draft20(
crate::draft20::data_stream::FetchObjectWriter::new(match group_order {
AnyFetchGroupOrder::Ascending => {
crate::draft20::data_stream::GroupOrder::Ascending
}
AnyFetchGroupOrder::Descending => {
crate::draft20::data_stream::GroupOrder::Descending
}
}),
),
#[allow(unreachable_patterns)]
_ => {
return Err(CodecError::UnsupportedDraft(format!(
"draft {:?} not enabled via feature flag",
header.draft()
)));
}
};
Ok(Self { state })
}
#[must_use]
#[allow(unreachable_code)]
pub fn draft(&self) -> DraftVersion {
match &self.state {
#[cfg(any(
feature = "draft07",
feature = "draft08",
feature = "draft09",
feature = "draft10",
feature = "draft11",
feature = "draft12",
feature = "draft13",
feature = "draft14"
))]
FetchWriterState::Absolute(draft) => *draft,
#[cfg(feature = "draft15")]
FetchWriterState::Draft15(_) => DraftVersion::Draft15,
#[cfg(feature = "draft16")]
FetchWriterState::Draft16(_) => DraftVersion::Draft16,
#[cfg(feature = "draft17")]
FetchWriterState::Draft17(_) => DraftVersion::Draft17,
#[cfg(feature = "draft18")]
FetchWriterState::Draft18(_) => DraftVersion::Draft18,
#[cfg(feature = "draft19")]
FetchWriterState::Draft19(_) => DraftVersion::Draft19,
#[cfg(feature = "draft20")]
FetchWriterState::Draft20(_) => DraftVersion::Draft20,
#[allow(unreachable_patterns)]
_ => unreachable!("AnyFetchObjectWriter has no enabled variants"),
}
}
#[allow(unused_variables)]
pub fn reemit_object(
&mut self,
frame: &AnyFetchFrame,
raw: &[u8],
out: &mut impl BufMut,
) -> Result<FetchReemit, CodecError> {
if frame.draft != self.draft() {
return Err(CodecError::UnsupportedDraft(format!(
"a draft {:?} fetch frame cannot be written onto a draft {:?} stream",
frame.draft,
self.draft()
)));
}
let framing_len = frame.meta.wire_len.saturating_sub(frame.meta.payload_length);
let framing_len = usize::try_from(framing_len).map_err(|_| CodecError::InvalidField)?;
if framing_len > raw.len() {
return Err(CodecError::InvalidField);
}
let (framing, rest) = raw.split_at(framing_len);
match (&mut self.state, &frame.shape) {
#[cfg(any(
feature = "draft07",
feature = "draft08",
feature = "draft09",
feature = "draft10",
feature = "draft11",
feature = "draft12",
feature = "draft13",
feature = "draft14"
))]
(FetchWriterState::Absolute(_), FetchFrameShape::Absolute) => {
Ok(FetchReemit::Unchanged)
}
#[cfg(feature = "draft15")]
(FetchWriterState::Draft15(writer), FetchFrameShape::Draft15(original)) => {
let reframed = writer.header_for(original)?;
if reframed == *original {
writer.advance(original);
return Ok(FetchReemit::Unchanged);
}
let mut encoded = Vec::with_capacity(framing.len() + 16);
reframed.encode(&mut encoded)?;
writer.advance(&reframed);
Ok(put_reframed(&encoded, framing.len(), rest, out))
}
#[cfg(feature = "draft16")]
(FetchWriterState::Draft16(writer), FetchFrameShape::Draft16(original, location)) => {
let reframed = writer.header_for(original, location)?;
if reframed == *original {
writer.advance(original, location);
return Ok(FetchReemit::Unchanged);
}
let mut encoded = Vec::with_capacity(framing.len() + 16);
reframed.encode(&mut encoded)?;
writer.advance(&reframed, location);
Ok(put_reframed(&encoded, framing.len(), rest, out))
}
#[cfg(feature = "draft17")]
(FetchWriterState::Draft17(writer), FetchFrameShape::Draft17(original)) => {
let reframed = writer.header_for(original)?;
if reframed == original.header {
writer.advance(original);
return Ok(FetchReemit::Unchanged);
}
let mut encoded = Vec::with_capacity(framing.len() + 16);
reframed.encode(&mut encoded)?;
writer.advance(original);
Ok(put_reframed(&encoded, framing.len(), rest, out))
}
#[cfg(feature = "draft18")]
(FetchWriterState::Draft18(writer), FetchFrameShape::Draft18(original)) => {
let reframed = writer.header_for(original)?;
if reframed == original.header {
writer.advance(original);
return Ok(FetchReemit::Unchanged);
}
let mut encoded = Vec::with_capacity(framing.len() + 16);
reframed.encode(&mut encoded)?;
writer.advance(original);
Ok(put_reframed(&encoded, framing.len(), rest, out))
}
#[cfg(feature = "draft19")]
(FetchWriterState::Draft19(writer), FetchFrameShape::Draft19(original)) => {
let reframed = writer.header_for(original)?;
if reframed == original.header {
writer.advance(original);
return Ok(FetchReemit::Unchanged);
}
let mut encoded = Vec::with_capacity(framing.len() + 16);
reframed.encode(&mut encoded)?;
writer.advance(original);
Ok(put_reframed(&encoded, framing.len(), rest, out))
}
#[cfg(feature = "draft20")]
(FetchWriterState::Draft20(writer), FetchFrameShape::Draft20(original)) => {
let reframed = writer.header_for(original)?;
if reframed == original.header {
writer.advance(original);
return Ok(FetchReemit::Unchanged);
}
let mut encoded = Vec::with_capacity(framing.len() + 16);
reframed.encode(&mut encoded)?;
writer.advance(original);
Ok(put_reframed(&encoded, framing.len(), rest, out))
}
#[allow(unreachable_patterns)]
_ => Err(CodecError::UnsupportedDraft(format!(
"no fetch writer for draft {:?}",
frame.draft
))),
}
}
}
#[cfg(any(
feature = "draft15",
feature = "draft16",
feature = "draft17",
feature = "draft18",
feature = "draft19",
feature = "draft20"
))]
fn put_reframed(
encoded: &[u8],
framing_bytes_before: usize,
rest: &[u8],
out: &mut impl BufMut,
) -> FetchReemit {
out.put_slice(encoded);
out.put_slice(rest);
FetchReemit::Reframed { framing_bytes_before, framing_bytes_after: encoded.len() }
}
macro_rules! subgroup_extensions_fn {
($name:ident, $feat:literal, $draft:ident) => {
#[cfg(feature = $feat)]
fn $name(header: &crate::$draft::data_stream::SubgroupHeader) -> Result<bool, CodecError> {
if !header.stream_type.is_subgroup() {
return Err(CodecError::InvalidField);
}
Ok(header.stream_type.has_extensions())
}
};
}
subgroup_extensions_fn!(subgroup_extensions_11, "draft11", draft11);
subgroup_extensions_fn!(subgroup_extensions_12, "draft12", draft12);
subgroup_extensions_fn!(subgroup_extensions_13, "draft13", draft13);