1use std::collections::BTreeMap;
2use std::sync::Arc;
3
4use bytes::Bytes;
5use tokio::sync::{Mutex, OnceCell};
6use xet_client::cas_client::{Client, ProgressCallback};
7use xet_client::cas_types::ChunkRange;
8#[cfg(not(target_family = "wasm"))]
9use xet_client::cas_types::Key;
10use xet_client::chunk_cache::ChunkCache;
11use xet_core_structures::merklehash::MerkleHash;
12use xet_runtime::core::XetContext;
13use xet_runtime::utils::UniqueId;
14
15use super::super::error::Result;
16use super::retrieval_urls::{TermBlockRetrievalURLs, XorbURLProvider};
17use crate::progress_tracking::ItemProgressUpdater;
18
19pub struct XorbBlockData {
25 pub chunk_offsets: Vec<(usize, usize)>,
29
30 pub data: Bytes,
32}
33
34#[derive(Debug)]
38pub struct XorbReference {
39 pub term_chunks: ChunkRange,
41 pub uncompressed_size: usize,
43}
44
45pub struct XorbBlock {
52 pub xorb_hash: MerkleHash,
53 pub chunk_ranges: Vec<ChunkRange>,
56 pub xorb_block_index: usize,
58 pub references: Vec<XorbReference>,
61 pub uncompressed_size_if_known: Option<usize>,
64 pub data: OnceCell<Arc<XorbBlockData>>,
65}
66
67impl PartialEq for XorbBlock {
68 fn eq(&self, other: &Self) -> bool {
69 self.xorb_hash == other.xorb_hash
70 && self.chunk_ranges == other.chunk_ranges
71 && self.xorb_block_index == other.xorb_block_index
72 }
73}
74
75impl Eq for XorbBlock {}
76
77fn build_chunk_offsets(chunk_ranges: &[ChunkRange], byte_offsets: &[u32]) -> Vec<(usize, usize)> {
79 let mut chunk_offsets = Vec::new();
80 let mut offset_idx = 0;
81 for range in chunk_ranges {
82 for chunk_idx in range.start..range.end {
83 chunk_offsets.push((chunk_idx as usize, byte_offsets[offset_idx] as usize));
84 offset_idx += 1;
85 }
86 }
87 chunk_offsets
88}
89
90impl XorbBlock {
91 pub async fn retrieve_data(
98 self: Arc<Self>,
99 ctx: XetContext,
100 client: Arc<dyn Client>,
101 url_info: Arc<TermBlockRetrievalURLs>,
102 progress_updater: Option<Arc<ItemProgressUpdater>>,
103 #[cfg_attr(target_family = "wasm", allow(unused_variables))] chunk_cache: Option<Arc<dyn ChunkCache>>,
104 ) -> Result<Arc<XorbBlockData>> {
105 let xorb_block_index = self.xorb_block_index;
106 let uncompressed_size_if_known = self.uncompressed_size_if_known;
107 let chunk_ranges = self.chunk_ranges.clone();
108
109 self.data
110 .get_or_try_init(|| async {
111 #[cfg(not(target_family = "wasm"))]
117 if let Some(ref cache) = chunk_cache {
118 let cache_key = Key {
119 prefix: ctx.config.data.default_prefix.clone(),
120 hash: self.xorb_hash,
121 };
122 let chunk_range = chunk_ranges.first().copied().unwrap_or_default();
123
124 if let Ok(Some(cache_range)) = cache.get(&cache_key, &chunk_range).await {
125 if let Some(ref updater) = progress_updater {
127 let (_, _, http_ranges) = url_info.get_retrieval_url(xorb_block_index).await;
128 let transfer_bytes: u64 = http_ranges.iter().map(|r| r.length()).sum();
129 updater.report_transfer_progress(transfer_bytes);
130 }
131 let chunk_offsets = build_chunk_offsets(&chunk_ranges, &cache_range.offsets);
132 let data = Bytes::from(cache_range.data);
133 return Ok(Arc::new(XorbBlockData { chunk_offsets, data }));
134 }
135 }
136
137 let permit = client.acquire_download_permit().await?;
139
140 let url_provider = XorbURLProvider {
141 ctx: ctx.clone(),
142 client: client.clone(),
143 url_info,
144 xorb_block_index,
145 last_acquisition_id: Mutex::new(UniqueId::null()),
146 };
147
148 let progress_callback: Option<ProgressCallback> = progress_updater.as_ref().map(|updater| {
151 let updater = updater.clone();
152 Arc::new(move |delta: u64, _completed: u64, _total: u64| {
153 updater.report_transfer_progress(delta);
154 }) as ProgressCallback
155 });
156
157 let (data, chunk_byte_offsets) = client
158 .get_file_term_data(Box::new(url_provider), permit, progress_callback, uncompressed_size_if_known)
159 .await?;
160
161 #[cfg(not(target_family = "wasm"))]
163 if let Some(cache) = chunk_cache {
164 let cache_key = Key {
165 prefix: ctx.config.data.default_prefix.clone(),
166 hash: self.xorb_hash,
167 };
168 let chunk_range = chunk_ranges.first().copied().unwrap_or_default();
169 let data = data.clone();
170 let chunk_byte_offsets = chunk_byte_offsets.clone();
171 tokio::spawn(async move {
172 if let Err(err) = cache.put(&cache_key, &chunk_range, &chunk_byte_offsets, &data).await {
173 tracing::warn!("chunk cache put failed: {err}");
174 }
175 });
176 }
177
178 let chunk_offsets = build_chunk_offsets(&chunk_ranges, &chunk_byte_offsets);
179
180 Ok(Arc::new(XorbBlockData { chunk_offsets, data }))
181 })
182 .await
183 .cloned()
184 }
185
186 pub fn determine_size_if_possible(xorb_ranges: &[ChunkRange], terms: &[XorbReference]) -> Option<usize> {
202 debug_assert!(
203 terms.windows(2).all(|w| w[0].term_chunks.start <= w[1].term_chunks.start),
204 "terms must be sorted by chunk range start"
205 );
206
207 debug_assert!(
208 terms.iter().all(|term| xorb_ranges
209 .iter()
210 .any(|r| term.term_chunks.start >= r.start && term.term_chunks.end <= r.end)),
211 "all terms must fall within one of the xorb ranges"
212 );
213
214 if xorb_ranges.is_empty() {
215 return Some(0);
216 }
217
218 let gap_bridges: BTreeMap<u32, u32> = xorb_ranges
222 .windows(2)
223 .filter(|pair| pair[0].end < pair[1].start)
224 .map(|pair| (pair[0].end, pair[1].start))
225 .collect();
226
227 let mut reachable: BTreeMap<u32, usize> = BTreeMap::new();
230 reachable.insert(xorb_ranges[0].start, 0);
231
232 for term in terms {
234 if let Some(&accumulated) = reachable.get(&term.term_chunks.start) {
235 let new_end = term.term_chunks.end;
236 let new_size = accumulated + term.uncompressed_size;
237
238 reachable.entry(new_end).or_insert(new_size);
239
240 if let Some(&bridge_target) = gap_bridges.get(&new_end) {
243 reachable.entry(bridge_target).or_insert(new_size);
244 }
245 }
246 }
247
248 reachable.get(&xorb_ranges.last().unwrap().end).copied()
250 }
251}
252
253#[cfg(test)]
254mod tests {
255 use xet_client::cas_types::ChunkRange;
256
257 use super::*;
258
259 fn build_refs(pairs: &[(ChunkRange, usize)]) -> Vec<XorbReference> {
260 pairs
261 .iter()
262 .map(|(range, size)| XorbReference {
263 term_chunks: *range,
264 uncompressed_size: *size,
265 })
266 .collect()
267 }
268
269 #[test]
270 fn test_single_term_exact_match() {
271 let ranges = &[ChunkRange::new(0, 5)];
272 let terms = build_refs(&[(ChunkRange::new(0, 5), 1000)]);
273 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), Some(1000));
274 }
275
276 #[test]
277 fn test_two_terms_chained() {
278 let ranges = &[ChunkRange::new(0, 5)];
279 let terms = build_refs(&[(ChunkRange::new(0, 3), 600), (ChunkRange::new(3, 5), 400)]);
280 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), Some(1000));
281 }
282
283 #[test]
284 fn test_three_terms_chained() {
285 let ranges = &[ChunkRange::new(0, 6)];
286 let terms = build_refs(&[
287 (ChunkRange::new(0, 2), 200),
288 (ChunkRange::new(2, 4), 300),
289 (ChunkRange::new(4, 6), 500),
290 ]);
291 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), Some(1000));
292 }
293
294 #[test]
295 fn test_gap_in_chain() {
296 let ranges = &[ChunkRange::new(0, 6)];
297 let terms = build_refs(&[(ChunkRange::new(0, 2), 200), (ChunkRange::new(4, 6), 500)]);
298 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), None);
299 }
300
301 #[test]
302 fn test_does_not_start_at_xorb_start() {
303 let ranges = &[ChunkRange::new(0, 5)];
304 let terms = build_refs(&[(ChunkRange::new(1, 5), 800)]);
305 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), None);
306 }
307
308 #[test]
309 fn test_does_not_end_at_xorb_end() {
310 let ranges = &[ChunkRange::new(0, 5)];
311 let terms = build_refs(&[(ChunkRange::new(0, 3), 600)]);
312 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), None);
313 }
314
315 #[test]
316 fn test_empty_terms() {
317 let ranges = &[ChunkRange::new(0, 5)];
318 let terms: Vec<XorbReference> = vec![];
319 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), None);
320 }
321
322 #[test]
323 fn test_overlapping_terms_with_exact_cover() {
324 let ranges = &[ChunkRange::new(0, 5)];
327 let terms = build_refs(&[
328 (ChunkRange::new(0, 3), 600),
329 (ChunkRange::new(1, 4), 700),
330 (ChunkRange::new(3, 5), 400),
331 ]);
332 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), Some(1000));
333 }
334
335 #[test]
336 fn test_duplicate_terms_first_covers() {
337 let ranges = &[ChunkRange::new(0, 5)];
339 let terms = build_refs(&[(ChunkRange::new(0, 5), 1000), (ChunkRange::new(0, 5), 1000)]);
340 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), Some(1000));
341 }
342
343 #[test]
344 fn test_nonzero_xorb_start() {
345 let ranges = &[ChunkRange::new(3, 8)];
346 let terms = build_refs(&[(ChunkRange::new(3, 5), 400), (ChunkRange::new(5, 8), 600)]);
347 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), Some(1000));
348 }
349
350 #[test]
351 fn test_nonzero_xorb_start_no_match() {
352 let ranges = &[ChunkRange::new(3, 8)];
353 let terms = build_refs(&[(ChunkRange::new(3, 5), 400)]);
354 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), None);
355 }
356
357 #[test]
358 fn test_single_chunk_range() {
359 let ranges = &[ChunkRange::new(0, 1)];
360 let terms = build_refs(&[(ChunkRange::new(0, 1), 42)]);
361 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), Some(42));
362 }
363
364 #[test]
365 fn test_chain_with_overlapping_inner_terms() {
366 let ranges = &[ChunkRange::new(2, 8)];
367 let terms = build_refs(&[
370 (ChunkRange::new(2, 5), 500),
371 (ChunkRange::new(3, 6), 999),
372 (ChunkRange::new(5, 8), 300),
373 ]);
374 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), Some(800));
375 }
376
377 #[test]
378 fn test_partial_overlap_no_cover() {
379 let ranges = &[ChunkRange::new(0, 10)];
381 let terms = build_refs(&[
382 (ChunkRange::new(0, 4), 400),
383 (ChunkRange::new(3, 7), 400),
384 (ChunkRange::new(6, 10), 400),
385 ]);
386 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), None);
387 }
388
389 #[test]
390 fn test_same_start_short_then_long_covering_full() {
391 let ranges = &[ChunkRange::new(0, 5)];
393 let terms = build_refs(&[(ChunkRange::new(0, 3), 300), (ChunkRange::new(0, 5), 500)]);
394 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), Some(500));
395 }
396
397 #[test]
398 fn test_same_start_short_then_long_with_chain() {
399 let ranges = &[ChunkRange::new(0, 6)];
402 let terms = build_refs(&[
403 (ChunkRange::new(0, 2), 200),
404 (ChunkRange::new(0, 3), 300),
405 (ChunkRange::new(3, 6), 300),
406 ]);
407 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), Some(600));
408 }
409
410 #[test]
411 fn test_same_start_multiple_duplicates_chain_through_second() {
412 let ranges = &[ChunkRange::new(0, 6)];
415 let terms = build_refs(&[
416 (ChunkRange::new(0, 2), 200),
417 (ChunkRange::new(0, 4), 400),
418 (ChunkRange::new(0, 5), 500),
419 (ChunkRange::new(4, 6), 200),
420 ]);
421 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), Some(600));
422 }
423
424 #[test]
425 fn test_same_start_at_midpoint() {
426 let ranges = &[ChunkRange::new(0, 8)];
429 let terms = build_refs(&[
430 (ChunkRange::new(0, 3), 300),
431 (ChunkRange::new(3, 5), 200),
432 (ChunkRange::new(3, 6), 300),
433 (ChunkRange::new(6, 8), 200),
434 ]);
435 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), Some(800));
436 }
437
438 #[test]
439 fn test_same_start_none_covers() {
440 let ranges = &[ChunkRange::new(0, 10)];
442 let terms = build_refs(&[
443 (ChunkRange::new(0, 2), 200),
444 (ChunkRange::new(0, 4), 400),
445 (ChunkRange::new(0, 6), 600),
446 ]);
447 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), None);
448 }
449
450 #[test]
451 fn test_same_start_two_groups_chained() {
452 let ranges = &[ChunkRange::new(0, 6)];
455 let terms = build_refs(&[
456 (ChunkRange::new(0, 2), 200),
457 (ChunkRange::new(0, 3), 300),
458 (ChunkRange::new(3, 5), 200),
459 (ChunkRange::new(3, 6), 300),
460 ]);
461 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), Some(600));
462 }
463
464 #[test]
465 fn test_multiple_disjoint_ranges_both_covered() {
466 let ranges = &[ChunkRange::new(0, 3), ChunkRange::new(5, 8)];
467 let terms = build_refs(&[(ChunkRange::new(0, 3), 300), (ChunkRange::new(5, 8), 400)]);
468 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), Some(700));
469 }
470
471 #[test]
472 fn test_multiple_disjoint_ranges_one_uncovered() {
473 let ranges = &[ChunkRange::new(0, 3), ChunkRange::new(5, 8)];
474 let terms = build_refs(&[(ChunkRange::new(0, 3), 300)]);
475 assert_eq!(XorbBlock::determine_size_if_possible(ranges, &terms), None);
476 }
477}