crabka_protocol/legacy_compat/
fetch.rs1use crate::{
2 kafka_3_6_2,
3 owned::{fetch_request::FetchRequest, fetch_response::FetchResponse},
4};
5
6impl From<kafka_3_6_2::owned::fetch_request::FetchRequest> for FetchRequest {
9 fn from(l: kafka_3_6_2::owned::fetch_request::FetchRequest) -> Self {
10 Self {
11 replica_id: l.replica_id,
12 max_wait_ms: l.max_wait_ms,
13 min_bytes: l.min_bytes,
14 max_bytes: l.max_bytes,
15 isolation_level: l.isolation_level,
16 session_id: l.session_id,
17 session_epoch: l.session_epoch,
18 topics: l.topics.into_iter().map(Into::into).collect(),
19 forgotten_topics_data: l
20 .forgotten_topics_data
21 .into_iter()
22 .map(Into::into)
23 .collect(),
24 rack_id: l.rack_id,
25 cluster_id: l.cluster_id,
26 replica_state: l.replica_state.into(),
27 ..Default::default()
28 }
29 }
30}
31
32impl From<kafka_3_6_2::owned::fetch_request::ReplicaState>
33 for crate::owned::fetch_request::ReplicaState
34{
35 fn from(l: kafka_3_6_2::owned::fetch_request::ReplicaState) -> Self {
36 Self {
37 replica_id: l.replica_id,
38 replica_epoch: l.replica_epoch,
39 ..Default::default()
40 }
41 }
42}
43
44impl From<kafka_3_6_2::owned::fetch_request::FetchTopic>
45 for crate::owned::fetch_request::FetchTopic
46{
47 fn from(l: kafka_3_6_2::owned::fetch_request::FetchTopic) -> Self {
48 Self {
49 topic: l.topic,
50 partitions: l.partitions.into_iter().map(Into::into).collect(),
52 ..Default::default()
53 }
54 }
55}
56
57impl From<kafka_3_6_2::owned::fetch_request::FetchPartition>
58 for crate::owned::fetch_request::FetchPartition
59{
60 fn from(l: kafka_3_6_2::owned::fetch_request::FetchPartition) -> Self {
61 Self {
62 partition: l.partition,
63 current_leader_epoch: l.current_leader_epoch,
64 fetch_offset: l.fetch_offset,
65 last_fetched_epoch: l.last_fetched_epoch,
66 log_start_offset: l.log_start_offset,
67 partition_max_bytes: l.partition_max_bytes,
68 ..Default::default()
70 }
71 }
72}
73
74impl From<kafka_3_6_2::owned::fetch_request::ForgottenTopic>
75 for crate::owned::fetch_request::ForgottenTopic
76{
77 fn from(l: kafka_3_6_2::owned::fetch_request::ForgottenTopic) -> Self {
78 Self {
79 topic: l.topic,
80 partitions: l.partitions,
82 ..Default::default()
83 }
84 }
85}
86
87impl From<FetchResponse> for kafka_3_6_2::owned::fetch_response::FetchResponse {
90 fn from(c: FetchResponse) -> Self {
91 Self {
92 throttle_time_ms: c.throttle_time_ms,
93 error_code: c.error_code,
94 session_id: c.session_id,
95 responses: c.responses.into_iter().map(Into::into).collect(),
96 ..Default::default()
98 }
99 }
100}
101
102impl From<crate::owned::fetch_response::FetchableTopicResponse>
103 for kafka_3_6_2::owned::fetch_response::FetchableTopicResponse
104{
105 fn from(c: crate::owned::fetch_response::FetchableTopicResponse) -> Self {
106 Self {
107 topic: c.topic,
108 topic_id: c.topic_id,
109 partitions: c.partitions.into_iter().map(Into::into).collect(),
110 ..Default::default()
111 }
112 }
113}
114
115impl From<crate::owned::fetch_response::PartitionData>
116 for kafka_3_6_2::owned::fetch_response::PartitionData
117{
118 fn from(c: crate::owned::fetch_response::PartitionData) -> Self {
119 Self {
120 partition_index: c.partition_index,
121 error_code: c.error_code,
122 high_watermark: c.high_watermark,
123 last_stable_offset: c.last_stable_offset,
124 log_start_offset: c.log_start_offset,
125 aborted_transactions: c
126 .aborted_transactions
127 .map(|v| v.into_iter().map(Into::into).collect()),
128 preferred_read_replica: c.preferred_read_replica,
129 records: c.records,
130 diverging_epoch: c.diverging_epoch.into(),
131 current_leader: c.current_leader.into(),
132 snapshot_id: c.snapshot_id.into(),
133 ..Default::default()
134 }
135 }
136}
137
138impl From<crate::owned::fetch_response::AbortedTransaction>
139 for kafka_3_6_2::owned::fetch_response::AbortedTransaction
140{
141 fn from(c: crate::owned::fetch_response::AbortedTransaction) -> Self {
142 Self {
143 producer_id: c.producer_id,
144 first_offset: c.first_offset,
145 ..Default::default()
146 }
147 }
148}
149
150impl From<crate::owned::fetch_response::EpochEndOffset>
151 for kafka_3_6_2::owned::fetch_response::EpochEndOffset
152{
153 fn from(c: crate::owned::fetch_response::EpochEndOffset) -> Self {
154 Self {
155 epoch: c.epoch,
156 end_offset: c.end_offset,
157 ..Default::default()
158 }
159 }
160}
161
162impl From<crate::owned::fetch_response::LeaderIdAndEpoch>
163 for kafka_3_6_2::owned::fetch_response::LeaderIdAndEpoch
164{
165 fn from(c: crate::owned::fetch_response::LeaderIdAndEpoch) -> Self {
166 Self {
167 leader_id: c.leader_id,
168 leader_epoch: c.leader_epoch,
169 ..Default::default()
170 }
171 }
172}
173
174impl From<crate::owned::fetch_response::SnapshotId>
175 for kafka_3_6_2::owned::fetch_response::SnapshotId
176{
177 fn from(c: crate::owned::fetch_response::SnapshotId) -> Self {
178 Self {
179 end_offset: c.end_offset,
180 epoch: c.epoch,
181 ..Default::default()
182 }
183 }
184}
185
186#[cfg(test)]
187mod tests {
188 use assert2::assert;
189 use bytes::Bytes;
190
191 use super::*;
192 use crate::{UnknownTaggedFields, primitives::uuid::Uuid, records::RecordsPayload};
193
194 #[test]
195 fn legacy_fetch_request_conversion_preserves_mapped_fields() {
196 let mut legacy = kafka_3_6_2::owned::fetch_request::FetchRequest::populated(15);
197 legacy.replica_id = 42;
198 legacy.max_wait_ms = 43;
199 legacy.min_bytes = 44;
200 legacy.max_bytes = 45;
201 legacy.isolation_level = 1;
202 legacy.session_id = 46;
203 legacy.session_epoch = 47;
204 legacy.rack_id = "rack-a".into();
205 legacy.cluster_id = Some("cluster-a".into());
206 legacy.replica_state.replica_id = 48;
207 legacy.replica_state.replica_epoch = 49;
208 legacy.topics[0].topic = "fetch-topic".into();
209 legacy.topics[0].partitions[0].partition = 3;
210 legacy.topics[0].partitions[0].current_leader_epoch = 4;
211 legacy.topics[0].partitions[0].fetch_offset = 5;
212 legacy.topics[0].partitions[0].last_fetched_epoch = 6;
213 legacy.topics[0].partitions[0].log_start_offset = 7;
214 legacy.topics[0].partitions[0].partition_max_bytes = 8;
215 legacy.forgotten_topics_data[0].topic = "forgotten-topic".into();
216 legacy.forgotten_topics_data[0].partitions = vec![9, 10];
217 let converted = FetchRequest::from(legacy);
218
219 let expected = FetchRequest {
220 replica_id: 42,
221 max_wait_ms: 43,
222 min_bytes: 44,
223 max_bytes: 45,
224 isolation_level: 1,
225 session_id: 46,
226 session_epoch: 47,
227 topics: vec![crate::owned::fetch_request::FetchTopic {
228 topic: "fetch-topic".to_string(),
229 topic_id: Uuid::ZERO,
230 partitions: vec![crate::owned::fetch_request::FetchPartition {
231 partition: 3,
232 current_leader_epoch: 4,
233 fetch_offset: 5,
234 last_fetched_epoch: 6,
235 log_start_offset: 7,
236 partition_max_bytes: 8,
237 replica_directory_id: Uuid::ZERO,
238 high_watermark: 9_223_372_036_854_775_807,
239 unknown_tagged_fields: UnknownTaggedFields(vec![]),
240 }],
241 unknown_tagged_fields: UnknownTaggedFields(vec![]),
242 }],
243 forgotten_topics_data: vec![crate::owned::fetch_request::ForgottenTopic {
244 topic: "forgotten-topic".to_string(),
245 topic_id: Uuid::ZERO,
246 partitions: vec![9, 10],
247 unknown_tagged_fields: UnknownTaggedFields(vec![]),
248 }],
249 rack_id: "rack-a".to_string(),
250 cluster_id: Some("cluster-a".to_string()),
251 replica_state: crate::owned::fetch_request::ReplicaState {
252 replica_id: 48,
253 replica_epoch: 49,
254 unknown_tagged_fields: UnknownTaggedFields(vec![]),
255 },
256 unknown_tagged_fields: UnknownTaggedFields(vec![]),
257 };
258 assert!(converted == expected);
259 }
260
261 #[test]
262 fn fetch_response_conversion_preserves_mapped_fields() {
263 let mut canonical = FetchResponse::populated(18);
264 canonical.throttle_time_ms = 51;
265 canonical.error_code = 52;
266 canonical.session_id = 53;
267 canonical.responses[0].topic = "fetch-response-topic".into();
268 canonical.responses[0].topic_id = Uuid([0x12; 16]);
269 canonical.responses[0].partitions[0].partition_index = 54;
270 canonical.responses[0].partitions[0].error_code = 55;
271 canonical.responses[0].partitions[0].high_watermark = 56;
272 canonical.responses[0].partitions[0].last_stable_offset = 57;
273 canonical.responses[0].partitions[0].log_start_offset = 58;
274 canonical.responses[0].partitions[0].preferred_read_replica = 59;
275 canonical.responses[0].partitions[0].records =
276 Some(RecordsPayload::Legacy(Bytes::from_static(&[1, 2, 3])));
277 canonical.responses[0].partitions[0].diverging_epoch.epoch = 60;
278 canonical.responses[0].partitions[0]
279 .diverging_epoch
280 .end_offset = 61;
281 canonical.responses[0].partitions[0]
282 .current_leader
283 .leader_id = 62;
284 canonical.responses[0].partitions[0]
285 .current_leader
286 .leader_epoch = 63;
287 canonical.responses[0].partitions[0].snapshot_id.end_offset = 64;
288 canonical.responses[0].partitions[0].snapshot_id.epoch = 65;
289 canonical.responses[0].partitions[0]
290 .aborted_transactions
291 .as_mut()
292 .expect("aborted transactions")[0]
293 .producer_id = 66;
294 canonical.responses[0].partitions[0]
295 .aborted_transactions
296 .as_mut()
297 .expect("aborted transactions")[0]
298 .first_offset = 67;
299 let converted = kafka_3_6_2::owned::fetch_response::FetchResponse::from(canonical);
300
301 let expected = kafka_3_6_2::owned::fetch_response::FetchResponse {
302 throttle_time_ms: 51,
303 error_code: 52,
304 session_id: 53,
305 responses: vec![kafka_3_6_2::owned::fetch_response::FetchableTopicResponse {
306 topic: "fetch-response-topic".to_string(),
307 topic_id: Uuid([0x12; 16]),
308 partitions: vec![kafka_3_6_2::owned::fetch_response::PartitionData {
309 partition_index: 54,
310 error_code: 55,
311 high_watermark: 56,
312 last_stable_offset: 57,
313 log_start_offset: 58,
314 aborted_transactions: Some(vec![
315 kafka_3_6_2::owned::fetch_response::AbortedTransaction {
316 producer_id: 66,
317 first_offset: 67,
318 unknown_tagged_fields: UnknownTaggedFields(vec![]),
319 },
320 ]),
321 preferred_read_replica: 59,
322 records: Some(RecordsPayload::Legacy(Bytes::from_static(&[1, 2, 3]))),
323 diverging_epoch: kafka_3_6_2::owned::fetch_response::EpochEndOffset {
324 epoch: 60,
325 end_offset: 61,
326 unknown_tagged_fields: UnknownTaggedFields(vec![]),
327 },
328 current_leader: kafka_3_6_2::owned::fetch_response::LeaderIdAndEpoch {
329 leader_id: 62,
330 leader_epoch: 63,
331 unknown_tagged_fields: UnknownTaggedFields(vec![]),
332 },
333 snapshot_id: kafka_3_6_2::owned::fetch_response::SnapshotId {
334 end_offset: 64,
335 epoch: 65,
336 unknown_tagged_fields: UnknownTaggedFields(vec![]),
337 },
338 unknown_tagged_fields: UnknownTaggedFields(vec![]),
339 }],
340 unknown_tagged_fields: UnknownTaggedFields(vec![]),
341 }],
342 unknown_tagged_fields: UnknownTaggedFields(vec![]),
343 };
344 assert!(converted == expected);
345 }
346}