Skip to main content

crabka_protocol/legacy_compat/
fetch.rs

1use crate::{
2    kafka_3_6_2,
3    owned::{fetch_request::FetchRequest, fetch_response::FetchResponse},
4};
5
6// ── Request: legacy → canonical ──────────────────────────────────────────────
7
8impl 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            // topic_id (v13+) defaults to Uuid::nil()
51            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            // replica_directory_id (v17+ tagged) and high_watermark (v18+ tagged) default
69            ..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            // topic_id (v13+) defaults to Uuid::nil()
81            partitions: l.partitions,
82            ..Default::default()
83        }
84    }
85}
86
87// ── Response: canonical → legacy ─────────────────────────────────────────────
88
89impl 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            // node_endpoints (canonical v16+ tagged field) dropped — not present in 3.6.2 schema
97            ..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}