1use bytes::{Buf, BufMut};
36
37use crate::dispatch::{AnyFetchHeader, AnySubgroupHeader};
38use crate::error::CodecError;
39use crate::varint::VarInt;
40use crate::version::DraftVersion;
41
42#[derive(Debug, Clone, PartialEq, Eq)]
51pub struct AnySubgroupObject {
52 pub object_id: u64,
55 pub extension_headers: Vec<u8>,
68 pub extension_count: Option<u64>,
73 pub status: Option<u64>,
85 pub payload: Vec<u8>,
87}
88
89#[derive(Debug, Clone, Copy, PartialEq, Eq)]
94pub struct AnySubgroupObjectMeta {
95 pub object_id: u64,
97 pub payload_length: u64,
99 pub status: Option<u64>,
101 pub extension_headers_len: u64,
106 pub wire_len: u64,
109}
110
111#[derive(Debug, Clone, Copy, PartialEq, Eq)]
123pub enum AnyFetchEndOfRange {
124 NonExistent,
126 Unknown,
128}
129
130#[derive(Debug, Clone, Copy, PartialEq, Eq)]
141pub enum AnyFetchGroupOrder {
142 Ascending,
144 Descending,
146}
147
148#[derive(Debug, Clone, PartialEq, Eq)]
154pub struct AnyFetchObject {
155 pub group_id: u64,
159 pub subgroup_id: u64,
162 pub has_subgroup_id: bool,
175 pub object_id: u64,
177 pub publisher_priority: u8,
185 pub extension_headers: Vec<u8>,
188 pub extension_count: Option<u64>,
191 pub status: Option<u64>,
199 pub end_of_range: Option<AnyFetchEndOfRange>,
206 pub payload: Vec<u8>,
208}
209
210#[derive(Debug, Clone, Copy, PartialEq, Eq)]
214pub struct AnyFetchObjectMeta {
215 pub group_id: u64,
217 pub subgroup_id: u64,
219 pub has_subgroup_id: bool,
222 pub object_id: u64,
224 pub publisher_priority: u8,
227 pub payload_length: u64,
229 pub status: Option<u64>,
231 pub end_of_range: Option<AnyFetchEndOfRange>,
233 pub extension_headers_len: u64,
235 pub wire_len: u64,
237}
238
239#[cfg(any(feature = "draft16", feature = "draft17", feature = "draft18", feature = "draft19"))]
263const DEFAULT_PUBLISHER_PRIORITY: u8 = 128;
264
265#[allow(dead_code)]
268mod conv {
269 use super::{AnySubgroupObject, Buf, CodecError};
270 use crate::varint::VarInt;
271
272 pub fn skip(buf: &mut impl Buf, len: u64) -> Result<(), CodecError> {
274 let len = usize::try_from(len).map_err(|_| CodecError::UnexpectedEnd)?;
275 if buf.remaining() < len {
276 return Err(CodecError::UnexpectedEnd);
277 }
278 buf.advance(len);
279 Ok(())
280 }
281
282 pub fn take(buf: &mut impl Buf, len: u64) -> Result<Vec<u8>, CodecError> {
284 let len = usize::try_from(len).map_err(|_| CodecError::UnexpectedEnd)?;
285 crate::types::read_bytes(buf, len)
286 }
287
288 pub fn varint(v: u64) -> Result<VarInt, CodecError> {
290 VarInt::from_u64(v).map_err(|_| CodecError::InvalidField)
291 }
292
293 pub fn status_to_write(object: &AnySubgroupObject) -> Result<Option<u64>, CodecError> {
308 match (object.status, object.payload.is_empty()) {
309 (Some(0), false) => Ok(None),
310 (Some(_), false) => Err(CodecError::InvalidField),
311 (Some(code), true) => Ok(Some(code)),
312 (None, true) => Ok(Some(0)),
313 (None, false) => Ok(None),
314 }
315 }
316}
317
318macro_rules! legacy_subgroup_glue {
328 (no_extensions $name:ident, $feat:literal, $draft:ident) => {
329 #[cfg(feature = $feat)]
330 mod $name {
331 use super::conv;
332 use super::{AnySubgroupObject, AnySubgroupObjectMeta};
333 use crate::error::CodecError;
334 use crate::$draft::data_stream::ObjectHeader;
335 use crate::$draft::types::ObjectStatus;
336 use bytes::{Buf, BufMut};
337
338 pub fn read_object(buf: &mut impl Buf) -> Result<AnySubgroupObject, CodecError> {
339 let header = ObjectHeader::decode(buf)?;
340 let payload_length = header.payload_length.into_inner();
341 let (status, payload) = if payload_length == 0 {
342 (Some(header.object_status as u64), Vec::new())
343 } else {
344 (None, conv::take(buf, payload_length)?)
345 };
346 Ok(AnySubgroupObject {
347 object_id: header.object_id.into_inner(),
348 extension_headers: Vec::new(),
349 extension_count: None,
350 status,
351 payload,
352 })
353 }
354
355 pub fn read_object_meta(
356 buf: &mut impl Buf,
357 ) -> Result<AnySubgroupObjectMeta, CodecError> {
358 let start = buf.remaining();
359 let header = ObjectHeader::decode(buf)?;
360 let payload_length = header.payload_length.into_inner();
361 let status = if payload_length == 0 {
362 Some(header.object_status as u64)
363 } else {
364 conv::skip(buf, payload_length)?;
365 None
366 };
367 Ok(AnySubgroupObjectMeta {
368 object_id: header.object_id.into_inner(),
369 payload_length,
370 status,
371 extension_headers_len: 0,
372 wire_len: (start - buf.remaining()) as u64,
373 })
374 }
375
376 pub fn write_object(
377 object: &AnySubgroupObject,
378 buf: &mut impl BufMut,
379 ) -> Result<(), CodecError> {
380 if !object.extension_headers.is_empty() {
381 return Err(CodecError::InvalidField);
382 }
383 let object_status = match conv::status_to_write(object)? {
384 Some(code) => ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?,
385 None => ObjectStatus::Normal,
386 };
387 ObjectHeader {
388 object_id: conv::varint(object.object_id)?,
389 payload_length: conv::varint(object.payload.len() as u64)?,
390 object_status,
391 }
392 .encode(buf);
393 buf.put_slice(&object.payload);
394 Ok(())
395 }
396 }
397 };
398
399 (count_extensions $name:ident, $feat:literal, $draft:ident) => {
400 #[cfg(feature = $feat)]
401 mod $name {
402 use super::conv;
403 use super::{AnySubgroupObject, AnySubgroupObjectMeta};
404 use crate::error::CodecError;
405 use crate::$draft::data_stream::ObjectHeader;
406 use crate::$draft::types::ObjectStatus;
407 use bytes::{Buf, BufMut};
408
409 pub fn read_object(buf: &mut impl Buf) -> Result<AnySubgroupObject, CodecError> {
410 let header = ObjectHeader::decode(buf)?;
411 let payload_length = header.payload_length.into_inner();
412 let (status, payload) = if payload_length == 0 {
413 (Some(header.object_status as u64), Vec::new())
414 } else {
415 (None, conv::take(buf, payload_length)?)
416 };
417 Ok(AnySubgroupObject {
418 object_id: header.object_id.into_inner(),
419 extension_headers: header.extensions,
420 extension_count: Some(header.extension_count.into_inner()),
421 status,
422 payload,
423 })
424 }
425
426 pub fn read_object_meta(
427 buf: &mut impl Buf,
428 ) -> Result<AnySubgroupObjectMeta, CodecError> {
429 let start = buf.remaining();
430 let header = ObjectHeader::decode(buf)?;
431 let payload_length = header.payload_length.into_inner();
432 let status = if payload_length == 0 {
433 Some(header.object_status as u64)
434 } else {
435 conv::skip(buf, payload_length)?;
436 None
437 };
438 Ok(AnySubgroupObjectMeta {
439 object_id: header.object_id.into_inner(),
440 payload_length,
441 status,
442 extension_headers_len: header.extensions.len() as u64,
443 wire_len: (start - buf.remaining()) as u64,
444 })
445 }
446
447 pub fn write_object(
448 object: &AnySubgroupObject,
449 buf: &mut impl BufMut,
450 ) -> Result<(), CodecError> {
451 let extension_count = match object.extension_count {
452 Some(count) => count,
453 None if object.extension_headers.is_empty() => 0,
454 None => return Err(CodecError::InvalidField),
455 };
456 let object_status = match conv::status_to_write(object)? {
457 Some(code) => ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?,
458 None => ObjectStatus::Normal,
459 };
460 ObjectHeader {
461 object_id: conv::varint(object.object_id)?,
462 extension_count: conv::varint(extension_count)?,
463 extensions: object.extension_headers.clone(),
464 payload_length: conv::varint(object.payload.len() as u64)?,
465 object_status,
466 }
467 .encode(buf);
468 buf.put_slice(&object.payload);
469 Ok(())
470 }
471 }
472 };
473
474 (length_extensions $name:ident, $feat:literal, $draft:ident) => {
475 #[cfg(feature = $feat)]
476 mod $name {
477 use super::conv;
478 use super::{AnySubgroupObject, AnySubgroupObjectMeta};
479 use crate::error::CodecError;
480 use crate::$draft::data_stream::ObjectHeader;
481 use crate::$draft::types::ObjectStatus;
482 use bytes::{Buf, BufMut};
483
484 pub fn read_object(buf: &mut impl Buf) -> Result<AnySubgroupObject, CodecError> {
485 let header = ObjectHeader::decode(buf)?;
486 let payload_length = header.payload_length.into_inner();
487 let (status, payload) = if payload_length == 0 {
488 (Some(header.object_status as u64), Vec::new())
489 } else {
490 (None, conv::take(buf, payload_length)?)
491 };
492 Ok(AnySubgroupObject {
493 object_id: header.object_id.into_inner(),
494 extension_headers: header.extensions,
495 extension_count: None,
496 status,
497 payload,
498 })
499 }
500
501 pub fn read_object_meta(
502 buf: &mut impl Buf,
503 ) -> Result<AnySubgroupObjectMeta, CodecError> {
504 let start = buf.remaining();
505 let header = ObjectHeader::decode(buf)?;
506 let payload_length = header.payload_length.into_inner();
507 let status = if payload_length == 0 {
508 Some(header.object_status as u64)
509 } else {
510 conv::skip(buf, payload_length)?;
511 None
512 };
513 Ok(AnySubgroupObjectMeta {
514 object_id: header.object_id.into_inner(),
515 payload_length,
516 status,
517 extension_headers_len: header.extension_headers_length.into_inner(),
518 wire_len: (start - buf.remaining()) as u64,
519 })
520 }
521
522 pub fn write_object(
523 object: &AnySubgroupObject,
524 buf: &mut impl BufMut,
525 ) -> Result<(), CodecError> {
526 let object_status = match conv::status_to_write(object)? {
527 Some(code) => ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?,
528 None => ObjectStatus::Normal,
529 };
530 ObjectHeader {
531 object_id: conv::varint(object.object_id)?,
532 extension_headers_length: conv::varint(object.extension_headers.len() as u64)?,
533 extensions: object.extension_headers.clone(),
534 payload_length: conv::varint(object.payload.len() as u64)?,
535 object_status,
536 }
537 .encode(buf);
538 buf.put_slice(&object.payload);
539 Ok(())
540 }
541 }
542 };
543
544 (gated_extensions $name:ident, $feat:literal, $draft:ident) => {
545 #[cfg(feature = $feat)]
546 mod $name {
547 use super::conv;
548 use super::{AnySubgroupObject, AnySubgroupObjectMeta};
549 use crate::error::CodecError;
550 use crate::$draft::data_stream::ObjectHeader;
551 use crate::$draft::types::ObjectStatus;
552 use bytes::{Buf, BufMut};
553
554 pub fn read_object(
555 extensions: bool,
556 buf: &mut impl Buf,
557 ) -> Result<AnySubgroupObject, CodecError> {
558 let header = ObjectHeader::decode_with_extensions(extensions, buf)?;
559 let payload_length = header.payload_length.into_inner();
560 let (status, payload) = if payload_length == 0 {
561 (Some(header.object_status as u64), Vec::new())
562 } else {
563 (None, conv::take(buf, payload_length)?)
564 };
565 Ok(AnySubgroupObject {
566 object_id: header.object_id.into_inner(),
567 extension_headers: header.extensions,
568 extension_count: None,
569 status,
570 payload,
571 })
572 }
573
574 pub fn read_object_meta(
575 extensions: bool,
576 buf: &mut impl Buf,
577 ) -> Result<AnySubgroupObjectMeta, CodecError> {
578 let start = buf.remaining();
579 let header = ObjectHeader::decode_with_extensions(extensions, buf)?;
580 let payload_length = header.payload_length.into_inner();
581 let status = if payload_length == 0 {
582 Some(header.object_status as u64)
583 } else {
584 conv::skip(buf, payload_length)?;
585 None
586 };
587 Ok(AnySubgroupObjectMeta {
588 object_id: header.object_id.into_inner(),
589 payload_length,
590 status,
591 extension_headers_len: header.extension_headers_length.into_inner(),
592 wire_len: (start - buf.remaining()) as u64,
593 })
594 }
595
596 pub fn write_object(
597 extensions: bool,
598 object: &AnySubgroupObject,
599 buf: &mut impl BufMut,
600 ) -> Result<(), CodecError> {
601 if !extensions && !object.extension_headers.is_empty() {
602 return Err(CodecError::InvalidField);
603 }
604 let object_status = match conv::status_to_write(object)? {
605 Some(code) => ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?,
606 None => ObjectStatus::Normal,
607 };
608 ObjectHeader {
609 object_id: conv::varint(object.object_id)?,
610 extension_headers_length: conv::varint(object.extension_headers.len() as u64)?,
611 extensions: object.extension_headers.clone(),
612 payload_length: conv::varint(object.payload.len() as u64)?,
613 object_status,
614 }
615 .encode_with_extensions(extensions, buf);
616 buf.put_slice(&object.payload);
617 Ok(())
618 }
619 }
620 };
621}
622
623legacy_subgroup_glue!(no_extensions sg07, "draft07", draft07);
624legacy_subgroup_glue!(count_extensions sg08, "draft08", draft08);
625legacy_subgroup_glue!(length_extensions sg09, "draft09", draft09);
626legacy_subgroup_glue!(length_extensions sg10, "draft10", draft10);
627legacy_subgroup_glue!(gated_extensions sg11, "draft11", draft11);
628legacy_subgroup_glue!(gated_extensions sg12, "draft12", draft12);
629legacy_subgroup_glue!(gated_extensions sg13, "draft13", draft13);
630
631macro_rules! modern_subgroup_glue {
645 (derived_length $name:ident, $feat:literal, $draft:ident) => {
646 #[cfg(feature = $feat)]
647 mod $name {
648 use super::conv;
649 use super::{AnySubgroupObject, AnySubgroupObjectMeta};
650 use crate::error::CodecError;
651 use crate::$draft::data_stream::{SubgroupObject, SubgroupObjectReader};
652 use crate::$draft::types::ObjectStatus;
653 use bytes::{Buf, BufMut};
654
655 pub fn read_object(
656 reader: &mut SubgroupObjectReader,
657 buf: &mut impl Buf,
658 ) -> Result<AnySubgroupObject, CodecError> {
659 let object = reader.read_object(buf)?;
660 Ok(AnySubgroupObject {
661 object_id: object.object_id.into_inner(),
662 extension_headers: object.extension_headers,
663 extension_count: None,
664 status: object.status.map(ObjectStatus::as_u64),
665 payload: object.payload,
666 })
667 }
668
669 pub fn read_object_meta(
670 reader: &mut SubgroupObjectReader,
671 buf: &mut impl Buf,
672 ) -> Result<AnySubgroupObjectMeta, CodecError> {
673 let meta = reader.read_object_meta(buf)?;
674 Ok(AnySubgroupObjectMeta {
675 object_id: meta.object_id,
676 payload_length: meta.payload_length,
677 status: meta.status,
678 extension_headers_len: meta.extension_headers_len,
679 wire_len: meta.wire_len,
680 })
681 }
682
683 pub fn write_object(
684 writer: &mut SubgroupObjectReader,
685 object: &AnySubgroupObject,
686 buf: &mut impl BufMut,
687 ) -> Result<(), CodecError> {
688 let status = match conv::status_to_write(object)? {
689 Some(code) => {
690 Some(ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?)
691 }
692 None => None,
693 };
694 writer.write_object(
695 &SubgroupObject {
696 object_id: conv::varint(object.object_id)?,
697 extension_headers: object.extension_headers.clone(),
698 status,
699 payload: object.payload.clone(),
700 },
701 buf,
702 )
703 }
704 }
705 };
706
707 (explicit_length $name:ident, $feat:literal, $draft:ident) => {
708 #[cfg(feature = $feat)]
709 mod $name {
710 use super::conv;
711 use super::{AnySubgroupObject, AnySubgroupObjectMeta};
712 use crate::error::CodecError;
713 use crate::$draft::data_stream::{SubgroupObject, SubgroupObjectReader};
714 use crate::$draft::types::ObjectStatus;
715 use bytes::{Buf, BufMut};
716
717 pub fn read_object(
718 reader: &mut SubgroupObjectReader,
719 buf: &mut impl Buf,
720 ) -> Result<AnySubgroupObject, CodecError> {
721 let object = reader.read_object(buf)?;
722 Ok(AnySubgroupObject {
723 object_id: object.object_id.into_inner(),
724 extension_headers: object.extension_headers,
725 extension_count: None,
726 status: object.object_status.map(ObjectStatus::as_u64),
727 payload: object.payload,
728 })
729 }
730
731 pub fn read_object_meta(
732 reader: &mut SubgroupObjectReader,
733 buf: &mut impl Buf,
734 ) -> Result<AnySubgroupObjectMeta, CodecError> {
735 let meta = reader.read_object_meta(buf)?;
736 Ok(AnySubgroupObjectMeta {
737 object_id: meta.object_id,
738 payload_length: meta.payload_length,
739 status: meta.status,
740 extension_headers_len: meta.extension_headers_len,
741 wire_len: meta.wire_len,
742 })
743 }
744
745 pub fn write_object(
746 writer: &mut SubgroupObjectReader,
747 object: &AnySubgroupObject,
748 buf: &mut impl BufMut,
749 ) -> Result<(), CodecError> {
750 let object_status = match conv::status_to_write(object)? {
751 Some(code) => {
752 Some(ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?)
753 }
754 None => None,
755 };
756 writer.write_object(
757 &SubgroupObject {
758 object_id: conv::varint(object.object_id)?,
759 extension_headers: object.extension_headers.clone(),
760 payload_length: conv::varint(object.payload.len() as u64)?,
761 object_status,
762 payload: object.payload.clone(),
763 },
764 buf,
765 )
766 }
767 }
768 };
769}
770
771modern_subgroup_glue!(derived_length sg14, "draft14", draft14);
772modern_subgroup_glue!(explicit_length sg15, "draft15", draft15);
773modern_subgroup_glue!(explicit_length sg16, "draft16", draft16);
774modern_subgroup_glue!(explicit_length sg17, "draft17", draft17);
775modern_subgroup_glue!(explicit_length sg18, "draft18", draft18);
776modern_subgroup_glue!(explicit_length sg19, "draft19", draft19);
777
778macro_rules! fetch_glue {
785 (no_extensions $name:ident, $feat:literal, $draft:ident) => {
786 #[cfg(feature = $feat)]
787 mod $name {
788 use super::conv;
789 use super::{AnyFetchObject, AnyFetchObjectMeta};
790 use crate::error::CodecError;
791 use crate::$draft::data_stream::FetchObjectHeader;
792 use bytes::Buf;
793
794 pub fn read_object(buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
795 let header = FetchObjectHeader::decode(buf)?;
796 let payload_length = header.payload_length.into_inner();
797 let (status, payload) = if payload_length == 0 {
798 (Some(header.object_status as u64), Vec::new())
799 } else {
800 (None, conv::take(buf, payload_length)?)
801 };
802 Ok(AnyFetchObject {
803 group_id: header.group_id.into_inner(),
804 subgroup_id: header.subgroup_id.into_inner(),
805 has_subgroup_id: true,
806 object_id: header.object_id.into_inner(),
807 publisher_priority: header.publisher_priority,
808 extension_headers: Vec::new(),
809 extension_count: None,
810 status,
811 end_of_range: None,
812 payload,
813 })
814 }
815
816 pub fn read_object_meta(buf: &mut impl Buf) -> Result<AnyFetchObjectMeta, CodecError> {
817 let start = buf.remaining();
818 let header = FetchObjectHeader::decode(buf)?;
819 let payload_length = header.payload_length.into_inner();
820 let status = if payload_length == 0 {
821 Some(header.object_status as u64)
822 } else {
823 conv::skip(buf, payload_length)?;
824 None
825 };
826 Ok(AnyFetchObjectMeta {
827 group_id: header.group_id.into_inner(),
828 subgroup_id: header.subgroup_id.into_inner(),
829 has_subgroup_id: true,
830 object_id: header.object_id.into_inner(),
831 publisher_priority: header.publisher_priority,
832 payload_length,
833 status,
834 end_of_range: None,
835 extension_headers_len: 0,
836 wire_len: (start - buf.remaining()) as u64,
837 })
838 }
839 }
840 };
841
842 (count_extensions $name:ident, $feat:literal, $draft:ident) => {
843 #[cfg(feature = $feat)]
844 mod $name {
845 use super::conv;
846 use super::{AnyFetchObject, AnyFetchObjectMeta};
847 use crate::error::CodecError;
848 use crate::$draft::data_stream::FetchObjectHeader;
849 use bytes::Buf;
850
851 pub fn read_object(buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
852 let header = FetchObjectHeader::decode(buf)?;
853 let payload_length = header.payload_length.into_inner();
854 let (status, payload) = if payload_length == 0 {
855 (Some(header.object_status as u64), Vec::new())
856 } else {
857 (None, conv::take(buf, payload_length)?)
858 };
859 Ok(AnyFetchObject {
860 group_id: header.group_id.into_inner(),
861 subgroup_id: header.subgroup_id.into_inner(),
862 has_subgroup_id: true,
863 object_id: header.object_id.into_inner(),
864 publisher_priority: header.publisher_priority,
865 extension_headers: header.extensions,
866 extension_count: Some(header.extension_count.into_inner()),
867 status,
868 end_of_range: None,
869 payload,
870 })
871 }
872
873 pub fn read_object_meta(buf: &mut impl Buf) -> Result<AnyFetchObjectMeta, CodecError> {
874 let start = buf.remaining();
875 let header = FetchObjectHeader::decode(buf)?;
876 let payload_length = header.payload_length.into_inner();
877 let status = if payload_length == 0 {
878 Some(header.object_status as u64)
879 } else {
880 conv::skip(buf, payload_length)?;
881 None
882 };
883 Ok(AnyFetchObjectMeta {
884 group_id: header.group_id.into_inner(),
885 subgroup_id: header.subgroup_id.into_inner(),
886 has_subgroup_id: true,
887 object_id: header.object_id.into_inner(),
888 publisher_priority: header.publisher_priority,
889 payload_length,
890 status,
891 end_of_range: None,
892 extension_headers_len: header.extensions.len() as u64,
893 wire_len: (start - buf.remaining()) as u64,
894 })
895 }
896 }
897 };
898
899 (length_extensions $name:ident, $feat:literal, $draft:ident) => {
900 #[cfg(feature = $feat)]
901 mod $name {
902 use super::conv;
903 use super::{AnyFetchObject, AnyFetchObjectMeta};
904 use crate::error::CodecError;
905 use crate::$draft::data_stream::FetchObjectHeader;
906 use bytes::Buf;
907
908 pub fn read_object(buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
909 let header = FetchObjectHeader::decode(buf)?;
910 let payload_length = header.payload_length.into_inner();
911 let (status, payload) = if payload_length == 0 {
912 (Some(header.object_status as u64), Vec::new())
913 } else {
914 (None, conv::take(buf, payload_length)?)
915 };
916 Ok(AnyFetchObject {
917 group_id: header.group_id.into_inner(),
918 subgroup_id: header.subgroup_id.into_inner(),
919 has_subgroup_id: true,
920 object_id: header.object_id.into_inner(),
921 publisher_priority: header.publisher_priority,
922 extension_headers: header.extensions,
923 extension_count: None,
924 status,
925 end_of_range: None,
926 payload,
927 })
928 }
929
930 pub fn read_object_meta(buf: &mut impl Buf) -> Result<AnyFetchObjectMeta, CodecError> {
931 let start = buf.remaining();
932 let header = FetchObjectHeader::decode(buf)?;
933 let payload_length = header.payload_length.into_inner();
934 let status = if payload_length == 0 {
935 Some(header.object_status as u64)
936 } else {
937 conv::skip(buf, payload_length)?;
938 None
939 };
940 Ok(AnyFetchObjectMeta {
941 group_id: header.group_id.into_inner(),
942 subgroup_id: header.subgroup_id.into_inner(),
943 has_subgroup_id: true,
944 object_id: header.object_id.into_inner(),
945 publisher_priority: header.publisher_priority,
946 payload_length,
947 status,
948 end_of_range: None,
949 extension_headers_len: header.extension_headers_length.into_inner(),
950 wire_len: (start - buf.remaining()) as u64,
951 })
952 }
953 }
954 };
955}
956
957fetch_glue!(no_extensions fo07, "draft07", draft07);
958fetch_glue!(count_extensions fo08, "draft08", draft08);
959fetch_glue!(length_extensions fo09, "draft09", draft09);
960fetch_glue!(length_extensions fo10, "draft10", draft10);
961fetch_glue!(length_extensions fo11, "draft11", draft11);
962fetch_glue!(length_extensions fo12, "draft12", draft12);
963fetch_glue!(length_extensions fo13, "draft13", draft13);
964
965#[cfg(feature = "draft14")]
966mod fo14 {
967 use super::{AnyFetchObject, AnyFetchObjectMeta};
968 use crate::draft14::data_stream::FetchObject;
969 use crate::draft14::types::ObjectStatus;
970 use crate::error::CodecError;
971 use bytes::Buf;
972
973 pub fn read_object(buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
974 let object = FetchObject::decode(buf)?;
975 Ok(AnyFetchObject {
976 group_id: object.group_id.into_inner(),
977 subgroup_id: object.subgroup_id.into_inner(),
978 has_subgroup_id: true,
979 object_id: object.object_id.into_inner(),
980 publisher_priority: object.publisher_priority,
981 extension_headers: object.extension_headers,
982 extension_count: None,
983 status: object.status.map(ObjectStatus::as_u64),
984 end_of_range: None,
985 payload: object.payload,
986 })
987 }
988
989 pub fn read_object_meta(buf: &mut impl Buf) -> Result<AnyFetchObjectMeta, CodecError> {
990 let meta = FetchObject::decode_meta(buf)?;
991 Ok(AnyFetchObjectMeta {
992 group_id: meta.group_id,
993 subgroup_id: meta.subgroup_id,
994 has_subgroup_id: true,
995 object_id: meta.object_id,
996 publisher_priority: meta.publisher_priority,
997 payload_length: meta.payload_length,
998 status: meta.status,
999 end_of_range: None,
1000 extension_headers_len: meta.extension_headers_len,
1001 wire_len: meta.wire_len,
1002 })
1003 }
1004}
1005
1006#[cfg(feature = "draft15")]
1007mod fo15 {
1008 use super::{conv, AnyFetchObject, AnyFetchObjectMeta};
1009 use crate::draft15::data_stream::FetchObjectReader;
1010 use crate::draft15::types::ObjectStatus;
1011 use crate::error::CodecError;
1012 use bytes::Buf;
1013
1014 pub fn read_object(
1015 reader: &mut FetchObjectReader,
1016 buf: &mut impl Buf,
1017 ) -> Result<AnyFetchObject, CodecError> {
1018 let header = reader.read_object_header(buf)?;
1019 let payload = conv::take(buf, header.payload_length.into_inner())?;
1023 Ok(AnyFetchObject {
1024 group_id: header.group_id.into_inner(),
1025 subgroup_id: header.subgroup_id.into_inner(),
1026 has_subgroup_id: true,
1027 object_id: header.object_id.into_inner(),
1028 publisher_priority: header.publisher_priority,
1029 extension_headers: header.extension_headers,
1030 extension_count: None,
1031 status: header.object_status.map(ObjectStatus::as_u64),
1034 end_of_range: None,
1035 payload,
1036 })
1037 }
1038
1039 pub fn read_object_frame(
1040 reader: &mut FetchObjectReader,
1041 buf: &mut impl Buf,
1042 ) -> Result<super::AnyFetchFrame, CodecError> {
1043 let start = buf.remaining();
1044 let header = reader.read_object_header(buf)?;
1045 let payload_length = header.payload_length.into_inner();
1046 conv::skip(buf, payload_length)?;
1047 let meta = AnyFetchObjectMeta {
1048 group_id: header.group_id.into_inner(),
1049 subgroup_id: header.subgroup_id.into_inner(),
1050 has_subgroup_id: true,
1051 object_id: header.object_id.into_inner(),
1052 publisher_priority: header.publisher_priority,
1053 payload_length,
1054 status: header.object_status.map(ObjectStatus::as_u64),
1055 end_of_range: None,
1056 extension_headers_len: header.extension_headers.len() as u64,
1057 wire_len: (start - buf.remaining()) as u64,
1058 };
1059 Ok(super::AnyFetchFrame {
1060 meta,
1061 draft: crate::version::DraftVersion::Draft15,
1062 shape: super::FetchFrameShape::Draft15(header),
1063 })
1064 }
1065
1066 pub fn read_object_meta(
1067 reader: &mut FetchObjectReader,
1068 buf: &mut impl Buf,
1069 ) -> Result<AnyFetchObjectMeta, CodecError> {
1070 read_object_frame(reader, buf).map(|frame| frame.meta)
1071 }
1072}
1073
1074#[cfg(feature = "draft16")]
1075mod fo16 {
1076 use super::DEFAULT_PUBLISHER_PRIORITY;
1077 use super::{conv, AnyFetchEndOfRange, AnyFetchObject, AnyFetchObjectMeta};
1078 use crate::draft16::data_stream::{
1079 FetchEndOfRange, FetchObjectHeader, FetchObjectLocation, FetchObjectReader,
1080 };
1081 use crate::error::CodecError;
1082 use bytes::Buf;
1083
1084 fn parts(
1088 reader: &mut FetchObjectReader,
1089 buf: &mut impl Buf,
1090 ) -> Result<(FetchObjectHeader, FetchObjectLocation), CodecError> {
1091 let header = FetchObjectHeader::decode(buf)?;
1092 let location = reader.resolve(&header)?;
1093 Ok((header, location))
1094 }
1095
1096 fn resolved(location: &FetchObjectLocation) -> super::Resolved {
1097 let end_of_range = location.end_of_range;
1098 super::Resolved {
1099 group_id: location.group_id,
1100 subgroup_id: location.subgroup_id.filter(|_| end_of_range.is_none()),
1107 object_id: location.object_id,
1108 publisher_priority: location.publisher_priority.unwrap_or(DEFAULT_PUBLISHER_PRIORITY),
1109 end_of_range: end_of_range.map(|r| match r {
1110 FetchEndOfRange::NonExistent => AnyFetchEndOfRange::NonExistent,
1111 FetchEndOfRange::Unknown => AnyFetchEndOfRange::Unknown,
1112 }),
1113 }
1114 }
1115
1116 pub fn read_object(
1117 reader: &mut FetchObjectReader,
1118 buf: &mut impl Buf,
1119 ) -> Result<AnyFetchObject, CodecError> {
1120 let (header, location) = parts(reader, buf)?;
1121 let payload = conv::take(buf, header.payload_length.into_inner())?;
1122 Ok(resolved(&location).into_object(header.extensions.unwrap_or_default(), payload))
1123 }
1124
1125 pub fn read_object_frame(
1126 reader: &mut FetchObjectReader,
1127 buf: &mut impl Buf,
1128 ) -> Result<super::AnyFetchFrame, CodecError> {
1129 let start = buf.remaining();
1130 let (header, location) = parts(reader, buf)?;
1131 let payload_length = header.payload_length.into_inner();
1132 conv::skip(buf, payload_length)?;
1133 let meta = resolved(&location).into_meta(
1134 header.extensions.as_ref().map_or(0, |e| e.len() as u64),
1135 payload_length,
1136 (start - buf.remaining()) as u64,
1137 );
1138 Ok(super::AnyFetchFrame {
1139 meta,
1140 draft: crate::version::DraftVersion::Draft16,
1141 shape: super::FetchFrameShape::Draft16(header, location),
1142 })
1143 }
1144
1145 pub fn read_object_meta(
1146 reader: &mut FetchObjectReader,
1147 buf: &mut impl Buf,
1148 ) -> Result<AnyFetchObjectMeta, CodecError> {
1149 read_object_frame(reader, buf).map(|frame| frame.meta)
1150 }
1151}
1152
1153#[cfg(feature = "draft17")]
1154mod fo17 {
1155 use super::{conv, AnyFetchEndOfRange, AnyFetchObject, AnyFetchObjectMeta};
1156 use crate::draft17::data_stream::{EndOfRange, FetchObject, FetchObjectReader};
1157 use crate::error::CodecError;
1158 use bytes::Buf;
1159
1160 fn resolved(object: &FetchObject) -> super::Resolved {
1161 super::Resolved {
1162 group_id: object.group_id,
1163 subgroup_id: object.subgroup_id,
1164 object_id: object.object_id,
1165 publisher_priority: object
1166 .publisher_priority
1167 .unwrap_or(super::DEFAULT_PUBLISHER_PRIORITY),
1168 end_of_range: object.header.end_of_range().map(|r| match r {
1169 EndOfRange::NonExistent => AnyFetchEndOfRange::NonExistent,
1170 EndOfRange::Unknown => AnyFetchEndOfRange::Unknown,
1171 }),
1172 }
1173 }
1174
1175 pub fn read_object(
1176 reader: &mut FetchObjectReader,
1177 buf: &mut impl Buf,
1178 ) -> Result<AnyFetchObject, CodecError> {
1179 let object = reader.read_object_header(buf)?;
1180 let resolved = resolved(&object);
1181 let payload = conv::take(buf, object.header.payload_length.into_inner())?;
1182 Ok(resolved.into_object(object.header.properties, payload))
1183 }
1184
1185 pub fn read_object_frame(
1186 reader: &mut FetchObjectReader,
1187 buf: &mut impl Buf,
1188 ) -> Result<super::AnyFetchFrame, CodecError> {
1189 let start = buf.remaining();
1190 let object = reader.read_object_header(buf)?;
1191 let resolved = resolved(&object);
1192 let payload_length = object.header.payload_length.into_inner();
1193 conv::skip(buf, payload_length)?;
1194 let meta = resolved.into_meta(
1195 object.header.properties.len() as u64,
1196 payload_length,
1197 (start - buf.remaining()) as u64,
1198 );
1199 Ok(super::AnyFetchFrame {
1200 meta,
1201 draft: crate::version::DraftVersion::Draft17,
1202 shape: super::FetchFrameShape::Draft17(object),
1203 })
1204 }
1205
1206 pub fn read_object_meta(
1207 reader: &mut FetchObjectReader,
1208 buf: &mut impl Buf,
1209 ) -> Result<AnyFetchObjectMeta, CodecError> {
1210 read_object_frame(reader, buf).map(|frame| frame.meta)
1211 }
1212}
1213
1214#[cfg(feature = "draft18")]
1215mod fo18 {
1216 use super::{conv, AnyFetchEndOfRange, AnyFetchObject, AnyFetchObjectMeta};
1217 use crate::draft18::data_stream::{EndOfRange, FetchObject, FetchObjectReader};
1218 use crate::error::CodecError;
1219 use bytes::Buf;
1220
1221 fn resolved(object: &FetchObject) -> super::Resolved {
1222 super::Resolved {
1223 group_id: object.group_id,
1224 subgroup_id: object.subgroup_id,
1225 object_id: object.object_id,
1226 publisher_priority: object
1227 .publisher_priority
1228 .unwrap_or(super::DEFAULT_PUBLISHER_PRIORITY),
1229 end_of_range: object.header.end_of_range().map(|r| match r {
1230 EndOfRange::NonExistent => AnyFetchEndOfRange::NonExistent,
1231 EndOfRange::Unknown => AnyFetchEndOfRange::Unknown,
1232 }),
1233 }
1234 }
1235
1236 pub fn read_object(
1237 reader: &mut FetchObjectReader,
1238 buf: &mut impl Buf,
1239 ) -> Result<AnyFetchObject, CodecError> {
1240 let object = reader.read_object_header(buf)?;
1241 let resolved = resolved(&object);
1242 let payload = conv::take(buf, object.header.payload_length.into_inner())?;
1243 Ok(resolved.into_object(object.header.properties, payload))
1244 }
1245
1246 pub fn read_object_frame(
1247 reader: &mut FetchObjectReader,
1248 buf: &mut impl Buf,
1249 ) -> Result<super::AnyFetchFrame, CodecError> {
1250 let start = buf.remaining();
1251 let object = reader.read_object_header(buf)?;
1252 let resolved = resolved(&object);
1253 let payload_length = object.header.payload_length.into_inner();
1254 conv::skip(buf, payload_length)?;
1255 let meta = resolved.into_meta(
1256 object.header.properties.len() as u64,
1257 payload_length,
1258 (start - buf.remaining()) as u64,
1259 );
1260 Ok(super::AnyFetchFrame {
1261 meta,
1262 draft: crate::version::DraftVersion::Draft18,
1263 shape: super::FetchFrameShape::Draft18(object),
1264 })
1265 }
1266
1267 pub fn read_object_meta(
1268 reader: &mut FetchObjectReader,
1269 buf: &mut impl Buf,
1270 ) -> Result<AnyFetchObjectMeta, CodecError> {
1271 read_object_frame(reader, buf).map(|frame| frame.meta)
1272 }
1273}
1274
1275#[cfg(feature = "draft19")]
1276mod fo19 {
1277 use super::{conv, AnyFetchEndOfRange, AnyFetchObject, AnyFetchObjectMeta};
1278 use crate::draft19::data_stream::{FetchEndOfRange, FetchObject, FetchObjectReader};
1279 use crate::error::CodecError;
1280 use bytes::Buf;
1281
1282 fn resolved(object: &FetchObject) -> super::Resolved {
1283 super::Resolved {
1284 group_id: object.group_id,
1285 subgroup_id: object.subgroup_id,
1286 object_id: object.object_id,
1287 publisher_priority: object
1288 .publisher_priority
1289 .unwrap_or(super::DEFAULT_PUBLISHER_PRIORITY),
1290 end_of_range: object.header.end_of_range().map(|r| match r {
1291 FetchEndOfRange::NonExistent => AnyFetchEndOfRange::NonExistent,
1292 FetchEndOfRange::Unknown => AnyFetchEndOfRange::Unknown,
1293 }),
1294 }
1295 }
1296
1297 pub fn read_object(
1298 reader: &mut FetchObjectReader,
1299 buf: &mut impl Buf,
1300 ) -> Result<AnyFetchObject, CodecError> {
1301 let object = reader.read_object_header(buf)?;
1302 let resolved = resolved(&object);
1303 let payload = conv::take(buf, object.header.payload_length.into_inner())?;
1304 Ok(resolved.into_object(object.header.properties.unwrap_or_default(), payload))
1305 }
1306
1307 pub fn read_object_frame(
1308 reader: &mut FetchObjectReader,
1309 buf: &mut impl Buf,
1310 ) -> Result<super::AnyFetchFrame, CodecError> {
1311 let start = buf.remaining();
1312 let object = reader.read_object_header(buf)?;
1313 let resolved = resolved(&object);
1314 let payload_length = object.header.payload_length.into_inner();
1315 conv::skip(buf, payload_length)?;
1316 let meta = resolved.into_meta(
1317 object.header.properties.as_ref().map_or(0, |p| p.len() as u64),
1318 payload_length,
1319 (start - buf.remaining()) as u64,
1320 );
1321 Ok(super::AnyFetchFrame {
1322 meta,
1323 draft: crate::version::DraftVersion::Draft19,
1324 shape: super::FetchFrameShape::Draft19(object),
1325 })
1326 }
1327
1328 pub fn read_object_meta(
1329 reader: &mut FetchObjectReader,
1330 buf: &mut impl Buf,
1331 ) -> Result<AnyFetchObjectMeta, CodecError> {
1332 read_object_frame(reader, buf).map(|frame| frame.meta)
1333 }
1334}
1335
1336#[cfg(any(feature = "draft16", feature = "draft17", feature = "draft18", feature = "draft19"))]
1345struct Resolved {
1346 group_id: u64,
1347 subgroup_id: Option<u64>,
1348 object_id: u64,
1349 publisher_priority: u8,
1350 end_of_range: Option<AnyFetchEndOfRange>,
1351}
1352
1353#[cfg(any(feature = "draft16", feature = "draft17", feature = "draft18", feature = "draft19"))]
1354impl Resolved {
1355 fn into_object(self, extension_headers: Vec<u8>, payload: Vec<u8>) -> AnyFetchObject {
1356 AnyFetchObject {
1357 group_id: self.group_id,
1358 subgroup_id: self.subgroup_id.unwrap_or(0),
1359 has_subgroup_id: self.subgroup_id.is_some(),
1360 object_id: self.object_id,
1361 publisher_priority: self.publisher_priority,
1362 extension_headers,
1363 extension_count: None,
1364 status: None,
1367 end_of_range: self.end_of_range,
1368 payload,
1369 }
1370 }
1371
1372 fn into_meta(
1373 self,
1374 extension_headers_len: u64,
1375 payload_length: u64,
1376 wire_len: u64,
1377 ) -> AnyFetchObjectMeta {
1378 AnyFetchObjectMeta {
1379 group_id: self.group_id,
1380 subgroup_id: self.subgroup_id.unwrap_or(0),
1381 has_subgroup_id: self.subgroup_id.is_some(),
1382 object_id: self.object_id,
1383 publisher_priority: self.publisher_priority,
1384 payload_length,
1385 status: None,
1386 end_of_range: self.end_of_range,
1387 extension_headers_len,
1388 wire_len,
1389 }
1390 }
1391}
1392
1393#[derive(Debug, Clone)]
1399enum SubgroupReaderState {
1400 #[cfg(feature = "draft07")]
1401 Draft07,
1402 #[cfg(feature = "draft08")]
1403 Draft08,
1404 #[cfg(feature = "draft09")]
1405 Draft09,
1406 #[cfg(feature = "draft10")]
1407 Draft10,
1408 #[cfg(feature = "draft11")]
1409 Draft11 { extensions: bool },
1410 #[cfg(feature = "draft12")]
1411 Draft12 { extensions: bool },
1412 #[cfg(feature = "draft13")]
1413 Draft13 { extensions: bool },
1414 #[cfg(feature = "draft14")]
1415 Draft14(crate::draft14::data_stream::SubgroupObjectReader),
1416 #[cfg(feature = "draft15")]
1417 Draft15(crate::draft15::data_stream::SubgroupObjectReader),
1418 #[cfg(feature = "draft16")]
1419 Draft16(crate::draft16::data_stream::SubgroupObjectReader),
1420 #[cfg(feature = "draft17")]
1421 Draft17(crate::draft17::data_stream::SubgroupObjectReader),
1422 #[cfg(feature = "draft18")]
1423 Draft18(crate::draft18::data_stream::SubgroupObjectReader),
1424 #[cfg(feature = "draft19")]
1425 Draft19(crate::draft19::data_stream::SubgroupObjectReader),
1426}
1427
1428#[derive(Debug, Clone)]
1439pub struct AnySubgroupObjectReader {
1440 state: SubgroupReaderState,
1441}
1442
1443impl AnySubgroupObjectReader {
1444 #[allow(unused_variables, unreachable_code)]
1450 pub fn new(header: &AnySubgroupHeader) -> Result<Self, CodecError> {
1451 let state = match header {
1452 #[cfg(feature = "draft07")]
1453 AnySubgroupHeader::Draft07(_) => SubgroupReaderState::Draft07,
1454 #[cfg(feature = "draft08")]
1455 AnySubgroupHeader::Draft08(_) => SubgroupReaderState::Draft08,
1456 #[cfg(feature = "draft09")]
1457 AnySubgroupHeader::Draft09(_) => SubgroupReaderState::Draft09,
1458 #[cfg(feature = "draft10")]
1459 AnySubgroupHeader::Draft10(_) => SubgroupReaderState::Draft10,
1460 #[cfg(feature = "draft11")]
1461 AnySubgroupHeader::Draft11(h) => {
1462 SubgroupReaderState::Draft11 { extensions: subgroup_extensions_11(h)? }
1463 }
1464 #[cfg(feature = "draft12")]
1465 AnySubgroupHeader::Draft12(h) => {
1466 SubgroupReaderState::Draft12 { extensions: subgroup_extensions_12(h)? }
1467 }
1468 #[cfg(feature = "draft13")]
1469 AnySubgroupHeader::Draft13(h) => {
1470 SubgroupReaderState::Draft13 { extensions: subgroup_extensions_13(h)? }
1471 }
1472 #[cfg(feature = "draft14")]
1473 AnySubgroupHeader::Draft14(h) => SubgroupReaderState::Draft14(
1474 crate::draft14::data_stream::SubgroupObjectReader::new(h),
1475 ),
1476 #[cfg(feature = "draft15")]
1477 AnySubgroupHeader::Draft15(h) => SubgroupReaderState::Draft15(
1478 crate::draft15::data_stream::SubgroupObjectReader::new(h),
1479 ),
1480 #[cfg(feature = "draft16")]
1481 AnySubgroupHeader::Draft16(h) => SubgroupReaderState::Draft16(
1482 crate::draft16::data_stream::SubgroupObjectReader::new(h),
1483 ),
1484 #[cfg(feature = "draft17")]
1485 AnySubgroupHeader::Draft17(h) => SubgroupReaderState::Draft17(
1486 crate::draft17::data_stream::SubgroupObjectReader::new(h),
1487 ),
1488 #[cfg(feature = "draft18")]
1489 AnySubgroupHeader::Draft18(h) => SubgroupReaderState::Draft18(
1490 crate::draft18::data_stream::SubgroupObjectReader::new(h),
1491 ),
1492 #[cfg(feature = "draft19")]
1493 AnySubgroupHeader::Draft19(h) => SubgroupReaderState::Draft19(
1494 crate::draft19::data_stream::SubgroupObjectReader::new(h),
1495 ),
1496 #[allow(unreachable_patterns)]
1497 _ => {
1498 return Err(CodecError::UnsupportedDraft(format!(
1499 "draft {:?} not enabled via feature flag",
1500 header.draft()
1501 )));
1502 }
1503 };
1504 Ok(Self { state })
1505 }
1506
1507 #[allow(unreachable_code)]
1509 pub fn draft(&self) -> DraftVersion {
1510 match &self.state {
1511 #[cfg(feature = "draft07")]
1512 SubgroupReaderState::Draft07 => DraftVersion::Draft07,
1513 #[cfg(feature = "draft08")]
1514 SubgroupReaderState::Draft08 => DraftVersion::Draft08,
1515 #[cfg(feature = "draft09")]
1516 SubgroupReaderState::Draft09 => DraftVersion::Draft09,
1517 #[cfg(feature = "draft10")]
1518 SubgroupReaderState::Draft10 => DraftVersion::Draft10,
1519 #[cfg(feature = "draft11")]
1520 SubgroupReaderState::Draft11 { .. } => DraftVersion::Draft11,
1521 #[cfg(feature = "draft12")]
1522 SubgroupReaderState::Draft12 { .. } => DraftVersion::Draft12,
1523 #[cfg(feature = "draft13")]
1524 SubgroupReaderState::Draft13 { .. } => DraftVersion::Draft13,
1525 #[cfg(feature = "draft14")]
1526 SubgroupReaderState::Draft14(_) => DraftVersion::Draft14,
1527 #[cfg(feature = "draft15")]
1528 SubgroupReaderState::Draft15(_) => DraftVersion::Draft15,
1529 #[cfg(feature = "draft16")]
1530 SubgroupReaderState::Draft16(_) => DraftVersion::Draft16,
1531 #[cfg(feature = "draft17")]
1532 SubgroupReaderState::Draft17(_) => DraftVersion::Draft17,
1533 #[cfg(feature = "draft18")]
1534 SubgroupReaderState::Draft18(_) => DraftVersion::Draft18,
1535 #[cfg(feature = "draft19")]
1536 SubgroupReaderState::Draft19(_) => DraftVersion::Draft19,
1537 #[allow(unreachable_patterns)]
1538 _ => unreachable!("AnySubgroupObjectReader has no enabled variants"),
1539 }
1540 }
1541
1542 #[allow(unused_variables, unreachable_code)]
1548 pub fn read_object(&mut self, buf: &mut impl Buf) -> Result<AnySubgroupObject, CodecError> {
1549 match &mut self.state {
1550 #[cfg(feature = "draft07")]
1551 SubgroupReaderState::Draft07 => sg07::read_object(buf),
1552 #[cfg(feature = "draft08")]
1553 SubgroupReaderState::Draft08 => sg08::read_object(buf),
1554 #[cfg(feature = "draft09")]
1555 SubgroupReaderState::Draft09 => sg09::read_object(buf),
1556 #[cfg(feature = "draft10")]
1557 SubgroupReaderState::Draft10 => sg10::read_object(buf),
1558 #[cfg(feature = "draft11")]
1559 SubgroupReaderState::Draft11 { extensions } => sg11::read_object(*extensions, buf),
1560 #[cfg(feature = "draft12")]
1561 SubgroupReaderState::Draft12 { extensions } => sg12::read_object(*extensions, buf),
1562 #[cfg(feature = "draft13")]
1563 SubgroupReaderState::Draft13 { extensions } => sg13::read_object(*extensions, buf),
1564 #[cfg(feature = "draft14")]
1565 SubgroupReaderState::Draft14(inner) => sg14::read_object(inner, buf),
1566 #[cfg(feature = "draft15")]
1567 SubgroupReaderState::Draft15(inner) => sg15::read_object(inner, buf),
1568 #[cfg(feature = "draft16")]
1569 SubgroupReaderState::Draft16(inner) => sg16::read_object(inner, buf),
1570 #[cfg(feature = "draft17")]
1571 SubgroupReaderState::Draft17(inner) => sg17::read_object(inner, buf),
1572 #[cfg(feature = "draft18")]
1573 SubgroupReaderState::Draft18(inner) => sg18::read_object(inner, buf),
1574 #[cfg(feature = "draft19")]
1575 SubgroupReaderState::Draft19(inner) => sg19::read_object(inner, buf),
1576 #[allow(unreachable_patterns)]
1577 _ => unreachable!("AnySubgroupObjectReader has no enabled variants"),
1578 }
1579 }
1580
1581 #[allow(unused_variables, unreachable_code)]
1588 pub fn read_object_meta(
1589 &mut self,
1590 buf: &mut impl Buf,
1591 ) -> Result<AnySubgroupObjectMeta, CodecError> {
1592 match &mut self.state {
1593 #[cfg(feature = "draft07")]
1594 SubgroupReaderState::Draft07 => sg07::read_object_meta(buf),
1595 #[cfg(feature = "draft08")]
1596 SubgroupReaderState::Draft08 => sg08::read_object_meta(buf),
1597 #[cfg(feature = "draft09")]
1598 SubgroupReaderState::Draft09 => sg09::read_object_meta(buf),
1599 #[cfg(feature = "draft10")]
1600 SubgroupReaderState::Draft10 => sg10::read_object_meta(buf),
1601 #[cfg(feature = "draft11")]
1602 SubgroupReaderState::Draft11 { extensions } => sg11::read_object_meta(*extensions, buf),
1603 #[cfg(feature = "draft12")]
1604 SubgroupReaderState::Draft12 { extensions } => sg12::read_object_meta(*extensions, buf),
1605 #[cfg(feature = "draft13")]
1606 SubgroupReaderState::Draft13 { extensions } => sg13::read_object_meta(*extensions, buf),
1607 #[cfg(feature = "draft14")]
1608 SubgroupReaderState::Draft14(inner) => sg14::read_object_meta(inner, buf),
1609 #[cfg(feature = "draft15")]
1610 SubgroupReaderState::Draft15(inner) => sg15::read_object_meta(inner, buf),
1611 #[cfg(feature = "draft16")]
1612 SubgroupReaderState::Draft16(inner) => sg16::read_object_meta(inner, buf),
1613 #[cfg(feature = "draft17")]
1614 SubgroupReaderState::Draft17(inner) => sg17::read_object_meta(inner, buf),
1615 #[cfg(feature = "draft18")]
1616 SubgroupReaderState::Draft18(inner) => sg18::read_object_meta(inner, buf),
1617 #[cfg(feature = "draft19")]
1618 SubgroupReaderState::Draft19(inner) => sg19::read_object_meta(inner, buf),
1619 #[allow(unreachable_patterns)]
1620 _ => unreachable!("AnySubgroupObjectReader has no enabled variants"),
1621 }
1622 }
1623}
1624
1625#[derive(Debug, Clone)]
1633enum SubgroupWriterState {
1634 #[cfg(feature = "draft07")]
1635 Draft07 { prev_object_id: Option<u64> },
1636 #[cfg(feature = "draft08")]
1637 Draft08 { prev_object_id: Option<u64> },
1638 #[cfg(feature = "draft09")]
1639 Draft09 { prev_object_id: Option<u64> },
1640 #[cfg(feature = "draft10")]
1641 Draft10 { prev_object_id: Option<u64> },
1642 #[cfg(feature = "draft11")]
1643 Draft11 { extensions: bool, prev_object_id: Option<u64> },
1644 #[cfg(feature = "draft12")]
1645 Draft12 { extensions: bool, prev_object_id: Option<u64> },
1646 #[cfg(feature = "draft13")]
1647 Draft13 { extensions: bool, prev_object_id: Option<u64> },
1648 #[cfg(feature = "draft14")]
1649 Draft14 { inner: crate::draft14::data_stream::SubgroupObjectReader, extensions: bool },
1650 #[cfg(feature = "draft15")]
1651 Draft15 { inner: crate::draft15::data_stream::SubgroupObjectReader, extensions: bool },
1652 #[cfg(feature = "draft16")]
1653 Draft16 { inner: crate::draft16::data_stream::SubgroupObjectReader, extensions: bool },
1654 #[cfg(feature = "draft17")]
1655 Draft17 { inner: crate::draft17::data_stream::SubgroupObjectReader, extensions: bool },
1656 #[cfg(feature = "draft18")]
1657 Draft18 { inner: crate::draft18::data_stream::SubgroupObjectReader, extensions: bool },
1658 #[cfg(feature = "draft19")]
1659 Draft19 { inner: crate::draft19::data_stream::SubgroupObjectReader, extensions: bool },
1660}
1661
1662#[derive(Debug, Clone)]
1678pub struct AnySubgroupObjectWriter {
1679 state: SubgroupWriterState,
1680}
1681
1682impl AnySubgroupObjectWriter {
1683 #[allow(unused_variables, unreachable_code)]
1689 pub fn new(header: &AnySubgroupHeader) -> Result<Self, CodecError> {
1690 let state = match header {
1691 #[cfg(feature = "draft07")]
1692 AnySubgroupHeader::Draft07(_) => SubgroupWriterState::Draft07 { prev_object_id: None },
1693 #[cfg(feature = "draft08")]
1694 AnySubgroupHeader::Draft08(_) => SubgroupWriterState::Draft08 { prev_object_id: None },
1695 #[cfg(feature = "draft09")]
1696 AnySubgroupHeader::Draft09(_) => SubgroupWriterState::Draft09 { prev_object_id: None },
1697 #[cfg(feature = "draft10")]
1698 AnySubgroupHeader::Draft10(_) => SubgroupWriterState::Draft10 { prev_object_id: None },
1699 #[cfg(feature = "draft11")]
1700 AnySubgroupHeader::Draft11(h) => SubgroupWriterState::Draft11 {
1701 extensions: subgroup_extensions_11(h)?,
1702 prev_object_id: None,
1703 },
1704 #[cfg(feature = "draft12")]
1705 AnySubgroupHeader::Draft12(h) => SubgroupWriterState::Draft12 {
1706 extensions: subgroup_extensions_12(h)?,
1707 prev_object_id: None,
1708 },
1709 #[cfg(feature = "draft13")]
1710 AnySubgroupHeader::Draft13(h) => SubgroupWriterState::Draft13 {
1711 extensions: subgroup_extensions_13(h)?,
1712 prev_object_id: None,
1713 },
1714 #[cfg(feature = "draft14")]
1715 AnySubgroupHeader::Draft14(h) => SubgroupWriterState::Draft14 {
1716 inner: crate::draft14::data_stream::SubgroupObjectReader::new(h),
1717 extensions: h.stream_type.extensions_present(),
1718 },
1719 #[cfg(feature = "draft15")]
1720 AnySubgroupHeader::Draft15(h) => SubgroupWriterState::Draft15 {
1721 inner: crate::draft15::data_stream::SubgroupObjectReader::new(h),
1722 extensions: h.has_extensions(),
1723 },
1724 #[cfg(feature = "draft16")]
1725 AnySubgroupHeader::Draft16(h) => SubgroupWriterState::Draft16 {
1726 inner: crate::draft16::data_stream::SubgroupObjectReader::new(h),
1727 extensions: h.has_extensions(),
1728 },
1729 #[cfg(feature = "draft17")]
1730 AnySubgroupHeader::Draft17(h) => SubgroupWriterState::Draft17 {
1731 inner: crate::draft17::data_stream::SubgroupObjectReader::new(h),
1732 extensions: h.has_properties(),
1733 },
1734 #[cfg(feature = "draft18")]
1735 AnySubgroupHeader::Draft18(h) => SubgroupWriterState::Draft18 {
1736 inner: crate::draft18::data_stream::SubgroupObjectReader::new(h),
1737 extensions: h.has_properties(),
1738 },
1739 #[cfg(feature = "draft19")]
1740 AnySubgroupHeader::Draft19(h) => SubgroupWriterState::Draft19 {
1741 inner: crate::draft19::data_stream::SubgroupObjectReader::new(h),
1742 extensions: h.has_properties(),
1743 },
1744 #[allow(unreachable_patterns)]
1745 _ => {
1746 return Err(CodecError::UnsupportedDraft(format!(
1747 "draft {:?} not enabled via feature flag",
1748 header.draft()
1749 )));
1750 }
1751 };
1752 Ok(Self { state })
1753 }
1754
1755 #[allow(unreachable_code)]
1757 pub fn draft(&self) -> DraftVersion {
1758 match &self.state {
1759 #[cfg(feature = "draft07")]
1760 SubgroupWriterState::Draft07 { .. } => DraftVersion::Draft07,
1761 #[cfg(feature = "draft08")]
1762 SubgroupWriterState::Draft08 { .. } => DraftVersion::Draft08,
1763 #[cfg(feature = "draft09")]
1764 SubgroupWriterState::Draft09 { .. } => DraftVersion::Draft09,
1765 #[cfg(feature = "draft10")]
1766 SubgroupWriterState::Draft10 { .. } => DraftVersion::Draft10,
1767 #[cfg(feature = "draft11")]
1768 SubgroupWriterState::Draft11 { .. } => DraftVersion::Draft11,
1769 #[cfg(feature = "draft12")]
1770 SubgroupWriterState::Draft12 { .. } => DraftVersion::Draft12,
1771 #[cfg(feature = "draft13")]
1772 SubgroupWriterState::Draft13 { .. } => DraftVersion::Draft13,
1773 #[cfg(feature = "draft14")]
1774 SubgroupWriterState::Draft14 { .. } => DraftVersion::Draft14,
1775 #[cfg(feature = "draft15")]
1776 SubgroupWriterState::Draft15 { .. } => DraftVersion::Draft15,
1777 #[cfg(feature = "draft16")]
1778 SubgroupWriterState::Draft16 { .. } => DraftVersion::Draft16,
1779 #[cfg(feature = "draft17")]
1780 SubgroupWriterState::Draft17 { .. } => DraftVersion::Draft17,
1781 #[cfg(feature = "draft18")]
1782 SubgroupWriterState::Draft18 { .. } => DraftVersion::Draft18,
1783 #[cfg(feature = "draft19")]
1784 SubgroupWriterState::Draft19 { .. } => DraftVersion::Draft19,
1785 #[allow(unreachable_patterns)]
1786 _ => unreachable!("AnySubgroupObjectWriter has no enabled variants"),
1787 }
1788 }
1789
1790 #[allow(unused_variables, unreachable_code)]
1819 pub fn write_object(
1820 &mut self,
1821 object: &AnySubgroupObject,
1822 buf: &mut impl BufMut,
1823 ) -> Result<(), CodecError> {
1824 match &mut self.state {
1825 #[cfg(feature = "draft07")]
1826 SubgroupWriterState::Draft07 { prev_object_id } => {
1827 advance_absolute_id(prev_object_id, object, |o| sg07::write_object(o, buf))
1828 }
1829 #[cfg(feature = "draft08")]
1830 SubgroupWriterState::Draft08 { prev_object_id } => {
1831 advance_absolute_id(prev_object_id, object, |o| sg08::write_object(o, buf))
1832 }
1833 #[cfg(feature = "draft09")]
1834 SubgroupWriterState::Draft09 { prev_object_id } => {
1835 advance_absolute_id(prev_object_id, object, |o| sg09::write_object(o, buf))
1836 }
1837 #[cfg(feature = "draft10")]
1838 SubgroupWriterState::Draft10 { prev_object_id } => {
1839 advance_absolute_id(prev_object_id, object, |o| sg10::write_object(o, buf))
1840 }
1841 #[cfg(feature = "draft11")]
1842 SubgroupWriterState::Draft11 { extensions, prev_object_id } => {
1843 let extensions = *extensions;
1844 advance_absolute_id(prev_object_id, object, |o| {
1845 sg11::write_object(extensions, o, buf)
1846 })
1847 }
1848 #[cfg(feature = "draft12")]
1849 SubgroupWriterState::Draft12 { extensions, prev_object_id } => {
1850 let extensions = *extensions;
1851 advance_absolute_id(prev_object_id, object, |o| {
1852 sg12::write_object(extensions, o, buf)
1853 })
1854 }
1855 #[cfg(feature = "draft13")]
1856 SubgroupWriterState::Draft13 { extensions, prev_object_id } => {
1857 let extensions = *extensions;
1858 advance_absolute_id(prev_object_id, object, |o| {
1859 sg13::write_object(extensions, o, buf)
1860 })
1861 }
1862 #[cfg(feature = "draft14")]
1863 SubgroupWriterState::Draft14 { inner, extensions } => {
1864 reject_unrepresentable_extensions(*extensions, object)?;
1865 sg14::write_object(inner, object, buf)
1866 }
1867 #[cfg(feature = "draft15")]
1868 SubgroupWriterState::Draft15 { inner, extensions } => {
1869 reject_unrepresentable_extensions(*extensions, object)?;
1870 sg15::write_object(inner, object, buf)
1871 }
1872 #[cfg(feature = "draft16")]
1873 SubgroupWriterState::Draft16 { inner, extensions } => {
1874 reject_unrepresentable_extensions(*extensions, object)?;
1875 sg16::write_object(inner, object, buf)
1876 }
1877 #[cfg(feature = "draft17")]
1878 SubgroupWriterState::Draft17 { inner, extensions } => {
1879 reject_unrepresentable_extensions(*extensions, object)?;
1880 sg17::write_object(inner, object, buf)
1881 }
1882 #[cfg(feature = "draft18")]
1883 SubgroupWriterState::Draft18 { inner, extensions } => {
1884 reject_unrepresentable_extensions(*extensions, object)?;
1885 sg18::write_object(inner, object, buf)
1886 }
1887 #[cfg(feature = "draft19")]
1888 SubgroupWriterState::Draft19 { inner, extensions } => {
1889 reject_unrepresentable_extensions(*extensions, object)?;
1890 sg19::write_object(inner, object, buf)
1891 }
1892 #[allow(unreachable_patterns)]
1893 _ => unreachable!("AnySubgroupObjectWriter has no enabled variants"),
1894 }
1895 }
1896}
1897
1898#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1902pub enum Reemit {
1903 Verbatim,
1905 Reencoded {
1907 id_bytes_before: usize,
1909 id_bytes_after: usize,
1911 },
1912}
1913
1914pub fn reemit_subgroup_object(
1976 draft: DraftVersion,
1977 prev_forwarded: Option<u64>,
1978 object_id: u64,
1979 raw: &[u8],
1980 out: &mut impl BufMut,
1981) -> Result<Reemit, CodecError> {
1982 if matches!(prev_forwarded, Some(prev) if object_id <= prev) {
1983 return Err(CodecError::InvalidField);
1984 }
1985
1986 let mut cursor: &[u8] = raw;
1990 draft.decode_varint(&mut cursor).map_err(|_| CodecError::InvalidField)?;
1991 let id_bytes_before = raw.len() - cursor.len();
1992
1993 if !delta_encodes_object_ids(draft) {
1996 out.put_slice(raw);
1997 return Ok(Reemit::Verbatim);
1998 }
1999
2000 let delta = match prev_forwarded {
2001 None => object_id,
2002 Some(prev) => object_id
2003 .checked_sub(prev)
2004 .and_then(|v| v.checked_sub(1))
2005 .ok_or(CodecError::InvalidField)?,
2006 };
2007
2008 let field = if draft.uses_moqt_varint() {
2011 VarInt::from_u64_moqt(delta)
2012 } else {
2013 VarInt::from_u64(delta).map_err(|_| CodecError::InvalidField)?
2014 };
2015
2016 let mut encoded = [0u8; 9];
2018 let mut slot: &mut [u8] = &mut encoded;
2019 draft.encode_varint(field, &mut slot);
2020 let id_bytes_after = 9 - slot.len();
2021 let encoded = &encoded[..id_bytes_after];
2022
2023 if encoded == &raw[..id_bytes_before] {
2024 out.put_slice(raw);
2025 return Ok(Reemit::Verbatim);
2026 }
2027
2028 out.put_slice(encoded);
2029 out.put_slice(&raw[id_bytes_before..]);
2030 Ok(Reemit::Reencoded { id_bytes_before, id_bytes_after })
2031}
2032
2033fn delta_encodes_object_ids(draft: DraftVersion) -> bool {
2039 matches!(
2040 draft,
2041 DraftVersion::Draft14
2042 | DraftVersion::Draft15
2043 | DraftVersion::Draft16
2044 | DraftVersion::Draft17
2045 | DraftVersion::Draft18
2046 | DraftVersion::Draft19
2047 )
2048}
2049
2050#[cfg(any(
2060 feature = "draft07",
2061 feature = "draft08",
2062 feature = "draft09",
2063 feature = "draft10",
2064 feature = "draft11",
2065 feature = "draft12",
2066 feature = "draft13"
2067))]
2068fn advance_absolute_id(
2069 prev_object_id: &mut Option<u64>,
2070 object: &AnySubgroupObject,
2071 write: impl FnOnce(&AnySubgroupObject) -> Result<(), CodecError>,
2072) -> Result<(), CodecError> {
2073 if matches!(*prev_object_id, Some(prev) if object.object_id <= prev) {
2074 return Err(CodecError::InvalidField);
2075 }
2076 write(object)?;
2077 *prev_object_id = Some(object.object_id);
2078 Ok(())
2079}
2080
2081#[cfg(any(
2084 feature = "draft14",
2085 feature = "draft15",
2086 feature = "draft16",
2087 feature = "draft17",
2088 feature = "draft18",
2089 feature = "draft19"
2090))]
2091fn reject_unrepresentable_extensions(
2092 extensions: bool,
2093 object: &AnySubgroupObject,
2094) -> Result<(), CodecError> {
2095 if !extensions && !object.extension_headers.is_empty() {
2096 return Err(CodecError::InvalidField);
2097 }
2098 Ok(())
2099}
2100
2101#[derive(Debug, Clone)]
2107enum FetchReaderState {
2108 #[cfg(feature = "draft07")]
2109 Draft07,
2110 #[cfg(feature = "draft08")]
2111 Draft08,
2112 #[cfg(feature = "draft09")]
2113 Draft09,
2114 #[cfg(feature = "draft10")]
2115 Draft10,
2116 #[cfg(feature = "draft11")]
2117 Draft11,
2118 #[cfg(feature = "draft12")]
2119 Draft12,
2120 #[cfg(feature = "draft13")]
2121 Draft13,
2122 #[cfg(feature = "draft14")]
2123 Draft14,
2124 #[cfg(feature = "draft15")]
2125 Draft15(crate::draft15::data_stream::FetchObjectReader),
2126 #[cfg(feature = "draft16")]
2127 Draft16(crate::draft16::data_stream::FetchObjectReader),
2128 #[cfg(feature = "draft17")]
2129 Draft17(crate::draft17::data_stream::FetchObjectReader),
2130 #[cfg(feature = "draft18")]
2131 Draft18(crate::draft18::data_stream::FetchObjectReader),
2132 #[cfg(feature = "draft19")]
2133 Draft19(crate::draft19::data_stream::FetchObjectReader),
2134}
2135
2136#[derive(Debug, Clone)]
2161pub struct AnyFetchObjectReader {
2162 state: FetchReaderState,
2163}
2164
2165impl AnyFetchObjectReader {
2166 #[allow(unused_variables, unreachable_code)]
2184 pub fn new(
2185 header: &AnyFetchHeader,
2186 group_order: AnyFetchGroupOrder,
2187 ) -> Result<Self, CodecError> {
2188 let state = match header {
2189 #[cfg(feature = "draft07")]
2190 AnyFetchHeader::Draft07(_) => FetchReaderState::Draft07,
2191 #[cfg(feature = "draft08")]
2192 AnyFetchHeader::Draft08(_) => FetchReaderState::Draft08,
2193 #[cfg(feature = "draft09")]
2194 AnyFetchHeader::Draft09(_) => FetchReaderState::Draft09,
2195 #[cfg(feature = "draft10")]
2196 AnyFetchHeader::Draft10(_) => FetchReaderState::Draft10,
2197 #[cfg(feature = "draft11")]
2198 AnyFetchHeader::Draft11(_) => FetchReaderState::Draft11,
2199 #[cfg(feature = "draft12")]
2200 AnyFetchHeader::Draft12(_) => FetchReaderState::Draft12,
2201 #[cfg(feature = "draft13")]
2202 AnyFetchHeader::Draft13(_) => FetchReaderState::Draft13,
2203 #[cfg(feature = "draft14")]
2204 AnyFetchHeader::Draft14(_) => FetchReaderState::Draft14,
2205 #[cfg(feature = "draft15")]
2208 AnyFetchHeader::Draft15(_) => {
2209 FetchReaderState::Draft15(crate::draft15::data_stream::FetchObjectReader::new())
2210 }
2211 #[cfg(feature = "draft16")]
2212 AnyFetchHeader::Draft16(_) => {
2213 FetchReaderState::Draft16(crate::draft16::data_stream::FetchObjectReader::new())
2214 }
2215 #[cfg(feature = "draft17")]
2216 AnyFetchHeader::Draft17(_) => {
2217 FetchReaderState::Draft17(crate::draft17::data_stream::FetchObjectReader::new())
2218 }
2219 #[cfg(feature = "draft18")]
2220 AnyFetchHeader::Draft18(_) => FetchReaderState::Draft18(
2221 crate::draft18::data_stream::FetchObjectReader::new(match group_order {
2222 AnyFetchGroupOrder::Ascending => {
2223 crate::draft18::data_stream::GroupOrder::Ascending
2224 }
2225 AnyFetchGroupOrder::Descending => {
2226 crate::draft18::data_stream::GroupOrder::Descending
2227 }
2228 }),
2229 ),
2230 #[cfg(feature = "draft19")]
2231 AnyFetchHeader::Draft19(_) => FetchReaderState::Draft19(
2232 crate::draft19::data_stream::FetchObjectReader::new(match group_order {
2233 AnyFetchGroupOrder::Ascending => {
2234 crate::draft19::data_stream::GroupOrder::Ascending
2235 }
2236 AnyFetchGroupOrder::Descending => {
2237 crate::draft19::data_stream::GroupOrder::Descending
2238 }
2239 }),
2240 ),
2241 #[allow(unreachable_patterns)]
2242 _ => {
2243 return Err(CodecError::UnsupportedDraft(format!(
2244 "draft {:?} not enabled via feature flag",
2245 header.draft()
2246 )));
2247 }
2248 };
2249 Ok(Self { state })
2250 }
2251
2252 #[allow(unreachable_code)]
2254 pub fn draft(&self) -> DraftVersion {
2255 match &self.state {
2256 #[cfg(feature = "draft07")]
2257 FetchReaderState::Draft07 => DraftVersion::Draft07,
2258 #[cfg(feature = "draft08")]
2259 FetchReaderState::Draft08 => DraftVersion::Draft08,
2260 #[cfg(feature = "draft09")]
2261 FetchReaderState::Draft09 => DraftVersion::Draft09,
2262 #[cfg(feature = "draft10")]
2263 FetchReaderState::Draft10 => DraftVersion::Draft10,
2264 #[cfg(feature = "draft11")]
2265 FetchReaderState::Draft11 => DraftVersion::Draft11,
2266 #[cfg(feature = "draft12")]
2267 FetchReaderState::Draft12 => DraftVersion::Draft12,
2268 #[cfg(feature = "draft13")]
2269 FetchReaderState::Draft13 => DraftVersion::Draft13,
2270 #[cfg(feature = "draft14")]
2271 FetchReaderState::Draft14 => DraftVersion::Draft14,
2272 #[cfg(feature = "draft15")]
2273 FetchReaderState::Draft15(_) => DraftVersion::Draft15,
2274 #[cfg(feature = "draft16")]
2275 FetchReaderState::Draft16(_) => DraftVersion::Draft16,
2276 #[cfg(feature = "draft17")]
2277 FetchReaderState::Draft17(_) => DraftVersion::Draft17,
2278 #[cfg(feature = "draft18")]
2279 FetchReaderState::Draft18(_) => DraftVersion::Draft18,
2280 #[cfg(feature = "draft19")]
2281 FetchReaderState::Draft19(_) => DraftVersion::Draft19,
2282 #[allow(unreachable_patterns)]
2283 _ => unreachable!("AnyFetchObjectReader has no enabled variants"),
2284 }
2285 }
2286
2287 #[allow(unused_variables, unreachable_code)]
2299 pub fn read_object(&mut self, buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
2300 match &mut self.state {
2301 #[cfg(feature = "draft07")]
2302 FetchReaderState::Draft07 => fo07::read_object(buf),
2303 #[cfg(feature = "draft08")]
2304 FetchReaderState::Draft08 => fo08::read_object(buf),
2305 #[cfg(feature = "draft09")]
2306 FetchReaderState::Draft09 => fo09::read_object(buf),
2307 #[cfg(feature = "draft10")]
2308 FetchReaderState::Draft10 => fo10::read_object(buf),
2309 #[cfg(feature = "draft11")]
2310 FetchReaderState::Draft11 => fo11::read_object(buf),
2311 #[cfg(feature = "draft12")]
2312 FetchReaderState::Draft12 => fo12::read_object(buf),
2313 #[cfg(feature = "draft13")]
2314 FetchReaderState::Draft13 => fo13::read_object(buf),
2315 #[cfg(feature = "draft14")]
2316 FetchReaderState::Draft14 => fo14::read_object(buf),
2317 #[cfg(feature = "draft15")]
2318 FetchReaderState::Draft15(inner) => fo15::read_object(inner, buf),
2319 #[cfg(feature = "draft16")]
2320 FetchReaderState::Draft16(inner) => fo16::read_object(inner, buf),
2321 #[cfg(feature = "draft17")]
2322 FetchReaderState::Draft17(inner) => fo17::read_object(inner, buf),
2323 #[cfg(feature = "draft18")]
2324 FetchReaderState::Draft18(inner) => fo18::read_object(inner, buf),
2325 #[cfg(feature = "draft19")]
2326 FetchReaderState::Draft19(inner) => fo19::read_object(inner, buf),
2327 #[allow(unreachable_patterns)]
2328 _ => unreachable!("AnyFetchObjectReader has no enabled variants"),
2329 }
2330 }
2331
2332 #[allow(unused_variables, unreachable_code)]
2345 pub fn read_object_frame(&mut self, buf: &mut impl Buf) -> Result<AnyFetchFrame, CodecError> {
2346 match &mut self.state {
2347 #[cfg(feature = "draft07")]
2348 FetchReaderState::Draft07 => fo07::read_object_meta(buf)
2349 .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft07, meta)),
2350 #[cfg(feature = "draft08")]
2351 FetchReaderState::Draft08 => fo08::read_object_meta(buf)
2352 .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft08, meta)),
2353 #[cfg(feature = "draft09")]
2354 FetchReaderState::Draft09 => fo09::read_object_meta(buf)
2355 .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft09, meta)),
2356 #[cfg(feature = "draft10")]
2357 FetchReaderState::Draft10 => fo10::read_object_meta(buf)
2358 .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft10, meta)),
2359 #[cfg(feature = "draft11")]
2360 FetchReaderState::Draft11 => fo11::read_object_meta(buf)
2361 .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft11, meta)),
2362 #[cfg(feature = "draft12")]
2363 FetchReaderState::Draft12 => fo12::read_object_meta(buf)
2364 .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft12, meta)),
2365 #[cfg(feature = "draft13")]
2366 FetchReaderState::Draft13 => fo13::read_object_meta(buf)
2367 .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft13, meta)),
2368 #[cfg(feature = "draft14")]
2369 FetchReaderState::Draft14 => fo14::read_object_meta(buf)
2370 .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft14, meta)),
2371 #[cfg(feature = "draft15")]
2372 FetchReaderState::Draft15(inner) => fo15::read_object_frame(inner, buf),
2373 #[cfg(feature = "draft16")]
2374 FetchReaderState::Draft16(inner) => fo16::read_object_frame(inner, buf),
2375 #[cfg(feature = "draft17")]
2376 FetchReaderState::Draft17(inner) => fo17::read_object_frame(inner, buf),
2377 #[cfg(feature = "draft18")]
2378 FetchReaderState::Draft18(inner) => fo18::read_object_frame(inner, buf),
2379 #[cfg(feature = "draft19")]
2380 FetchReaderState::Draft19(inner) => fo19::read_object_frame(inner, buf),
2381 #[allow(unreachable_patterns)]
2382 _ => unreachable!("AnyFetchObjectReader has no enabled variants"),
2383 }
2384 }
2385
2386 #[allow(unused_variables, unreachable_code)]
2392 pub fn read_object_meta(
2393 &mut self,
2394 buf: &mut impl Buf,
2395 ) -> Result<AnyFetchObjectMeta, CodecError> {
2396 match &mut self.state {
2397 #[cfg(feature = "draft07")]
2398 FetchReaderState::Draft07 => fo07::read_object_meta(buf),
2399 #[cfg(feature = "draft08")]
2400 FetchReaderState::Draft08 => fo08::read_object_meta(buf),
2401 #[cfg(feature = "draft09")]
2402 FetchReaderState::Draft09 => fo09::read_object_meta(buf),
2403 #[cfg(feature = "draft10")]
2404 FetchReaderState::Draft10 => fo10::read_object_meta(buf),
2405 #[cfg(feature = "draft11")]
2406 FetchReaderState::Draft11 => fo11::read_object_meta(buf),
2407 #[cfg(feature = "draft12")]
2408 FetchReaderState::Draft12 => fo12::read_object_meta(buf),
2409 #[cfg(feature = "draft13")]
2410 FetchReaderState::Draft13 => fo13::read_object_meta(buf),
2411 #[cfg(feature = "draft14")]
2412 FetchReaderState::Draft14 => fo14::read_object_meta(buf),
2413 #[cfg(feature = "draft15")]
2414 FetchReaderState::Draft15(inner) => fo15::read_object_meta(inner, buf),
2415 #[cfg(feature = "draft16")]
2416 FetchReaderState::Draft16(inner) => fo16::read_object_meta(inner, buf),
2417 #[cfg(feature = "draft17")]
2418 FetchReaderState::Draft17(inner) => fo17::read_object_meta(inner, buf),
2419 #[cfg(feature = "draft18")]
2420 FetchReaderState::Draft18(inner) => fo18::read_object_meta(inner, buf),
2421 #[cfg(feature = "draft19")]
2422 FetchReaderState::Draft19(inner) => fo19::read_object_meta(inner, buf),
2423 #[allow(unreachable_patterns)]
2424 _ => unreachable!("AnyFetchObjectReader has no enabled variants"),
2425 }
2426 }
2427}
2428
2429#[derive(Debug, Clone)]
2439enum FetchFrameShape {
2440 #[cfg(any(
2442 feature = "draft07",
2443 feature = "draft08",
2444 feature = "draft09",
2445 feature = "draft10",
2446 feature = "draft11",
2447 feature = "draft12",
2448 feature = "draft13",
2449 feature = "draft14"
2450 ))]
2451 Absolute,
2452 #[cfg(feature = "draft15")]
2453 Draft15(crate::draft15::data_stream::FetchObjectHeader),
2454 #[cfg(feature = "draft16")]
2455 Draft16(
2456 crate::draft16::data_stream::FetchObjectHeader,
2457 crate::draft16::data_stream::FetchObjectLocation,
2458 ),
2459 #[cfg(feature = "draft17")]
2460 Draft17(crate::draft17::data_stream::FetchObject),
2461 #[cfg(feature = "draft18")]
2462 Draft18(crate::draft18::data_stream::FetchObject),
2463 #[cfg(feature = "draft19")]
2464 Draft19(crate::draft19::data_stream::FetchObject),
2465}
2466
2467#[derive(Debug, Clone)]
2476pub struct AnyFetchFrame {
2477 pub meta: AnyFetchObjectMeta,
2480 draft: DraftVersion,
2481 shape: FetchFrameShape,
2482}
2483
2484impl AnyFetchFrame {
2485 #[must_use]
2491 pub fn draft(&self) -> DraftVersion {
2492 self.draft
2493 }
2494
2495 #[cfg(any(
2497 feature = "draft07",
2498 feature = "draft08",
2499 feature = "draft09",
2500 feature = "draft10",
2501 feature = "draft11",
2502 feature = "draft12",
2503 feature = "draft13",
2504 feature = "draft14"
2505 ))]
2506 fn absolute(draft: DraftVersion, meta: AnyFetchObjectMeta) -> Self {
2507 Self { meta, draft, shape: FetchFrameShape::Absolute }
2508 }
2509}
2510
2511#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2521pub enum FetchReemit {
2522 Unchanged,
2529 Reframed {
2533 framing_bytes_before: usize,
2535 framing_bytes_after: usize,
2537 },
2538}
2539
2540#[derive(Debug, Clone)]
2542enum FetchWriterState {
2543 #[cfg(any(
2547 feature = "draft07",
2548 feature = "draft08",
2549 feature = "draft09",
2550 feature = "draft10",
2551 feature = "draft11",
2552 feature = "draft12",
2553 feature = "draft13",
2554 feature = "draft14"
2555 ))]
2556 Absolute(DraftVersion),
2557 #[cfg(feature = "draft15")]
2558 Draft15(crate::draft15::data_stream::FetchObjectWriter),
2559 #[cfg(feature = "draft16")]
2560 Draft16(crate::draft16::data_stream::FetchObjectWriter),
2561 #[cfg(feature = "draft17")]
2562 Draft17(crate::draft17::data_stream::FetchObjectWriter),
2563 #[cfg(feature = "draft18")]
2564 Draft18(crate::draft18::data_stream::FetchObjectWriter),
2565 #[cfg(feature = "draft19")]
2566 Draft19(crate::draft19::data_stream::FetchObjectWriter),
2567}
2568
2569#[derive(Debug, Clone)]
2603pub struct AnyFetchObjectWriter {
2604 state: FetchWriterState,
2605}
2606
2607impl AnyFetchObjectWriter {
2608 #[allow(unused_variables, unreachable_code)]
2616 pub fn new(
2617 header: &AnyFetchHeader,
2618 group_order: AnyFetchGroupOrder,
2619 ) -> Result<Self, CodecError> {
2620 let state = match header {
2621 #[cfg(feature = "draft07")]
2622 AnyFetchHeader::Draft07(_) => FetchWriterState::Absolute(DraftVersion::Draft07),
2623 #[cfg(feature = "draft08")]
2624 AnyFetchHeader::Draft08(_) => FetchWriterState::Absolute(DraftVersion::Draft08),
2625 #[cfg(feature = "draft09")]
2626 AnyFetchHeader::Draft09(_) => FetchWriterState::Absolute(DraftVersion::Draft09),
2627 #[cfg(feature = "draft10")]
2628 AnyFetchHeader::Draft10(_) => FetchWriterState::Absolute(DraftVersion::Draft10),
2629 #[cfg(feature = "draft11")]
2630 AnyFetchHeader::Draft11(_) => FetchWriterState::Absolute(DraftVersion::Draft11),
2631 #[cfg(feature = "draft12")]
2632 AnyFetchHeader::Draft12(_) => FetchWriterState::Absolute(DraftVersion::Draft12),
2633 #[cfg(feature = "draft13")]
2634 AnyFetchHeader::Draft13(_) => FetchWriterState::Absolute(DraftVersion::Draft13),
2635 #[cfg(feature = "draft14")]
2636 AnyFetchHeader::Draft14(_) => FetchWriterState::Absolute(DraftVersion::Draft14),
2637 #[cfg(feature = "draft15")]
2638 AnyFetchHeader::Draft15(_) => {
2639 FetchWriterState::Draft15(crate::draft15::data_stream::FetchObjectWriter::new())
2640 }
2641 #[cfg(feature = "draft16")]
2642 AnyFetchHeader::Draft16(_) => {
2643 FetchWriterState::Draft16(crate::draft16::data_stream::FetchObjectWriter::new())
2644 }
2645 #[cfg(feature = "draft17")]
2646 AnyFetchHeader::Draft17(_) => {
2647 FetchWriterState::Draft17(crate::draft17::data_stream::FetchObjectWriter::new())
2648 }
2649 #[cfg(feature = "draft18")]
2650 AnyFetchHeader::Draft18(_) => FetchWriterState::Draft18(
2651 crate::draft18::data_stream::FetchObjectWriter::new(match group_order {
2652 AnyFetchGroupOrder::Ascending => {
2653 crate::draft18::data_stream::GroupOrder::Ascending
2654 }
2655 AnyFetchGroupOrder::Descending => {
2656 crate::draft18::data_stream::GroupOrder::Descending
2657 }
2658 }),
2659 ),
2660 #[cfg(feature = "draft19")]
2661 AnyFetchHeader::Draft19(_) => FetchWriterState::Draft19(
2662 crate::draft19::data_stream::FetchObjectWriter::new(match group_order {
2663 AnyFetchGroupOrder::Ascending => {
2664 crate::draft19::data_stream::GroupOrder::Ascending
2665 }
2666 AnyFetchGroupOrder::Descending => {
2667 crate::draft19::data_stream::GroupOrder::Descending
2668 }
2669 }),
2670 ),
2671 #[allow(unreachable_patterns)]
2672 _ => {
2673 return Err(CodecError::UnsupportedDraft(format!(
2674 "draft {:?} not enabled via feature flag",
2675 header.draft()
2676 )));
2677 }
2678 };
2679 Ok(Self { state })
2680 }
2681
2682 #[must_use]
2684 #[allow(unreachable_code)]
2685 pub fn draft(&self) -> DraftVersion {
2686 match &self.state {
2687 #[cfg(any(
2688 feature = "draft07",
2689 feature = "draft08",
2690 feature = "draft09",
2691 feature = "draft10",
2692 feature = "draft11",
2693 feature = "draft12",
2694 feature = "draft13",
2695 feature = "draft14"
2696 ))]
2697 FetchWriterState::Absolute(draft) => *draft,
2698 #[cfg(feature = "draft15")]
2699 FetchWriterState::Draft15(_) => DraftVersion::Draft15,
2700 #[cfg(feature = "draft16")]
2701 FetchWriterState::Draft16(_) => DraftVersion::Draft16,
2702 #[cfg(feature = "draft17")]
2703 FetchWriterState::Draft17(_) => DraftVersion::Draft17,
2704 #[cfg(feature = "draft18")]
2705 FetchWriterState::Draft18(_) => DraftVersion::Draft18,
2706 #[cfg(feature = "draft19")]
2707 FetchWriterState::Draft19(_) => DraftVersion::Draft19,
2708 #[allow(unreachable_patterns)]
2709 _ => unreachable!("AnyFetchObjectWriter has no enabled variants"),
2710 }
2711 }
2712
2713 #[allow(unused_variables)]
2745 pub fn reemit_object(
2746 &mut self,
2747 frame: &AnyFetchFrame,
2748 raw: &[u8],
2749 out: &mut impl BufMut,
2750 ) -> Result<FetchReemit, CodecError> {
2751 if frame.draft != self.draft() {
2752 return Err(CodecError::UnsupportedDraft(format!(
2753 "a draft {:?} fetch frame cannot be written onto a draft {:?} stream",
2754 frame.draft,
2755 self.draft()
2756 )));
2757 }
2758
2759 let framing_len = frame.meta.wire_len.saturating_sub(frame.meta.payload_length);
2760 let framing_len = usize::try_from(framing_len).map_err(|_| CodecError::InvalidField)?;
2761 if framing_len > raw.len() {
2762 return Err(CodecError::InvalidField);
2763 }
2764 let (framing, rest) = raw.split_at(framing_len);
2765
2766 match (&mut self.state, &frame.shape) {
2767 #[cfg(any(
2770 feature = "draft07",
2771 feature = "draft08",
2772 feature = "draft09",
2773 feature = "draft10",
2774 feature = "draft11",
2775 feature = "draft12",
2776 feature = "draft13",
2777 feature = "draft14"
2778 ))]
2779 (FetchWriterState::Absolute(_), FetchFrameShape::Absolute) => {
2780 Ok(FetchReemit::Unchanged)
2781 }
2782 #[cfg(feature = "draft15")]
2783 (FetchWriterState::Draft15(writer), FetchFrameShape::Draft15(original)) => {
2784 let reframed = writer.header_for(original)?;
2785 if reframed == *original {
2786 writer.advance(original);
2787 return Ok(FetchReemit::Unchanged);
2788 }
2789 let mut encoded = Vec::with_capacity(framing.len() + 16);
2790 reframed.encode(&mut encoded)?;
2791 writer.advance(&reframed);
2792 Ok(put_reframed(&encoded, framing.len(), rest, out))
2793 }
2794 #[cfg(feature = "draft16")]
2795 (FetchWriterState::Draft16(writer), FetchFrameShape::Draft16(original, location)) => {
2796 let reframed = writer.header_for(original, location)?;
2797 if reframed == *original {
2798 writer.advance(original, location);
2799 return Ok(FetchReemit::Unchanged);
2800 }
2801 let mut encoded = Vec::with_capacity(framing.len() + 16);
2802 reframed.encode(&mut encoded)?;
2803 writer.advance(&reframed, location);
2804 Ok(put_reframed(&encoded, framing.len(), rest, out))
2805 }
2806 #[cfg(feature = "draft17")]
2807 (FetchWriterState::Draft17(writer), FetchFrameShape::Draft17(original)) => {
2808 let reframed = writer.header_for(original)?;
2809 if reframed == original.header {
2810 writer.advance(original);
2811 return Ok(FetchReemit::Unchanged);
2812 }
2813 let mut encoded = Vec::with_capacity(framing.len() + 16);
2814 reframed.encode(&mut encoded)?;
2815 writer.advance(original);
2816 Ok(put_reframed(&encoded, framing.len(), rest, out))
2817 }
2818 #[cfg(feature = "draft18")]
2819 (FetchWriterState::Draft18(writer), FetchFrameShape::Draft18(original)) => {
2820 let reframed = writer.header_for(original)?;
2821 if reframed == original.header {
2822 writer.advance(original);
2823 return Ok(FetchReemit::Unchanged);
2824 }
2825 let mut encoded = Vec::with_capacity(framing.len() + 16);
2826 reframed.encode(&mut encoded)?;
2827 writer.advance(original);
2828 Ok(put_reframed(&encoded, framing.len(), rest, out))
2829 }
2830 #[cfg(feature = "draft19")]
2831 (FetchWriterState::Draft19(writer), FetchFrameShape::Draft19(original)) => {
2832 let reframed = writer.header_for(original)?;
2833 if reframed == original.header {
2834 writer.advance(original);
2835 return Ok(FetchReemit::Unchanged);
2836 }
2837 let mut encoded = Vec::with_capacity(framing.len() + 16);
2838 reframed.encode(&mut encoded)?;
2839 writer.advance(original);
2840 Ok(put_reframed(&encoded, framing.len(), rest, out))
2841 }
2842 #[allow(unreachable_patterns)]
2845 _ => Err(CodecError::UnsupportedDraft(format!(
2846 "no fetch writer for draft {:?}",
2847 frame.draft
2848 ))),
2849 }
2850 }
2851}
2852
2853#[cfg(any(
2856 feature = "draft15",
2857 feature = "draft16",
2858 feature = "draft17",
2859 feature = "draft18",
2860 feature = "draft19"
2861))]
2862fn put_reframed(
2863 encoded: &[u8],
2864 framing_bytes_before: usize,
2865 rest: &[u8],
2866 out: &mut impl BufMut,
2867) -> FetchReemit {
2868 out.put_slice(encoded);
2869 out.put_slice(rest);
2870 FetchReemit::Reframed { framing_bytes_before, framing_bytes_after: encoded.len() }
2871}
2872
2873macro_rules! subgroup_extensions_fn {
2879 ($name:ident, $feat:literal, $draft:ident) => {
2880 #[cfg(feature = $feat)]
2881 fn $name(header: &crate::$draft::data_stream::SubgroupHeader) -> Result<bool, CodecError> {
2882 if !header.stream_type.is_subgroup() {
2883 return Err(CodecError::InvalidField);
2884 }
2885 Ok(header.stream_type.has_extensions())
2886 }
2887 };
2888}
2889
2890subgroup_extensions_fn!(subgroup_extensions_11, "draft11", draft11);
2891subgroup_extensions_fn!(subgroup_extensions_12, "draft12", draft12);
2892subgroup_extensions_fn!(subgroup_extensions_13, "draft13", draft13);