use crate::auth_token::{AuthorizationToken, AUTH_TOKEN_PARAMETER};
use crate::error::{
CodecError, MAX_FULL_TRACK_NAME_LENGTH, MAX_GOAWAY_URI_LENGTH, MAX_MESSAGE_LENGTH,
MAX_REASON_PHRASE_LENGTH,
};
use crate::kvp::{KeyValuePair, KvpValue};
use crate::types::*;
pub use crate::types::{check_group_range, check_location_range, check_open_ended_group_range};
use crate::varint::VarInt;
use bytes::{Buf, BufMut};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(u64)]
pub enum MessageType {
SubscribeUpdate = 0x02,
Subscribe = 0x03,
SubscribeOk = 0x04,
SubscribeError = 0x05,
Announce = 0x06,
AnnounceOk = 0x07,
AnnounceError = 0x08,
Unannounce = 0x09,
Unsubscribe = 0x0A,
SubscribeDone = 0x0B,
AnnounceCancel = 0x0C,
TrackStatus = 0x0D,
TrackStatusOk = 0x0E,
TrackStatusError = 0x0F,
GoAway = 0x10,
SubscribeNamespace = 0x11,
SubscribeNamespaceOk = 0x12,
SubscribeNamespaceError = 0x13,
UnsubscribeNamespace = 0x14,
MaxRequestId = 0x15,
Fetch = 0x16,
FetchCancel = 0x17,
FetchOk = 0x18,
FetchError = 0x19,
RequestsBlocked = 0x1A,
Publish = 0x1D,
PublishOk = 0x1E,
PublishError = 0x1F,
ClientSetup = 0x20,
ServerSetup = 0x21,
}
impl MessageType {
pub fn from_id(id: u64) -> Option<Self> {
match id {
0x02 => Some(MessageType::SubscribeUpdate),
0x03 => Some(MessageType::Subscribe),
0x04 => Some(MessageType::SubscribeOk),
0x05 => Some(MessageType::SubscribeError),
0x06 => Some(MessageType::Announce),
0x07 => Some(MessageType::AnnounceOk),
0x08 => Some(MessageType::AnnounceError),
0x09 => Some(MessageType::Unannounce),
0x0A => Some(MessageType::Unsubscribe),
0x0B => Some(MessageType::SubscribeDone),
0x0C => Some(MessageType::AnnounceCancel),
0x0D => Some(MessageType::TrackStatus),
0x0E => Some(MessageType::TrackStatusOk),
0x0F => Some(MessageType::TrackStatusError),
0x10 => Some(MessageType::GoAway),
0x11 => Some(MessageType::SubscribeNamespace),
0x12 => Some(MessageType::SubscribeNamespaceOk),
0x13 => Some(MessageType::SubscribeNamespaceError),
0x14 => Some(MessageType::UnsubscribeNamespace),
0x15 => Some(MessageType::MaxRequestId),
0x16 => Some(MessageType::Fetch),
0x17 => Some(MessageType::FetchCancel),
0x18 => Some(MessageType::FetchOk),
0x19 => Some(MessageType::FetchError),
0x1A => Some(MessageType::RequestsBlocked),
0x1D => Some(MessageType::Publish),
0x1E => Some(MessageType::PublishOk),
0x1F => Some(MessageType::PublishError),
0x20 => Some(MessageType::ClientSetup),
0x21 => Some(MessageType::ServerSetup),
_ => None,
}
}
pub fn id(&self) -> u64 {
*self as u64
}
pub fn name(&self) -> &'static str {
match self {
MessageType::SubscribeUpdate => "subscribe_update",
MessageType::Subscribe => "subscribe",
MessageType::SubscribeOk => "subscribe_ok",
MessageType::SubscribeError => "subscribe_error",
MessageType::Announce => "announce",
MessageType::AnnounceOk => "announce_ok",
MessageType::AnnounceError => "announce_error",
MessageType::Unannounce => "unannounce",
MessageType::Unsubscribe => "unsubscribe",
MessageType::SubscribeDone => "subscribe_done",
MessageType::AnnounceCancel => "announce_cancel",
MessageType::TrackStatus => "track_status",
MessageType::TrackStatusOk => "track_status_ok",
MessageType::TrackStatusError => "track_status_error",
MessageType::GoAway => "goaway",
MessageType::SubscribeNamespace => "subscribe_namespace",
MessageType::SubscribeNamespaceOk => "subscribe_namespace_ok",
MessageType::SubscribeNamespaceError => "subscribe_namespace_error",
MessageType::UnsubscribeNamespace => "unsubscribe_namespace",
MessageType::MaxRequestId => "max_request_id",
MessageType::Fetch => "fetch",
MessageType::FetchCancel => "fetch_cancel",
MessageType::FetchOk => "fetch_ok",
MessageType::FetchError => "fetch_error",
MessageType::RequestsBlocked => "requests_blocked",
MessageType::Publish => "publish",
MessageType::PublishOk => "publish_ok",
MessageType::PublishError => "publish_error",
MessageType::ClientSetup => "client_setup",
MessageType::ServerSetup => "server_setup",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ClientSetup {
pub supported_versions: Vec<VarInt>,
pub parameters: Vec<KeyValuePair>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ServerSetup {
pub selected_version: VarInt,
pub parameters: Vec<KeyValuePair>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct GoAway {
pub new_session_uri: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct MaxRequestId {
pub request_id: VarInt,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RequestsBlocked {
pub maximum_request_id: VarInt,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Subscribe {
pub request_id: VarInt,
pub track_namespace: TrackNamespace,
pub track_name: Vec<u8>,
pub subscriber_priority: u8,
pub group_order: GroupOrder,
pub forward: Forward,
pub filter_type: FilterType,
pub start_group: Option<VarInt>,
pub start_object: Option<VarInt>,
pub end_group: Option<VarInt>,
pub parameters: Vec<KeyValuePair>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SubscribeOk {
pub request_id: VarInt,
pub track_alias: VarInt,
pub expires: VarInt,
pub group_order: GroupOrder,
pub content_exists: ContentExists,
pub largest_location: Option<Location>,
pub parameters: Vec<KeyValuePair>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SubscribeError {
pub request_id: VarInt,
pub error_code: VarInt,
pub reason_phrase: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SubscribeUpdate {
pub request_id: VarInt,
pub start_group: VarInt,
pub start_object: VarInt,
pub end_group: VarInt,
pub subscriber_priority: u8,
pub forward: Forward,
pub parameters: Vec<KeyValuePair>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SubscribeDone {
pub request_id: VarInt,
pub status_code: VarInt,
pub stream_count: VarInt,
pub reason_phrase: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Unsubscribe {
pub request_id: VarInt,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Announce {
pub request_id: VarInt,
pub track_namespace: TrackNamespace,
pub parameters: Vec<KeyValuePair>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AnnounceOk {
pub request_id: VarInt,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AnnounceError {
pub request_id: VarInt,
pub error_code: VarInt,
pub reason_phrase: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct AnnounceCancel {
pub track_namespace: TrackNamespace,
pub error_code: VarInt,
pub reason_phrase: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Unannounce {
pub track_namespace: TrackNamespace,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SubscribeNamespace {
pub request_id: VarInt,
pub track_namespace_prefix: TrackNamespace,
pub parameters: Vec<KeyValuePair>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SubscribeNamespaceOk {
pub request_id: VarInt,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct SubscribeNamespaceError {
pub request_id: VarInt,
pub error_code: VarInt,
pub reason_phrase: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct UnsubscribeNamespace {
pub track_namespace_prefix: TrackNamespace,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TrackStatus {
pub request_id: VarInt,
pub track_namespace: TrackNamespace,
pub track_name: Vec<u8>,
pub subscriber_priority: u8,
pub group_order: GroupOrder,
pub forward: Forward,
pub filter_type: FilterType,
pub start_group: Option<VarInt>,
pub start_object: Option<VarInt>,
pub end_group: Option<VarInt>,
pub parameters: Vec<KeyValuePair>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TrackStatusOk {
pub request_id: VarInt,
pub track_alias: VarInt,
pub expires: VarInt,
pub group_order: GroupOrder,
pub content_exists: ContentExists,
pub largest_location: Option<Location>,
pub parameters: Vec<KeyValuePair>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TrackStatusError {
pub request_id: VarInt,
pub error_code: VarInt,
pub reason_phrase: Vec<u8>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[repr(u64)]
pub enum FetchType {
Standalone = 1,
RelativeJoining = 2,
AbsoluteJoining = 3,
}
impl FetchType {
pub fn from_u64(v: u64) -> Option<Self> {
match v {
1 => Some(FetchType::Standalone),
2 => Some(FetchType::RelativeJoining),
3 => Some(FetchType::AbsoluteJoining),
_ => None,
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Fetch {
pub request_id: VarInt,
pub subscriber_priority: u8,
pub group_order: GroupOrder,
pub fetch_type: FetchType,
pub fetch_payload: FetchPayload,
pub parameters: Vec<KeyValuePair>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum FetchPayload {
Standalone {
track_namespace: TrackNamespace,
track_name: Vec<u8>,
start_group: VarInt,
start_object: VarInt,
end_group: VarInt,
end_object: VarInt,
},
Joining {
joining_request_id: VarInt,
joining_start: VarInt,
},
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct FetchOk {
pub request_id: VarInt,
pub group_order: GroupOrder,
pub end_of_track: u8,
pub end_location: Location,
pub parameters: Vec<KeyValuePair>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct FetchError {
pub request_id: VarInt,
pub error_code: VarInt,
pub reason_phrase: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct FetchCancel {
pub request_id: VarInt,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Publish {
pub request_id: VarInt,
pub track_namespace: TrackNamespace,
pub track_name: Vec<u8>,
pub track_alias: VarInt,
pub group_order: GroupOrder,
pub content_exists: ContentExists,
pub largest_location: Option<Location>,
pub forward: Forward,
pub parameters: Vec<KeyValuePair>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PublishOk {
pub request_id: VarInt,
pub forward: Forward,
pub subscriber_priority: u8,
pub group_order: GroupOrder,
pub filter_type: FilterType,
pub start_group: Option<VarInt>,
pub start_object: Option<VarInt>,
pub end_group: Option<VarInt>,
pub parameters: Vec<KeyValuePair>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PublishError {
pub request_id: VarInt,
pub error_code: VarInt,
pub reason_phrase: Vec<u8>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ControlMessage {
ClientSetup(ClientSetup),
ServerSetup(ServerSetup),
GoAway(GoAway),
MaxRequestId(MaxRequestId),
RequestsBlocked(RequestsBlocked),
Subscribe(Subscribe),
SubscribeOk(SubscribeOk),
SubscribeError(SubscribeError),
SubscribeUpdate(SubscribeUpdate),
SubscribeDone(SubscribeDone),
Unsubscribe(Unsubscribe),
Announce(Announce),
AnnounceOk(AnnounceOk),
AnnounceError(AnnounceError),
AnnounceCancel(AnnounceCancel),
Unannounce(Unannounce),
SubscribeNamespace(SubscribeNamespace),
SubscribeNamespaceOk(SubscribeNamespaceOk),
SubscribeNamespaceError(SubscribeNamespaceError),
UnsubscribeNamespace(UnsubscribeNamespace),
TrackStatus(TrackStatus),
TrackStatusOk(TrackStatusOk),
TrackStatusError(TrackStatusError),
Fetch(Fetch),
FetchOk(FetchOk),
FetchError(FetchError),
FetchCancel(FetchCancel),
Publish(Publish),
PublishOk(PublishOk),
PublishError(PublishError),
}
fn check_full_track_name(namespace: &TrackNamespace, track_name: &[u8]) -> Result<(), CodecError> {
let total = namespace.field_bytes_len().saturating_add(track_name.len());
if total > MAX_FULL_TRACK_NAME_LENGTH {
return Err(CodecError::TrackNameTooLong);
}
Ok(())
}
fn read_reason_phrase(buf: &mut impl Buf) -> Result<Vec<u8>, CodecError> {
let len = VarInt::decode(buf)?.into_inner() as usize;
if len > MAX_REASON_PHRASE_LENGTH {
return Err(CodecError::ReasonPhraseTooLong);
}
read_bytes(buf, len)
}
fn read_u8(buf: &mut impl Buf) -> Result<u8, CodecError> {
if !buf.has_remaining() {
return Err(CodecError::UnexpectedEnd);
}
Ok(buf.get_u8())
}
fn read_group_order(buf: &mut impl Buf) -> Result<GroupOrder, CodecError> {
GroupOrder::from_u8(read_u8(buf)?).ok_or(CodecError::InvalidField)
}
fn read_group_order_response(buf: &mut impl Buf) -> Result<GroupOrder, CodecError> {
match read_group_order(buf)? {
GroupOrder::Publisher => Err(CodecError::InvalidField),
order => Ok(order),
}
}
fn read_forward(buf: &mut impl Buf) -> Result<Forward, CodecError> {
match read_u8(buf)? {
0 => Ok(Forward::DontForward),
1 => Ok(Forward::Forward),
other => Err(CodecError::InvalidForward(other)),
}
}
fn read_content_exists(buf: &mut impl Buf) -> Result<ContentExists, CodecError> {
match read_u8(buf)? {
0 => Ok(ContentExists::NoLargestLocation),
1 => Ok(ContentExists::HasLargestLocation),
other => Err(CodecError::InvalidContentExists(other)),
}
}
fn read_filter_type(buf: &mut impl Buf) -> Result<FilterType, CodecError> {
let value = VarInt::decode(buf)?.into_inner();
FilterType::from_u64(value).ok_or(CodecError::InvalidFilterType(value))
}
type StartLocation = (Option<VarInt>, Option<VarInt>);
fn read_start_location(
filter_type: FilterType,
buf: &mut impl Buf,
) -> Result<StartLocation, CodecError> {
match filter_type {
FilterType::AbsoluteStart | FilterType::AbsoluteRange => {
Ok((Some(VarInt::decode(buf)?), Some(VarInt::decode(buf)?)))
}
_ => Ok((None, None)),
}
}
fn read_end_group(
filter_type: FilterType,
buf: &mut impl Buf,
) -> Result<Option<VarInt>, CodecError> {
match filter_type {
FilterType::AbsoluteRange => Ok(Some(VarInt::decode(buf)?)),
_ => Ok(None),
}
}
fn read_largest_location(
content_exists: ContentExists,
buf: &mut impl Buf,
) -> Result<Option<Location>, CodecError> {
match content_exists {
ContentExists::HasLargestLocation => Ok(Some(Location::decode(buf)?)),
ContentExists::NoLargestLocation => Ok(None),
}
}
fn check_group_order(message: &ControlMessage) -> Result<(), CodecError> {
let order = match message {
ControlMessage::SubscribeOk(m) => m.group_order,
ControlMessage::TrackStatusOk(m) => m.group_order,
ControlMessage::FetchOk(m) => m.group_order,
ControlMessage::Publish(m) => m.group_order,
ControlMessage::PublishOk(m) => m.group_order,
_ => return Ok(()),
};
if order == GroupOrder::Publisher {
return Err(CodecError::InvalidField);
}
Ok(())
}
fn check_discriminators(message: &ControlMessage) -> Result<(), CodecError> {
fn check_filter(
filter_type: FilterType,
start_group: &Option<VarInt>,
start_object: &Option<VarInt>,
end_group: &Option<VarInt>,
) -> Result<(), CodecError> {
let wants_start =
matches!(filter_type, FilterType::AbsoluteStart | FilterType::AbsoluteRange);
if wants_start != start_group.is_some() || wants_start != start_object.is_some() {
return Err(CodecError::InvalidField);
}
if (filter_type == FilterType::AbsoluteRange) != end_group.is_some() {
return Err(CodecError::InvalidField);
}
Ok(())
}
fn check_content(
content_exists: ContentExists,
largest_location: &Option<Location>,
) -> Result<(), CodecError> {
if (content_exists == ContentExists::HasLargestLocation) != largest_location.is_some() {
return Err(CodecError::InvalidField);
}
Ok(())
}
match message {
ControlMessage::Subscribe(m) => {
check_filter(m.filter_type, &m.start_group, &m.start_object, &m.end_group)
}
ControlMessage::TrackStatus(m) => {
check_filter(m.filter_type, &m.start_group, &m.start_object, &m.end_group)
}
ControlMessage::PublishOk(m) => {
check_filter(m.filter_type, &m.start_group, &m.start_object, &m.end_group)
}
ControlMessage::SubscribeOk(m) => check_content(m.content_exists, &m.largest_location),
ControlMessage::TrackStatusOk(m) => check_content(m.content_exists, &m.largest_location),
ControlMessage::Publish(m) => check_content(m.content_exists, &m.largest_location),
ControlMessage::Fetch(m) => {
let body_is_standalone = matches!(m.fetch_payload, FetchPayload::Standalone { .. });
if body_is_standalone != (m.fetch_type == FetchType::Standalone) {
return Err(CodecError::InvalidField);
}
Ok(())
}
_ => Ok(()),
}
}
fn check_ranges(message: &ControlMessage) -> Result<(), CodecError> {
match message {
ControlMessage::Subscribe(m) => match (&m.start_group, &m.end_group) {
(Some(start_group), Some(end_group)) => {
check_group_range(start_group.into_inner(), end_group.into_inner())
}
_ => Ok(()),
},
ControlMessage::SubscribeUpdate(m) => {
check_open_ended_group_range(m.start_group.into_inner(), m.end_group.into_inner())
}
ControlMessage::Fetch(m) => match &m.fetch_payload {
FetchPayload::Standalone {
start_group, start_object, end_group, end_object, ..
} => check_location_range(
start_group.into_inner(),
start_object.into_inner(),
end_group.into_inner(),
end_object.into_inner(),
),
FetchPayload::Joining { .. } => Ok(()),
},
_ => Ok(()),
}
}
const AUTHORIZATION_TOKEN: u64 = 0x03;
const DELIVERY_TIMEOUT: u64 = 0x02;
const MAX_CACHE_DURATION: u64 = 0x04;
const SETUP_PATH: u64 = 0x01;
const SETUP_MAX_REQUEST_ID: u64 = 0x02;
const SETUP_MAX_AUTH_TOKEN_CACHE_SIZE: u64 = 0x04;
const SETUP_AUTHORIZATION_TOKEN: u64 = AUTHORIZATION_TOKEN;
const KNOWN_PARAMETERS: &[u64] = &[AUTHORIZATION_TOKEN, DELIVERY_TIMEOUT, MAX_CACHE_DURATION];
const REPEATABLE_PARAMETERS: &[u64] = &[AUTHORIZATION_TOKEN];
const KNOWN_SETUP_PARAMETERS: &[u64] =
&[SETUP_PATH, SETUP_MAX_REQUEST_ID, SETUP_AUTHORIZATION_TOKEN, SETUP_MAX_AUTH_TOKEN_CACHE_SIZE];
const REPEATABLE_SETUP_PARAMETERS: &[u64] = &[SETUP_AUTHORIZATION_TOKEN];
fn check_sender_parameters(
parameters: &[KeyValuePair],
repeatable: &[u64],
) -> Result<(), CodecError> {
for (i, parameter) in parameters.iter().enumerate() {
let key = parameter.key.into_inner();
if repeatable.contains(&key) {
continue;
}
if parameters[..i].iter().any(|earlier| earlier.key == parameter.key) {
return Err(CodecError::DuplicateParameter(key));
}
}
Ok(())
}
fn check_receiver_parameters(
parameters: &[KeyValuePair],
known: &[u64],
repeatable: &[u64],
) -> Result<(), CodecError> {
for (i, parameter) in parameters.iter().enumerate() {
let key = parameter.key.into_inner();
if repeatable.contains(&key) || !known.contains(&key) {
continue;
}
if parameters[..i].iter().any(|earlier| earlier.key == parameter.key) {
return Err(CodecError::DuplicateParameter(key));
}
}
Ok(())
}
fn check_authorization_tokens(parameters: &[KeyValuePair]) -> Result<(), CodecError> {
for parameter in parameters {
let key = parameter.key.into_inner();
if key != AUTH_TOKEN_PARAMETER {
continue;
}
match ¶meter.value {
KvpValue::Bytes(value) => {
AuthorizationToken::decode(key, value)?;
}
KvpValue::Varint(_) => {
return Err(CodecError::KeyValueFormatting {
key,
detail: "its value is a bare varint where the type defines a Token structure",
});
}
}
}
Ok(())
}
fn decode_parameters(buf: &mut impl Buf) -> Result<Vec<KeyValuePair>, CodecError> {
let parameters = KeyValuePair::decode_list(buf)?;
check_receiver_parameters(¶meters, KNOWN_PARAMETERS, REPEATABLE_PARAMETERS)?;
check_authorization_tokens(¶meters)?;
Ok(parameters)
}
fn encode_parameters(parameters: &[KeyValuePair], buf: &mut impl BufMut) -> Result<(), CodecError> {
check_sender_parameters(parameters, REPEATABLE_PARAMETERS)?;
check_authorization_tokens(parameters)?;
KeyValuePair::encode_list_checked(parameters, buf)?;
Ok(())
}
fn decode_setup_parameters(buf: &mut impl Buf) -> Result<Vec<KeyValuePair>, CodecError> {
let parameters = KeyValuePair::decode_list(buf)?;
check_receiver_parameters(¶meters, KNOWN_SETUP_PARAMETERS, REPEATABLE_SETUP_PARAMETERS)?;
check_authorization_tokens(¶meters)?;
Ok(parameters)
}
fn encode_setup_parameters(
parameters: &[KeyValuePair],
buf: &mut impl BufMut,
) -> Result<(), CodecError> {
check_sender_parameters(parameters, REPEATABLE_SETUP_PARAMETERS)?;
check_authorization_tokens(parameters)?;
KeyValuePair::encode_list_checked(parameters, buf)?;
Ok(())
}
impl ControlMessage {
pub fn encode(&self, buf: &mut impl BufMut) -> Result<(), CodecError> {
check_discriminators(self)?;
check_group_order(self)?;
check_ranges(self)?;
let mut payload = Vec::with_capacity(256);
self.encode_payload(&mut payload)?;
if payload.len() > MAX_MESSAGE_LENGTH {
return Err(CodecError::MessageTooLong(payload.len()));
}
VarInt::from_usize(self.message_type().id() as usize).encode(buf);
buf.put_u16(payload.len() as u16);
buf.put_slice(&payload);
Ok(())
}
pub fn decode(buf: &mut impl Buf) -> Result<Self, CodecError> {
let type_id = VarInt::decode(buf)?.into_inner();
let msg_type =
MessageType::from_id(type_id).ok_or(CodecError::UnknownMessageType(type_id))?;
if buf.remaining() < 2 {
return Err(CodecError::UnexpectedEnd);
}
let payload_len = buf.get_u16() as usize;
if buf.remaining() < payload_len {
return Err(CodecError::UnexpectedEnd);
}
let payload_bytes = buf.copy_to_bytes(payload_len);
let mut payload = &payload_bytes[..];
let msg = match Self::decode_payload(msg_type, &mut payload) {
Ok(msg) => msg,
Err(
CodecError::UnexpectedEnd
| CodecError::Kvp(crate::kvp::KvpError::UnexpectedEnd)
| CodecError::Kvp(crate::kvp::KvpError::VarInt(
crate::varint::VarIntError::UnexpectedEnd,
))
| CodecError::VarInt(crate::varint::VarIntError::UnexpectedEnd),
) => {
return Err(CodecError::ControlMessageLengthMismatch {
declared: payload_len,
detail: "its fields ran past the end",
});
}
Err(e) => return Err(e),
};
check_ranges(&msg)?;
if payload.has_remaining() {
return Err(CodecError::ControlMessageLengthMismatch {
declared: payload_len,
detail: "its fields left bytes unread",
});
}
Ok(msg)
}
fn encode_payload(&self, buf: &mut impl BufMut) -> Result<(), CodecError> {
match self {
ControlMessage::ClientSetup(m) => {
VarInt::from_usize(m.supported_versions.len()).encode(buf);
for v in &m.supported_versions {
v.encode(buf);
}
encode_setup_parameters(&m.parameters, buf)?;
}
ControlMessage::ServerSetup(m) => {
m.selected_version.encode(buf);
encode_setup_parameters(&m.parameters, buf)?;
}
ControlMessage::GoAway(m) => {
if m.new_session_uri.len() > MAX_GOAWAY_URI_LENGTH {
return Err(CodecError::GoAwayUriTooLong);
}
VarInt::from_usize(m.new_session_uri.len()).encode(buf);
buf.put_slice(&m.new_session_uri);
}
ControlMessage::MaxRequestId(m) => {
m.request_id.encode(buf);
}
ControlMessage::RequestsBlocked(m) => {
m.maximum_request_id.encode(buf);
}
ControlMessage::Subscribe(m) => {
m.request_id.encode(buf);
m.track_namespace.validate(TrackNamespaceRules::for_draft(13))?;
m.track_namespace.encode(buf);
check_full_track_name(&m.track_namespace, &m.track_name)?;
VarInt::from_usize(m.track_name.len()).encode(buf);
buf.put_slice(&m.track_name);
buf.put_u8(m.subscriber_priority);
buf.put_u8(m.group_order as u8);
buf.put_u8(m.forward as u8);
VarInt::from_usize(m.filter_type as usize).encode(buf);
if let Some(sg) = &m.start_group {
sg.encode(buf);
}
if let Some(so) = &m.start_object {
so.encode(buf);
}
if let Some(eg) = &m.end_group {
eg.encode(buf);
}
encode_parameters(&m.parameters, buf)?;
}
ControlMessage::SubscribeOk(m) => {
m.request_id.encode(buf);
m.track_alias.encode(buf);
m.expires.encode(buf);
buf.put_u8(m.group_order as u8);
buf.put_u8(m.content_exists as u8);
if let Some(loc) = &m.largest_location {
loc.encode(buf);
}
encode_parameters(&m.parameters, buf)?;
}
ControlMessage::SubscribeError(m) => {
if m.reason_phrase.len() > MAX_REASON_PHRASE_LENGTH {
return Err(CodecError::ReasonPhraseTooLong);
}
m.request_id.encode(buf);
m.error_code.encode(buf);
VarInt::from_usize(m.reason_phrase.len()).encode(buf);
buf.put_slice(&m.reason_phrase);
}
ControlMessage::SubscribeUpdate(m) => {
m.request_id.encode(buf);
m.start_group.encode(buf);
m.start_object.encode(buf);
m.end_group.encode(buf);
buf.put_u8(m.subscriber_priority);
buf.put_u8(m.forward as u8);
encode_parameters(&m.parameters, buf)?;
}
ControlMessage::SubscribeDone(m) => {
if m.reason_phrase.len() > MAX_REASON_PHRASE_LENGTH {
return Err(CodecError::ReasonPhraseTooLong);
}
m.request_id.encode(buf);
m.status_code.encode(buf);
m.stream_count.encode(buf);
VarInt::from_usize(m.reason_phrase.len()).encode(buf);
buf.put_slice(&m.reason_phrase);
}
ControlMessage::Unsubscribe(m) => {
m.request_id.encode(buf);
}
ControlMessage::Announce(m) => {
m.request_id.encode(buf);
m.track_namespace.validate(TrackNamespaceRules::for_draft(13))?;
m.track_namespace.encode(buf);
encode_parameters(&m.parameters, buf)?;
}
ControlMessage::AnnounceOk(m) => {
m.request_id.encode(buf);
}
ControlMessage::AnnounceError(m) => {
if m.reason_phrase.len() > MAX_REASON_PHRASE_LENGTH {
return Err(CodecError::ReasonPhraseTooLong);
}
m.request_id.encode(buf);
m.error_code.encode(buf);
VarInt::from_usize(m.reason_phrase.len()).encode(buf);
buf.put_slice(&m.reason_phrase);
}
ControlMessage::AnnounceCancel(m) => {
if m.reason_phrase.len() > MAX_REASON_PHRASE_LENGTH {
return Err(CodecError::ReasonPhraseTooLong);
}
m.track_namespace.validate(TrackNamespaceRules::for_draft(13))?;
m.track_namespace.encode(buf);
m.error_code.encode(buf);
VarInt::from_usize(m.reason_phrase.len()).encode(buf);
buf.put_slice(&m.reason_phrase);
}
ControlMessage::Unannounce(m) => {
m.track_namespace.validate(TrackNamespaceRules::for_draft(13))?;
m.track_namespace.encode(buf);
}
ControlMessage::SubscribeNamespace(m) => {
m.request_id.encode(buf);
m.track_namespace_prefix.validate(TrackNamespaceRules::for_draft(13))?;
m.track_namespace_prefix.encode(buf);
encode_parameters(&m.parameters, buf)?;
}
ControlMessage::SubscribeNamespaceOk(m) => {
m.request_id.encode(buf);
}
ControlMessage::SubscribeNamespaceError(m) => {
if m.reason_phrase.len() > MAX_REASON_PHRASE_LENGTH {
return Err(CodecError::ReasonPhraseTooLong);
}
m.request_id.encode(buf);
m.error_code.encode(buf);
VarInt::from_usize(m.reason_phrase.len()).encode(buf);
buf.put_slice(&m.reason_phrase);
}
ControlMessage::UnsubscribeNamespace(m) => {
m.track_namespace_prefix.validate(TrackNamespaceRules::for_draft(13))?;
m.track_namespace_prefix.encode(buf);
}
ControlMessage::TrackStatus(m) => {
m.request_id.encode(buf);
m.track_namespace.validate(TrackNamespaceRules::for_draft(13))?;
m.track_namespace.encode(buf);
check_full_track_name(&m.track_namespace, &m.track_name)?;
VarInt::from_usize(m.track_name.len()).encode(buf);
buf.put_slice(&m.track_name);
buf.put_u8(m.subscriber_priority);
buf.put_u8(m.group_order as u8);
buf.put_u8(m.forward as u8);
VarInt::from_usize(m.filter_type as usize).encode(buf);
if let Some(sg) = &m.start_group {
sg.encode(buf);
}
if let Some(so) = &m.start_object {
so.encode(buf);
}
if let Some(eg) = &m.end_group {
eg.encode(buf);
}
encode_parameters(&m.parameters, buf)?;
}
ControlMessage::TrackStatusOk(m) => {
m.request_id.encode(buf);
m.track_alias.encode(buf);
m.expires.encode(buf);
buf.put_u8(m.group_order as u8);
buf.put_u8(m.content_exists as u8);
if let Some(loc) = &m.largest_location {
loc.encode(buf);
}
encode_parameters(&m.parameters, buf)?;
}
ControlMessage::TrackStatusError(m) => {
if m.reason_phrase.len() > MAX_REASON_PHRASE_LENGTH {
return Err(CodecError::ReasonPhraseTooLong);
}
m.request_id.encode(buf);
m.error_code.encode(buf);
VarInt::from_usize(m.reason_phrase.len()).encode(buf);
buf.put_slice(&m.reason_phrase);
}
ControlMessage::Fetch(m) => {
m.request_id.encode(buf);
buf.put_u8(m.subscriber_priority);
buf.put_u8(m.group_order as u8);
VarInt::from_usize(m.fetch_type as usize).encode(buf);
match &m.fetch_payload {
FetchPayload::Standalone {
track_namespace,
track_name,
start_group,
start_object,
end_group,
end_object,
} => {
track_namespace.validate(TrackNamespaceRules::for_draft(13))?;
track_namespace.encode(buf);
check_full_track_name(track_namespace, track_name)?;
VarInt::from_usize(track_name.len()).encode(buf);
buf.put_slice(track_name);
start_group.encode(buf);
start_object.encode(buf);
end_group.encode(buf);
end_object.encode(buf);
}
FetchPayload::Joining { joining_request_id, joining_start } => {
joining_request_id.encode(buf);
joining_start.encode(buf);
}
}
encode_parameters(&m.parameters, buf)?;
}
ControlMessage::FetchOk(m) => {
m.request_id.encode(buf);
buf.put_u8(m.group_order as u8);
buf.put_u8(m.end_of_track);
m.end_location.encode(buf);
encode_parameters(&m.parameters, buf)?;
}
ControlMessage::FetchError(m) => {
if m.reason_phrase.len() > MAX_REASON_PHRASE_LENGTH {
return Err(CodecError::ReasonPhraseTooLong);
}
m.request_id.encode(buf);
m.error_code.encode(buf);
VarInt::from_usize(m.reason_phrase.len()).encode(buf);
buf.put_slice(&m.reason_phrase);
}
ControlMessage::FetchCancel(m) => {
m.request_id.encode(buf);
}
ControlMessage::Publish(m) => {
m.request_id.encode(buf);
m.track_namespace.validate(TrackNamespaceRules::for_draft(13))?;
m.track_namespace.encode(buf);
check_full_track_name(&m.track_namespace, &m.track_name)?;
VarInt::from_usize(m.track_name.len()).encode(buf);
buf.put_slice(&m.track_name);
m.track_alias.encode(buf);
buf.put_u8(m.group_order as u8);
buf.put_u8(m.content_exists as u8);
if let Some(loc) = &m.largest_location {
loc.encode(buf);
}
buf.put_u8(m.forward as u8);
encode_parameters(&m.parameters, buf)?;
}
ControlMessage::PublishOk(m) => {
m.request_id.encode(buf);
buf.put_u8(m.forward as u8);
buf.put_u8(m.subscriber_priority);
buf.put_u8(m.group_order as u8);
VarInt::from_usize(m.filter_type as usize).encode(buf);
if let Some(sg) = &m.start_group {
sg.encode(buf);
}
if let Some(so) = &m.start_object {
so.encode(buf);
}
if let Some(eg) = &m.end_group {
eg.encode(buf);
}
encode_parameters(&m.parameters, buf)?;
}
ControlMessage::PublishError(m) => {
if m.reason_phrase.len() > MAX_REASON_PHRASE_LENGTH {
return Err(CodecError::ReasonPhraseTooLong);
}
m.request_id.encode(buf);
m.error_code.encode(buf);
VarInt::from_usize(m.reason_phrase.len()).encode(buf);
buf.put_slice(&m.reason_phrase);
}
}
Ok(())
}
fn decode_payload(msg_type: MessageType, buf: &mut impl Buf) -> Result<Self, CodecError> {
match msg_type {
MessageType::ClientSetup => {
let num_versions = VarInt::decode(buf)?.into_inner() as usize;
if num_versions == 0 {
return Err(CodecError::InvalidField);
}
let mut supported_versions = crate::types::reserve_bounded(num_versions, buf);
for _ in 0..num_versions {
supported_versions.push(VarInt::decode(buf)?);
}
let parameters = decode_setup_parameters(buf)?;
Ok(ControlMessage::ClientSetup(ClientSetup { supported_versions, parameters }))
}
MessageType::ServerSetup => {
let selected_version = VarInt::decode(buf)?;
let parameters = decode_setup_parameters(buf)?;
Ok(ControlMessage::ServerSetup(ServerSetup { selected_version, parameters }))
}
MessageType::GoAway => {
let uri_len = VarInt::decode(buf)?.into_inner() as usize;
if uri_len > MAX_GOAWAY_URI_LENGTH {
return Err(CodecError::GoAwayUriTooLong);
}
let uri = read_bytes(buf, uri_len)?;
Ok(ControlMessage::GoAway(GoAway { new_session_uri: uri }))
}
MessageType::MaxRequestId => {
let request_id = VarInt::decode(buf)?;
Ok(ControlMessage::MaxRequestId(MaxRequestId { request_id }))
}
MessageType::RequestsBlocked => {
let maximum_request_id = VarInt::decode(buf)?;
Ok(ControlMessage::RequestsBlocked(RequestsBlocked { maximum_request_id }))
}
MessageType::Subscribe => {
let request_id = VarInt::decode(buf)?;
let track_namespace = TrackNamespace::decode(buf)?;
let track_name_len = VarInt::decode(buf)?.into_inner() as usize;
let track_name = read_bytes(buf, track_name_len)?;
check_full_track_name(&track_namespace, &track_name)?;
let subscriber_priority = read_u8(buf)?;
let group_order = read_group_order(buf)?;
let forward = read_forward(buf)?;
let filter_type = read_filter_type(buf)?;
let (start_group, start_object) = read_start_location(filter_type, buf)?;
let end_group = read_end_group(filter_type, buf)?;
let parameters = decode_parameters(buf)?;
Ok(ControlMessage::Subscribe(Subscribe {
request_id,
track_namespace,
track_name,
subscriber_priority,
group_order,
forward,
filter_type,
start_group,
start_object,
end_group,
parameters,
}))
}
MessageType::SubscribeOk => {
let request_id = VarInt::decode(buf)?;
let track_alias = VarInt::decode(buf)?;
let expires = VarInt::decode(buf)?;
let group_order = read_group_order_response(buf)?;
let content_exists = read_content_exists(buf)?;
let largest_location = read_largest_location(content_exists, buf)?;
let parameters = decode_parameters(buf)?;
Ok(ControlMessage::SubscribeOk(SubscribeOk {
request_id,
track_alias,
expires,
group_order,
content_exists,
largest_location,
parameters,
}))
}
MessageType::SubscribeError => {
let request_id = VarInt::decode(buf)?;
let error_code = VarInt::decode(buf)?;
let reason_phrase = read_reason_phrase(buf)?;
Ok(ControlMessage::SubscribeError(SubscribeError {
request_id,
error_code,
reason_phrase,
}))
}
MessageType::SubscribeUpdate => {
let request_id = VarInt::decode(buf)?;
let start_group = VarInt::decode(buf)?;
let start_object = VarInt::decode(buf)?;
let end_group = VarInt::decode(buf)?;
let subscriber_priority = read_u8(buf)?;
let forward = read_forward(buf)?;
let parameters = decode_parameters(buf)?;
Ok(ControlMessage::SubscribeUpdate(SubscribeUpdate {
request_id,
start_group,
start_object,
end_group,
subscriber_priority,
forward,
parameters,
}))
}
MessageType::SubscribeDone => {
let request_id = VarInt::decode(buf)?;
let status_code = VarInt::decode(buf)?;
let stream_count = VarInt::decode(buf)?;
let reason_phrase = read_reason_phrase(buf)?;
Ok(ControlMessage::SubscribeDone(SubscribeDone {
request_id,
status_code,
stream_count,
reason_phrase,
}))
}
MessageType::Unsubscribe => {
let request_id = VarInt::decode(buf)?;
Ok(ControlMessage::Unsubscribe(Unsubscribe { request_id }))
}
MessageType::Announce => {
let request_id = VarInt::decode(buf)?;
let track_namespace = TrackNamespace::decode(buf)?;
let parameters = decode_parameters(buf)?;
Ok(ControlMessage::Announce(Announce { request_id, track_namespace, parameters }))
}
MessageType::AnnounceOk => {
let request_id = VarInt::decode(buf)?;
Ok(ControlMessage::AnnounceOk(AnnounceOk { request_id }))
}
MessageType::AnnounceError => {
let request_id = VarInt::decode(buf)?;
let error_code = VarInt::decode(buf)?;
let reason_phrase = read_reason_phrase(buf)?;
Ok(ControlMessage::AnnounceError(AnnounceError {
request_id,
error_code,
reason_phrase,
}))
}
MessageType::AnnounceCancel => {
let track_namespace = TrackNamespace::decode(buf)?;
let error_code = VarInt::decode(buf)?;
let reason_phrase = read_reason_phrase(buf)?;
Ok(ControlMessage::AnnounceCancel(AnnounceCancel {
track_namespace,
error_code,
reason_phrase,
}))
}
MessageType::Unannounce => {
let track_namespace = TrackNamespace::decode(buf)?;
Ok(ControlMessage::Unannounce(Unannounce { track_namespace }))
}
MessageType::SubscribeNamespace => {
let request_id = VarInt::decode(buf)?;
let track_namespace_prefix = TrackNamespace::decode(buf)?;
let parameters = decode_parameters(buf)?;
Ok(ControlMessage::SubscribeNamespace(SubscribeNamespace {
request_id,
track_namespace_prefix,
parameters,
}))
}
MessageType::SubscribeNamespaceOk => {
let request_id = VarInt::decode(buf)?;
Ok(ControlMessage::SubscribeNamespaceOk(SubscribeNamespaceOk { request_id }))
}
MessageType::SubscribeNamespaceError => {
let request_id = VarInt::decode(buf)?;
let error_code = VarInt::decode(buf)?;
let reason_phrase = read_reason_phrase(buf)?;
Ok(ControlMessage::SubscribeNamespaceError(SubscribeNamespaceError {
request_id,
error_code,
reason_phrase,
}))
}
MessageType::UnsubscribeNamespace => {
let track_namespace_prefix = TrackNamespace::decode(buf)?;
Ok(ControlMessage::UnsubscribeNamespace(UnsubscribeNamespace {
track_namespace_prefix,
}))
}
MessageType::TrackStatus => {
let request_id = VarInt::decode(buf)?;
let track_namespace = TrackNamespace::decode(buf)?;
let track_name_len = VarInt::decode(buf)?.into_inner() as usize;
let track_name = read_bytes(buf, track_name_len)?;
check_full_track_name(&track_namespace, &track_name)?;
let subscriber_priority = read_u8(buf)?;
let group_order = read_group_order(buf)?;
let forward = read_forward(buf)?;
let filter_type = read_filter_type(buf)?;
let (start_group, start_object) = read_start_location(filter_type, buf)?;
let end_group = read_end_group(filter_type, buf)?;
let parameters = decode_parameters(buf)?;
Ok(ControlMessage::TrackStatus(TrackStatus {
request_id,
track_namespace,
track_name,
subscriber_priority,
group_order,
forward,
filter_type,
start_group,
start_object,
end_group,
parameters,
}))
}
MessageType::TrackStatusOk => {
let request_id = VarInt::decode(buf)?;
let track_alias = VarInt::decode(buf)?;
let expires = VarInt::decode(buf)?;
let group_order = read_group_order_response(buf)?;
let content_exists = read_content_exists(buf)?;
let largest_location = read_largest_location(content_exists, buf)?;
let parameters = decode_parameters(buf)?;
Ok(ControlMessage::TrackStatusOk(TrackStatusOk {
request_id,
track_alias,
expires,
group_order,
content_exists,
largest_location,
parameters,
}))
}
MessageType::TrackStatusError => {
let request_id = VarInt::decode(buf)?;
let error_code = VarInt::decode(buf)?;
let reason_phrase = read_reason_phrase(buf)?;
Ok(ControlMessage::TrackStatusError(TrackStatusError {
request_id,
error_code,
reason_phrase,
}))
}
MessageType::Fetch => {
let request_id = VarInt::decode(buf)?;
let subscriber_priority = read_u8(buf)?;
let group_order = read_group_order(buf)?;
let fetch_type_val = VarInt::decode(buf)?.into_inner();
let fetch_type = FetchType::from_u64(fetch_type_val)
.ok_or(CodecError::InvalidFetchType(fetch_type_val))?;
let fetch_payload = match fetch_type {
FetchType::Standalone => {
let track_namespace = TrackNamespace::decode(buf)?;
let track_name_len = VarInt::decode(buf)?.into_inner() as usize;
let track_name = read_bytes(buf, track_name_len)?;
check_full_track_name(&track_namespace, &track_name)?;
let start_group = VarInt::decode(buf)?;
let start_object = VarInt::decode(buf)?;
let end_group = VarInt::decode(buf)?;
let end_object = VarInt::decode(buf)?;
FetchPayload::Standalone {
track_namespace,
track_name,
start_group,
start_object,
end_group,
end_object,
}
}
FetchType::RelativeJoining | FetchType::AbsoluteJoining => {
let joining_request_id = VarInt::decode(buf)?;
let joining_start = VarInt::decode(buf)?;
FetchPayload::Joining { joining_request_id, joining_start }
}
};
let parameters = decode_parameters(buf)?;
Ok(ControlMessage::Fetch(Fetch {
request_id,
subscriber_priority,
group_order,
fetch_type,
fetch_payload,
parameters,
}))
}
MessageType::FetchOk => {
let request_id = VarInt::decode(buf)?;
let group_order = read_group_order_response(buf)?;
let end_of_track = read_u8(buf)?;
let end_location = Location::decode(buf)?;
let parameters = decode_parameters(buf)?;
Ok(ControlMessage::FetchOk(FetchOk {
request_id,
group_order,
end_of_track,
end_location,
parameters,
}))
}
MessageType::FetchError => {
let request_id = VarInt::decode(buf)?;
let error_code = VarInt::decode(buf)?;
let reason_phrase = read_reason_phrase(buf)?;
Ok(ControlMessage::FetchError(FetchError { request_id, error_code, reason_phrase }))
}
MessageType::FetchCancel => {
let request_id = VarInt::decode(buf)?;
Ok(ControlMessage::FetchCancel(FetchCancel { request_id }))
}
MessageType::Publish => {
let request_id = VarInt::decode(buf)?;
let track_namespace = TrackNamespace::decode(buf)?;
let track_name_len = VarInt::decode(buf)?.into_inner() as usize;
let track_name = read_bytes(buf, track_name_len)?;
check_full_track_name(&track_namespace, &track_name)?;
let track_alias = VarInt::decode(buf)?;
let group_order = read_group_order_response(buf)?;
let content_exists = read_content_exists(buf)?;
let largest_location = read_largest_location(content_exists, buf)?;
let forward = read_forward(buf)?;
let parameters = decode_parameters(buf)?;
Ok(ControlMessage::Publish(Publish {
request_id,
track_namespace,
track_name,
track_alias,
group_order,
content_exists,
largest_location,
forward,
parameters,
}))
}
MessageType::PublishOk => {
let request_id = VarInt::decode(buf)?;
let forward = read_forward(buf)?;
let subscriber_priority = read_u8(buf)?;
let group_order = read_group_order_response(buf)?;
let filter_type = read_filter_type(buf)?;
let (start_group, start_object) = read_start_location(filter_type, buf)?;
let end_group = read_end_group(filter_type, buf)?;
let parameters = decode_parameters(buf)?;
Ok(ControlMessage::PublishOk(PublishOk {
request_id,
forward,
subscriber_priority,
group_order,
filter_type,
start_group,
start_object,
end_group,
parameters,
}))
}
MessageType::PublishError => {
let request_id = VarInt::decode(buf)?;
let error_code = VarInt::decode(buf)?;
let reason_phrase = read_reason_phrase(buf)?;
Ok(ControlMessage::PublishError(PublishError {
request_id,
error_code,
reason_phrase,
}))
}
}
}
pub fn message_type(&self) -> MessageType {
match self {
ControlMessage::ClientSetup(_) => MessageType::ClientSetup,
ControlMessage::ServerSetup(_) => MessageType::ServerSetup,
ControlMessage::GoAway(_) => MessageType::GoAway,
ControlMessage::MaxRequestId(_) => MessageType::MaxRequestId,
ControlMessage::RequestsBlocked(_) => MessageType::RequestsBlocked,
ControlMessage::Subscribe(_) => MessageType::Subscribe,
ControlMessage::SubscribeOk(_) => MessageType::SubscribeOk,
ControlMessage::SubscribeError(_) => MessageType::SubscribeError,
ControlMessage::SubscribeUpdate(_) => MessageType::SubscribeUpdate,
ControlMessage::SubscribeDone(_) => MessageType::SubscribeDone,
ControlMessage::Unsubscribe(_) => MessageType::Unsubscribe,
ControlMessage::Announce(_) => MessageType::Announce,
ControlMessage::AnnounceOk(_) => MessageType::AnnounceOk,
ControlMessage::AnnounceError(_) => MessageType::AnnounceError,
ControlMessage::AnnounceCancel(_) => MessageType::AnnounceCancel,
ControlMessage::Unannounce(_) => MessageType::Unannounce,
ControlMessage::SubscribeNamespace(_) => MessageType::SubscribeNamespace,
ControlMessage::SubscribeNamespaceOk(_) => MessageType::SubscribeNamespaceOk,
ControlMessage::SubscribeNamespaceError(_) => MessageType::SubscribeNamespaceError,
ControlMessage::UnsubscribeNamespace(_) => MessageType::UnsubscribeNamespace,
ControlMessage::TrackStatus(_) => MessageType::TrackStatus,
ControlMessage::TrackStatusOk(_) => MessageType::TrackStatusOk,
ControlMessage::TrackStatusError(_) => MessageType::TrackStatusError,
ControlMessage::Fetch(_) => MessageType::Fetch,
ControlMessage::FetchOk(_) => MessageType::FetchOk,
ControlMessage::FetchError(_) => MessageType::FetchError,
ControlMessage::FetchCancel(_) => MessageType::FetchCancel,
ControlMessage::Publish(_) => MessageType::Publish,
ControlMessage::PublishOk(_) => MessageType::PublishOk,
ControlMessage::PublishError(_) => MessageType::PublishError,
}
}
}