Skip to main content

crabka_protocol/legacy_compat/
fetch.rs

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