pulsar_binary_protocol_spec/
command.rs

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
            }
        }
    }
}