1
  2
  3
  4
  5
  6
  7
  8
  9
 10
 11
 12
 13
 14
 15
 16
 17
 18
 19
 20
 21
 22
 23
 24
 25
 26
 27
 28
 29
 30
 31
 32
 33
 34
 35
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
use crate::{
    frame::{FrameParseBatchPayloadError, FrameParseSinglePayloadError},
    protos::protobuf::pulsar_api::{BaseCommand, MessageMetadata, SingleMessageMetadata},
    types::AckValidationError,
};

#[derive(Debug, Clone)]
pub enum Command {
    Simple(SimpleCommand),
    Payload(Box<PayloadCommand>),
}

#[derive(Debug, Clone)]
pub enum CommandWithParsed {
    Simple(SimpleCommand),
    Payload(Box<PayloadCommandWithParsed>),
}

// http://pulsar.apache.org/docs/en/develop-binary-protocol/#simple-commands
#[derive(Debug, Clone)]
pub struct SimpleCommand {
    pub message: BaseCommand,
}

// http://pulsar.apache.org/docs/en/develop-binary-protocol/#payload-commands
#[derive(Debug, Clone)]
pub struct PayloadCommand {
    pub message: BaseCommand,
    pub metadata: MessageMetadata,
    pub payload: PayloadCommandPayload,
}

#[derive(Debug, Clone)]
pub enum PayloadCommandPayload {
    Single(Vec<u8>),
    // http://pulsar.apache.org/docs/en/develop-binary-protocol/#batch-messages
    Batch(Vec<(SingleMessageMetadata, Vec<u8>)>),
}

#[derive(Debug, Clone)]
pub struct PayloadCommandWithParsed {
    pub message: BaseCommand,
    pub metadata: MessageMetadata,
    pub payload: PayloadCommandPayloadWithParsed,
    pub is_checksum_match: Option<bool>,
}

#[derive(Debug, Clone)]
pub enum PayloadCommandPayloadWithParsed {
    Single(Result<Vec<u8>, PayloadCommandPayloadErrorWithParsed>),
    Batch(Result<Vec<(SingleMessageMetadata, Vec<u8>)>, PayloadCommandPayloadErrorWithParsed>),
}

#[derive(PartialEq, Eq, Debug, Clone)]
pub enum PayloadCommandPayloadErrorWithParsed {
    DecompressionError,
    UncompressedSizeCorruption,
    BatchDeSerializeError,
}
impl From<FrameParseSinglePayloadError> for PayloadCommandPayloadErrorWithParsed {
    fn from(err: FrameParseSinglePayloadError) -> Self {
        match err {
            FrameParseSinglePayloadError::CompressionUnsupported { type_code: _ } => {
                Self::DecompressionError
            }
            #[cfg(feature = "with-compression-lz4")]
            FrameParseSinglePayloadError::CompressionLZ4DecompressError(_) => {
                Self::DecompressionError
            }
            #[cfg(feature = "with-compression-zlib")]
            FrameParseSinglePayloadError::CompressionZlibDecompressError(_) => {
                Self::DecompressionError
            }
            FrameParseSinglePayloadError::UncompressedSizeCorruption => {
                Self::UncompressedSizeCorruption
            }
        }
    }
}
impl From<FrameParseBatchPayloadError> for PayloadCommandPayloadErrorWithParsed {
    fn from(err: FrameParseBatchPayloadError) -> Self {
        match err {
            FrameParseBatchPayloadError::CompressionUnsupported { type_code: _ } => {
                Self::DecompressionError
            }
            #[cfg(feature = "with-compression-lz4")]
            FrameParseBatchPayloadError::CompressionLZ4DecompressError(_) => {
                Self::DecompressionError
            }
            #[cfg(feature = "with-compression-zlib")]
            FrameParseBatchPayloadError::CompressionZlibDecompressError(_) => {
                Self::DecompressionError
            }
            FrameParseBatchPayloadError::UncompressedSizeCorruption => {
                Self::UncompressedSizeCorruption
            }
            FrameParseBatchPayloadError::GetSingleMessageMetadataSizeFailed => {
                Self::BatchDeSerializeError
            }
            FrameParseBatchPayloadError::GetSingleMessageMetadataFailed => {
                Self::BatchDeSerializeError
            }
            FrameParseBatchPayloadError::DeserializeSingleMessageMetadataError(_) => {
                Self::BatchDeSerializeError
            }
            FrameParseBatchPayloadError::GetSingleMessagePayloadFailed => {
                Self::BatchDeSerializeError
            }
        }
    }
}
impl From<&PayloadCommandPayloadErrorWithParsed> for AckValidationError {
    fn from(err: &PayloadCommandPayloadErrorWithParsed) -> Self {
        match err {
            PayloadCommandPayloadErrorWithParsed::DecompressionError => Self::DecompressionError,
            PayloadCommandPayloadErrorWithParsed::UncompressedSizeCorruption => {
                Self::UncompressedSizeCorruption
            }
            PayloadCommandPayloadErrorWithParsed::BatchDeSerializeError => {
                Self::BatchDeSerializeError
            }
        }
    }
}