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