1use crate::error::CodecError;
2use crate::kvp::KeyValuePair;
3use crate::types::read_bytes;
4use crate::types::*;
5use crate::types::{check_location_range, check_open_ended_group_range};
6use crate::varint::VarInt;
7use bytes::{Buf, BufMut};
8
9#[derive(Debug, Clone, Copy, PartialEq, Eq)]
11#[repr(u64)]
12pub enum MessageType {
13 SubscribeUpdate = 0x02,
15 Subscribe = 0x03,
17 SubscribeOk = 0x04,
19 SubscribeError = 0x05,
21 Announce = 0x06,
23 AnnounceOk = 0x07,
25 AnnounceError = 0x08,
27 Unannounce = 0x09,
29 Unsubscribe = 0x0A,
31 SubscribeDone = 0x0B,
33 AnnounceCancel = 0x0C,
35 TrackStatusRequest = 0x0D,
37 TrackStatus = 0x0E,
39 GoAway = 0x10,
41 SubscribeAnnounces = 0x11,
43 SubscribeAnnouncesOk = 0x12,
45 SubscribeAnnouncesError = 0x13,
47 UnsubscribeAnnounces = 0x14,
49 MaxSubscribeId = 0x15,
51 Fetch = 0x16,
53 FetchCancel = 0x17,
55 FetchOk = 0x18,
57 FetchError = 0x19,
59 ClientSetup = 0x40,
61 ServerSetup = 0x41,
63}
64
65impl MessageType {
66 pub fn from_id(id: u64) -> Option<Self> {
68 match id {
69 0x02 => Some(MessageType::SubscribeUpdate),
70 0x03 => Some(MessageType::Subscribe),
71 0x04 => Some(MessageType::SubscribeOk),
72 0x05 => Some(MessageType::SubscribeError),
73 0x06 => Some(MessageType::Announce),
74 0x07 => Some(MessageType::AnnounceOk),
75 0x08 => Some(MessageType::AnnounceError),
76 0x09 => Some(MessageType::Unannounce),
77 0x0A => Some(MessageType::Unsubscribe),
78 0x0B => Some(MessageType::SubscribeDone),
79 0x0C => Some(MessageType::AnnounceCancel),
80 0x0D => Some(MessageType::TrackStatusRequest),
81 0x0E => Some(MessageType::TrackStatus),
82 0x10 => Some(MessageType::GoAway),
83 0x11 => Some(MessageType::SubscribeAnnounces),
84 0x12 => Some(MessageType::SubscribeAnnouncesOk),
85 0x13 => Some(MessageType::SubscribeAnnouncesError),
86 0x14 => Some(MessageType::UnsubscribeAnnounces),
87 0x15 => Some(MessageType::MaxSubscribeId),
88 0x16 => Some(MessageType::Fetch),
89 0x17 => Some(MessageType::FetchCancel),
90 0x18 => Some(MessageType::FetchOk),
91 0x19 => Some(MessageType::FetchError),
92 0x40 => Some(MessageType::ClientSetup),
93 0x41 => Some(MessageType::ServerSetup),
94 _ => None,
95 }
96 }
97
98 pub fn id(&self) -> u64 {
100 *self as u64
101 }
102}
103
104#[derive(Debug, Clone, PartialEq, Eq)]
110pub struct ClientSetup {
111 pub supported_versions: Vec<VarInt>,
113 pub parameters: Vec<KeyValuePair>,
115}
116
117#[derive(Debug, Clone, PartialEq, Eq)]
119pub struct ServerSetup {
120 pub selected_version: VarInt,
122 pub parameters: Vec<KeyValuePair>,
124}
125
126#[derive(Debug, Clone, PartialEq, Eq)]
128pub struct GoAway {
129 pub new_session_uri: Vec<u8>,
131}
132
133#[derive(Debug, Clone, PartialEq, Eq)]
135pub struct MaxSubscribeId {
136 pub subscribe_id: VarInt,
138}
139
140#[derive(Debug, Clone, PartialEq, Eq)]
146pub struct Subscribe {
147 pub subscribe_id: VarInt,
149 pub track_alias: VarInt,
151 pub track_namespace: TrackNamespace,
153 pub track_name: Vec<u8>,
155 pub subscriber_priority: u8,
157 pub group_order: GroupOrder,
159 pub filter_type: FilterType,
161 pub start_location: Option<Location>,
163 pub end_group: Option<VarInt>,
165 pub end_object: Option<VarInt>,
167 pub parameters: Vec<KeyValuePair>,
169}
170
171#[derive(Debug, Clone, PartialEq, Eq)]
173pub struct SubscribeOk {
174 pub subscribe_id: VarInt,
176 pub expires: VarInt,
178 pub group_order: GroupOrder,
180 pub content_exists: ContentExists,
182 pub largest_group_id: Option<VarInt>,
184 pub largest_object_id: Option<VarInt>,
186 pub parameters: Vec<KeyValuePair>,
188}
189
190#[derive(Debug, Clone, PartialEq, Eq)]
192pub struct SubscribeError {
193 pub subscribe_id: VarInt,
195 pub error_code: VarInt,
197 pub reason_phrase: Vec<u8>,
199 pub track_alias: VarInt,
201}
202
203#[derive(Debug, Clone, PartialEq, Eq)]
205pub struct SubscribeUpdate {
206 pub subscribe_id: VarInt,
208 pub start_group: VarInt,
210 pub start_object: VarInt,
212 pub end_group: VarInt,
214 pub end_object: VarInt,
216 pub subscriber_priority: u8,
218 pub parameters: Vec<KeyValuePair>,
220}
221
222#[derive(Debug, Clone, PartialEq, Eq)]
224pub struct SubscribeDone {
225 pub subscribe_id: VarInt,
227 pub status_code: VarInt,
229 pub reason_phrase: Vec<u8>,
231 pub content_exists: ContentExists,
233 pub final_group: Option<VarInt>,
235 pub final_object: Option<VarInt>,
237}
238
239#[derive(Debug, Clone, PartialEq, Eq)]
241pub struct Unsubscribe {
242 pub subscribe_id: VarInt,
244}
245
246#[derive(Debug, Clone, PartialEq, Eq)]
252pub struct Announce {
253 pub track_namespace: TrackNamespace,
255 pub parameters: Vec<KeyValuePair>,
257}
258
259#[derive(Debug, Clone, PartialEq, Eq)]
261pub struct AnnounceOk {
262 pub track_namespace: TrackNamespace,
264}
265
266#[derive(Debug, Clone, PartialEq, Eq)]
268pub struct AnnounceError {
269 pub track_namespace: TrackNamespace,
271 pub error_code: VarInt,
273 pub reason_phrase: Vec<u8>,
275}
276
277#[derive(Debug, Clone, PartialEq, Eq)]
279pub struct AnnounceCancel {
280 pub track_namespace: TrackNamespace,
282 pub error_code: VarInt,
284 pub reason_phrase: Vec<u8>,
286}
287
288#[derive(Debug, Clone, PartialEq, Eq)]
290pub struct Unannounce {
291 pub track_namespace: TrackNamespace,
293}
294
295#[derive(Debug, Clone, PartialEq, Eq)]
301pub struct SubscribeAnnounces {
302 pub track_namespace_prefix: TrackNamespace,
304 pub parameters: Vec<KeyValuePair>,
306}
307
308#[derive(Debug, Clone, PartialEq, Eq)]
310pub struct SubscribeAnnouncesOk {
311 pub track_namespace_prefix: TrackNamespace,
313}
314
315#[derive(Debug, Clone, PartialEq, Eq)]
317pub struct SubscribeAnnouncesError {
318 pub track_namespace_prefix: TrackNamespace,
320 pub error_code: VarInt,
322 pub reason_phrase: Vec<u8>,
324}
325
326#[derive(Debug, Clone, PartialEq, Eq)]
328pub struct UnsubscribeAnnounces {
329 pub track_namespace_prefix: TrackNamespace,
331}
332
333#[derive(Debug, Clone, PartialEq, Eq)]
339pub struct TrackStatusRequest {
340 pub track_namespace: TrackNamespace,
342 pub track_name: Vec<u8>,
344}
345
346#[derive(Debug, Clone, PartialEq, Eq)]
348pub struct TrackStatus {
349 pub track_namespace: TrackNamespace,
351 pub track_name: Vec<u8>,
353 pub status_code: VarInt,
355 pub last_group_id: VarInt,
357 pub last_object_id: VarInt,
359}
360
361#[derive(Debug, Clone, PartialEq, Eq)]
367pub struct Fetch {
368 pub subscribe_id: VarInt,
370 pub track_namespace: TrackNamespace,
372 pub track_name: Vec<u8>,
374 pub subscriber_priority: u8,
376 pub group_order: GroupOrder,
378 pub start_group: VarInt,
380 pub start_object: VarInt,
382 pub end_group: VarInt,
384 pub end_object: VarInt,
386 pub parameters: Vec<KeyValuePair>,
388}
389
390#[derive(Debug, Clone, PartialEq, Eq)]
392pub struct FetchOk {
393 pub subscribe_id: VarInt,
395 pub group_order: GroupOrder,
397 pub end_of_track: u8,
399 pub largest_group_id: Option<VarInt>,
401 pub largest_object_id: Option<VarInt>,
403 pub parameters: Vec<KeyValuePair>,
405}
406
407#[derive(Debug, Clone, PartialEq, Eq)]
409pub struct FetchError {
410 pub subscribe_id: VarInt,
412 pub error_code: VarInt,
414 pub reason_phrase: Vec<u8>,
416}
417
418#[derive(Debug, Clone, PartialEq, Eq)]
420pub struct FetchCancel {
421 pub subscribe_id: VarInt,
423}
424
425fn read_group_order_response(buf: &mut impl Buf) -> Result<GroupOrder, CodecError> {
440 if !buf.has_remaining() {
441 return Err(CodecError::UnexpectedEnd);
442 }
443 match GroupOrder::from_u8(buf.get_u8()).ok_or(CodecError::InvalidField)? {
444 GroupOrder::Publisher => Err(CodecError::InvalidField),
445 order => Ok(order),
446 }
447}
448
449fn check_track_status(
460 status_code: VarInt,
461 last_group_id: VarInt,
462 last_object_id: VarInt,
463) -> Result<(), CodecError> {
464 let code = crate::draft07::error_codes::TrackStatusCode::from_u64(status_code.into_inner())
465 .ok_or(CodecError::InvalidField)?;
466 if code.requires_zero_location()
467 && (last_group_id.into_inner() != 0 || last_object_id.into_inner() != 0)
468 {
469 return Err(CodecError::InvalidField);
470 }
471 Ok(())
472}
473
474fn check_discriminators(message: &ControlMessage) -> Result<(), CodecError> {
485 match message {
486 ControlMessage::Subscribe(m) => {
487 let wants_start =
488 matches!(m.filter_type, FilterType::AbsoluteStart | FilterType::AbsoluteRange);
489 if wants_start != m.start_location.is_some() {
490 return Err(CodecError::InvalidField);
491 }
492 let wants_end = m.filter_type == FilterType::AbsoluteRange;
493 if wants_end != m.end_group.is_some() || wants_end != m.end_object.is_some() {
494 return Err(CodecError::InvalidField);
495 }
496 Ok(())
497 }
498 ControlMessage::SubscribeOk(m) => {
499 let has = m.content_exists == ContentExists::HasLargestLocation;
500 if has != m.largest_group_id.is_some() || has != m.largest_object_id.is_some() {
501 return Err(CodecError::InvalidField);
502 }
503 Ok(())
504 }
505 _ => Ok(()),
506 }
507}
508
509fn check_group_order(message: &ControlMessage) -> Result<(), CodecError> {
514 let order = match message {
515 ControlMessage::SubscribeOk(m) => m.group_order,
516 ControlMessage::FetchOk(m) => m.group_order,
517 _ => return Ok(()),
518 };
519 if order == GroupOrder::Publisher {
520 return Err(CodecError::InvalidField);
521 }
522 Ok(())
523}
524
525#[derive(Debug, Clone, PartialEq, Eq)]
527pub enum ControlMessage {
528 ClientSetup(ClientSetup),
530 ServerSetup(ServerSetup),
532 GoAway(GoAway),
534 MaxSubscribeId(MaxSubscribeId),
536 Subscribe(Subscribe),
538 SubscribeOk(SubscribeOk),
540 SubscribeError(SubscribeError),
542 SubscribeUpdate(SubscribeUpdate),
544 SubscribeDone(SubscribeDone),
546 Unsubscribe(Unsubscribe),
548 Announce(Announce),
550 AnnounceOk(AnnounceOk),
552 AnnounceError(AnnounceError),
554 AnnounceCancel(AnnounceCancel),
556 Unannounce(Unannounce),
558 SubscribeAnnounces(SubscribeAnnounces),
560 SubscribeAnnouncesOk(SubscribeAnnouncesOk),
562 SubscribeAnnouncesError(SubscribeAnnouncesError),
564 UnsubscribeAnnounces(UnsubscribeAnnounces),
566 TrackStatusRequest(TrackStatusRequest),
568 TrackStatus(TrackStatus),
570 Fetch(Fetch),
572 FetchOk(FetchOk),
574 FetchError(FetchError),
576 FetchCancel(FetchCancel),
578}
579
580fn check_ranges(message: &ControlMessage) -> Result<(), CodecError> {
594 match message {
595 ControlMessage::Subscribe(m) => match (&m.start_location, &m.end_group, &m.end_object) {
596 (Some(start), Some(end_group), Some(end_object)) => check_location_range(
597 start.group.into_inner(),
598 start.object.into_inner(),
599 end_group.into_inner(),
600 end_object.into_inner(),
601 ),
602 _ => Ok(()),
603 },
604 ControlMessage::SubscribeUpdate(m) => {
605 check_open_ended_group_range(m.start_group.into_inner(), m.end_group.into_inner())
606 }
607 ControlMessage::Fetch(m) => check_location_range(
608 m.start_group.into_inner(),
609 m.start_object.into_inner(),
610 m.end_group.into_inner(),
611 m.end_object.into_inner(),
612 ),
613 _ => Ok(()),
614 }
615}
616
617fn check_no_duplicate_parameters(parameters: &[KeyValuePair]) -> Result<(), CodecError> {
628 for (i, parameter) in parameters.iter().enumerate() {
629 if parameters[..i].iter().any(|earlier| earlier.key == parameter.key) {
630 return Err(CodecError::DuplicateParameter(parameter.key.into_inner()));
631 }
632 }
633 Ok(())
634}
635
636const SETUP_VARINT_PARAMETERS: &[u64] = &[0x00, 0x02];
643
644const VERSION_VARINT_PARAMETERS: &[u64] = &[0x03, 0x04];
657
658fn check_parameter_lengths(
676 parameters: &[KeyValuePair],
677 varint_typed: &[u64],
678) -> Result<(), CodecError> {
679 for parameter in parameters {
680 let key = parameter.key.into_inner();
681 if !varint_typed.contains(&key) {
682 continue;
683 }
684 if let crate::kvp::KvpValue::Bytes(bytes) = ¶meter.value {
685 let mut cursor = &bytes[..];
686 let one_varint = VarInt::decode(&mut cursor).is_ok() && !cursor.has_remaining();
687 if !one_varint {
688 return Err(CodecError::ParameterLengthMismatch(key));
689 }
690 }
691 }
692 Ok(())
693}
694
695fn decode_parameters(buf: &mut impl Buf) -> Result<Vec<KeyValuePair>, CodecError> {
698 let parameters = KeyValuePair::decode_list_d07(buf)?;
699 check_no_duplicate_parameters(¶meters)?;
700 check_parameter_lengths(¶meters, VERSION_VARINT_PARAMETERS)?;
701 Ok(parameters)
702}
703
704fn encode_parameters(parameters: &[KeyValuePair], buf: &mut impl BufMut) -> Result<(), CodecError> {
707 check_no_duplicate_parameters(parameters)?;
708 check_parameter_lengths(parameters, VERSION_VARINT_PARAMETERS)?;
709 KeyValuePair::encode_list_d07(parameters, buf);
710 Ok(())
711}
712
713fn decode_setup_parameters(buf: &mut impl Buf) -> Result<Vec<KeyValuePair>, CodecError> {
719 let parameters = KeyValuePair::decode_list_d07(buf)?;
720 check_no_duplicate_parameters(¶meters)?;
721 check_parameter_lengths(¶meters, SETUP_VARINT_PARAMETERS)?;
722 Ok(parameters)
723}
724
725fn encode_setup_parameters(
727 parameters: &[KeyValuePair],
728 buf: &mut impl BufMut,
729) -> Result<(), CodecError> {
730 check_no_duplicate_parameters(parameters)?;
731 check_parameter_lengths(parameters, SETUP_VARINT_PARAMETERS)?;
732 KeyValuePair::encode_list_d07(parameters, buf);
733 Ok(())
734}
735
736impl ControlMessage {
737 pub fn message_type(&self) -> MessageType {
739 match self {
740 ControlMessage::ClientSetup(_) => MessageType::ClientSetup,
741 ControlMessage::ServerSetup(_) => MessageType::ServerSetup,
742 ControlMessage::GoAway(_) => MessageType::GoAway,
743 ControlMessage::MaxSubscribeId(_) => MessageType::MaxSubscribeId,
744 ControlMessage::Subscribe(_) => MessageType::Subscribe,
745 ControlMessage::SubscribeOk(_) => MessageType::SubscribeOk,
746 ControlMessage::SubscribeError(_) => MessageType::SubscribeError,
747 ControlMessage::SubscribeUpdate(_) => MessageType::SubscribeUpdate,
748 ControlMessage::SubscribeDone(_) => MessageType::SubscribeDone,
749 ControlMessage::Unsubscribe(_) => MessageType::Unsubscribe,
750 ControlMessage::Announce(_) => MessageType::Announce,
751 ControlMessage::AnnounceOk(_) => MessageType::AnnounceOk,
752 ControlMessage::AnnounceError(_) => MessageType::AnnounceError,
753 ControlMessage::AnnounceCancel(_) => MessageType::AnnounceCancel,
754 ControlMessage::Unannounce(_) => MessageType::Unannounce,
755 ControlMessage::SubscribeAnnounces(_) => MessageType::SubscribeAnnounces,
756 ControlMessage::SubscribeAnnouncesOk(_) => MessageType::SubscribeAnnouncesOk,
757 ControlMessage::SubscribeAnnouncesError(_) => MessageType::SubscribeAnnouncesError,
758 ControlMessage::UnsubscribeAnnounces(_) => MessageType::UnsubscribeAnnounces,
759 ControlMessage::TrackStatusRequest(_) => MessageType::TrackStatusRequest,
760 ControlMessage::TrackStatus(_) => MessageType::TrackStatus,
761 ControlMessage::Fetch(_) => MessageType::Fetch,
762 ControlMessage::FetchOk(_) => MessageType::FetchOk,
763 ControlMessage::FetchError(_) => MessageType::FetchError,
764 ControlMessage::FetchCancel(_) => MessageType::FetchCancel,
765 }
766 }
767
768 pub fn encode(&self, buf: &mut impl BufMut) -> Result<(), CodecError> {
770 check_discriminators(self)?;
771 check_group_order(self)?;
772 check_ranges(self)?;
773 let mut payload = Vec::with_capacity(256);
774 self.encode_payload(&mut payload)?;
775
776 VarInt::from_usize(self.message_type().id() as usize).encode(buf);
784 VarInt::from_usize(payload.len()).encode(buf);
785 buf.put_slice(&payload);
786 Ok(())
787 }
788
789 pub fn decode(buf: &mut impl Buf) -> Result<Self, CodecError> {
791 let type_id = VarInt::decode(buf)?.into_inner();
792 let msg_type =
793 MessageType::from_id(type_id).ok_or(CodecError::UnknownMessageType(type_id))?;
794 let payload_len = VarInt::decode(buf)?.into_inner() as usize;
795 if buf.remaining() < payload_len {
796 return Err(CodecError::UnexpectedEnd);
797 }
798 let payload_bytes = buf.copy_to_bytes(payload_len);
799 let mut payload = &payload_bytes[..];
800 let msg = match Self::decode_payload(msg_type, &mut payload) {
801 Ok(msg) => msg,
802 Err(
808 CodecError::UnexpectedEnd
809 | CodecError::Kvp(crate::kvp::KvpError::UnexpectedEnd)
810 | CodecError::Kvp(crate::kvp::KvpError::VarInt(
811 crate::varint::VarIntError::UnexpectedEnd,
812 ))
813 | CodecError::VarInt(crate::varint::VarIntError::UnexpectedEnd),
814 ) => {
815 return Err(CodecError::ControlMessageLengthMismatch {
816 declared: payload_len,
817 detail: "its fields ran past the end",
818 });
819 }
820 Err(e) => return Err(e),
821 };
822 check_ranges(&msg)?;
823 if payload.has_remaining() {
828 return Err(CodecError::ControlMessageLengthMismatch {
829 declared: payload_len,
830 detail: "its fields left bytes unread",
831 });
832 }
833 Ok(msg)
834 }
835
836 fn encode_payload(&self, buf: &mut impl BufMut) -> Result<(), CodecError> {
837 match self {
838 ControlMessage::ClientSetup(m) => {
839 VarInt::from_usize(m.supported_versions.len()).encode(buf);
840 for v in &m.supported_versions {
841 v.encode(buf);
842 }
843 encode_setup_parameters(&m.parameters, buf)?;
844 }
845 ControlMessage::ServerSetup(m) => {
846 m.selected_version.encode(buf);
847 encode_setup_parameters(&m.parameters, buf)?;
848 }
849 ControlMessage::GoAway(m) => {
850 VarInt::from_usize(m.new_session_uri.len()).encode(buf);
851 buf.put_slice(&m.new_session_uri);
852 }
853 ControlMessage::MaxSubscribeId(m) => {
854 m.subscribe_id.encode(buf);
855 }
856 ControlMessage::Subscribe(m) => {
857 m.subscribe_id.encode(buf);
858 m.track_alias.encode(buf);
859 m.track_namespace.validate(TrackNamespaceRules::for_draft(7))?;
860 m.track_namespace.encode(buf);
861 VarInt::from_usize(m.track_name.len()).encode(buf);
862 buf.put_slice(&m.track_name);
863 buf.put_u8(m.subscriber_priority);
864 buf.put_u8(m.group_order as u8);
865 buf.put_u8(m.filter_type as u8);
866 if let Some(loc) = &m.start_location {
867 loc.encode(buf);
868 }
869 if let Some(eg) = &m.end_group {
870 eg.encode(buf);
871 }
872 if let Some(eo) = &m.end_object {
873 eo.encode(buf);
874 }
875 encode_parameters(&m.parameters, buf)?;
876 }
877 ControlMessage::SubscribeOk(m) => {
878 m.subscribe_id.encode(buf);
879 m.expires.encode(buf);
880 buf.put_u8(m.group_order as u8);
881 buf.put_u8(m.content_exists as u8);
882 if let Some(gid) = &m.largest_group_id {
883 gid.encode(buf);
884 }
885 if let Some(oid) = &m.largest_object_id {
886 oid.encode(buf);
887 }
888 encode_parameters(&m.parameters, buf)?;
889 }
890 ControlMessage::SubscribeError(m) => {
891 m.subscribe_id.encode(buf);
892 m.error_code.encode(buf);
893 VarInt::from_usize(m.reason_phrase.len()).encode(buf);
894 buf.put_slice(&m.reason_phrase);
895 m.track_alias.encode(buf);
896 }
897 ControlMessage::SubscribeUpdate(m) => {
898 m.subscribe_id.encode(buf);
899 m.start_group.encode(buf);
900 m.start_object.encode(buf);
901 m.end_group.encode(buf);
902 m.end_object.encode(buf);
903 buf.put_u8(m.subscriber_priority);
904 encode_parameters(&m.parameters, buf)?;
905 }
906 ControlMessage::SubscribeDone(m) => {
907 m.subscribe_id.encode(buf);
908 m.status_code.encode(buf);
909 VarInt::from_usize(m.reason_phrase.len()).encode(buf);
910 buf.put_slice(&m.reason_phrase);
911 buf.put_u8(m.content_exists as u8);
912 if let Some(fg) = &m.final_group {
913 fg.encode(buf);
914 }
915 if let Some(fo) = &m.final_object {
916 fo.encode(buf);
917 }
918 }
919 ControlMessage::Unsubscribe(m) => {
920 m.subscribe_id.encode(buf);
921 }
922 ControlMessage::Announce(m) => {
923 m.track_namespace.validate(TrackNamespaceRules::for_draft(7))?;
924 m.track_namespace.encode(buf);
925 encode_parameters(&m.parameters, buf)?;
926 }
927 ControlMessage::AnnounceOk(m) => {
928 m.track_namespace.validate(TrackNamespaceRules::for_draft(7))?;
929 m.track_namespace.encode(buf);
930 }
931 ControlMessage::AnnounceError(m) => {
932 m.track_namespace.validate(TrackNamespaceRules::for_draft(7))?;
933 m.track_namespace.encode(buf);
934 m.error_code.encode(buf);
935 VarInt::from_usize(m.reason_phrase.len()).encode(buf);
936 buf.put_slice(&m.reason_phrase);
937 }
938 ControlMessage::AnnounceCancel(m) => {
939 m.track_namespace.validate(TrackNamespaceRules::for_draft(7))?;
940 m.track_namespace.encode(buf);
941 m.error_code.encode(buf);
942 VarInt::from_usize(m.reason_phrase.len()).encode(buf);
943 buf.put_slice(&m.reason_phrase);
944 }
945 ControlMessage::Unannounce(m) => {
946 m.track_namespace.validate(TrackNamespaceRules::for_draft(7))?;
947 m.track_namespace.encode(buf);
948 }
949 ControlMessage::SubscribeAnnounces(m) => {
950 m.track_namespace_prefix.validate(TrackNamespaceRules::for_draft(7))?;
951 m.track_namespace_prefix.encode(buf);
952 encode_parameters(&m.parameters, buf)?;
953 }
954 ControlMessage::SubscribeAnnouncesOk(m) => {
955 m.track_namespace_prefix.validate(TrackNamespaceRules::for_draft(7))?;
956 m.track_namespace_prefix.encode(buf);
957 }
958 ControlMessage::SubscribeAnnouncesError(m) => {
959 m.track_namespace_prefix.validate(TrackNamespaceRules::for_draft(7))?;
960 m.track_namespace_prefix.encode(buf);
961 m.error_code.encode(buf);
962 VarInt::from_usize(m.reason_phrase.len()).encode(buf);
963 buf.put_slice(&m.reason_phrase);
964 }
965 ControlMessage::UnsubscribeAnnounces(m) => {
966 m.track_namespace_prefix.validate(TrackNamespaceRules::for_draft(7))?;
967 m.track_namespace_prefix.encode(buf);
968 }
969 ControlMessage::TrackStatusRequest(m) => {
970 m.track_namespace.validate(TrackNamespaceRules::for_draft(7))?;
971 m.track_namespace.encode(buf);
972 VarInt::from_usize(m.track_name.len()).encode(buf);
973 buf.put_slice(&m.track_name);
974 }
975 ControlMessage::TrackStatus(m) => {
976 m.track_namespace.validate(TrackNamespaceRules::for_draft(7))?;
977 m.track_namespace.encode(buf);
978 VarInt::from_usize(m.track_name.len()).encode(buf);
979 buf.put_slice(&m.track_name);
980 check_track_status(m.status_code, m.last_group_id, m.last_object_id)?;
981 m.status_code.encode(buf);
982 m.last_group_id.encode(buf);
983 m.last_object_id.encode(buf);
984 }
985 ControlMessage::Fetch(m) => {
986 m.subscribe_id.encode(buf);
987 m.track_namespace.validate(TrackNamespaceRules::for_draft(7))?;
988 m.track_namespace.encode(buf);
989 VarInt::from_usize(m.track_name.len()).encode(buf);
990 buf.put_slice(&m.track_name);
991 buf.put_u8(m.subscriber_priority);
992 buf.put_u8(m.group_order as u8);
993 m.start_group.encode(buf);
994 m.start_object.encode(buf);
995 m.end_group.encode(buf);
996 m.end_object.encode(buf);
997 encode_parameters(&m.parameters, buf)?;
998 }
999 ControlMessage::FetchOk(m) => {
1000 m.subscribe_id.encode(buf);
1001 buf.put_u8(m.group_order as u8);
1002 buf.put_u8(m.end_of_track);
1003 if let Some(gid) = &m.largest_group_id {
1004 gid.encode(buf);
1005 }
1006 if let Some(oid) = &m.largest_object_id {
1007 oid.encode(buf);
1008 }
1009 encode_parameters(&m.parameters, buf)?;
1010 }
1011 ControlMessage::FetchError(m) => {
1012 m.subscribe_id.encode(buf);
1013 m.error_code.encode(buf);
1014 VarInt::from_usize(m.reason_phrase.len()).encode(buf);
1015 buf.put_slice(&m.reason_phrase);
1016 }
1017 ControlMessage::FetchCancel(m) => {
1018 m.subscribe_id.encode(buf);
1019 }
1020 }
1021 Ok(())
1022 }
1023
1024 fn decode_payload(msg_type: MessageType, buf: &mut impl Buf) -> Result<Self, CodecError> {
1025 match msg_type {
1026 MessageType::ClientSetup => {
1027 let num_versions = VarInt::decode(buf)?.into_inner() as usize;
1028 if num_versions == 0 {
1038 return Err(CodecError::InvalidField);
1039 }
1040 let mut supported_versions = crate::types::reserve_bounded(num_versions, buf);
1041 for _ in 0..num_versions {
1042 supported_versions.push(VarInt::decode(buf)?);
1043 }
1044 let parameters = decode_setup_parameters(buf)?;
1045 Ok(ControlMessage::ClientSetup(ClientSetup { supported_versions, parameters }))
1046 }
1047 MessageType::ServerSetup => {
1048 let selected_version = VarInt::decode(buf)?;
1049 let parameters = decode_setup_parameters(buf)?;
1050 Ok(ControlMessage::ServerSetup(ServerSetup { selected_version, parameters }))
1051 }
1052 MessageType::GoAway => {
1053 let uri_len = VarInt::decode(buf)?.into_inner() as usize;
1054 let uri = read_bytes(buf, uri_len)?;
1055 Ok(ControlMessage::GoAway(GoAway { new_session_uri: uri }))
1056 }
1057 MessageType::MaxSubscribeId => {
1058 let subscribe_id = VarInt::decode(buf)?;
1059 Ok(ControlMessage::MaxSubscribeId(MaxSubscribeId { subscribe_id }))
1060 }
1061 MessageType::Subscribe => {
1062 let subscribe_id = VarInt::decode(buf)?;
1063 let track_alias = VarInt::decode(buf)?;
1064 let track_namespace = TrackNamespace::decode(buf)?;
1065 let track_name_len = VarInt::decode(buf)?.into_inner() as usize;
1066 let track_name = read_bytes(buf, track_name_len)?;
1067 if buf.remaining() < 3 {
1068 return Err(CodecError::UnexpectedEnd);
1069 }
1070 let subscriber_priority = buf.get_u8();
1071 let group_order =
1072 GroupOrder::from_u8(buf.get_u8()).ok_or(CodecError::InvalidField)?;
1073 let filter_val = buf.get_u8();
1074 let filter_type = FilterType::from_u8(filter_val)
1075 .ok_or(CodecError::InvalidFilterType(filter_val as u64))?;
1076 let start_location = match filter_type {
1077 FilterType::AbsoluteStart | FilterType::AbsoluteRange => {
1078 Some(Location::decode(buf)?)
1079 }
1080 _ => None,
1081 };
1082 let (end_group, end_object) = match filter_type {
1083 FilterType::AbsoluteRange => {
1084 let eg = VarInt::decode(buf)?;
1085 let eo = VarInt::decode(buf)?;
1086 (Some(eg), Some(eo))
1087 }
1088 _ => (None, None),
1089 };
1090 let parameters = decode_parameters(buf)?;
1091 Ok(ControlMessage::Subscribe(Subscribe {
1092 subscribe_id,
1093 track_alias,
1094 track_namespace,
1095 track_name,
1096 subscriber_priority,
1097 group_order,
1098 filter_type,
1099 start_location,
1100 end_group,
1101 end_object,
1102 parameters,
1103 }))
1104 }
1105 MessageType::SubscribeOk => {
1106 let subscribe_id = VarInt::decode(buf)?;
1107 let expires = VarInt::decode(buf)?;
1108 if buf.remaining() < 2 {
1109 return Err(CodecError::UnexpectedEnd);
1110 }
1111 let group_order = read_group_order_response(buf)?;
1112 let content_exists_val = buf.get_u8();
1113 let content_exists = match content_exists_val {
1114 0 => ContentExists::NoLargestLocation,
1115 1 => ContentExists::HasLargestLocation,
1116 other => return Err(CodecError::InvalidContentExists(other)),
1117 };
1118 let (largest_group_id, largest_object_id) =
1119 if content_exists == ContentExists::HasLargestLocation {
1120 let gid = VarInt::decode(buf)?;
1121 let oid = VarInt::decode(buf)?;
1122 (Some(gid), Some(oid))
1123 } else {
1124 (None, None)
1125 };
1126 let parameters = decode_parameters(buf)?;
1127 Ok(ControlMessage::SubscribeOk(SubscribeOk {
1128 subscribe_id,
1129 expires,
1130 group_order,
1131 content_exists,
1132 largest_group_id,
1133 largest_object_id,
1134 parameters,
1135 }))
1136 }
1137 MessageType::SubscribeError => {
1138 let subscribe_id = VarInt::decode(buf)?;
1139 let error_code = VarInt::decode(buf)?;
1140 let reason_len = VarInt::decode(buf)?.into_inner() as usize;
1141 let reason_phrase = read_bytes(buf, reason_len)?;
1142 let track_alias = VarInt::decode(buf)?;
1143 Ok(ControlMessage::SubscribeError(SubscribeError {
1144 subscribe_id,
1145 error_code,
1146 reason_phrase,
1147 track_alias,
1148 }))
1149 }
1150 MessageType::SubscribeUpdate => {
1151 let subscribe_id = VarInt::decode(buf)?;
1152 let start_group = VarInt::decode(buf)?;
1153 let start_object = VarInt::decode(buf)?;
1154 let end_group = VarInt::decode(buf)?;
1155 let end_object = VarInt::decode(buf)?;
1156 if buf.remaining() < 1 {
1157 return Err(CodecError::UnexpectedEnd);
1158 }
1159 let subscriber_priority = buf.get_u8();
1160 let parameters = decode_parameters(buf)?;
1161 Ok(ControlMessage::SubscribeUpdate(SubscribeUpdate {
1162 subscribe_id,
1163 start_group,
1164 start_object,
1165 end_group,
1166 end_object,
1167 subscriber_priority,
1168 parameters,
1169 }))
1170 }
1171 MessageType::SubscribeDone => {
1172 let subscribe_id = VarInt::decode(buf)?;
1173 let status_code = VarInt::decode(buf)?;
1174 let reason_len = VarInt::decode(buf)?.into_inner() as usize;
1175 let reason_phrase = read_bytes(buf, reason_len)?;
1176 if buf.remaining() < 1 {
1177 return Err(CodecError::UnexpectedEnd);
1178 }
1179 let content_exists_val = buf.get_u8();
1180 let content_exists = match content_exists_val {
1181 0 => ContentExists::NoLargestLocation,
1182 1 => ContentExists::HasLargestLocation,
1183 other => return Err(CodecError::InvalidContentExists(other)),
1184 };
1185 let (final_group, final_object) =
1186 if content_exists == ContentExists::HasLargestLocation {
1187 let fg = VarInt::decode(buf)?;
1188 let fo = VarInt::decode(buf)?;
1189 (Some(fg), Some(fo))
1190 } else {
1191 (None, None)
1192 };
1193 Ok(ControlMessage::SubscribeDone(SubscribeDone {
1194 subscribe_id,
1195 status_code,
1196 reason_phrase,
1197 content_exists,
1198 final_group,
1199 final_object,
1200 }))
1201 }
1202 MessageType::Unsubscribe => {
1203 let subscribe_id = VarInt::decode(buf)?;
1204 Ok(ControlMessage::Unsubscribe(Unsubscribe { subscribe_id }))
1205 }
1206 MessageType::Announce => {
1207 let track_namespace = TrackNamespace::decode(buf)?;
1208 let parameters = decode_parameters(buf)?;
1209 Ok(ControlMessage::Announce(Announce { track_namespace, parameters }))
1210 }
1211 MessageType::AnnounceOk => {
1212 let track_namespace = TrackNamespace::decode(buf)?;
1213 Ok(ControlMessage::AnnounceOk(AnnounceOk { track_namespace }))
1214 }
1215 MessageType::AnnounceError => {
1216 let track_namespace = TrackNamespace::decode(buf)?;
1217 let error_code = VarInt::decode(buf)?;
1218 let reason_len = VarInt::decode(buf)?.into_inner() as usize;
1219 let reason_phrase = read_bytes(buf, reason_len)?;
1220 Ok(ControlMessage::AnnounceError(AnnounceError {
1221 track_namespace,
1222 error_code,
1223 reason_phrase,
1224 }))
1225 }
1226 MessageType::AnnounceCancel => {
1227 let track_namespace = TrackNamespace::decode(buf)?;
1228 let error_code = VarInt::decode(buf)?;
1229 let reason_len = VarInt::decode(buf)?.into_inner() as usize;
1230 let reason_phrase = read_bytes(buf, reason_len)?;
1231 Ok(ControlMessage::AnnounceCancel(AnnounceCancel {
1232 track_namespace,
1233 error_code,
1234 reason_phrase,
1235 }))
1236 }
1237 MessageType::Unannounce => {
1238 let track_namespace = TrackNamespace::decode(buf)?;
1239 Ok(ControlMessage::Unannounce(Unannounce { track_namespace }))
1240 }
1241 MessageType::SubscribeAnnounces => {
1242 let track_namespace_prefix = TrackNamespace::decode(buf)?;
1243 let parameters = decode_parameters(buf)?;
1244 Ok(ControlMessage::SubscribeAnnounces(SubscribeAnnounces {
1245 track_namespace_prefix,
1246 parameters,
1247 }))
1248 }
1249 MessageType::SubscribeAnnouncesOk => {
1250 let track_namespace_prefix = TrackNamespace::decode(buf)?;
1251 Ok(ControlMessage::SubscribeAnnouncesOk(SubscribeAnnouncesOk {
1252 track_namespace_prefix,
1253 }))
1254 }
1255 MessageType::SubscribeAnnouncesError => {
1256 let track_namespace_prefix = TrackNamespace::decode(buf)?;
1257 let error_code = VarInt::decode(buf)?;
1258 let reason_len = VarInt::decode(buf)?.into_inner() as usize;
1259 let reason_phrase = read_bytes(buf, reason_len)?;
1260 Ok(ControlMessage::SubscribeAnnouncesError(SubscribeAnnouncesError {
1261 track_namespace_prefix,
1262 error_code,
1263 reason_phrase,
1264 }))
1265 }
1266 MessageType::UnsubscribeAnnounces => {
1267 let track_namespace_prefix = TrackNamespace::decode(buf)?;
1268 Ok(ControlMessage::UnsubscribeAnnounces(UnsubscribeAnnounces {
1269 track_namespace_prefix,
1270 }))
1271 }
1272 MessageType::TrackStatusRequest => {
1273 let track_namespace = TrackNamespace::decode(buf)?;
1274 let track_name_len = VarInt::decode(buf)?.into_inner() as usize;
1275 let track_name = read_bytes(buf, track_name_len)?;
1276 Ok(ControlMessage::TrackStatusRequest(TrackStatusRequest {
1277 track_namespace,
1278 track_name,
1279 }))
1280 }
1281 MessageType::TrackStatus => {
1282 let track_namespace = TrackNamespace::decode(buf)?;
1283 let track_name_len = VarInt::decode(buf)?.into_inner() as usize;
1284 let track_name = read_bytes(buf, track_name_len)?;
1285 let status_code = VarInt::decode(buf)?;
1286 let last_group_id = VarInt::decode(buf)?;
1287 let last_object_id = VarInt::decode(buf)?;
1288 check_track_status(status_code, last_group_id, last_object_id)?;
1289 Ok(ControlMessage::TrackStatus(TrackStatus {
1290 track_namespace,
1291 track_name,
1292 status_code,
1293 last_group_id,
1294 last_object_id,
1295 }))
1296 }
1297 MessageType::Fetch => {
1298 let subscribe_id = VarInt::decode(buf)?;
1299 let track_namespace = TrackNamespace::decode(buf)?;
1300 let track_name_len = VarInt::decode(buf)?.into_inner() as usize;
1301 let track_name = read_bytes(buf, track_name_len)?;
1302 if buf.remaining() < 2 {
1303 return Err(CodecError::UnexpectedEnd);
1304 }
1305 let subscriber_priority = buf.get_u8();
1306 let group_order =
1307 GroupOrder::from_u8(buf.get_u8()).ok_or(CodecError::InvalidField)?;
1308 let start_group = VarInt::decode(buf)?;
1309 let start_object = VarInt::decode(buf)?;
1310 let end_group = VarInt::decode(buf)?;
1311 let end_object = VarInt::decode(buf)?;
1312 let parameters = decode_parameters(buf)?;
1313 Ok(ControlMessage::Fetch(Fetch {
1314 subscribe_id,
1315 track_namespace,
1316 track_name,
1317 subscriber_priority,
1318 group_order,
1319 start_group,
1320 start_object,
1321 end_group,
1322 end_object,
1323 parameters,
1324 }))
1325 }
1326 MessageType::FetchOk => {
1327 let subscribe_id = VarInt::decode(buf)?;
1328 if buf.remaining() < 2 {
1329 return Err(CodecError::UnexpectedEnd);
1330 }
1331 let group_order = read_group_order_response(buf)?;
1332 let end_of_track = buf.get_u8();
1333 let largest_group_id = Some(VarInt::decode(buf)?);
1334 let largest_object_id = Some(VarInt::decode(buf)?);
1335 let parameters = decode_parameters(buf)?;
1336 Ok(ControlMessage::FetchOk(FetchOk {
1337 subscribe_id,
1338 group_order,
1339 end_of_track,
1340 largest_group_id,
1341 largest_object_id,
1342 parameters,
1343 }))
1344 }
1345 MessageType::FetchError => {
1346 let subscribe_id = VarInt::decode(buf)?;
1347 let error_code = VarInt::decode(buf)?;
1348 let reason_len = VarInt::decode(buf)?.into_inner() as usize;
1349 let reason_phrase = read_bytes(buf, reason_len)?;
1350 Ok(ControlMessage::FetchError(FetchError {
1351 subscribe_id,
1352 error_code,
1353 reason_phrase,
1354 }))
1355 }
1356 MessageType::FetchCancel => {
1357 let subscribe_id = VarInt::decode(buf)?;
1358 Ok(ControlMessage::FetchCancel(FetchCancel { subscribe_id }))
1359 }
1360 }
1361 }
1362}