1use bytes::{Buf, BufMut};
37
38use crate::dispatch::{AnyFetchHeader, AnySubgroupHeader};
39use crate::error::CodecError;
40use crate::varint::VarInt;
41use crate::version::DraftVersion;
42
43#[derive(Debug, Clone, PartialEq, Eq)]
52pub struct AnySubgroupObject {
53 pub object_id: u64,
56 pub extension_headers: Vec<u8>,
69 pub extension_count: Option<u64>,
74 pub status: Option<u64>,
86 pub payload: Vec<u8>,
88}
89
90#[derive(Debug, Clone, Copy, PartialEq, Eq)]
95pub struct AnySubgroupObjectMeta {
96 pub object_id: u64,
98 pub payload_length: u64,
100 pub status: Option<u64>,
102 pub extension_headers_len: u64,
107 pub wire_len: u64,
110}
111
112#[derive(Debug, Clone, Copy, PartialEq, Eq)]
124pub enum AnyFetchEndOfRange {
125 NonExistent,
127 Unknown,
129 TimedOut,
138}
139
140#[derive(Debug, Clone, Copy, PartialEq, Eq)]
151pub enum AnyFetchGroupOrder {
152 Ascending,
154 Descending,
156}
157
158#[derive(Debug, Clone, PartialEq, Eq)]
164pub struct AnyFetchObject {
165 pub group_id: u64,
169 pub subgroup_id: u64,
172 pub has_subgroup_id: bool,
185 pub object_id: u64,
187 pub publisher_priority: u8,
195 pub extension_headers: Vec<u8>,
198 pub extension_count: Option<u64>,
201 pub status: Option<u64>,
209 pub end_of_range: Option<AnyFetchEndOfRange>,
216 pub payload: Vec<u8>,
218}
219
220#[derive(Debug, Clone, Copy, PartialEq, Eq)]
224pub struct AnyFetchObjectMeta {
225 pub group_id: u64,
227 pub subgroup_id: u64,
229 pub has_subgroup_id: bool,
232 pub object_id: u64,
234 pub publisher_priority: u8,
237 pub payload_length: u64,
239 pub status: Option<u64>,
241 pub end_of_range: Option<AnyFetchEndOfRange>,
243 pub extension_headers_len: u64,
245 pub wire_len: u64,
247}
248
249#[cfg(any(
273 feature = "draft16",
274 feature = "draft17",
275 feature = "draft18",
276 feature = "draft19",
277 feature = "draft20"
278))]
279const DEFAULT_PUBLISHER_PRIORITY: u8 = 128;
280
281#[allow(dead_code)]
284mod conv {
285 use super::{AnySubgroupObject, Buf, CodecError};
286 use crate::varint::VarInt;
287
288 pub fn skip(buf: &mut impl Buf, len: u64) -> Result<(), CodecError> {
290 let len = usize::try_from(len).map_err(|_| CodecError::UnexpectedEnd)?;
291 if buf.remaining() < len {
292 return Err(CodecError::UnexpectedEnd);
293 }
294 buf.advance(len);
295 Ok(())
296 }
297
298 pub fn take(buf: &mut impl Buf, len: u64) -> Result<Vec<u8>, CodecError> {
300 let len = usize::try_from(len).map_err(|_| CodecError::UnexpectedEnd)?;
301 crate::types::read_bytes(buf, len)
302 }
303
304 pub fn varint(v: u64) -> Result<VarInt, CodecError> {
306 VarInt::from_u64(v).map_err(|_| CodecError::InvalidField)
307 }
308
309 pub fn status_to_write(object: &AnySubgroupObject) -> Result<Option<u64>, CodecError> {
324 match (object.status, object.payload.is_empty()) {
325 (Some(0), false) => Ok(None),
326 (Some(_), false) => Err(CodecError::InvalidField),
327 (Some(code), true) => Ok(Some(code)),
328 (None, true) => Ok(Some(0)),
329 (None, false) => Ok(None),
330 }
331 }
332}
333
334macro_rules! legacy_subgroup_glue {
344 (no_extensions $name:ident, $feat:literal, $draft:ident) => {
345 #[cfg(feature = $feat)]
346 mod $name {
347 use super::conv;
348 use super::{AnySubgroupObject, AnySubgroupObjectMeta};
349 use crate::error::CodecError;
350 use crate::$draft::data_stream::ObjectHeader;
351 use crate::$draft::types::ObjectStatus;
352 use bytes::{Buf, BufMut};
353
354 pub fn read_object(buf: &mut impl Buf) -> Result<AnySubgroupObject, CodecError> {
355 let header = ObjectHeader::decode(buf)?;
356 let payload_length = header.payload_length.into_inner();
357 let (status, payload) = if payload_length == 0 {
358 (Some(header.object_status as u64), Vec::new())
359 } else {
360 (None, conv::take(buf, payload_length)?)
361 };
362 Ok(AnySubgroupObject {
363 object_id: header.object_id.into_inner(),
364 extension_headers: Vec::new(),
365 extension_count: None,
366 status,
367 payload,
368 })
369 }
370
371 pub fn read_object_meta(
372 buf: &mut impl Buf,
373 ) -> Result<AnySubgroupObjectMeta, CodecError> {
374 let start = buf.remaining();
375 let header = ObjectHeader::decode(buf)?;
376 let payload_length = header.payload_length.into_inner();
377 let status = if payload_length == 0 {
378 Some(header.object_status as u64)
379 } else {
380 conv::skip(buf, payload_length)?;
381 None
382 };
383 Ok(AnySubgroupObjectMeta {
384 object_id: header.object_id.into_inner(),
385 payload_length,
386 status,
387 extension_headers_len: 0,
388 wire_len: (start - buf.remaining()) as u64,
389 })
390 }
391
392 pub fn write_object(
393 object: &AnySubgroupObject,
394 buf: &mut impl BufMut,
395 ) -> Result<(), CodecError> {
396 if !object.extension_headers.is_empty() {
397 return Err(CodecError::InvalidField);
398 }
399 let object_status = match conv::status_to_write(object)? {
400 Some(code) => ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?,
401 None => ObjectStatus::Normal,
402 };
403 ObjectHeader {
404 object_id: conv::varint(object.object_id)?,
405 payload_length: conv::varint(object.payload.len() as u64)?,
406 object_status,
407 }
408 .encode(buf);
409 buf.put_slice(&object.payload);
410 Ok(())
411 }
412 }
413 };
414
415 (count_extensions $name:ident, $feat:literal, $draft:ident) => {
416 #[cfg(feature = $feat)]
417 mod $name {
418 use super::conv;
419 use super::{AnySubgroupObject, AnySubgroupObjectMeta};
420 use crate::error::CodecError;
421 use crate::$draft::data_stream::ObjectHeader;
422 use crate::$draft::types::ObjectStatus;
423 use bytes::{Buf, BufMut};
424
425 pub fn read_object(buf: &mut impl Buf) -> Result<AnySubgroupObject, CodecError> {
426 let header = ObjectHeader::decode(buf)?;
427 let payload_length = header.payload_length.into_inner();
428 let (status, payload) = if payload_length == 0 {
429 (Some(header.object_status as u64), Vec::new())
430 } else {
431 (None, conv::take(buf, payload_length)?)
432 };
433 Ok(AnySubgroupObject {
434 object_id: header.object_id.into_inner(),
435 extension_headers: header.extensions,
436 extension_count: Some(header.extension_count.into_inner()),
437 status,
438 payload,
439 })
440 }
441
442 pub fn read_object_meta(
443 buf: &mut impl Buf,
444 ) -> Result<AnySubgroupObjectMeta, CodecError> {
445 let start = buf.remaining();
446 let header = ObjectHeader::decode(buf)?;
447 let payload_length = header.payload_length.into_inner();
448 let status = if payload_length == 0 {
449 Some(header.object_status as u64)
450 } else {
451 conv::skip(buf, payload_length)?;
452 None
453 };
454 Ok(AnySubgroupObjectMeta {
455 object_id: header.object_id.into_inner(),
456 payload_length,
457 status,
458 extension_headers_len: header.extensions.len() as u64,
459 wire_len: (start - buf.remaining()) as u64,
460 })
461 }
462
463 pub fn write_object(
464 object: &AnySubgroupObject,
465 buf: &mut impl BufMut,
466 ) -> Result<(), CodecError> {
467 let extension_count = match object.extension_count {
468 Some(count) => count,
469 None if object.extension_headers.is_empty() => 0,
470 None => return Err(CodecError::InvalidField),
471 };
472 let object_status = match conv::status_to_write(object)? {
473 Some(code) => ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?,
474 None => ObjectStatus::Normal,
475 };
476 ObjectHeader {
477 object_id: conv::varint(object.object_id)?,
478 extension_count: conv::varint(extension_count)?,
479 extensions: object.extension_headers.clone(),
480 payload_length: conv::varint(object.payload.len() as u64)?,
481 object_status,
482 }
483 .encode(buf);
484 buf.put_slice(&object.payload);
485 Ok(())
486 }
487 }
488 };
489
490 (length_extensions $name:ident, $feat:literal, $draft:ident) => {
491 #[cfg(feature = $feat)]
492 mod $name {
493 use super::conv;
494 use super::{AnySubgroupObject, AnySubgroupObjectMeta};
495 use crate::error::CodecError;
496 use crate::$draft::data_stream::ObjectHeader;
497 use crate::$draft::types::ObjectStatus;
498 use bytes::{Buf, BufMut};
499
500 pub fn read_object(buf: &mut impl Buf) -> Result<AnySubgroupObject, CodecError> {
501 let header = ObjectHeader::decode(buf)?;
502 let payload_length = header.payload_length.into_inner();
503 let (status, payload) = if payload_length == 0 {
504 (Some(header.object_status as u64), Vec::new())
505 } else {
506 (None, conv::take(buf, payload_length)?)
507 };
508 Ok(AnySubgroupObject {
509 object_id: header.object_id.into_inner(),
510 extension_headers: header.extensions,
511 extension_count: None,
512 status,
513 payload,
514 })
515 }
516
517 pub fn read_object_meta(
518 buf: &mut impl Buf,
519 ) -> Result<AnySubgroupObjectMeta, CodecError> {
520 let start = buf.remaining();
521 let header = ObjectHeader::decode(buf)?;
522 let payload_length = header.payload_length.into_inner();
523 let status = if payload_length == 0 {
524 Some(header.object_status as u64)
525 } else {
526 conv::skip(buf, payload_length)?;
527 None
528 };
529 Ok(AnySubgroupObjectMeta {
530 object_id: header.object_id.into_inner(),
531 payload_length,
532 status,
533 extension_headers_len: header.extension_headers_length.into_inner(),
534 wire_len: (start - buf.remaining()) as u64,
535 })
536 }
537
538 pub fn write_object(
539 object: &AnySubgroupObject,
540 buf: &mut impl BufMut,
541 ) -> Result<(), CodecError> {
542 let object_status = match conv::status_to_write(object)? {
543 Some(code) => ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?,
544 None => ObjectStatus::Normal,
545 };
546 ObjectHeader {
547 object_id: conv::varint(object.object_id)?,
548 extension_headers_length: conv::varint(object.extension_headers.len() as u64)?,
549 extensions: object.extension_headers.clone(),
550 payload_length: conv::varint(object.payload.len() as u64)?,
551 object_status,
552 }
553 .encode(buf);
554 buf.put_slice(&object.payload);
555 Ok(())
556 }
557 }
558 };
559
560 (gated_extensions $name:ident, $feat:literal, $draft:ident) => {
561 #[cfg(feature = $feat)]
562 mod $name {
563 use super::conv;
564 use super::{AnySubgroupObject, AnySubgroupObjectMeta};
565 use crate::error::CodecError;
566 use crate::$draft::data_stream::ObjectHeader;
567 use crate::$draft::types::ObjectStatus;
568 use bytes::{Buf, BufMut};
569
570 pub fn read_object(
571 extensions: bool,
572 buf: &mut impl Buf,
573 ) -> Result<AnySubgroupObject, CodecError> {
574 let header = ObjectHeader::decode_with_extensions(extensions, buf)?;
575 let payload_length = header.payload_length.into_inner();
576 let (status, payload) = if payload_length == 0 {
577 (Some(header.object_status as u64), Vec::new())
578 } else {
579 (None, conv::take(buf, payload_length)?)
580 };
581 Ok(AnySubgroupObject {
582 object_id: header.object_id.into_inner(),
583 extension_headers: header.extensions,
584 extension_count: None,
585 status,
586 payload,
587 })
588 }
589
590 pub fn read_object_meta(
591 extensions: bool,
592 buf: &mut impl Buf,
593 ) -> Result<AnySubgroupObjectMeta, CodecError> {
594 let start = buf.remaining();
595 let header = ObjectHeader::decode_with_extensions(extensions, buf)?;
596 let payload_length = header.payload_length.into_inner();
597 let status = if payload_length == 0 {
598 Some(header.object_status as u64)
599 } else {
600 conv::skip(buf, payload_length)?;
601 None
602 };
603 Ok(AnySubgroupObjectMeta {
604 object_id: header.object_id.into_inner(),
605 payload_length,
606 status,
607 extension_headers_len: header.extension_headers_length.into_inner(),
608 wire_len: (start - buf.remaining()) as u64,
609 })
610 }
611
612 pub fn write_object(
613 extensions: bool,
614 object: &AnySubgroupObject,
615 buf: &mut impl BufMut,
616 ) -> Result<(), CodecError> {
617 if !extensions && !object.extension_headers.is_empty() {
618 return Err(CodecError::InvalidField);
619 }
620 let object_status = match conv::status_to_write(object)? {
621 Some(code) => ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?,
622 None => ObjectStatus::Normal,
623 };
624 ObjectHeader {
625 object_id: conv::varint(object.object_id)?,
626 extension_headers_length: conv::varint(object.extension_headers.len() as u64)?,
627 extensions: object.extension_headers.clone(),
628 payload_length: conv::varint(object.payload.len() as u64)?,
629 object_status,
630 }
631 .encode_with_extensions(extensions, buf);
632 buf.put_slice(&object.payload);
633 Ok(())
634 }
635 }
636 };
637}
638
639legacy_subgroup_glue!(no_extensions sg07, "draft07", draft07);
640legacy_subgroup_glue!(count_extensions sg08, "draft08", draft08);
641legacy_subgroup_glue!(length_extensions sg09, "draft09", draft09);
642legacy_subgroup_glue!(length_extensions sg10, "draft10", draft10);
643legacy_subgroup_glue!(gated_extensions sg11, "draft11", draft11);
644legacy_subgroup_glue!(gated_extensions sg12, "draft12", draft12);
645legacy_subgroup_glue!(gated_extensions sg13, "draft13", draft13);
646
647macro_rules! modern_subgroup_glue {
661 (derived_length $name:ident, $feat:literal, $draft:ident) => {
662 #[cfg(feature = $feat)]
663 mod $name {
664 use super::conv;
665 use super::{AnySubgroupObject, AnySubgroupObjectMeta};
666 use crate::error::CodecError;
667 use crate::$draft::data_stream::{SubgroupObject, SubgroupObjectReader};
668 use crate::$draft::types::ObjectStatus;
669 use bytes::{Buf, BufMut};
670
671 pub fn read_object(
672 reader: &mut SubgroupObjectReader,
673 buf: &mut impl Buf,
674 ) -> Result<AnySubgroupObject, CodecError> {
675 let object = reader.read_object(buf)?;
676 Ok(AnySubgroupObject {
677 object_id: object.object_id.into_inner(),
678 extension_headers: object.extension_headers,
679 extension_count: None,
680 status: object.status.map(ObjectStatus::as_u64),
681 payload: object.payload,
682 })
683 }
684
685 pub fn read_object_meta(
686 reader: &mut SubgroupObjectReader,
687 buf: &mut impl Buf,
688 ) -> Result<AnySubgroupObjectMeta, CodecError> {
689 let meta = reader.read_object_meta(buf)?;
690 Ok(AnySubgroupObjectMeta {
691 object_id: meta.object_id,
692 payload_length: meta.payload_length,
693 status: meta.status,
694 extension_headers_len: meta.extension_headers_len,
695 wire_len: meta.wire_len,
696 })
697 }
698
699 pub fn write_object(
700 writer: &mut SubgroupObjectReader,
701 object: &AnySubgroupObject,
702 buf: &mut impl BufMut,
703 ) -> Result<(), CodecError> {
704 let status = match conv::status_to_write(object)? {
705 Some(code) => {
706 Some(ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?)
707 }
708 None => None,
709 };
710 writer.write_object(
711 &SubgroupObject {
712 object_id: conv::varint(object.object_id)?,
713 extension_headers: object.extension_headers.clone(),
714 status,
715 payload: object.payload.clone(),
716 },
717 buf,
718 )
719 }
720 }
721 };
722
723 (explicit_length $name:ident, $feat:literal, $draft:ident) => {
724 #[cfg(feature = $feat)]
725 mod $name {
726 use super::conv;
727 use super::{AnySubgroupObject, AnySubgroupObjectMeta};
728 use crate::error::CodecError;
729 use crate::$draft::data_stream::{SubgroupObject, SubgroupObjectReader};
730 use crate::$draft::types::ObjectStatus;
731 use bytes::{Buf, BufMut};
732
733 pub fn read_object(
734 reader: &mut SubgroupObjectReader,
735 buf: &mut impl Buf,
736 ) -> Result<AnySubgroupObject, CodecError> {
737 let object = reader.read_object(buf)?;
738 Ok(AnySubgroupObject {
739 object_id: object.object_id.into_inner(),
740 extension_headers: object.extension_headers,
741 extension_count: None,
742 status: object.object_status.map(ObjectStatus::as_u64),
743 payload: object.payload,
744 })
745 }
746
747 pub fn read_object_meta(
748 reader: &mut SubgroupObjectReader,
749 buf: &mut impl Buf,
750 ) -> Result<AnySubgroupObjectMeta, CodecError> {
751 let meta = reader.read_object_meta(buf)?;
752 Ok(AnySubgroupObjectMeta {
753 object_id: meta.object_id,
754 payload_length: meta.payload_length,
755 status: meta.status,
756 extension_headers_len: meta.extension_headers_len,
757 wire_len: meta.wire_len,
758 })
759 }
760
761 pub fn write_object(
762 writer: &mut SubgroupObjectReader,
763 object: &AnySubgroupObject,
764 buf: &mut impl BufMut,
765 ) -> Result<(), CodecError> {
766 let object_status = match conv::status_to_write(object)? {
767 Some(code) => {
768 Some(ObjectStatus::from_u64(code).ok_or(CodecError::InvalidField)?)
769 }
770 None => None,
771 };
772 writer.write_object(
773 &SubgroupObject {
774 object_id: conv::varint(object.object_id)?,
775 extension_headers: object.extension_headers.clone(),
776 payload_length: conv::varint(object.payload.len() as u64)?,
777 object_status,
778 payload: object.payload.clone(),
779 },
780 buf,
781 )
782 }
783 }
784 };
785}
786
787modern_subgroup_glue!(derived_length sg14, "draft14", draft14);
788modern_subgroup_glue!(explicit_length sg15, "draft15", draft15);
789modern_subgroup_glue!(explicit_length sg16, "draft16", draft16);
790modern_subgroup_glue!(explicit_length sg17, "draft17", draft17);
791modern_subgroup_glue!(explicit_length sg18, "draft18", draft18);
792modern_subgroup_glue!(explicit_length sg19, "draft19", draft19);
793modern_subgroup_glue!(explicit_length sg20, "draft20", draft20);
794
795macro_rules! fetch_glue {
802 (no_extensions $name:ident, $feat:literal, $draft:ident) => {
803 #[cfg(feature = $feat)]
804 mod $name {
805 use super::conv;
806 use super::{AnyFetchObject, AnyFetchObjectMeta};
807 use crate::error::CodecError;
808 use crate::$draft::data_stream::FetchObjectHeader;
809 use bytes::Buf;
810
811 pub fn read_object(buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
812 let header = FetchObjectHeader::decode(buf)?;
813 let payload_length = header.payload_length.into_inner();
814 let (status, payload) = if payload_length == 0 {
815 (Some(header.object_status as u64), Vec::new())
816 } else {
817 (None, conv::take(buf, payload_length)?)
818 };
819 Ok(AnyFetchObject {
820 group_id: header.group_id.into_inner(),
821 subgroup_id: header.subgroup_id.into_inner(),
822 has_subgroup_id: true,
823 object_id: header.object_id.into_inner(),
824 publisher_priority: header.publisher_priority,
825 extension_headers: Vec::new(),
826 extension_count: None,
827 status,
828 end_of_range: None,
829 payload,
830 })
831 }
832
833 pub fn read_object_meta(buf: &mut impl Buf) -> Result<AnyFetchObjectMeta, CodecError> {
834 let start = buf.remaining();
835 let header = FetchObjectHeader::decode(buf)?;
836 let payload_length = header.payload_length.into_inner();
837 let status = if payload_length == 0 {
838 Some(header.object_status as u64)
839 } else {
840 conv::skip(buf, payload_length)?;
841 None
842 };
843 Ok(AnyFetchObjectMeta {
844 group_id: header.group_id.into_inner(),
845 subgroup_id: header.subgroup_id.into_inner(),
846 has_subgroup_id: true,
847 object_id: header.object_id.into_inner(),
848 publisher_priority: header.publisher_priority,
849 payload_length,
850 status,
851 end_of_range: None,
852 extension_headers_len: 0,
853 wire_len: (start - buf.remaining()) as u64,
854 })
855 }
856 }
857 };
858
859 (count_extensions $name:ident, $feat:literal, $draft:ident) => {
860 #[cfg(feature = $feat)]
861 mod $name {
862 use super::conv;
863 use super::{AnyFetchObject, AnyFetchObjectMeta};
864 use crate::error::CodecError;
865 use crate::$draft::data_stream::FetchObjectHeader;
866 use bytes::Buf;
867
868 pub fn read_object(buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
869 let header = FetchObjectHeader::decode(buf)?;
870 let payload_length = header.payload_length.into_inner();
871 let (status, payload) = if payload_length == 0 {
872 (Some(header.object_status as u64), Vec::new())
873 } else {
874 (None, conv::take(buf, payload_length)?)
875 };
876 Ok(AnyFetchObject {
877 group_id: header.group_id.into_inner(),
878 subgroup_id: header.subgroup_id.into_inner(),
879 has_subgroup_id: true,
880 object_id: header.object_id.into_inner(),
881 publisher_priority: header.publisher_priority,
882 extension_headers: header.extensions,
883 extension_count: Some(header.extension_count.into_inner()),
884 status,
885 end_of_range: None,
886 payload,
887 })
888 }
889
890 pub fn read_object_meta(buf: &mut impl Buf) -> Result<AnyFetchObjectMeta, CodecError> {
891 let start = buf.remaining();
892 let header = FetchObjectHeader::decode(buf)?;
893 let payload_length = header.payload_length.into_inner();
894 let status = if payload_length == 0 {
895 Some(header.object_status as u64)
896 } else {
897 conv::skip(buf, payload_length)?;
898 None
899 };
900 Ok(AnyFetchObjectMeta {
901 group_id: header.group_id.into_inner(),
902 subgroup_id: header.subgroup_id.into_inner(),
903 has_subgroup_id: true,
904 object_id: header.object_id.into_inner(),
905 publisher_priority: header.publisher_priority,
906 payload_length,
907 status,
908 end_of_range: None,
909 extension_headers_len: header.extensions.len() as u64,
910 wire_len: (start - buf.remaining()) as u64,
911 })
912 }
913 }
914 };
915
916 (length_extensions $name:ident, $feat:literal, $draft:ident) => {
917 #[cfg(feature = $feat)]
918 mod $name {
919 use super::conv;
920 use super::{AnyFetchObject, AnyFetchObjectMeta};
921 use crate::error::CodecError;
922 use crate::$draft::data_stream::FetchObjectHeader;
923 use bytes::Buf;
924
925 pub fn read_object(buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
926 let header = FetchObjectHeader::decode(buf)?;
927 let payload_length = header.payload_length.into_inner();
928 let (status, payload) = if payload_length == 0 {
929 (Some(header.object_status as u64), Vec::new())
930 } else {
931 (None, conv::take(buf, payload_length)?)
932 };
933 Ok(AnyFetchObject {
934 group_id: header.group_id.into_inner(),
935 subgroup_id: header.subgroup_id.into_inner(),
936 has_subgroup_id: true,
937 object_id: header.object_id.into_inner(),
938 publisher_priority: header.publisher_priority,
939 extension_headers: header.extensions,
940 extension_count: None,
941 status,
942 end_of_range: None,
943 payload,
944 })
945 }
946
947 pub fn read_object_meta(buf: &mut impl Buf) -> Result<AnyFetchObjectMeta, CodecError> {
948 let start = buf.remaining();
949 let header = FetchObjectHeader::decode(buf)?;
950 let payload_length = header.payload_length.into_inner();
951 let status = if payload_length == 0 {
952 Some(header.object_status as u64)
953 } else {
954 conv::skip(buf, payload_length)?;
955 None
956 };
957 Ok(AnyFetchObjectMeta {
958 group_id: header.group_id.into_inner(),
959 subgroup_id: header.subgroup_id.into_inner(),
960 has_subgroup_id: true,
961 object_id: header.object_id.into_inner(),
962 publisher_priority: header.publisher_priority,
963 payload_length,
964 status,
965 end_of_range: None,
966 extension_headers_len: header.extension_headers_length.into_inner(),
967 wire_len: (start - buf.remaining()) as u64,
968 })
969 }
970 }
971 };
972}
973
974fetch_glue!(no_extensions fo07, "draft07", draft07);
975fetch_glue!(count_extensions fo08, "draft08", draft08);
976fetch_glue!(length_extensions fo09, "draft09", draft09);
977fetch_glue!(length_extensions fo10, "draft10", draft10);
978fetch_glue!(length_extensions fo11, "draft11", draft11);
979fetch_glue!(length_extensions fo12, "draft12", draft12);
980fetch_glue!(length_extensions fo13, "draft13", draft13);
981
982#[cfg(feature = "draft14")]
983mod fo14 {
984 use super::{AnyFetchObject, AnyFetchObjectMeta};
985 use crate::draft14::data_stream::FetchObject;
986 use crate::draft14::types::ObjectStatus;
987 use crate::error::CodecError;
988 use bytes::Buf;
989
990 pub fn read_object(buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
991 let object = FetchObject::decode(buf)?;
992 Ok(AnyFetchObject {
993 group_id: object.group_id.into_inner(),
994 subgroup_id: object.subgroup_id.into_inner(),
995 has_subgroup_id: true,
996 object_id: object.object_id.into_inner(),
997 publisher_priority: object.publisher_priority,
998 extension_headers: object.extension_headers,
999 extension_count: None,
1000 status: object.status.map(ObjectStatus::as_u64),
1001 end_of_range: None,
1002 payload: object.payload,
1003 })
1004 }
1005
1006 pub fn read_object_meta(buf: &mut impl Buf) -> Result<AnyFetchObjectMeta, CodecError> {
1007 let meta = FetchObject::decode_meta(buf)?;
1008 Ok(AnyFetchObjectMeta {
1009 group_id: meta.group_id,
1010 subgroup_id: meta.subgroup_id,
1011 has_subgroup_id: true,
1012 object_id: meta.object_id,
1013 publisher_priority: meta.publisher_priority,
1014 payload_length: meta.payload_length,
1015 status: meta.status,
1016 end_of_range: None,
1017 extension_headers_len: meta.extension_headers_len,
1018 wire_len: meta.wire_len,
1019 })
1020 }
1021}
1022
1023#[cfg(feature = "draft15")]
1024mod fo15 {
1025 use super::{conv, AnyFetchObject, AnyFetchObjectMeta};
1026 use crate::draft15::data_stream::FetchObjectReader;
1027 use crate::draft15::types::ObjectStatus;
1028 use crate::error::CodecError;
1029 use bytes::Buf;
1030
1031 pub fn read_object(
1032 reader: &mut FetchObjectReader,
1033 buf: &mut impl Buf,
1034 ) -> Result<AnyFetchObject, CodecError> {
1035 let header = reader.read_object_header(buf)?;
1036 let payload = conv::take(buf, header.payload_length.into_inner())?;
1040 Ok(AnyFetchObject {
1041 group_id: header.group_id.into_inner(),
1042 subgroup_id: header.subgroup_id.into_inner(),
1043 has_subgroup_id: true,
1044 object_id: header.object_id.into_inner(),
1045 publisher_priority: header.publisher_priority,
1046 extension_headers: header.extension_headers,
1047 extension_count: None,
1048 status: header.object_status.map(ObjectStatus::as_u64),
1051 end_of_range: None,
1052 payload,
1053 })
1054 }
1055
1056 pub fn read_object_frame(
1057 reader: &mut FetchObjectReader,
1058 buf: &mut impl Buf,
1059 ) -> Result<super::AnyFetchFrame, CodecError> {
1060 let start = buf.remaining();
1061 let header = reader.read_object_header(buf)?;
1062 let payload_length = header.payload_length.into_inner();
1063 conv::skip(buf, payload_length)?;
1064 let meta = AnyFetchObjectMeta {
1065 group_id: header.group_id.into_inner(),
1066 subgroup_id: header.subgroup_id.into_inner(),
1067 has_subgroup_id: true,
1068 object_id: header.object_id.into_inner(),
1069 publisher_priority: header.publisher_priority,
1070 payload_length,
1071 status: header.object_status.map(ObjectStatus::as_u64),
1072 end_of_range: None,
1073 extension_headers_len: header.extension_headers.len() as u64,
1074 wire_len: (start - buf.remaining()) as u64,
1075 };
1076 Ok(super::AnyFetchFrame {
1077 meta,
1078 draft: crate::version::DraftVersion::Draft15,
1079 shape: super::FetchFrameShape::Draft15(header),
1080 })
1081 }
1082
1083 pub fn read_object_meta(
1084 reader: &mut FetchObjectReader,
1085 buf: &mut impl Buf,
1086 ) -> Result<AnyFetchObjectMeta, CodecError> {
1087 read_object_frame(reader, buf).map(|frame| frame.meta)
1088 }
1089}
1090
1091#[cfg(feature = "draft16")]
1092mod fo16 {
1093 use super::DEFAULT_PUBLISHER_PRIORITY;
1094 use super::{conv, AnyFetchEndOfRange, AnyFetchObject, AnyFetchObjectMeta};
1095 use crate::draft16::data_stream::{
1096 FetchEndOfRange, FetchObjectHeader, FetchObjectLocation, FetchObjectReader,
1097 };
1098 use crate::error::CodecError;
1099 use bytes::Buf;
1100
1101 fn parts(
1105 reader: &mut FetchObjectReader,
1106 buf: &mut impl Buf,
1107 ) -> Result<(FetchObjectHeader, FetchObjectLocation), CodecError> {
1108 let header = FetchObjectHeader::decode(buf)?;
1109 let location = reader.resolve(&header)?;
1110 Ok((header, location))
1111 }
1112
1113 fn resolved(location: &FetchObjectLocation) -> super::Resolved {
1114 let end_of_range = location.end_of_range;
1115 super::Resolved {
1116 group_id: location.group_id,
1117 subgroup_id: location.subgroup_id.filter(|_| end_of_range.is_none()),
1124 object_id: location.object_id,
1125 publisher_priority: location.publisher_priority.unwrap_or(DEFAULT_PUBLISHER_PRIORITY),
1126 end_of_range: end_of_range.map(|r| match r {
1127 FetchEndOfRange::NonExistent => AnyFetchEndOfRange::NonExistent,
1128 FetchEndOfRange::Unknown => AnyFetchEndOfRange::Unknown,
1129 }),
1130 }
1131 }
1132
1133 pub fn read_object(
1134 reader: &mut FetchObjectReader,
1135 buf: &mut impl Buf,
1136 ) -> Result<AnyFetchObject, CodecError> {
1137 let (header, location) = parts(reader, buf)?;
1138 let payload = conv::take(buf, header.payload_length.into_inner())?;
1139 Ok(resolved(&location).into_object(header.extensions.unwrap_or_default(), payload))
1140 }
1141
1142 pub fn read_object_frame(
1143 reader: &mut FetchObjectReader,
1144 buf: &mut impl Buf,
1145 ) -> Result<super::AnyFetchFrame, CodecError> {
1146 let start = buf.remaining();
1147 let (header, location) = parts(reader, buf)?;
1148 let payload_length = header.payload_length.into_inner();
1149 conv::skip(buf, payload_length)?;
1150 let meta = resolved(&location).into_meta(
1151 header.extensions.as_ref().map_or(0, |e| e.len() as u64),
1152 payload_length,
1153 (start - buf.remaining()) as u64,
1154 );
1155 Ok(super::AnyFetchFrame {
1156 meta,
1157 draft: crate::version::DraftVersion::Draft16,
1158 shape: super::FetchFrameShape::Draft16(header, location),
1159 })
1160 }
1161
1162 pub fn read_object_meta(
1163 reader: &mut FetchObjectReader,
1164 buf: &mut impl Buf,
1165 ) -> Result<AnyFetchObjectMeta, CodecError> {
1166 read_object_frame(reader, buf).map(|frame| frame.meta)
1167 }
1168}
1169
1170#[cfg(feature = "draft17")]
1171mod fo17 {
1172 use super::{conv, AnyFetchEndOfRange, AnyFetchObject, AnyFetchObjectMeta};
1173 use crate::draft17::data_stream::{EndOfRange, FetchObject, FetchObjectReader};
1174 use crate::error::CodecError;
1175 use bytes::Buf;
1176
1177 fn resolved(object: &FetchObject) -> super::Resolved {
1178 super::Resolved {
1179 group_id: object.group_id,
1180 subgroup_id: object.subgroup_id,
1181 object_id: object.object_id,
1182 publisher_priority: object
1183 .publisher_priority
1184 .unwrap_or(super::DEFAULT_PUBLISHER_PRIORITY),
1185 end_of_range: object.header.end_of_range().map(|r| match r {
1186 EndOfRange::NonExistent => AnyFetchEndOfRange::NonExistent,
1187 EndOfRange::Unknown => AnyFetchEndOfRange::Unknown,
1188 }),
1189 }
1190 }
1191
1192 pub fn read_object(
1193 reader: &mut FetchObjectReader,
1194 buf: &mut impl Buf,
1195 ) -> Result<AnyFetchObject, CodecError> {
1196 let object = reader.read_object_header(buf)?;
1197 let resolved = resolved(&object);
1198 let payload = conv::take(buf, object.header.payload_length.into_inner())?;
1199 Ok(resolved.into_object(object.header.properties, payload))
1200 }
1201
1202 pub fn read_object_frame(
1203 reader: &mut FetchObjectReader,
1204 buf: &mut impl Buf,
1205 ) -> Result<super::AnyFetchFrame, CodecError> {
1206 let start = buf.remaining();
1207 let object = reader.read_object_header(buf)?;
1208 let resolved = resolved(&object);
1209 let payload_length = object.header.payload_length.into_inner();
1210 conv::skip(buf, payload_length)?;
1211 let meta = resolved.into_meta(
1212 object.header.properties.len() as u64,
1213 payload_length,
1214 (start - buf.remaining()) as u64,
1215 );
1216 Ok(super::AnyFetchFrame {
1217 meta,
1218 draft: crate::version::DraftVersion::Draft17,
1219 shape: super::FetchFrameShape::Draft17(object),
1220 })
1221 }
1222
1223 pub fn read_object_meta(
1224 reader: &mut FetchObjectReader,
1225 buf: &mut impl Buf,
1226 ) -> Result<AnyFetchObjectMeta, CodecError> {
1227 read_object_frame(reader, buf).map(|frame| frame.meta)
1228 }
1229}
1230
1231#[cfg(feature = "draft18")]
1232mod fo18 {
1233 use super::{conv, AnyFetchEndOfRange, AnyFetchObject, AnyFetchObjectMeta};
1234 use crate::draft18::data_stream::{EndOfRange, FetchObject, FetchObjectReader};
1235 use crate::error::CodecError;
1236 use bytes::Buf;
1237
1238 fn resolved(object: &FetchObject) -> super::Resolved {
1239 super::Resolved {
1240 group_id: object.group_id,
1241 subgroup_id: object.subgroup_id,
1242 object_id: object.object_id,
1243 publisher_priority: object
1244 .publisher_priority
1245 .unwrap_or(super::DEFAULT_PUBLISHER_PRIORITY),
1246 end_of_range: object.header.end_of_range().map(|r| match r {
1247 EndOfRange::NonExistent => AnyFetchEndOfRange::NonExistent,
1248 EndOfRange::Unknown => AnyFetchEndOfRange::Unknown,
1249 }),
1250 }
1251 }
1252
1253 pub fn read_object(
1254 reader: &mut FetchObjectReader,
1255 buf: &mut impl Buf,
1256 ) -> Result<AnyFetchObject, CodecError> {
1257 let object = reader.read_object_header(buf)?;
1258 let resolved = resolved(&object);
1259 let payload = conv::take(buf, object.header.payload_length.into_inner())?;
1260 Ok(resolved.into_object(object.header.properties, payload))
1261 }
1262
1263 pub fn read_object_frame(
1264 reader: &mut FetchObjectReader,
1265 buf: &mut impl Buf,
1266 ) -> Result<super::AnyFetchFrame, CodecError> {
1267 let start = buf.remaining();
1268 let object = reader.read_object_header(buf)?;
1269 let resolved = resolved(&object);
1270 let payload_length = object.header.payload_length.into_inner();
1271 conv::skip(buf, payload_length)?;
1272 let meta = resolved.into_meta(
1273 object.header.properties.len() as u64,
1274 payload_length,
1275 (start - buf.remaining()) as u64,
1276 );
1277 Ok(super::AnyFetchFrame {
1278 meta,
1279 draft: crate::version::DraftVersion::Draft18,
1280 shape: super::FetchFrameShape::Draft18(object),
1281 })
1282 }
1283
1284 pub fn read_object_meta(
1285 reader: &mut FetchObjectReader,
1286 buf: &mut impl Buf,
1287 ) -> Result<AnyFetchObjectMeta, CodecError> {
1288 read_object_frame(reader, buf).map(|frame| frame.meta)
1289 }
1290}
1291
1292#[cfg(feature = "draft19")]
1293mod fo19 {
1294 use super::{conv, AnyFetchEndOfRange, AnyFetchObject, AnyFetchObjectMeta};
1295 use crate::draft19::data_stream::{FetchEndOfRange, FetchObject, FetchObjectReader};
1296 use crate::error::CodecError;
1297 use bytes::Buf;
1298
1299 fn resolved(object: &FetchObject) -> super::Resolved {
1300 super::Resolved {
1301 group_id: object.group_id,
1302 subgroup_id: object.subgroup_id,
1303 object_id: object.object_id,
1304 publisher_priority: object
1305 .publisher_priority
1306 .unwrap_or(super::DEFAULT_PUBLISHER_PRIORITY),
1307 end_of_range: object.header.end_of_range().map(|r| match r {
1308 FetchEndOfRange::NonExistent => AnyFetchEndOfRange::NonExistent,
1309 FetchEndOfRange::Unknown => AnyFetchEndOfRange::Unknown,
1310 }),
1311 }
1312 }
1313
1314 pub fn read_object(
1315 reader: &mut FetchObjectReader,
1316 buf: &mut impl Buf,
1317 ) -> Result<AnyFetchObject, CodecError> {
1318 let object = reader.read_object_header(buf)?;
1319 let resolved = resolved(&object);
1320 let payload = conv::take(buf, object.header.payload_length.into_inner())?;
1321 Ok(resolved.into_object(object.header.properties.unwrap_or_default(), payload))
1322 }
1323
1324 pub fn read_object_frame(
1325 reader: &mut FetchObjectReader,
1326 buf: &mut impl Buf,
1327 ) -> Result<super::AnyFetchFrame, CodecError> {
1328 let start = buf.remaining();
1329 let object = reader.read_object_header(buf)?;
1330 let resolved = resolved(&object);
1331 let payload_length = object.header.payload_length.into_inner();
1332 conv::skip(buf, payload_length)?;
1333 let meta = resolved.into_meta(
1334 object.header.properties.as_ref().map_or(0, |p| p.len() as u64),
1335 payload_length,
1336 (start - buf.remaining()) as u64,
1337 );
1338 Ok(super::AnyFetchFrame {
1339 meta,
1340 draft: crate::version::DraftVersion::Draft19,
1341 shape: super::FetchFrameShape::Draft19(object),
1342 })
1343 }
1344
1345 pub fn read_object_meta(
1346 reader: &mut FetchObjectReader,
1347 buf: &mut impl Buf,
1348 ) -> Result<AnyFetchObjectMeta, CodecError> {
1349 read_object_frame(reader, buf).map(|frame| frame.meta)
1350 }
1351}
1352
1353#[cfg(feature = "draft20")]
1354mod fo20 {
1355 use super::{conv, AnyFetchEndOfRange, AnyFetchObject, AnyFetchObjectMeta};
1356 use crate::draft20::data_stream::{FetchEndOfRange, FetchObject, FetchObjectReader};
1357 use crate::error::CodecError;
1358 use bytes::Buf;
1359
1360 fn resolved(object: &FetchObject) -> super::Resolved {
1361 super::Resolved {
1362 group_id: object.group_id,
1363 subgroup_id: object.subgroup_id,
1364 object_id: object.object_id,
1365 publisher_priority: object
1366 .publisher_priority
1367 .unwrap_or(super::DEFAULT_PUBLISHER_PRIORITY),
1368 end_of_range: object.header.end_of_range().map(|r| match r {
1369 FetchEndOfRange::NonExistent => AnyFetchEndOfRange::NonExistent,
1370 FetchEndOfRange::Unknown => AnyFetchEndOfRange::Unknown,
1371 FetchEndOfRange::TimedOut => AnyFetchEndOfRange::TimedOut,
1374 }),
1375 }
1376 }
1377
1378 pub fn read_object(
1379 reader: &mut FetchObjectReader,
1380 buf: &mut impl Buf,
1381 ) -> Result<AnyFetchObject, CodecError> {
1382 let object = reader.read_object_header(buf)?;
1383 let resolved = resolved(&object);
1384 let payload = conv::take(buf, object.header.payload_length.into_inner())?;
1385 Ok(resolved.into_object(object.header.properties.unwrap_or_default(), payload))
1386 }
1387
1388 pub fn read_object_frame(
1389 reader: &mut FetchObjectReader,
1390 buf: &mut impl Buf,
1391 ) -> Result<super::AnyFetchFrame, CodecError> {
1392 let start = buf.remaining();
1393 let object = reader.read_object_header(buf)?;
1394 let resolved = resolved(&object);
1395 let payload_length = object.header.payload_length.into_inner();
1396 conv::skip(buf, payload_length)?;
1397 let meta = resolved.into_meta(
1398 object.header.properties.as_ref().map_or(0, |p| p.len() as u64),
1399 payload_length,
1400 (start - buf.remaining()) as u64,
1401 );
1402 Ok(super::AnyFetchFrame {
1403 meta,
1404 draft: crate::version::DraftVersion::Draft20,
1405 shape: super::FetchFrameShape::Draft20(object),
1406 })
1407 }
1408
1409 pub fn read_object_meta(
1410 reader: &mut FetchObjectReader,
1411 buf: &mut impl Buf,
1412 ) -> Result<AnyFetchObjectMeta, CodecError> {
1413 read_object_frame(reader, buf).map(|frame| frame.meta)
1414 }
1415}
1416
1417#[cfg(any(
1426 feature = "draft16",
1427 feature = "draft17",
1428 feature = "draft18",
1429 feature = "draft19",
1430 feature = "draft20"
1431))]
1432struct Resolved {
1433 group_id: u64,
1434 subgroup_id: Option<u64>,
1435 object_id: u64,
1436 publisher_priority: u8,
1437 end_of_range: Option<AnyFetchEndOfRange>,
1438}
1439
1440#[cfg(any(
1441 feature = "draft16",
1442 feature = "draft17",
1443 feature = "draft18",
1444 feature = "draft19",
1445 feature = "draft20"
1446))]
1447impl Resolved {
1448 fn into_object(self, extension_headers: Vec<u8>, payload: Vec<u8>) -> AnyFetchObject {
1449 AnyFetchObject {
1450 group_id: self.group_id,
1451 subgroup_id: self.subgroup_id.unwrap_or(0),
1452 has_subgroup_id: self.subgroup_id.is_some(),
1453 object_id: self.object_id,
1454 publisher_priority: self.publisher_priority,
1455 extension_headers,
1456 extension_count: None,
1457 status: None,
1460 end_of_range: self.end_of_range,
1461 payload,
1462 }
1463 }
1464
1465 fn into_meta(
1466 self,
1467 extension_headers_len: u64,
1468 payload_length: u64,
1469 wire_len: u64,
1470 ) -> AnyFetchObjectMeta {
1471 AnyFetchObjectMeta {
1472 group_id: self.group_id,
1473 subgroup_id: self.subgroup_id.unwrap_or(0),
1474 has_subgroup_id: self.subgroup_id.is_some(),
1475 object_id: self.object_id,
1476 publisher_priority: self.publisher_priority,
1477 payload_length,
1478 status: None,
1479 end_of_range: self.end_of_range,
1480 extension_headers_len,
1481 wire_len,
1482 }
1483 }
1484}
1485
1486#[derive(Debug, Clone)]
1492enum SubgroupReaderState {
1493 #[cfg(feature = "draft07")]
1494 Draft07,
1495 #[cfg(feature = "draft08")]
1496 Draft08,
1497 #[cfg(feature = "draft09")]
1498 Draft09,
1499 #[cfg(feature = "draft10")]
1500 Draft10,
1501 #[cfg(feature = "draft11")]
1502 Draft11 { extensions: bool },
1503 #[cfg(feature = "draft12")]
1504 Draft12 { extensions: bool },
1505 #[cfg(feature = "draft13")]
1506 Draft13 { extensions: bool },
1507 #[cfg(feature = "draft14")]
1508 Draft14(crate::draft14::data_stream::SubgroupObjectReader),
1509 #[cfg(feature = "draft15")]
1510 Draft15(crate::draft15::data_stream::SubgroupObjectReader),
1511 #[cfg(feature = "draft16")]
1512 Draft16(crate::draft16::data_stream::SubgroupObjectReader),
1513 #[cfg(feature = "draft17")]
1514 Draft17(crate::draft17::data_stream::SubgroupObjectReader),
1515 #[cfg(feature = "draft18")]
1516 Draft18(crate::draft18::data_stream::SubgroupObjectReader),
1517 #[cfg(feature = "draft19")]
1518 Draft19(crate::draft19::data_stream::SubgroupObjectReader),
1519 #[cfg(feature = "draft20")]
1520 Draft20(crate::draft20::data_stream::SubgroupObjectReader),
1521}
1522
1523#[derive(Debug, Clone)]
1534pub struct AnySubgroupObjectReader {
1535 state: SubgroupReaderState,
1536}
1537
1538impl AnySubgroupObjectReader {
1539 #[allow(unused_variables, unreachable_code)]
1545 pub fn new(header: &AnySubgroupHeader) -> Result<Self, CodecError> {
1546 let state = match header {
1547 #[cfg(feature = "draft07")]
1548 AnySubgroupHeader::Draft07(_) => SubgroupReaderState::Draft07,
1549 #[cfg(feature = "draft08")]
1550 AnySubgroupHeader::Draft08(_) => SubgroupReaderState::Draft08,
1551 #[cfg(feature = "draft09")]
1552 AnySubgroupHeader::Draft09(_) => SubgroupReaderState::Draft09,
1553 #[cfg(feature = "draft10")]
1554 AnySubgroupHeader::Draft10(_) => SubgroupReaderState::Draft10,
1555 #[cfg(feature = "draft11")]
1556 AnySubgroupHeader::Draft11(h) => {
1557 SubgroupReaderState::Draft11 { extensions: subgroup_extensions_11(h)? }
1558 }
1559 #[cfg(feature = "draft12")]
1560 AnySubgroupHeader::Draft12(h) => {
1561 SubgroupReaderState::Draft12 { extensions: subgroup_extensions_12(h)? }
1562 }
1563 #[cfg(feature = "draft13")]
1564 AnySubgroupHeader::Draft13(h) => {
1565 SubgroupReaderState::Draft13 { extensions: subgroup_extensions_13(h)? }
1566 }
1567 #[cfg(feature = "draft14")]
1568 AnySubgroupHeader::Draft14(h) => SubgroupReaderState::Draft14(
1569 crate::draft14::data_stream::SubgroupObjectReader::new(h),
1570 ),
1571 #[cfg(feature = "draft15")]
1572 AnySubgroupHeader::Draft15(h) => SubgroupReaderState::Draft15(
1573 crate::draft15::data_stream::SubgroupObjectReader::new(h),
1574 ),
1575 #[cfg(feature = "draft16")]
1576 AnySubgroupHeader::Draft16(h) => SubgroupReaderState::Draft16(
1577 crate::draft16::data_stream::SubgroupObjectReader::new(h),
1578 ),
1579 #[cfg(feature = "draft17")]
1580 AnySubgroupHeader::Draft17(h) => SubgroupReaderState::Draft17(
1581 crate::draft17::data_stream::SubgroupObjectReader::new(h),
1582 ),
1583 #[cfg(feature = "draft18")]
1584 AnySubgroupHeader::Draft18(h) => SubgroupReaderState::Draft18(
1585 crate::draft18::data_stream::SubgroupObjectReader::new(h),
1586 ),
1587 #[cfg(feature = "draft19")]
1588 AnySubgroupHeader::Draft19(h) => SubgroupReaderState::Draft19(
1589 crate::draft19::data_stream::SubgroupObjectReader::new(h),
1590 ),
1591 #[cfg(feature = "draft20")]
1592 AnySubgroupHeader::Draft20(h) => SubgroupReaderState::Draft20(
1593 crate::draft20::data_stream::SubgroupObjectReader::new(h),
1594 ),
1595 #[allow(unreachable_patterns)]
1596 _ => {
1597 return Err(CodecError::UnsupportedDraft(format!(
1598 "draft {:?} not enabled via feature flag",
1599 header.draft()
1600 )));
1601 }
1602 };
1603 Ok(Self { state })
1604 }
1605
1606 #[allow(unreachable_code)]
1608 pub fn draft(&self) -> DraftVersion {
1609 match &self.state {
1610 #[cfg(feature = "draft07")]
1611 SubgroupReaderState::Draft07 => DraftVersion::Draft07,
1612 #[cfg(feature = "draft08")]
1613 SubgroupReaderState::Draft08 => DraftVersion::Draft08,
1614 #[cfg(feature = "draft09")]
1615 SubgroupReaderState::Draft09 => DraftVersion::Draft09,
1616 #[cfg(feature = "draft10")]
1617 SubgroupReaderState::Draft10 => DraftVersion::Draft10,
1618 #[cfg(feature = "draft11")]
1619 SubgroupReaderState::Draft11 { .. } => DraftVersion::Draft11,
1620 #[cfg(feature = "draft12")]
1621 SubgroupReaderState::Draft12 { .. } => DraftVersion::Draft12,
1622 #[cfg(feature = "draft13")]
1623 SubgroupReaderState::Draft13 { .. } => DraftVersion::Draft13,
1624 #[cfg(feature = "draft14")]
1625 SubgroupReaderState::Draft14(_) => DraftVersion::Draft14,
1626 #[cfg(feature = "draft15")]
1627 SubgroupReaderState::Draft15(_) => DraftVersion::Draft15,
1628 #[cfg(feature = "draft16")]
1629 SubgroupReaderState::Draft16(_) => DraftVersion::Draft16,
1630 #[cfg(feature = "draft17")]
1631 SubgroupReaderState::Draft17(_) => DraftVersion::Draft17,
1632 #[cfg(feature = "draft18")]
1633 SubgroupReaderState::Draft18(_) => DraftVersion::Draft18,
1634 #[cfg(feature = "draft19")]
1635 SubgroupReaderState::Draft19(_) => DraftVersion::Draft19,
1636 #[cfg(feature = "draft20")]
1637 SubgroupReaderState::Draft20(_) => DraftVersion::Draft20,
1638 #[allow(unreachable_patterns)]
1639 _ => unreachable!("AnySubgroupObjectReader has no enabled variants"),
1640 }
1641 }
1642
1643 #[allow(unused_variables, unreachable_code)]
1649 pub fn read_object(&mut self, buf: &mut impl Buf) -> Result<AnySubgroupObject, CodecError> {
1650 match &mut self.state {
1651 #[cfg(feature = "draft07")]
1652 SubgroupReaderState::Draft07 => sg07::read_object(buf),
1653 #[cfg(feature = "draft08")]
1654 SubgroupReaderState::Draft08 => sg08::read_object(buf),
1655 #[cfg(feature = "draft09")]
1656 SubgroupReaderState::Draft09 => sg09::read_object(buf),
1657 #[cfg(feature = "draft10")]
1658 SubgroupReaderState::Draft10 => sg10::read_object(buf),
1659 #[cfg(feature = "draft11")]
1660 SubgroupReaderState::Draft11 { extensions } => sg11::read_object(*extensions, buf),
1661 #[cfg(feature = "draft12")]
1662 SubgroupReaderState::Draft12 { extensions } => sg12::read_object(*extensions, buf),
1663 #[cfg(feature = "draft13")]
1664 SubgroupReaderState::Draft13 { extensions } => sg13::read_object(*extensions, buf),
1665 #[cfg(feature = "draft14")]
1666 SubgroupReaderState::Draft14(inner) => sg14::read_object(inner, buf),
1667 #[cfg(feature = "draft15")]
1668 SubgroupReaderState::Draft15(inner) => sg15::read_object(inner, buf),
1669 #[cfg(feature = "draft16")]
1670 SubgroupReaderState::Draft16(inner) => sg16::read_object(inner, buf),
1671 #[cfg(feature = "draft17")]
1672 SubgroupReaderState::Draft17(inner) => sg17::read_object(inner, buf),
1673 #[cfg(feature = "draft18")]
1674 SubgroupReaderState::Draft18(inner) => sg18::read_object(inner, buf),
1675 #[cfg(feature = "draft19")]
1676 SubgroupReaderState::Draft19(inner) => sg19::read_object(inner, buf),
1677 #[cfg(feature = "draft20")]
1678 SubgroupReaderState::Draft20(inner) => sg20::read_object(inner, buf),
1679 #[allow(unreachable_patterns)]
1680 _ => unreachable!("AnySubgroupObjectReader has no enabled variants"),
1681 }
1682 }
1683
1684 #[allow(unused_variables, unreachable_code)]
1691 pub fn read_object_meta(
1692 &mut self,
1693 buf: &mut impl Buf,
1694 ) -> Result<AnySubgroupObjectMeta, CodecError> {
1695 match &mut self.state {
1696 #[cfg(feature = "draft07")]
1697 SubgroupReaderState::Draft07 => sg07::read_object_meta(buf),
1698 #[cfg(feature = "draft08")]
1699 SubgroupReaderState::Draft08 => sg08::read_object_meta(buf),
1700 #[cfg(feature = "draft09")]
1701 SubgroupReaderState::Draft09 => sg09::read_object_meta(buf),
1702 #[cfg(feature = "draft10")]
1703 SubgroupReaderState::Draft10 => sg10::read_object_meta(buf),
1704 #[cfg(feature = "draft11")]
1705 SubgroupReaderState::Draft11 { extensions } => sg11::read_object_meta(*extensions, buf),
1706 #[cfg(feature = "draft12")]
1707 SubgroupReaderState::Draft12 { extensions } => sg12::read_object_meta(*extensions, buf),
1708 #[cfg(feature = "draft13")]
1709 SubgroupReaderState::Draft13 { extensions } => sg13::read_object_meta(*extensions, buf),
1710 #[cfg(feature = "draft14")]
1711 SubgroupReaderState::Draft14(inner) => sg14::read_object_meta(inner, buf),
1712 #[cfg(feature = "draft15")]
1713 SubgroupReaderState::Draft15(inner) => sg15::read_object_meta(inner, buf),
1714 #[cfg(feature = "draft16")]
1715 SubgroupReaderState::Draft16(inner) => sg16::read_object_meta(inner, buf),
1716 #[cfg(feature = "draft17")]
1717 SubgroupReaderState::Draft17(inner) => sg17::read_object_meta(inner, buf),
1718 #[cfg(feature = "draft18")]
1719 SubgroupReaderState::Draft18(inner) => sg18::read_object_meta(inner, buf),
1720 #[cfg(feature = "draft19")]
1721 SubgroupReaderState::Draft19(inner) => sg19::read_object_meta(inner, buf),
1722 #[cfg(feature = "draft20")]
1723 SubgroupReaderState::Draft20(inner) => sg20::read_object_meta(inner, buf),
1724 #[allow(unreachable_patterns)]
1725 _ => unreachable!("AnySubgroupObjectReader has no enabled variants"),
1726 }
1727 }
1728}
1729
1730#[derive(Debug, Clone)]
1738enum SubgroupWriterState {
1739 #[cfg(feature = "draft07")]
1740 Draft07 { prev_object_id: Option<u64> },
1741 #[cfg(feature = "draft08")]
1742 Draft08 { prev_object_id: Option<u64> },
1743 #[cfg(feature = "draft09")]
1744 Draft09 { prev_object_id: Option<u64> },
1745 #[cfg(feature = "draft10")]
1746 Draft10 { prev_object_id: Option<u64> },
1747 #[cfg(feature = "draft11")]
1748 Draft11 { extensions: bool, prev_object_id: Option<u64> },
1749 #[cfg(feature = "draft12")]
1750 Draft12 { extensions: bool, prev_object_id: Option<u64> },
1751 #[cfg(feature = "draft13")]
1752 Draft13 { extensions: bool, prev_object_id: Option<u64> },
1753 #[cfg(feature = "draft14")]
1754 Draft14 { inner: crate::draft14::data_stream::SubgroupObjectReader, extensions: bool },
1755 #[cfg(feature = "draft15")]
1756 Draft15 { inner: crate::draft15::data_stream::SubgroupObjectReader, extensions: bool },
1757 #[cfg(feature = "draft16")]
1758 Draft16 { inner: crate::draft16::data_stream::SubgroupObjectReader, extensions: bool },
1759 #[cfg(feature = "draft17")]
1760 Draft17 { inner: crate::draft17::data_stream::SubgroupObjectReader, extensions: bool },
1761 #[cfg(feature = "draft18")]
1762 Draft18 { inner: crate::draft18::data_stream::SubgroupObjectReader, extensions: bool },
1763 #[cfg(feature = "draft19")]
1764 Draft19 { inner: crate::draft19::data_stream::SubgroupObjectReader, extensions: bool },
1765 #[cfg(feature = "draft20")]
1766 Draft20 { inner: crate::draft20::data_stream::SubgroupObjectReader, extensions: bool },
1767}
1768
1769#[derive(Debug, Clone)]
1785pub struct AnySubgroupObjectWriter {
1786 state: SubgroupWriterState,
1787}
1788
1789impl AnySubgroupObjectWriter {
1790 #[allow(unused_variables, unreachable_code)]
1796 pub fn new(header: &AnySubgroupHeader) -> Result<Self, CodecError> {
1797 let state = match header {
1798 #[cfg(feature = "draft07")]
1799 AnySubgroupHeader::Draft07(_) => SubgroupWriterState::Draft07 { prev_object_id: None },
1800 #[cfg(feature = "draft08")]
1801 AnySubgroupHeader::Draft08(_) => SubgroupWriterState::Draft08 { prev_object_id: None },
1802 #[cfg(feature = "draft09")]
1803 AnySubgroupHeader::Draft09(_) => SubgroupWriterState::Draft09 { prev_object_id: None },
1804 #[cfg(feature = "draft10")]
1805 AnySubgroupHeader::Draft10(_) => SubgroupWriterState::Draft10 { prev_object_id: None },
1806 #[cfg(feature = "draft11")]
1807 AnySubgroupHeader::Draft11(h) => SubgroupWriterState::Draft11 {
1808 extensions: subgroup_extensions_11(h)?,
1809 prev_object_id: None,
1810 },
1811 #[cfg(feature = "draft12")]
1812 AnySubgroupHeader::Draft12(h) => SubgroupWriterState::Draft12 {
1813 extensions: subgroup_extensions_12(h)?,
1814 prev_object_id: None,
1815 },
1816 #[cfg(feature = "draft13")]
1817 AnySubgroupHeader::Draft13(h) => SubgroupWriterState::Draft13 {
1818 extensions: subgroup_extensions_13(h)?,
1819 prev_object_id: None,
1820 },
1821 #[cfg(feature = "draft14")]
1822 AnySubgroupHeader::Draft14(h) => SubgroupWriterState::Draft14 {
1823 inner: crate::draft14::data_stream::SubgroupObjectReader::new(h),
1824 extensions: h.stream_type.extensions_present(),
1825 },
1826 #[cfg(feature = "draft15")]
1827 AnySubgroupHeader::Draft15(h) => SubgroupWriterState::Draft15 {
1828 inner: crate::draft15::data_stream::SubgroupObjectReader::new(h),
1829 extensions: h.has_extensions(),
1830 },
1831 #[cfg(feature = "draft16")]
1832 AnySubgroupHeader::Draft16(h) => SubgroupWriterState::Draft16 {
1833 inner: crate::draft16::data_stream::SubgroupObjectReader::new(h),
1834 extensions: h.has_extensions(),
1835 },
1836 #[cfg(feature = "draft17")]
1837 AnySubgroupHeader::Draft17(h) => SubgroupWriterState::Draft17 {
1838 inner: crate::draft17::data_stream::SubgroupObjectReader::new(h),
1839 extensions: h.has_properties(),
1840 },
1841 #[cfg(feature = "draft18")]
1842 AnySubgroupHeader::Draft18(h) => SubgroupWriterState::Draft18 {
1843 inner: crate::draft18::data_stream::SubgroupObjectReader::new(h),
1844 extensions: h.has_properties(),
1845 },
1846 #[cfg(feature = "draft19")]
1847 AnySubgroupHeader::Draft19(h) => SubgroupWriterState::Draft19 {
1848 inner: crate::draft19::data_stream::SubgroupObjectReader::new(h),
1849 extensions: h.has_properties(),
1850 },
1851 #[cfg(feature = "draft20")]
1852 AnySubgroupHeader::Draft20(h) => SubgroupWriterState::Draft20 {
1853 inner: crate::draft20::data_stream::SubgroupObjectReader::new(h),
1854 extensions: h.has_properties(),
1855 },
1856 #[allow(unreachable_patterns)]
1857 _ => {
1858 return Err(CodecError::UnsupportedDraft(format!(
1859 "draft {:?} not enabled via feature flag",
1860 header.draft()
1861 )));
1862 }
1863 };
1864 Ok(Self { state })
1865 }
1866
1867 #[allow(unreachable_code)]
1869 pub fn draft(&self) -> DraftVersion {
1870 match &self.state {
1871 #[cfg(feature = "draft07")]
1872 SubgroupWriterState::Draft07 { .. } => DraftVersion::Draft07,
1873 #[cfg(feature = "draft08")]
1874 SubgroupWriterState::Draft08 { .. } => DraftVersion::Draft08,
1875 #[cfg(feature = "draft09")]
1876 SubgroupWriterState::Draft09 { .. } => DraftVersion::Draft09,
1877 #[cfg(feature = "draft10")]
1878 SubgroupWriterState::Draft10 { .. } => DraftVersion::Draft10,
1879 #[cfg(feature = "draft11")]
1880 SubgroupWriterState::Draft11 { .. } => DraftVersion::Draft11,
1881 #[cfg(feature = "draft12")]
1882 SubgroupWriterState::Draft12 { .. } => DraftVersion::Draft12,
1883 #[cfg(feature = "draft13")]
1884 SubgroupWriterState::Draft13 { .. } => DraftVersion::Draft13,
1885 #[cfg(feature = "draft14")]
1886 SubgroupWriterState::Draft14 { .. } => DraftVersion::Draft14,
1887 #[cfg(feature = "draft15")]
1888 SubgroupWriterState::Draft15 { .. } => DraftVersion::Draft15,
1889 #[cfg(feature = "draft16")]
1890 SubgroupWriterState::Draft16 { .. } => DraftVersion::Draft16,
1891 #[cfg(feature = "draft17")]
1892 SubgroupWriterState::Draft17 { .. } => DraftVersion::Draft17,
1893 #[cfg(feature = "draft18")]
1894 SubgroupWriterState::Draft18 { .. } => DraftVersion::Draft18,
1895 #[cfg(feature = "draft19")]
1896 SubgroupWriterState::Draft19 { .. } => DraftVersion::Draft19,
1897 #[cfg(feature = "draft20")]
1898 SubgroupWriterState::Draft20 { .. } => DraftVersion::Draft20,
1899 #[allow(unreachable_patterns)]
1900 _ => unreachable!("AnySubgroupObjectWriter has no enabled variants"),
1901 }
1902 }
1903
1904 #[allow(unused_variables, unreachable_code)]
1933 pub fn write_object(
1934 &mut self,
1935 object: &AnySubgroupObject,
1936 buf: &mut impl BufMut,
1937 ) -> Result<(), CodecError> {
1938 match &mut self.state {
1939 #[cfg(feature = "draft07")]
1940 SubgroupWriterState::Draft07 { prev_object_id } => {
1941 advance_absolute_id(prev_object_id, object, |o| sg07::write_object(o, buf))
1942 }
1943 #[cfg(feature = "draft08")]
1944 SubgroupWriterState::Draft08 { prev_object_id } => {
1945 advance_absolute_id(prev_object_id, object, |o| sg08::write_object(o, buf))
1946 }
1947 #[cfg(feature = "draft09")]
1948 SubgroupWriterState::Draft09 { prev_object_id } => {
1949 advance_absolute_id(prev_object_id, object, |o| sg09::write_object(o, buf))
1950 }
1951 #[cfg(feature = "draft10")]
1952 SubgroupWriterState::Draft10 { prev_object_id } => {
1953 advance_absolute_id(prev_object_id, object, |o| sg10::write_object(o, buf))
1954 }
1955 #[cfg(feature = "draft11")]
1956 SubgroupWriterState::Draft11 { extensions, prev_object_id } => {
1957 let extensions = *extensions;
1958 advance_absolute_id(prev_object_id, object, |o| {
1959 sg11::write_object(extensions, o, buf)
1960 })
1961 }
1962 #[cfg(feature = "draft12")]
1963 SubgroupWriterState::Draft12 { extensions, prev_object_id } => {
1964 let extensions = *extensions;
1965 advance_absolute_id(prev_object_id, object, |o| {
1966 sg12::write_object(extensions, o, buf)
1967 })
1968 }
1969 #[cfg(feature = "draft13")]
1970 SubgroupWriterState::Draft13 { extensions, prev_object_id } => {
1971 let extensions = *extensions;
1972 advance_absolute_id(prev_object_id, object, |o| {
1973 sg13::write_object(extensions, o, buf)
1974 })
1975 }
1976 #[cfg(feature = "draft14")]
1977 SubgroupWriterState::Draft14 { inner, extensions } => {
1978 reject_unrepresentable_extensions(*extensions, object)?;
1979 sg14::write_object(inner, object, buf)
1980 }
1981 #[cfg(feature = "draft15")]
1982 SubgroupWriterState::Draft15 { inner, extensions } => {
1983 reject_unrepresentable_extensions(*extensions, object)?;
1984 sg15::write_object(inner, object, buf)
1985 }
1986 #[cfg(feature = "draft16")]
1987 SubgroupWriterState::Draft16 { inner, extensions } => {
1988 reject_unrepresentable_extensions(*extensions, object)?;
1989 sg16::write_object(inner, object, buf)
1990 }
1991 #[cfg(feature = "draft17")]
1992 SubgroupWriterState::Draft17 { inner, extensions } => {
1993 reject_unrepresentable_extensions(*extensions, object)?;
1994 sg17::write_object(inner, object, buf)
1995 }
1996 #[cfg(feature = "draft18")]
1997 SubgroupWriterState::Draft18 { inner, extensions } => {
1998 reject_unrepresentable_extensions(*extensions, object)?;
1999 sg18::write_object(inner, object, buf)
2000 }
2001 #[cfg(feature = "draft19")]
2002 SubgroupWriterState::Draft19 { inner, extensions } => {
2003 reject_unrepresentable_extensions(*extensions, object)?;
2004 sg19::write_object(inner, object, buf)
2005 }
2006 #[cfg(feature = "draft20")]
2007 SubgroupWriterState::Draft20 { inner, extensions } => {
2008 reject_unrepresentable_extensions(*extensions, object)?;
2009 sg20::write_object(inner, object, buf)
2010 }
2011 #[allow(unreachable_patterns)]
2012 _ => unreachable!("AnySubgroupObjectWriter has no enabled variants"),
2013 }
2014 }
2015}
2016
2017#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2021pub enum Reemit {
2022 Verbatim,
2024 Reencoded {
2026 id_bytes_before: usize,
2028 id_bytes_after: usize,
2030 },
2031}
2032
2033pub fn reemit_subgroup_object(
2095 draft: DraftVersion,
2096 prev_forwarded: Option<u64>,
2097 object_id: u64,
2098 raw: &[u8],
2099 out: &mut impl BufMut,
2100) -> Result<Reemit, CodecError> {
2101 if matches!(prev_forwarded, Some(prev) if object_id <= prev) {
2102 return Err(CodecError::InvalidField);
2103 }
2104
2105 let mut cursor: &[u8] = raw;
2109 draft.decode_varint(&mut cursor).map_err(|_| CodecError::InvalidField)?;
2110 let id_bytes_before = raw.len() - cursor.len();
2111
2112 if !delta_encodes_object_ids(draft) {
2115 out.put_slice(raw);
2116 return Ok(Reemit::Verbatim);
2117 }
2118
2119 let delta = match prev_forwarded {
2120 None => object_id,
2121 Some(prev) => object_id
2122 .checked_sub(prev)
2123 .and_then(|v| v.checked_sub(1))
2124 .ok_or(CodecError::InvalidField)?,
2125 };
2126
2127 let field = if draft.uses_moqt_varint() {
2130 VarInt::from_u64_moqt(delta)
2131 } else {
2132 VarInt::from_u64(delta).map_err(|_| CodecError::InvalidField)?
2133 };
2134
2135 let mut encoded = [0u8; 9];
2137 let mut slot: &mut [u8] = &mut encoded;
2138 draft.encode_varint(field, &mut slot);
2139 let id_bytes_after = 9 - slot.len();
2140 let encoded = &encoded[..id_bytes_after];
2141
2142 if encoded == &raw[..id_bytes_before] {
2143 out.put_slice(raw);
2144 return Ok(Reemit::Verbatim);
2145 }
2146
2147 out.put_slice(encoded);
2148 out.put_slice(&raw[id_bytes_before..]);
2149 Ok(Reemit::Reencoded { id_bytes_before, id_bytes_after })
2150}
2151
2152fn delta_encodes_object_ids(draft: DraftVersion) -> bool {
2158 matches!(
2159 draft,
2160 DraftVersion::Draft14
2161 | DraftVersion::Draft15
2162 | DraftVersion::Draft16
2163 | DraftVersion::Draft17
2164 | DraftVersion::Draft18
2165 | DraftVersion::Draft19
2166 | DraftVersion::Draft20
2167 )
2168}
2169
2170#[cfg(any(
2180 feature = "draft07",
2181 feature = "draft08",
2182 feature = "draft09",
2183 feature = "draft10",
2184 feature = "draft11",
2185 feature = "draft12",
2186 feature = "draft13"
2187))]
2188fn advance_absolute_id(
2189 prev_object_id: &mut Option<u64>,
2190 object: &AnySubgroupObject,
2191 write: impl FnOnce(&AnySubgroupObject) -> Result<(), CodecError>,
2192) -> Result<(), CodecError> {
2193 if matches!(*prev_object_id, Some(prev) if object.object_id <= prev) {
2194 return Err(CodecError::InvalidField);
2195 }
2196 write(object)?;
2197 *prev_object_id = Some(object.object_id);
2198 Ok(())
2199}
2200
2201#[cfg(any(
2204 feature = "draft14",
2205 feature = "draft15",
2206 feature = "draft16",
2207 feature = "draft17",
2208 feature = "draft18",
2209 feature = "draft19",
2210 feature = "draft20"
2211))]
2212fn reject_unrepresentable_extensions(
2213 extensions: bool,
2214 object: &AnySubgroupObject,
2215) -> Result<(), CodecError> {
2216 if !extensions && !object.extension_headers.is_empty() {
2217 return Err(CodecError::InvalidField);
2218 }
2219 Ok(())
2220}
2221
2222#[derive(Debug, Clone)]
2228enum FetchReaderState {
2229 #[cfg(feature = "draft07")]
2230 Draft07,
2231 #[cfg(feature = "draft08")]
2232 Draft08,
2233 #[cfg(feature = "draft09")]
2234 Draft09,
2235 #[cfg(feature = "draft10")]
2236 Draft10,
2237 #[cfg(feature = "draft11")]
2238 Draft11,
2239 #[cfg(feature = "draft12")]
2240 Draft12,
2241 #[cfg(feature = "draft13")]
2242 Draft13,
2243 #[cfg(feature = "draft14")]
2244 Draft14,
2245 #[cfg(feature = "draft15")]
2246 Draft15(crate::draft15::data_stream::FetchObjectReader),
2247 #[cfg(feature = "draft16")]
2248 Draft16(crate::draft16::data_stream::FetchObjectReader),
2249 #[cfg(feature = "draft17")]
2250 Draft17(crate::draft17::data_stream::FetchObjectReader),
2251 #[cfg(feature = "draft18")]
2252 Draft18(crate::draft18::data_stream::FetchObjectReader),
2253 #[cfg(feature = "draft19")]
2254 Draft19(crate::draft19::data_stream::FetchObjectReader),
2255 #[cfg(feature = "draft20")]
2256 Draft20(crate::draft20::data_stream::FetchObjectReader),
2257}
2258
2259#[derive(Debug, Clone)]
2284pub struct AnyFetchObjectReader {
2285 state: FetchReaderState,
2286}
2287
2288impl AnyFetchObjectReader {
2289 #[allow(unused_variables, unreachable_code)]
2307 pub fn new(
2308 header: &AnyFetchHeader,
2309 group_order: AnyFetchGroupOrder,
2310 ) -> Result<Self, CodecError> {
2311 let state = match header {
2312 #[cfg(feature = "draft07")]
2313 AnyFetchHeader::Draft07(_) => FetchReaderState::Draft07,
2314 #[cfg(feature = "draft08")]
2315 AnyFetchHeader::Draft08(_) => FetchReaderState::Draft08,
2316 #[cfg(feature = "draft09")]
2317 AnyFetchHeader::Draft09(_) => FetchReaderState::Draft09,
2318 #[cfg(feature = "draft10")]
2319 AnyFetchHeader::Draft10(_) => FetchReaderState::Draft10,
2320 #[cfg(feature = "draft11")]
2321 AnyFetchHeader::Draft11(_) => FetchReaderState::Draft11,
2322 #[cfg(feature = "draft12")]
2323 AnyFetchHeader::Draft12(_) => FetchReaderState::Draft12,
2324 #[cfg(feature = "draft13")]
2325 AnyFetchHeader::Draft13(_) => FetchReaderState::Draft13,
2326 #[cfg(feature = "draft14")]
2327 AnyFetchHeader::Draft14(_) => FetchReaderState::Draft14,
2328 #[cfg(feature = "draft15")]
2331 AnyFetchHeader::Draft15(_) => {
2332 FetchReaderState::Draft15(crate::draft15::data_stream::FetchObjectReader::new())
2333 }
2334 #[cfg(feature = "draft16")]
2335 AnyFetchHeader::Draft16(_) => {
2336 FetchReaderState::Draft16(crate::draft16::data_stream::FetchObjectReader::new())
2337 }
2338 #[cfg(feature = "draft17")]
2339 AnyFetchHeader::Draft17(_) => {
2340 FetchReaderState::Draft17(crate::draft17::data_stream::FetchObjectReader::new())
2341 }
2342 #[cfg(feature = "draft18")]
2343 AnyFetchHeader::Draft18(_) => FetchReaderState::Draft18(
2344 crate::draft18::data_stream::FetchObjectReader::new(match group_order {
2345 AnyFetchGroupOrder::Ascending => {
2346 crate::draft18::data_stream::GroupOrder::Ascending
2347 }
2348 AnyFetchGroupOrder::Descending => {
2349 crate::draft18::data_stream::GroupOrder::Descending
2350 }
2351 }),
2352 ),
2353 #[cfg(feature = "draft19")]
2354 AnyFetchHeader::Draft19(_) => FetchReaderState::Draft19(
2355 crate::draft19::data_stream::FetchObjectReader::new(match group_order {
2356 AnyFetchGroupOrder::Ascending => {
2357 crate::draft19::data_stream::GroupOrder::Ascending
2358 }
2359 AnyFetchGroupOrder::Descending => {
2360 crate::draft19::data_stream::GroupOrder::Descending
2361 }
2362 }),
2363 ),
2364 #[cfg(feature = "draft20")]
2365 AnyFetchHeader::Draft20(_) => FetchReaderState::Draft20(
2366 crate::draft20::data_stream::FetchObjectReader::new(match group_order {
2367 AnyFetchGroupOrder::Ascending => {
2368 crate::draft20::data_stream::GroupOrder::Ascending
2369 }
2370 AnyFetchGroupOrder::Descending => {
2371 crate::draft20::data_stream::GroupOrder::Descending
2372 }
2373 }),
2374 ),
2375 #[allow(unreachable_patterns)]
2376 _ => {
2377 return Err(CodecError::UnsupportedDraft(format!(
2378 "draft {:?} not enabled via feature flag",
2379 header.draft()
2380 )));
2381 }
2382 };
2383 Ok(Self { state })
2384 }
2385
2386 #[allow(unreachable_code)]
2388 pub fn draft(&self) -> DraftVersion {
2389 match &self.state {
2390 #[cfg(feature = "draft07")]
2391 FetchReaderState::Draft07 => DraftVersion::Draft07,
2392 #[cfg(feature = "draft08")]
2393 FetchReaderState::Draft08 => DraftVersion::Draft08,
2394 #[cfg(feature = "draft09")]
2395 FetchReaderState::Draft09 => DraftVersion::Draft09,
2396 #[cfg(feature = "draft10")]
2397 FetchReaderState::Draft10 => DraftVersion::Draft10,
2398 #[cfg(feature = "draft11")]
2399 FetchReaderState::Draft11 => DraftVersion::Draft11,
2400 #[cfg(feature = "draft12")]
2401 FetchReaderState::Draft12 => DraftVersion::Draft12,
2402 #[cfg(feature = "draft13")]
2403 FetchReaderState::Draft13 => DraftVersion::Draft13,
2404 #[cfg(feature = "draft14")]
2405 FetchReaderState::Draft14 => DraftVersion::Draft14,
2406 #[cfg(feature = "draft15")]
2407 FetchReaderState::Draft15(_) => DraftVersion::Draft15,
2408 #[cfg(feature = "draft16")]
2409 FetchReaderState::Draft16(_) => DraftVersion::Draft16,
2410 #[cfg(feature = "draft17")]
2411 FetchReaderState::Draft17(_) => DraftVersion::Draft17,
2412 #[cfg(feature = "draft18")]
2413 FetchReaderState::Draft18(_) => DraftVersion::Draft18,
2414 #[cfg(feature = "draft19")]
2415 FetchReaderState::Draft19(_) => DraftVersion::Draft19,
2416 #[cfg(feature = "draft20")]
2417 FetchReaderState::Draft20(_) => DraftVersion::Draft20,
2418 #[allow(unreachable_patterns)]
2419 _ => unreachable!("AnyFetchObjectReader has no enabled variants"),
2420 }
2421 }
2422
2423 #[allow(unused_variables, unreachable_code)]
2435 pub fn read_object(&mut self, buf: &mut impl Buf) -> Result<AnyFetchObject, CodecError> {
2436 match &mut self.state {
2437 #[cfg(feature = "draft07")]
2438 FetchReaderState::Draft07 => fo07::read_object(buf),
2439 #[cfg(feature = "draft08")]
2440 FetchReaderState::Draft08 => fo08::read_object(buf),
2441 #[cfg(feature = "draft09")]
2442 FetchReaderState::Draft09 => fo09::read_object(buf),
2443 #[cfg(feature = "draft10")]
2444 FetchReaderState::Draft10 => fo10::read_object(buf),
2445 #[cfg(feature = "draft11")]
2446 FetchReaderState::Draft11 => fo11::read_object(buf),
2447 #[cfg(feature = "draft12")]
2448 FetchReaderState::Draft12 => fo12::read_object(buf),
2449 #[cfg(feature = "draft13")]
2450 FetchReaderState::Draft13 => fo13::read_object(buf),
2451 #[cfg(feature = "draft14")]
2452 FetchReaderState::Draft14 => fo14::read_object(buf),
2453 #[cfg(feature = "draft15")]
2454 FetchReaderState::Draft15(inner) => fo15::read_object(inner, buf),
2455 #[cfg(feature = "draft16")]
2456 FetchReaderState::Draft16(inner) => fo16::read_object(inner, buf),
2457 #[cfg(feature = "draft17")]
2458 FetchReaderState::Draft17(inner) => fo17::read_object(inner, buf),
2459 #[cfg(feature = "draft18")]
2460 FetchReaderState::Draft18(inner) => fo18::read_object(inner, buf),
2461 #[cfg(feature = "draft19")]
2462 FetchReaderState::Draft19(inner) => fo19::read_object(inner, buf),
2463 #[cfg(feature = "draft20")]
2464 FetchReaderState::Draft20(inner) => fo20::read_object(inner, buf),
2465 #[allow(unreachable_patterns)]
2466 _ => unreachable!("AnyFetchObjectReader has no enabled variants"),
2467 }
2468 }
2469
2470 #[allow(unused_variables, unreachable_code)]
2483 pub fn read_object_frame(&mut self, buf: &mut impl Buf) -> Result<AnyFetchFrame, CodecError> {
2484 match &mut self.state {
2485 #[cfg(feature = "draft07")]
2486 FetchReaderState::Draft07 => fo07::read_object_meta(buf)
2487 .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft07, meta)),
2488 #[cfg(feature = "draft08")]
2489 FetchReaderState::Draft08 => fo08::read_object_meta(buf)
2490 .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft08, meta)),
2491 #[cfg(feature = "draft09")]
2492 FetchReaderState::Draft09 => fo09::read_object_meta(buf)
2493 .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft09, meta)),
2494 #[cfg(feature = "draft10")]
2495 FetchReaderState::Draft10 => fo10::read_object_meta(buf)
2496 .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft10, meta)),
2497 #[cfg(feature = "draft11")]
2498 FetchReaderState::Draft11 => fo11::read_object_meta(buf)
2499 .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft11, meta)),
2500 #[cfg(feature = "draft12")]
2501 FetchReaderState::Draft12 => fo12::read_object_meta(buf)
2502 .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft12, meta)),
2503 #[cfg(feature = "draft13")]
2504 FetchReaderState::Draft13 => fo13::read_object_meta(buf)
2505 .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft13, meta)),
2506 #[cfg(feature = "draft14")]
2507 FetchReaderState::Draft14 => fo14::read_object_meta(buf)
2508 .map(|meta| AnyFetchFrame::absolute(DraftVersion::Draft14, meta)),
2509 #[cfg(feature = "draft15")]
2510 FetchReaderState::Draft15(inner) => fo15::read_object_frame(inner, buf),
2511 #[cfg(feature = "draft16")]
2512 FetchReaderState::Draft16(inner) => fo16::read_object_frame(inner, buf),
2513 #[cfg(feature = "draft17")]
2514 FetchReaderState::Draft17(inner) => fo17::read_object_frame(inner, buf),
2515 #[cfg(feature = "draft18")]
2516 FetchReaderState::Draft18(inner) => fo18::read_object_frame(inner, buf),
2517 #[cfg(feature = "draft19")]
2518 FetchReaderState::Draft19(inner) => fo19::read_object_frame(inner, buf),
2519 #[cfg(feature = "draft20")]
2520 FetchReaderState::Draft20(inner) => fo20::read_object_frame(inner, buf),
2521 #[allow(unreachable_patterns)]
2522 _ => unreachable!("AnyFetchObjectReader has no enabled variants"),
2523 }
2524 }
2525
2526 #[allow(unused_variables, unreachable_code)]
2532 pub fn read_object_meta(
2533 &mut self,
2534 buf: &mut impl Buf,
2535 ) -> Result<AnyFetchObjectMeta, CodecError> {
2536 match &mut self.state {
2537 #[cfg(feature = "draft07")]
2538 FetchReaderState::Draft07 => fo07::read_object_meta(buf),
2539 #[cfg(feature = "draft08")]
2540 FetchReaderState::Draft08 => fo08::read_object_meta(buf),
2541 #[cfg(feature = "draft09")]
2542 FetchReaderState::Draft09 => fo09::read_object_meta(buf),
2543 #[cfg(feature = "draft10")]
2544 FetchReaderState::Draft10 => fo10::read_object_meta(buf),
2545 #[cfg(feature = "draft11")]
2546 FetchReaderState::Draft11 => fo11::read_object_meta(buf),
2547 #[cfg(feature = "draft12")]
2548 FetchReaderState::Draft12 => fo12::read_object_meta(buf),
2549 #[cfg(feature = "draft13")]
2550 FetchReaderState::Draft13 => fo13::read_object_meta(buf),
2551 #[cfg(feature = "draft14")]
2552 FetchReaderState::Draft14 => fo14::read_object_meta(buf),
2553 #[cfg(feature = "draft15")]
2554 FetchReaderState::Draft15(inner) => fo15::read_object_meta(inner, buf),
2555 #[cfg(feature = "draft16")]
2556 FetchReaderState::Draft16(inner) => fo16::read_object_meta(inner, buf),
2557 #[cfg(feature = "draft17")]
2558 FetchReaderState::Draft17(inner) => fo17::read_object_meta(inner, buf),
2559 #[cfg(feature = "draft18")]
2560 FetchReaderState::Draft18(inner) => fo18::read_object_meta(inner, buf),
2561 #[cfg(feature = "draft19")]
2562 FetchReaderState::Draft19(inner) => fo19::read_object_meta(inner, buf),
2563 #[cfg(feature = "draft20")]
2564 FetchReaderState::Draft20(inner) => fo20::read_object_meta(inner, buf),
2565 #[allow(unreachable_patterns)]
2566 _ => unreachable!("AnyFetchObjectReader has no enabled variants"),
2567 }
2568 }
2569}
2570
2571#[derive(Debug, Clone)]
2581enum FetchFrameShape {
2582 #[cfg(any(
2584 feature = "draft07",
2585 feature = "draft08",
2586 feature = "draft09",
2587 feature = "draft10",
2588 feature = "draft11",
2589 feature = "draft12",
2590 feature = "draft13",
2591 feature = "draft14"
2592 ))]
2593 Absolute,
2594 #[cfg(feature = "draft15")]
2595 Draft15(crate::draft15::data_stream::FetchObjectHeader),
2596 #[cfg(feature = "draft16")]
2597 Draft16(
2598 crate::draft16::data_stream::FetchObjectHeader,
2599 crate::draft16::data_stream::FetchObjectLocation,
2600 ),
2601 #[cfg(feature = "draft17")]
2602 Draft17(crate::draft17::data_stream::FetchObject),
2603 #[cfg(feature = "draft18")]
2604 Draft18(crate::draft18::data_stream::FetchObject),
2605 #[cfg(feature = "draft19")]
2606 Draft19(crate::draft19::data_stream::FetchObject),
2607 #[cfg(feature = "draft20")]
2608 Draft20(crate::draft20::data_stream::FetchObject),
2609}
2610
2611#[derive(Debug, Clone)]
2620pub struct AnyFetchFrame {
2621 pub meta: AnyFetchObjectMeta,
2624 draft: DraftVersion,
2625 shape: FetchFrameShape,
2626}
2627
2628impl AnyFetchFrame {
2629 #[must_use]
2635 pub fn draft(&self) -> DraftVersion {
2636 self.draft
2637 }
2638
2639 #[cfg(any(
2641 feature = "draft07",
2642 feature = "draft08",
2643 feature = "draft09",
2644 feature = "draft10",
2645 feature = "draft11",
2646 feature = "draft12",
2647 feature = "draft13",
2648 feature = "draft14"
2649 ))]
2650 fn absolute(draft: DraftVersion, meta: AnyFetchObjectMeta) -> Self {
2651 Self { meta, draft, shape: FetchFrameShape::Absolute }
2652 }
2653}
2654
2655#[derive(Debug, Clone, Copy, PartialEq, Eq)]
2665pub enum FetchReemit {
2666 Unchanged,
2673 Reframed {
2677 framing_bytes_before: usize,
2679 framing_bytes_after: usize,
2681 },
2682}
2683
2684#[derive(Debug, Clone)]
2686enum FetchWriterState {
2687 #[cfg(any(
2691 feature = "draft07",
2692 feature = "draft08",
2693 feature = "draft09",
2694 feature = "draft10",
2695 feature = "draft11",
2696 feature = "draft12",
2697 feature = "draft13",
2698 feature = "draft14"
2699 ))]
2700 Absolute(DraftVersion),
2701 #[cfg(feature = "draft15")]
2702 Draft15(crate::draft15::data_stream::FetchObjectWriter),
2703 #[cfg(feature = "draft16")]
2704 Draft16(crate::draft16::data_stream::FetchObjectWriter),
2705 #[cfg(feature = "draft17")]
2706 Draft17(crate::draft17::data_stream::FetchObjectWriter),
2707 #[cfg(feature = "draft18")]
2708 Draft18(crate::draft18::data_stream::FetchObjectWriter),
2709 #[cfg(feature = "draft19")]
2710 Draft19(crate::draft19::data_stream::FetchObjectWriter),
2711 #[cfg(feature = "draft20")]
2712 Draft20(crate::draft20::data_stream::FetchObjectWriter),
2713}
2714
2715#[derive(Debug, Clone)]
2749pub struct AnyFetchObjectWriter {
2750 state: FetchWriterState,
2751}
2752
2753impl AnyFetchObjectWriter {
2754 #[allow(unused_variables, unreachable_code)]
2762 pub fn new(
2763 header: &AnyFetchHeader,
2764 group_order: AnyFetchGroupOrder,
2765 ) -> Result<Self, CodecError> {
2766 let state = match header {
2767 #[cfg(feature = "draft07")]
2768 AnyFetchHeader::Draft07(_) => FetchWriterState::Absolute(DraftVersion::Draft07),
2769 #[cfg(feature = "draft08")]
2770 AnyFetchHeader::Draft08(_) => FetchWriterState::Absolute(DraftVersion::Draft08),
2771 #[cfg(feature = "draft09")]
2772 AnyFetchHeader::Draft09(_) => FetchWriterState::Absolute(DraftVersion::Draft09),
2773 #[cfg(feature = "draft10")]
2774 AnyFetchHeader::Draft10(_) => FetchWriterState::Absolute(DraftVersion::Draft10),
2775 #[cfg(feature = "draft11")]
2776 AnyFetchHeader::Draft11(_) => FetchWriterState::Absolute(DraftVersion::Draft11),
2777 #[cfg(feature = "draft12")]
2778 AnyFetchHeader::Draft12(_) => FetchWriterState::Absolute(DraftVersion::Draft12),
2779 #[cfg(feature = "draft13")]
2780 AnyFetchHeader::Draft13(_) => FetchWriterState::Absolute(DraftVersion::Draft13),
2781 #[cfg(feature = "draft14")]
2782 AnyFetchHeader::Draft14(_) => FetchWriterState::Absolute(DraftVersion::Draft14),
2783 #[cfg(feature = "draft15")]
2784 AnyFetchHeader::Draft15(_) => {
2785 FetchWriterState::Draft15(crate::draft15::data_stream::FetchObjectWriter::new())
2786 }
2787 #[cfg(feature = "draft16")]
2788 AnyFetchHeader::Draft16(_) => {
2789 FetchWriterState::Draft16(crate::draft16::data_stream::FetchObjectWriter::new())
2790 }
2791 #[cfg(feature = "draft17")]
2792 AnyFetchHeader::Draft17(_) => {
2793 FetchWriterState::Draft17(crate::draft17::data_stream::FetchObjectWriter::new())
2794 }
2795 #[cfg(feature = "draft18")]
2796 AnyFetchHeader::Draft18(_) => FetchWriterState::Draft18(
2797 crate::draft18::data_stream::FetchObjectWriter::new(match group_order {
2798 AnyFetchGroupOrder::Ascending => {
2799 crate::draft18::data_stream::GroupOrder::Ascending
2800 }
2801 AnyFetchGroupOrder::Descending => {
2802 crate::draft18::data_stream::GroupOrder::Descending
2803 }
2804 }),
2805 ),
2806 #[cfg(feature = "draft19")]
2807 AnyFetchHeader::Draft19(_) => FetchWriterState::Draft19(
2808 crate::draft19::data_stream::FetchObjectWriter::new(match group_order {
2809 AnyFetchGroupOrder::Ascending => {
2810 crate::draft19::data_stream::GroupOrder::Ascending
2811 }
2812 AnyFetchGroupOrder::Descending => {
2813 crate::draft19::data_stream::GroupOrder::Descending
2814 }
2815 }),
2816 ),
2817 #[cfg(feature = "draft20")]
2818 AnyFetchHeader::Draft20(_) => FetchWriterState::Draft20(
2819 crate::draft20::data_stream::FetchObjectWriter::new(match group_order {
2820 AnyFetchGroupOrder::Ascending => {
2821 crate::draft20::data_stream::GroupOrder::Ascending
2822 }
2823 AnyFetchGroupOrder::Descending => {
2824 crate::draft20::data_stream::GroupOrder::Descending
2825 }
2826 }),
2827 ),
2828 #[allow(unreachable_patterns)]
2829 _ => {
2830 return Err(CodecError::UnsupportedDraft(format!(
2831 "draft {:?} not enabled via feature flag",
2832 header.draft()
2833 )));
2834 }
2835 };
2836 Ok(Self { state })
2837 }
2838
2839 #[must_use]
2841 #[allow(unreachable_code)]
2842 pub fn draft(&self) -> DraftVersion {
2843 match &self.state {
2844 #[cfg(any(
2845 feature = "draft07",
2846 feature = "draft08",
2847 feature = "draft09",
2848 feature = "draft10",
2849 feature = "draft11",
2850 feature = "draft12",
2851 feature = "draft13",
2852 feature = "draft14"
2853 ))]
2854 FetchWriterState::Absolute(draft) => *draft,
2855 #[cfg(feature = "draft15")]
2856 FetchWriterState::Draft15(_) => DraftVersion::Draft15,
2857 #[cfg(feature = "draft16")]
2858 FetchWriterState::Draft16(_) => DraftVersion::Draft16,
2859 #[cfg(feature = "draft17")]
2860 FetchWriterState::Draft17(_) => DraftVersion::Draft17,
2861 #[cfg(feature = "draft18")]
2862 FetchWriterState::Draft18(_) => DraftVersion::Draft18,
2863 #[cfg(feature = "draft19")]
2864 FetchWriterState::Draft19(_) => DraftVersion::Draft19,
2865 #[cfg(feature = "draft20")]
2866 FetchWriterState::Draft20(_) => DraftVersion::Draft20,
2867 #[allow(unreachable_patterns)]
2868 _ => unreachable!("AnyFetchObjectWriter has no enabled variants"),
2869 }
2870 }
2871
2872 #[allow(unused_variables)]
2904 pub fn reemit_object(
2905 &mut self,
2906 frame: &AnyFetchFrame,
2907 raw: &[u8],
2908 out: &mut impl BufMut,
2909 ) -> Result<FetchReemit, CodecError> {
2910 if frame.draft != self.draft() {
2911 return Err(CodecError::UnsupportedDraft(format!(
2912 "a draft {:?} fetch frame cannot be written onto a draft {:?} stream",
2913 frame.draft,
2914 self.draft()
2915 )));
2916 }
2917
2918 let framing_len = frame.meta.wire_len.saturating_sub(frame.meta.payload_length);
2919 let framing_len = usize::try_from(framing_len).map_err(|_| CodecError::InvalidField)?;
2920 if framing_len > raw.len() {
2921 return Err(CodecError::InvalidField);
2922 }
2923 let (framing, rest) = raw.split_at(framing_len);
2924
2925 match (&mut self.state, &frame.shape) {
2926 #[cfg(any(
2929 feature = "draft07",
2930 feature = "draft08",
2931 feature = "draft09",
2932 feature = "draft10",
2933 feature = "draft11",
2934 feature = "draft12",
2935 feature = "draft13",
2936 feature = "draft14"
2937 ))]
2938 (FetchWriterState::Absolute(_), FetchFrameShape::Absolute) => {
2939 Ok(FetchReemit::Unchanged)
2940 }
2941 #[cfg(feature = "draft15")]
2942 (FetchWriterState::Draft15(writer), FetchFrameShape::Draft15(original)) => {
2943 let reframed = writer.header_for(original)?;
2944 if reframed == *original {
2945 writer.advance(original);
2946 return Ok(FetchReemit::Unchanged);
2947 }
2948 let mut encoded = Vec::with_capacity(framing.len() + 16);
2949 reframed.encode(&mut encoded)?;
2950 writer.advance(&reframed);
2951 Ok(put_reframed(&encoded, framing.len(), rest, out))
2952 }
2953 #[cfg(feature = "draft16")]
2954 (FetchWriterState::Draft16(writer), FetchFrameShape::Draft16(original, location)) => {
2955 let reframed = writer.header_for(original, location)?;
2956 if reframed == *original {
2957 writer.advance(original, location);
2958 return Ok(FetchReemit::Unchanged);
2959 }
2960 let mut encoded = Vec::with_capacity(framing.len() + 16);
2961 reframed.encode(&mut encoded)?;
2962 writer.advance(&reframed, location);
2963 Ok(put_reframed(&encoded, framing.len(), rest, out))
2964 }
2965 #[cfg(feature = "draft17")]
2966 (FetchWriterState::Draft17(writer), FetchFrameShape::Draft17(original)) => {
2967 let reframed = writer.header_for(original)?;
2968 if reframed == original.header {
2969 writer.advance(original);
2970 return Ok(FetchReemit::Unchanged);
2971 }
2972 let mut encoded = Vec::with_capacity(framing.len() + 16);
2973 reframed.encode(&mut encoded)?;
2974 writer.advance(original);
2975 Ok(put_reframed(&encoded, framing.len(), rest, out))
2976 }
2977 #[cfg(feature = "draft18")]
2978 (FetchWriterState::Draft18(writer), FetchFrameShape::Draft18(original)) => {
2979 let reframed = writer.header_for(original)?;
2980 if reframed == original.header {
2981 writer.advance(original);
2982 return Ok(FetchReemit::Unchanged);
2983 }
2984 let mut encoded = Vec::with_capacity(framing.len() + 16);
2985 reframed.encode(&mut encoded)?;
2986 writer.advance(original);
2987 Ok(put_reframed(&encoded, framing.len(), rest, out))
2988 }
2989 #[cfg(feature = "draft19")]
2990 (FetchWriterState::Draft19(writer), FetchFrameShape::Draft19(original)) => {
2991 let reframed = writer.header_for(original)?;
2992 if reframed == original.header {
2993 writer.advance(original);
2994 return Ok(FetchReemit::Unchanged);
2995 }
2996 let mut encoded = Vec::with_capacity(framing.len() + 16);
2997 reframed.encode(&mut encoded)?;
2998 writer.advance(original);
2999 Ok(put_reframed(&encoded, framing.len(), rest, out))
3000 }
3001 #[cfg(feature = "draft20")]
3002 (FetchWriterState::Draft20(writer), FetchFrameShape::Draft20(original)) => {
3003 let reframed = writer.header_for(original)?;
3004 if reframed == original.header {
3005 writer.advance(original);
3006 return Ok(FetchReemit::Unchanged);
3007 }
3008 let mut encoded = Vec::with_capacity(framing.len() + 16);
3009 reframed.encode(&mut encoded)?;
3010 writer.advance(original);
3011 Ok(put_reframed(&encoded, framing.len(), rest, out))
3012 }
3013 #[allow(unreachable_patterns)]
3016 _ => Err(CodecError::UnsupportedDraft(format!(
3017 "no fetch writer for draft {:?}",
3018 frame.draft
3019 ))),
3020 }
3021 }
3022}
3023
3024#[cfg(any(
3027 feature = "draft15",
3028 feature = "draft16",
3029 feature = "draft17",
3030 feature = "draft18",
3031 feature = "draft19",
3032 feature = "draft20"
3033))]
3034fn put_reframed(
3035 encoded: &[u8],
3036 framing_bytes_before: usize,
3037 rest: &[u8],
3038 out: &mut impl BufMut,
3039) -> FetchReemit {
3040 out.put_slice(encoded);
3041 out.put_slice(rest);
3042 FetchReemit::Reframed { framing_bytes_before, framing_bytes_after: encoded.len() }
3043}
3044
3045macro_rules! subgroup_extensions_fn {
3051 ($name:ident, $feat:literal, $draft:ident) => {
3052 #[cfg(feature = $feat)]
3053 fn $name(header: &crate::$draft::data_stream::SubgroupHeader) -> Result<bool, CodecError> {
3054 if !header.stream_type.is_subgroup() {
3055 return Err(CodecError::InvalidField);
3056 }
3057 Ok(header.stream_type.has_extensions())
3058 }
3059 };
3060}
3061
3062subgroup_extensions_fn!(subgroup_extensions_11, "draft11", draft11);
3063subgroup_extensions_fn!(subgroup_extensions_12, "draft12", draft12);
3064subgroup_extensions_fn!(subgroup_extensions_13, "draft13", draft13);