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>),
}
#[derive(Debug, Clone)]
pub struct SimpleCommand {
pub message: BaseCommand,
}
#[derive(Debug, Clone)]
pub struct PayloadCommand {
pub message: BaseCommand,
pub metadata: MessageMetadata,
pub payload: PayloadCommandPayload,
}
#[derive(Debug, Clone)]
pub enum PayloadCommandPayload {
Single(Vec<u8>),
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
}
}
}
}