1use std::collections::{HashMap, HashSet};
2use std::io::{BufReader, Cursor};
3use std::ops::Range;
4use std::sync::Arc;
5use std::sync::atomic::{AtomicBool, AtomicU16, AtomicU64, AtomicUsize, Ordering};
6
7use async_trait::async_trait;
8use bytes::Bytes;
9use rand::RngExt;
10use tokio::sync::RwLock;
11use tokio::time::{Duration, Instant};
12use tracing::{error, info};
13use xet_core_structures::MerkleHashMap;
14use xet_core_structures::merklehash::MerkleHash;
15#[cfg(not(target_family = "wasm"))]
16use xet_core_structures::merklehash::compute_data_hash;
17use xet_core_structures::metadata_shard::file_structs::MDBFileInfo;
18use xet_core_structures::metadata_shard::shard_in_memory::MDBInMemoryShard;
19use xet_core_structures::metadata_shard::streaming_shard::MDBMinimalShard;
20use xet_core_structures::metadata_shard::xorb_structs::MDBXorbInfo;
21use xet_core_structures::xorb_object::{SerializedXorbObject, XorbObject};
22use xet_runtime::core::XetContext;
23
24use super::super::Client;
25use super::super::adaptive_concurrency::AdaptiveConcurrencyController;
26use super::super::progress_tracked_streams::ProgressCallback;
27use super::client_testing_utils::{FileTermReference, RandomFileContents};
28#[cfg(not(target_family = "wasm"))]
29use super::deletion_controls::ObjectTag;
30use super::direct_access_client::DirectAccessClient;
31use super::random_xorb::RandomXorb;
32use super::xorb_utils::{self, REFERENCE_INSTANT, duration_to_expiration_secs_ceil};
33use crate::cas_client::chunk_window_builder::build_file_chunk_hashes_response;
34use crate::cas_types::{
35 BatchQueryReconstructionResponse, FileChunkHashesResponse, FileRange, HexMerkleHash, HttpRange,
36 QueryReconstructionResponse, QueryReconstructionResponseV2, XorbMultiRangeFetch, XorbRangeDescriptor,
37 XorbReconstructionFetchInfo,
38};
39use crate::error::{ClientError, Result};
40
41struct MaterializedXorb {
43 serialized_data: Bytes,
44 xorb_object: XorbObject,
45}
46
47enum XorbStorage {
52 Materialized { entry: MaterializedXorb, generation: u64 },
53 Random { xorb: RandomXorb, generation: u64 },
54}
55
56pub struct MemoryClient {
58 xorbs: RwLock<MerkleHashMap<XorbStorage>>,
60 shard: RwLock<MDBInMemoryShard>,
62 global_dedup: RwLock<MerkleHashMap<Bytes>>,
64 upload_concurrency_controller: Arc<AdaptiveConcurrencyController>,
66 xorb_generation: AtomicU64,
68 url_expiration_ms: AtomicU64,
70 global_dedup_expiration_secs: AtomicU64,
72 random_ms_delay_window: (AtomicU64, AtomicU64),
74 max_ranges_per_fetch: AtomicUsize,
76 v2_disabled_status: AtomicU16,
78 lifecycle_tag_deletion: AtomicBool,
85 gc_tagged_xorbs: RwLock<HashSet<MerkleHash>>,
88 gc_tagged_shard: RwLock<Option<MerkleHash>>,
91}
92
93impl MemoryClient {
94 pub fn new(ctx: XetContext) -> Arc<Self> {
96 Arc::new(Self {
97 xorbs: RwLock::new(MerkleHashMap::new()),
98 shard: RwLock::new(MDBInMemoryShard::default()),
99 global_dedup: RwLock::new(MerkleHashMap::new()),
100 upload_concurrency_controller: AdaptiveConcurrencyController::new_upload(ctx, "memory_uploads"),
101 xorb_generation: AtomicU64::new(0),
102 url_expiration_ms: AtomicU64::new(u64::MAX),
103 global_dedup_expiration_secs: AtomicU64::new(0),
104 random_ms_delay_window: (AtomicU64::new(0), AtomicU64::new(0)),
105 max_ranges_per_fetch: AtomicUsize::new(usize::MAX),
106 v2_disabled_status: AtomicU16::new(0),
107 lifecycle_tag_deletion: AtomicBool::new(false),
108 gc_tagged_xorbs: RwLock::new(HashSet::new()),
109 gc_tagged_shard: RwLock::new(None),
110 })
111 }
112
113 pub fn set_lifecycle_tag_deletion(&self, on: bool) {
115 self.lifecycle_tag_deletion.store(on, Ordering::Relaxed);
116 }
117
118 fn lifecycle_tag_deletion_enabled(&self) -> bool {
119 self.lifecycle_tag_deletion.load(Ordering::Relaxed)
120 }
121
122 async fn xorb_is_tagged(&self, hash: &MerkleHash) -> bool {
123 self.gc_tagged_xorbs.read().await.contains(hash)
124 }
125
126 async fn shard_is_tagged(&self, hash: &MerkleHash) -> bool {
127 self.gc_tagged_shard.read().await.is_some_and(|h| &h == hash)
128 }
129
130 pub async fn insert_random_xorb(&self, xorb: RandomXorb) -> Result<MerkleHash> {
136 use xet_core_structures::metadata_shard::xorb_structs::{
137 MDBXorbInfo, XorbChunkSequenceEntry, XorbChunkSequenceHeader,
138 };
139
140 let hash = xorb.xorb_hash();
141 let xorb_obj = xorb.get_xorb_object();
142
143 let mut chunk_entries = Vec::with_capacity(xorb.num_chunks() as usize);
145 let mut byte_offset = 0u32;
146
147 for i in 0..xorb.num_chunks() {
148 let chunk_size = xorb.chunk_size(i).unwrap();
149 chunk_entries.push(XorbChunkSequenceEntry {
150 chunk_hash: xorb.chunk_hash(i).unwrap(),
151 chunk_byte_range_start: byte_offset,
152 unpacked_segment_bytes: chunk_size,
153 flags: 0,
154 _unused: 0,
155 });
156 byte_offset += chunk_size;
157 }
158
159 let cas_info = MDBXorbInfo {
160 metadata: XorbChunkSequenceHeader::new(
161 hash,
162 xorb.num_chunks(),
163 xorb_obj.info.unpacked_chunk_offsets.last().copied().unwrap_or(0),
164 ),
165 chunks: chunk_entries,
166 };
167
168 {
169 let mut shard = self.shard.write().await;
170 shard.add_xorb_block(cas_info)?;
171 }
172
173 let generation = self.xorb_generation.fetch_add(1, Ordering::Relaxed);
174 self.xorbs.write().await.insert(hash, XorbStorage::Random { xorb, generation });
175 Ok(hash)
176 }
177
178 pub async fn insert_random_lazy_file(
187 &self,
188 term_spec: &[(u64, (u64, u64))],
189 chunk_size: usize,
190 ) -> Result<RandomFileContents> {
191 use xet_core_structures::metadata_shard::file_structs::{
192 FileDataSequenceEntry, FileDataSequenceHeader, MDBFileInfo,
193 };
194
195 let mut xorb_num_chunks = std::collections::HashMap::<u64, u64>::new();
197 for &(xorb_seed, (_, chunk_idx_end)) in term_spec {
198 let c = xorb_num_chunks.entry(xorb_seed).or_default();
199 *c = (*c).max(chunk_idx_end);
200 }
201
202 let mut random_xorbs = std::collections::HashMap::<u64, RandomXorb>::new();
204 for (&xorb_seed, &n_chunks) in &xorb_num_chunks {
205 let xorb = RandomXorb::from_seed(xorb_seed, n_chunks as u32, chunk_size as u32);
206 self.insert_random_xorb(xorb.clone()).await?;
207 random_xorbs.insert(xorb_seed, xorb);
208 }
209
210 let mut file_segments = Vec::new();
212 let mut chunk_file_hashes = Vec::new();
213 let mut term_infos = Vec::new();
214 let mut file_data = Vec::new();
215
216 for &(xorb_seed, (chunk_start, chunk_end)) in term_spec {
217 let xorb = random_xorbs.get(&xorb_seed).unwrap();
218 let (c_start, c_end) = (chunk_start as u32, chunk_end as u32);
219
220 chunk_file_hashes.extend(xorb.chunk_hash_sizes(c_start, c_end));
221
222 let term_data = xorb.get_chunk_range_data(c_start, c_end).unwrap_or_default();
223 file_data.extend_from_slice(&term_data);
224
225 file_segments.push(FileDataSequenceEntry::new(
226 xorb.xorb_hash(),
227 xorb.chunk_range_size(c_start, c_end) as usize,
228 chunk_start as usize,
229 chunk_end as usize,
230 ));
231
232 term_infos.push(FileTermReference {
233 xorb_hash: xorb.xorb_hash(),
234 chunk_start: c_start,
235 chunk_end: c_end,
236 data: term_data,
237 chunk_hashes: xorb.chunk_hashes_range(c_start, c_end),
238 });
239 }
240
241 let file_hash = xet_core_structures::merklehash::file_hash_with_salt(&chunk_file_hashes, &[0; 32]);
242
243 {
245 let mut shard = self.shard.write().await;
246 shard.add_file_reconstruction_info(MDBFileInfo {
247 metadata: FileDataSequenceHeader::new(file_hash, file_segments.len(), false, false),
248 segments: file_segments,
249 verification: vec![],
250 metadata_ext: None,
251 })?;
252 }
253
254 Ok(RandomFileContents {
255 file_hash,
256 data: Bytes::from(file_data),
257 xorbs: MerkleHashMap::new(),
258 terms: term_infos,
259 })
260 }
261
262 #[cfg(not(target_family = "wasm"))]
263 fn current_shard_hash_and_bytes(shard: &MDBInMemoryShard) -> Result<Option<(MerkleHash, Bytes)>> {
264 if shard.is_empty() {
265 return Ok(None);
266 }
267 let bytes = Bytes::from(shard.to_bytes()?);
268 let hash = compute_data_hash(bytes.as_ref());
269 Ok(Some((hash, bytes)))
270 }
271
272 #[cfg(not(target_family = "wasm"))]
273 fn object_tag_from_key_and_payload(prefix: &[u8], key: &MerkleHash, payload: &[u8]) -> ObjectTag {
274 let key_bytes: [u8; 32] = (*key).into();
275 let payload_hash: [u8; 32] = compute_data_hash(payload).into();
276 let mut entropy = Vec::with_capacity(prefix.len() + key_bytes.len() + payload_hash.len());
277 entropy.extend_from_slice(prefix);
278 entropy.extend_from_slice(&key_bytes);
279 entropy.extend_from_slice(&payload_hash);
280 compute_data_hash(&entropy).into()
281 }
282
283 #[cfg(not(target_family = "wasm"))]
284 fn xorb_tag(hash: &MerkleHash, storage: &XorbStorage) -> ObjectTag {
285 match storage {
286 XorbStorage::Materialized { entry, generation } => {
287 let mut payload = Vec::from(entry.serialized_data.as_ref());
288 payload.extend_from_slice(&generation.to_le_bytes());
289 Self::object_tag_from_key_and_payload(b"xorb", hash, &payload)
290 },
291 XorbStorage::Random { xorb, generation } => {
292 let mut entropy = Vec::with_capacity(16);
293 entropy.extend_from_slice(&xorb.num_chunks().to_le_bytes());
294 entropy.extend_from_slice(&generation.to_le_bytes());
295 Self::object_tag_from_key_and_payload(b"xorb", hash, &entropy)
296 },
297 }
298 }
299}
300
301#[cfg_attr(not(target_family = "wasm"), async_trait)]
302#[cfg_attr(target_family = "wasm", async_trait(?Send))]
303impl DirectAccessClient for MemoryClient {
304 fn set_fetch_term_url_expiration(&self, expiration: Duration) {
305 self.url_expiration_ms.store(expiration.as_millis() as u64, Ordering::Relaxed);
306 }
307
308 fn set_global_dedup_shard_expiration(&self, expiration: Option<Duration>) {
309 self.global_dedup_expiration_secs
310 .store(duration_to_expiration_secs_ceil(expiration), Ordering::Relaxed);
311 }
312
313 fn set_max_ranges_per_fetch(&self, max_ranges: usize) {
314 self.max_ranges_per_fetch.store(max_ranges, Ordering::Relaxed);
315 }
316
317 fn disable_v2_reconstruction(&self, status_code: u16) {
318 self.v2_disabled_status.store(status_code, Ordering::Relaxed);
319 }
320
321 fn v2_disabled_status_code(&self) -> u16 {
322 self.v2_disabled_status.load(Ordering::Relaxed)
323 }
324
325 async fn get_reconstruction_v1(
326 &self,
327 file_id: &MerkleHash,
328 bytes_range: Option<FileRange>,
329 ) -> Result<Option<QueryReconstructionResponse>> {
330 MemoryClient::get_reconstruction_v1(self, file_id, bytes_range).await
331 }
332
333 async fn get_reconstruction_v2(
334 &self,
335 file_id: &MerkleHash,
336 bytes_range: Option<FileRange>,
337 ) -> Result<Option<QueryReconstructionResponseV2>> {
338 MemoryClient::get_reconstruction_v2(self, file_id, bytes_range).await
339 }
340
341 fn set_api_delay_range(&self, delay_range: Option<Range<Duration>>) {
342 match delay_range {
343 Some(range) => {
344 self.random_ms_delay_window
345 .0
346 .store(range.start.as_millis() as u64, Ordering::Relaxed);
347 self.random_ms_delay_window
348 .1
349 .store(range.end.as_millis() as u64, Ordering::Relaxed);
350 },
351 None => {
352 self.random_ms_delay_window.0.store(0, Ordering::Relaxed);
353 self.random_ms_delay_window.1.store(0, Ordering::Relaxed);
354 },
355 }
356 }
357
358 async fn apply_api_delay(&self) {
359 let min_ms = self.random_ms_delay_window.0.load(Ordering::Relaxed);
360 let max_ms = self.random_ms_delay_window.1.load(Ordering::Relaxed);
361
362 if min_ms == 0 && max_ms == 0 {
363 return;
364 }
365
366 let delay_ms = if min_ms == max_ms {
367 min_ms
368 } else {
369 rand::rng().random_range(min_ms..max_ms)
370 };
371
372 tokio::time::sleep(Duration::from_millis(delay_ms)).await;
373 }
374
375 async fn list_xorbs(&self) -> Result<Vec<MerkleHash>> {
376 let tagged = self.gc_tagged_xorbs.read().await;
377 Ok(self
378 .xorbs
379 .read()
380 .await
381 .keys()
382 .copied()
383 .filter(|h| !tagged.contains(h))
384 .collect())
385 }
386
387 async fn get_full_xorb(&self, hash: &MerkleHash) -> Result<Bytes> {
388 if self.xorb_is_tagged(hash).await {
389 return Err(ClientError::XORBNotFound(*hash));
390 }
391 let xorbs = self.xorbs.read().await;
392 let storage = xorbs.get(hash).ok_or_else(|| {
393 error!("Unable to find xorb in memory CAS {:?}", hash);
394 ClientError::XORBNotFound(*hash)
395 })?;
396
397 match storage {
398 XorbStorage::Materialized { entry, .. } => {
399 let mut reader = BufReader::new(Cursor::new(&entry.serialized_data));
400 let xorb_obj = XorbObject::deserialize(&mut reader)?;
401 let result = xorb_obj.get_all_bytes(&mut reader)?;
402 Ok(Bytes::from(result))
403 },
404 XorbStorage::Random { xorb, .. } => xorb
405 .get_chunk_range_data(0, xorb.num_chunks())
406 .ok_or(ClientError::XORBNotFound(*hash)),
407 }
408 }
409
410 async fn get_xorb_ranges(&self, hash: &MerkleHash, chunk_ranges: Vec<(u32, u32)>) -> Result<Vec<Bytes>> {
411 if self.xorb_is_tagged(hash).await {
412 return Err(ClientError::XORBNotFound(*hash));
413 }
414 if chunk_ranges.is_empty() {
415 return Ok(vec![Bytes::new()]);
416 }
417
418 let xorbs = self.xorbs.read().await;
419 let storage = xorbs.get(hash).ok_or_else(|| {
420 error!("Unable to find xorb in memory CAS {:?}", hash);
421 ClientError::XORBNotFound(*hash)
422 })?;
423
424 match storage {
425 XorbStorage::Materialized { entry, .. } => {
426 let mut reader = BufReader::new(Cursor::new(&entry.serialized_data));
427 let xorb_obj = XorbObject::deserialize(&mut reader)?;
428
429 let mut ret: Vec<Bytes> = Vec::new();
430 for r in chunk_ranges {
431 if r.0 >= r.1 {
432 ret.push(Bytes::new());
433 continue;
434 }
435 let data = xorb_obj.get_bytes_by_chunk_range(&mut reader, r.0, r.1)?;
436 ret.push(Bytes::from(data));
437 }
438 Ok(ret)
439 },
440 XorbStorage::Random { xorb, .. } => {
441 let mut ret: Vec<Bytes> = Vec::new();
442 for r in chunk_ranges {
443 if r.0 >= r.1 {
444 ret.push(Bytes::new());
445 continue;
446 }
447 let data = xorb.get_chunk_range_data(r.0, r.1).ok_or(ClientError::XORBNotFound(*hash))?;
448 ret.push(data);
449 }
450 Ok(ret)
451 },
452 }
453 }
454
455 async fn xorb_length(&self, hash: &MerkleHash) -> Result<u32> {
456 if self.xorb_is_tagged(hash).await {
457 return Err(ClientError::XORBNotFound(*hash));
458 }
459 let data = self.get_full_xorb(hash).await?;
460 Ok(data.len() as u32)
461 }
462
463 async fn xorb_exists(&self, hash: &MerkleHash) -> Result<bool> {
464 if self.xorb_is_tagged(hash).await {
465 return Ok(false);
466 }
467 Ok(self.xorbs.read().await.contains_key(hash))
468 }
469
470 async fn xorb_footer(&self, hash: &MerkleHash) -> Result<XorbObject> {
471 if self.xorb_is_tagged(hash).await {
472 return Err(ClientError::XORBNotFound(*hash));
473 }
474 let xorbs = self.xorbs.read().await;
475 let storage = xorbs.get(hash).ok_or_else(|| {
476 error!("Unable to find xorb in memory CAS {:?}", hash);
477 ClientError::XORBNotFound(*hash)
478 })?;
479
480 match storage {
481 XorbStorage::Materialized { entry, .. } => Ok(entry.xorb_object.clone()),
482 XorbStorage::Random { xorb, .. } => Ok(xorb.get_xorb_object()),
483 }
484 }
485
486 async fn get_file_size(&self, hash: &MerkleHash) -> Result<u64> {
487 let shard = self.shard.read().await;
488 #[cfg(not(target_family = "wasm"))]
492 if let Some((shard_hash, _)) = Self::current_shard_hash_and_bytes(&shard)?
493 && self.shard_is_tagged(&shard_hash).await
494 {
495 return Err(ClientError::FileNotFound(*hash));
496 }
497 let file_info = shard
498 .get_file_reconstruction_info(hash)
499 .ok_or(ClientError::FileNotFound(*hash))?;
500 Ok(file_info.file_size())
501 }
502
503 async fn get_file_data(&self, hash: &MerkleHash, byte_range: Option<FileRange>) -> Result<Bytes> {
504 #[cfg(not(target_family = "wasm"))]
508 {
509 let shard = self.shard.read().await;
510 if let Some((shard_hash, _)) = Self::current_shard_hash_and_bytes(&shard)?
511 && self.shard_is_tagged(&shard_hash).await
512 {
513 return Err(ClientError::FileNotFound(*hash));
514 }
515 }
516 let file_info = {
517 let shard = self.shard.read().await;
518 shard
519 .get_file_reconstruction_info(hash)
520 .ok_or(ClientError::FileNotFound(*hash))?
521 };
522
523 let mut file_vec = Vec::new();
524 for entry in &file_info.segments {
525 let entry_bytes = self
526 .get_xorb_ranges(&entry.xorb_hash, vec![(entry.chunk_index_start, entry.chunk_index_end)])
527 .await?
528 .pop()
529 .unwrap();
530 file_vec.extend_from_slice(&entry_bytes);
531 }
532
533 let file_size = file_vec.len();
534
535 let start = byte_range.as_ref().map(|range| range.start as usize).unwrap_or(0);
536
537 if byte_range.is_some() && start >= file_size {
538 return Err(ClientError::InvalidRange);
539 }
540
541 let end = byte_range
542 .as_ref()
543 .map(|range| range.end as usize)
544 .unwrap_or(file_size)
545 .min(file_size);
546
547 Ok(Bytes::from(file_vec[start..end].to_vec()))
548 }
549
550 async fn get_xorb_raw_bytes(&self, hash: &MerkleHash, byte_range: Option<FileRange>) -> Result<Bytes> {
551 if self.xorb_is_tagged(hash).await {
552 return Err(ClientError::XORBNotFound(*hash));
553 }
554 let xorbs = self.xorbs.read().await;
555 let storage = xorbs.get(hash).ok_or(ClientError::XORBNotFound(*hash))?;
556
557 match storage {
558 XorbStorage::Materialized { entry, .. } => {
559 let data = &entry.serialized_data;
560
561 let start = byte_range.as_ref().map(|r| r.start as usize).unwrap_or(0);
562 let end = byte_range
563 .as_ref()
564 .map(|r| r.end as usize)
565 .unwrap_or(data.len())
566 .min(data.len());
567
568 if start >= data.len() {
569 return Err(ClientError::InvalidRange);
570 }
571
572 Ok(data.slice(start..end))
573 },
574 XorbStorage::Random { xorb, .. } => {
575 let total_len = xorb.serialized_length();
576 let start = byte_range.as_ref().map(|r| r.start).unwrap_or(0);
577 let end = byte_range.as_ref().map(|r| r.end).unwrap_or(total_len).min(total_len);
578
579 if start >= total_len {
580 return Err(ClientError::InvalidRange);
581 }
582
583 Ok(xorb.get_serialized_range(start, end))
584 },
585 }
586 }
587
588 async fn xorb_raw_length(&self, hash: &MerkleHash) -> Result<u64> {
589 if self.xorb_is_tagged(hash).await {
590 return Err(ClientError::XORBNotFound(*hash));
591 }
592 let xorbs = self.xorbs.read().await;
593 let storage = xorbs.get(hash).ok_or(ClientError::XORBNotFound(*hash))?;
594
595 match storage {
596 XorbStorage::Materialized { entry, .. } => Ok(entry.serialized_data.len() as u64),
597 XorbStorage::Random { xorb, .. } => Ok(xorb.serialized_length()),
598 }
599 }
600
601 async fn fetch_term_data(
602 &self,
603 hash: MerkleHash,
604 fetch_term: XorbReconstructionFetchInfo,
605 ) -> Result<(Bytes, Vec<u32>)> {
606 self.apply_api_delay().await;
607 let (xorb_hash, url_byte_range, url_timestamp) = parse_fetch_url(&fetch_term.url)?;
608
609 let expiration_ms = self.url_expiration_ms.load(Ordering::Relaxed);
611 let elapsed_ms = Instant::now().saturating_duration_since(url_timestamp).as_millis() as u64;
612 if elapsed_ms > expiration_ms {
613 return Err(ClientError::PresignedUrlExpirationError);
614 }
615
616 let fetch_byte_range = FileRange::from(fetch_term.url_range);
620 if url_byte_range.start != fetch_byte_range.start || url_byte_range.end != fetch_byte_range.end {
621 return Err(ClientError::InvalidArguments);
622 }
623
624 if self.xorb_is_tagged(&xorb_hash).await {
628 return Err(ClientError::XORBNotFound(xorb_hash));
629 }
630
631 let xorbs = self.xorbs.read().await;
632 let storage = xorbs.get(&xorb_hash).ok_or_else(|| {
633 error!("Unable to find xorb in memory CAS {:?}", hash);
634 ClientError::XORBNotFound(hash)
635 })?;
636
637 let (data, xorb_obj) = match storage {
638 XorbStorage::Materialized { entry, .. } => {
639 let mut reader = BufReader::new(Cursor::new(&entry.serialized_data));
640 let xorb_obj = XorbObject::deserialize(&mut reader)?;
641 let data =
642 xorb_obj.get_bytes_by_chunk_range(&mut reader, fetch_term.range.start, fetch_term.range.end)?;
643 (Bytes::from(data), xorb_obj)
644 },
645 XorbStorage::Random { xorb, .. } => {
646 let data = xorb
647 .get_chunk_range_data(fetch_term.range.start, fetch_term.range.end)
648 .ok_or(ClientError::XORBNotFound(hash))?;
649 let xorb_obj = xorb.get_xorb_object();
650 (data, xorb_obj)
651 },
652 };
653
654 let chunk_byte_indices = {
655 let mut indices = Vec::new();
656 let mut cumulative = 0u32;
657 indices.push(0);
658 for chunk_idx in fetch_term.range.start..fetch_term.range.end {
659 let chunk_len = xorb_obj
660 .uncompressed_chunk_length(chunk_idx)
661 .map_err(|e| ClientError::Other(format!("Failed to get chunk length: {e}")))?;
662 cumulative += chunk_len;
663 indices.push(cumulative);
664 }
665 indices
666 };
667
668 Ok((data, chunk_byte_indices))
669 }
670}
671
672impl MemoryClient {
673 async fn compute_reconstruction_ranges(
674 &self,
675 file_id: &MerkleHash,
676 bytes_range: Option<FileRange>,
677 ) -> Result<xorb_utils::ReconstructionRangesResult> {
678 let file_info = {
679 let shard = self.shard.read().await;
680 match shard.get_file_reconstruction_info(file_id) {
681 Some(fi) => fi,
682 None => return Ok(None),
683 }
684 };
685
686 let tagged = self.gc_tagged_xorbs.read().await.clone();
690 let xorbs = self.xorbs.read().await;
691 xorb_utils::compute_reconstruction_ranges(&file_info, bytes_range, &mut |hash| {
692 if tagged.contains(hash) {
693 return Err(ClientError::XORBNotFound(*hash));
694 }
695 let storage = xorbs.get(hash).ok_or_else(|| {
696 error!("Unable to find xorb in memory CAS {:?}", hash);
697 ClientError::XORBNotFound(*hash)
698 })?;
699 Ok(match storage {
700 XorbStorage::Materialized { entry, .. } => entry.xorb_object.clone(),
701 XorbStorage::Random { xorb, .. } => xorb.get_xorb_object(),
702 })
703 })
704 }
705
706 pub async fn get_reconstruction_v1(
708 &self,
709 file_id: &MerkleHash,
710 bytes_range: Option<FileRange>,
711 ) -> Result<Option<QueryReconstructionResponse>> {
712 self.apply_api_delay().await;
713
714 let result = self.compute_reconstruction_ranges(file_id, bytes_range).await?;
715 let Some((offset_into_first_range, terms, merged_ranges)) = result else {
716 return Ok(None);
717 };
718
719 if terms.is_empty() {
720 return Ok(Some(QueryReconstructionResponse {
721 offset_into_first_range,
722 terms,
723 fetch_info: HashMap::new(),
724 }));
725 }
726
727 let timestamp = Instant::now();
728 let mut fetch_info: HashMap<HexMerkleHash, Vec<XorbReconstructionFetchInfo>> = HashMap::new();
729 for (hash, ranges) in merged_ranges {
730 let entries = ranges
731 .into_iter()
732 .map(|r| XorbReconstructionFetchInfo {
733 range: r.chunk_range,
734 url: generate_fetch_url(&hash, &r.byte_range, timestamp),
735 url_range: HttpRange::from(r.byte_range),
736 })
737 .collect();
738 fetch_info.insert(hash.into(), entries);
739 }
740
741 Ok(Some(QueryReconstructionResponse {
742 offset_into_first_range,
743 terms,
744 fetch_info,
745 }))
746 }
747
748 pub async fn get_reconstruction_v2(
750 &self,
751 file_id: &MerkleHash,
752 bytes_range: Option<FileRange>,
753 ) -> Result<Option<QueryReconstructionResponseV2>> {
754 self.apply_api_delay().await;
755
756 let result = self.compute_reconstruction_ranges(file_id, bytes_range).await?;
757 let Some((offset_into_first_range, terms, merged_ranges)) = result else {
758 return Ok(None);
759 };
760
761 if terms.is_empty() {
762 return Ok(Some(QueryReconstructionResponseV2 {
763 offset_into_first_range,
764 terms,
765 xorbs: HashMap::new(),
766 }));
767 }
768
769 let timestamp = Instant::now();
770 let max_ranges = self.max_ranges_per_fetch.load(Ordering::Relaxed);
771
772 let mut xorbs: HashMap<HexMerkleHash, Vec<XorbMultiRangeFetch>> = HashMap::new();
773 for (hash, ranges) in merged_ranges {
774 let mut fetch_entries = Vec::new();
775
776 for chunk in ranges.chunks(max_ranges) {
777 let range_descriptors: Vec<XorbRangeDescriptor> = chunk
778 .iter()
779 .map(|r| XorbRangeDescriptor {
780 chunks: r.chunk_range,
781 bytes: HttpRange::from(r.byte_range),
782 })
783 .collect();
784
785 let url = generate_v2_fetch_url(&hash, &range_descriptors, timestamp);
786 fetch_entries.push(XorbMultiRangeFetch {
787 url,
788 ranges: range_descriptors,
789 });
790 }
791
792 xorbs.insert(hash.into(), fetch_entries);
793 }
794
795 Ok(Some(QueryReconstructionResponseV2 {
796 offset_into_first_range,
797 terms,
798 xorbs,
799 }))
800 }
801}
802
803#[cfg_attr(not(target_family = "wasm"), async_trait)]
804#[cfg_attr(target_family = "wasm", async_trait(?Send))]
805impl Client for MemoryClient {
806 async fn get_file_reconstruction_info(
807 &self,
808 file_hash: &MerkleHash,
809 ) -> Result<Option<(MDBFileInfo, Option<MerkleHash>)>> {
810 self.apply_api_delay().await;
811 let shard = self.shard.read().await;
812 Ok(shard.get_file_reconstruction_info(file_hash).map(|fi| (fi, None)))
813 }
814
815 async fn query_for_global_dedup_shard(&self, _prefix: &str, chunk_hash: &MerkleHash) -> Result<Option<Bytes>> {
816 self.apply_api_delay().await;
817 let shard_bytes = {
818 let dedup = self.global_dedup.read().await;
819 let Some(shard_bytes) = dedup.get(chunk_hash) else {
820 return Ok(None);
821 };
822 shard_bytes.clone()
823 };
824
825 let expiration_secs = self.global_dedup_expiration_secs.load(Ordering::Relaxed);
826 if expiration_secs == 0 {
827 return Ok(Some(shard_bytes));
828 }
829
830 let expiry = std::time::SystemTime::now() + Duration::from_secs(expiration_secs);
831
832 let mut reader = Cursor::new(shard_bytes.as_ref());
833 let minimal_shard = MDBMinimalShard::from_reader(&mut reader, true, true)?;
834
835 let mut out = Vec::new();
836 minimal_shard.serialize_xorb_subset_with_expiry(&mut out, Some(expiry), |_| true)?;
837 Ok(Some(out.into()))
838 }
839
840 async fn acquire_upload_permit(&self) -> Result<super::super::adaptive_concurrency::ConnectionPermit> {
841 self.apply_api_delay().await;
842 self.upload_concurrency_controller.acquire_connection_permit().await
843 }
844
845 async fn upload_shard(
846 &self,
847 shard_data: Bytes,
848 _permit: super::super::adaptive_concurrency::ConnectionPermit,
849 ) -> Result<bool> {
850 self.apply_api_delay().await;
851 let mut reader = Cursor::new(&shard_data);
853 let minimal_shard = MDBMinimalShard::from_reader(&mut reader, true, true)?;
854
855 {
857 let mut shard_lg = self.shard.write().await;
858
859 for i in 0..minimal_shard.num_files() {
861 let file_view = minimal_shard.file(i).unwrap();
862 shard_lg.add_file_reconstruction_info(MDBFileInfo::from(file_view))?;
863 }
864
865 for i in 0..minimal_shard.num_xorb() {
867 let xorb_view = minimal_shard.xorb(i).unwrap();
868 shard_lg.add_xorb_block(MDBXorbInfo::from(xorb_view))?;
869 }
870 }
871
872 let chunk_hashes = minimal_shard.global_dedup_eligible_chunks();
874
875 {
876 let mut dedup_lg = self.global_dedup.write().await;
877 for chunk in chunk_hashes {
878 dedup_lg.insert(chunk, shard_data.clone());
879 }
880 }
881
882 *self.gc_tagged_shard.write().await = None;
886
887 Ok(true)
888 }
889
890 async fn upload_xorb(
891 &self,
892 _prefix: &str,
893 serialized_xorb_object: SerializedXorbObject,
894 progress_callback: Option<ProgressCallback>,
895 _permit: super::super::adaptive_concurrency::ConnectionPermit,
896 ) -> Result<u64> {
897 self.apply_api_delay().await;
898 let hash = serialized_xorb_object.hash;
899 let footer_start = serialized_xorb_object.footer_start;
900 let serialized_data = serialized_xorb_object.serialized_data;
901
902 info!("Storing XORB {hash:?} in memory");
908
909 let (xorb_obj, serialized_data) = if footer_start.is_some() {
911 let mut reader = Cursor::new(&serialized_data);
912 let xorb_obj = XorbObject::deserialize(&mut reader)?;
913 (xorb_obj, serialized_data)
914 } else {
915 let mut data_with_footer = Vec::new();
916 let (xorb_obj, computed_hash) = xet_core_structures::xorb_object::reconstruct_xorb_with_footer(
917 &mut data_with_footer,
918 &serialized_data,
919 )?;
920 if computed_hash != hash {
921 return Err(ClientError::Other(format!(
922 "XORB hash mismatch: expected {}, got {}",
923 hash.hex(),
924 computed_hash.hex(),
925 )));
926 }
927 (xorb_obj, data_with_footer)
928 };
929
930 let bytes_written = serialized_data.len();
931 let generation = self.xorb_generation.fetch_add(1, Ordering::Relaxed);
932
933 {
934 let mut xorbs = self.xorbs.write().await;
935 xorbs.insert(
936 hash,
937 XorbStorage::Materialized {
938 entry: MaterializedXorb {
939 serialized_data: Bytes::from(serialized_data),
940 xorb_object: xorb_obj,
941 },
942 generation,
943 },
944 );
945 }
946
947 self.gc_tagged_xorbs.write().await.remove(&hash);
951
952 if let Some(ref cb) = progress_callback {
953 let n = bytes_written as u64;
954 cb(n, n, n);
955 }
956
957 info!("XORB {hash:?} successfully stored with {bytes_written} bytes.");
958
959 Ok(bytes_written as u64)
960 }
961
962 async fn get_reconstruction(
963 &self,
964 file_id: &MerkleHash,
965 bytes_range: Option<FileRange>,
966 ) -> Result<Option<QueryReconstructionResponseV2>> {
967 self.get_reconstruction_v2(file_id, bytes_range).await
968 }
969
970 async fn batch_get_reconstruction(&self, file_ids: &[MerkleHash]) -> Result<BatchQueryReconstructionResponse> {
971 self.apply_api_delay().await;
972 let mut files = HashMap::new();
973 let mut fetch_info_map: HashMap<HexMerkleHash, Vec<XorbReconstructionFetchInfo>> = HashMap::new();
974
975 for file_id in file_ids {
976 if let Some(response) = self.get_reconstruction_v1(file_id, None).await? {
977 let hex_hash: HexMerkleHash = (*file_id).into();
978 files.insert(hex_hash, response.terms);
979
980 for (hash, fetch_infos) in response.fetch_info {
981 fetch_info_map.entry(hash).or_default().extend(fetch_infos);
982 }
983 }
984 }
985
986 Ok(BatchQueryReconstructionResponse {
987 files,
988 fetch_info: fetch_info_map,
989 })
990 }
991
992 async fn acquire_download_permit(&self) -> Result<super::super::adaptive_concurrency::ConnectionPermit> {
993 self.apply_api_delay().await;
994 self.upload_concurrency_controller.acquire_connection_permit().await
995 }
996
997 async fn get_file_term_data(
998 &self,
999 url_info: Box<dyn super::super::interface::URLProvider>,
1000 _download_permit: super::super::adaptive_concurrency::ConnectionPermit,
1001 progress_callback: Option<ProgressCallback>,
1002 uncompressed_size_if_known: Option<usize>,
1003 ) -> Result<(Bytes, Vec<u32>)> {
1004 self.apply_api_delay().await;
1005 let (url, http_ranges) = url_info.retrieve_url().await?;
1006 let (xorb_hash, url_timestamp) = parse_any_fetch_url(&url)?;
1007
1008 let expiration_ms = self.url_expiration_ms.load(Ordering::Relaxed);
1010 let elapsed_ms = Instant::now().saturating_duration_since(url_timestamp).as_millis() as u64;
1011 if elapsed_ms > expiration_ms {
1012 return Err(ClientError::PresignedUrlExpirationError);
1013 }
1014
1015 if self.xorb_is_tagged(&xorb_hash).await {
1018 return Err(ClientError::XORBNotFound(xorb_hash));
1019 }
1020
1021 let xorbs = self.xorbs.read().await;
1022 let storage = xorbs.get(&xorb_hash).ok_or(ClientError::XORBNotFound(xorb_hash))?;
1023
1024 let mut all_decompressed = Vec::new();
1026 let mut all_chunk_indices = Vec::<u32>::new();
1027 let mut total_transfer = 0u64;
1028
1029 for http_range in &http_ranges {
1030 let start = http_range.start as usize;
1031 let end = http_range.end as usize + 1;
1032 total_transfer += http_range.length();
1033
1034 let (data, chunk_indices) = match storage {
1035 XorbStorage::Materialized { entry, .. } => {
1036 let range_data = &entry.serialized_data[start..end];
1037 xet_core_structures::xorb_object::deserialize_chunks(&mut Cursor::new(range_data))?
1038 },
1039 XorbStorage::Random { xorb, .. } => {
1040 let range_data = xorb.get_serialized_range(start as u64, end as u64);
1041 xet_core_structures::xorb_object::deserialize_chunks(&mut Cursor::new(range_data.as_ref()))?
1042 },
1043 };
1044
1045 xet_core_structures::xorb_object::append_chunk_segment(
1046 &mut all_decompressed,
1047 &mut all_chunk_indices,
1048 &data,
1049 &chunk_indices,
1050 );
1051 }
1052
1053 if let Some(expected) = uncompressed_size_if_known {
1054 debug_assert_eq!(
1055 all_decompressed.len(),
1056 expected,
1057 "get_file_term_data: expected {} bytes, got {}",
1058 expected,
1059 all_decompressed.len()
1060 );
1061 }
1062
1063 if let Some(ref cb) = progress_callback {
1064 cb(total_transfer, total_transfer, total_transfer);
1065 }
1066 Ok((Bytes::from(all_decompressed), all_chunk_indices))
1067 }
1068
1069 async fn get_file_chunk_hashes(
1070 &self,
1071 file_id: &MerkleHash,
1072 dirty_ranges: Vec<FileRange>,
1073 ) -> Result<FileChunkHashesResponse> {
1074 self.apply_api_delay().await;
1075
1076 let file_info = {
1077 let shard = self.shard.read().await;
1078 shard
1079 .get_file_reconstruction_info(file_id)
1080 .ok_or(ClientError::FileNotFound(*file_id))?
1081 };
1082
1083 let tagged = self.gc_tagged_xorbs.read().await.clone();
1086 let xorbs = self.xorbs.read().await;
1087 let mut chunks: Vec<(MerkleHash, u64)> = Vec::new();
1088 for segment in &file_info.segments {
1089 if tagged.contains(&segment.xorb_hash) {
1090 return Err(ClientError::XORBNotFound(segment.xorb_hash));
1091 }
1092 let storage = xorbs
1093 .get(&segment.xorb_hash)
1094 .ok_or(ClientError::XORBNotFound(segment.xorb_hash))?;
1095 let xorb_obj = match storage {
1096 XorbStorage::Materialized { entry, .. } => std::borrow::Cow::Borrowed(&entry.xorb_object),
1097 XorbStorage::Random { xorb, .. } => std::borrow::Cow::Owned(xorb.get_xorb_object()),
1098 };
1099 chunks.extend(
1100 xorb_obj
1101 .chunk_hash_sizes(segment.chunk_index_start, segment.chunk_index_end)
1102 .map_err(|err| ClientError::Other(format!("chunk_hash_sizes error: {err}")))?,
1103 );
1104 }
1105
1106 build_file_chunk_hashes_response(&file_info, dirty_ranges, chunks)
1107 }
1108}
1109
1110#[cfg(not(target_family = "wasm"))]
1111#[async_trait]
1112impl super::DeletionControlableClient for MemoryClient {
1113 async fn list_shard_entries(&self) -> Result<Vec<MerkleHash>> {
1114 let shard = self.shard.read().await;
1115 let Some((h, _)) = Self::current_shard_hash_and_bytes(&shard)? else {
1116 return Ok(Vec::new());
1117 };
1118 if self.shard_is_tagged(&h).await {
1119 return Ok(Vec::new());
1120 }
1121 Ok(vec![h])
1122 }
1123
1124 async fn get_shard_bytes(&self, hash: &MerkleHash) -> Result<Bytes> {
1125 if self.shard_is_tagged(hash).await {
1126 return Err(ClientError::Other(format!("Shard not found: {}", hash.hex())));
1127 }
1128 let shard = self.shard.read().await;
1129 let Some((current_hash, bytes)) = Self::current_shard_hash_and_bytes(&shard)? else {
1130 return Err(ClientError::Other(format!("Shard not found: {}", hash.hex())));
1131 };
1132 if ¤t_hash != hash {
1133 return Err(ClientError::Other(format!("Shard not found: {}", hash.hex())));
1134 }
1135 Ok(bytes)
1136 }
1137
1138 async fn delete_shard_entry(&self, hash: &MerkleHash) -> Result<()> {
1139 let mut shard = self.shard.write().await;
1140 let Some((current_hash, _)) = Self::current_shard_hash_and_bytes(&shard)? else {
1141 return Err(ClientError::Other(format!("Shard not found: {}", hash.hex())));
1142 };
1143 if ¤t_hash != hash {
1144 return Err(ClientError::Other(format!("Shard not found: {}", hash.hex())));
1145 }
1146 if self.lifecycle_tag_deletion_enabled() {
1147 *self.gc_tagged_shard.write().await = Some(current_hash);
1148 } else {
1149 *shard = MDBInMemoryShard::default();
1150 }
1151 Ok(())
1152 }
1153
1154 async fn list_file_shard_entries(&self) -> Result<Vec<(MerkleHash, MerkleHash)>> {
1155 let shard = self.shard.read().await;
1156 let Some((shard_hash, _)) = Self::current_shard_hash_and_bytes(&shard)? else {
1157 return Ok(Vec::new());
1158 };
1159 if self.shard_is_tagged(&shard_hash).await {
1160 return Ok(Vec::new());
1161 }
1162 Ok(shard
1163 .file_content
1164 .keys()
1165 .copied()
1166 .map(|file_hash| (file_hash, shard_hash))
1167 .collect())
1168 }
1169
1170 async fn delete_file_entry(&self, file_hash: &MerkleHash) -> Result<()> {
1171 let mut shard = self.shard.write().await;
1172 if shard.file_content.remove(file_hash).is_some() {
1173 shard.recalculate_shard_size();
1174 }
1175 Ok(())
1176 }
1177
1178 async fn remove_shard_dedup_entries(&self, shard_hash: &MerkleHash) -> Result<()> {
1179 let shard = self.shard.read().await;
1180 let Some((current_hash, _)) = Self::current_shard_hash_and_bytes(&shard)? else {
1181 return Ok(());
1182 };
1183 if ¤t_hash != shard_hash {
1184 return Ok(());
1185 }
1186 drop(shard);
1187
1188 self.global_dedup.write().await.clear();
1189 Ok(())
1190 }
1191
1192 async fn delete_xorb(&self, hash: &MerkleHash) {
1193 if self.lifecycle_tag_deletion_enabled() {
1194 self.gc_tagged_xorbs.write().await.insert(*hash);
1195 } else {
1196 self.xorbs.write().await.remove(hash);
1197 }
1198 }
1199
1200 async fn list_xorbs_and_tags(&self) -> Result<Vec<(MerkleHash, ObjectTag)>> {
1201 let tagged = self.gc_tagged_xorbs.read().await;
1202 let xorbs = self.xorbs.read().await;
1203 Ok(xorbs
1204 .iter()
1205 .filter(|(hash, _)| !tagged.contains(hash))
1206 .map(|(hash, storage)| (*hash, Self::xorb_tag(hash, storage)))
1207 .collect())
1208 }
1209
1210 async fn delete_xorb_if_tag_matches(&self, hash: &MerkleHash, tag: &ObjectTag) -> Result<bool> {
1211 let current_tag = {
1212 let xorbs = self.xorbs.read().await;
1213 let Some(storage) = xorbs.get(hash) else {
1214 return Err(ClientError::XORBNotFound(*hash));
1215 };
1216 Self::xorb_tag(hash, storage)
1217 };
1218 if ¤t_tag != tag {
1219 return Ok(false);
1220 }
1221 if self.lifecycle_tag_deletion_enabled() {
1222 self.gc_tagged_xorbs.write().await.insert(*hash);
1223 } else {
1224 self.xorbs.write().await.remove(hash);
1225 }
1226 Ok(true)
1227 }
1228
1229 async fn list_shards_with_tags(&self) -> Result<Vec<(MerkleHash, ObjectTag)>> {
1230 let shard = self.shard.read().await;
1231 let Some((shard_hash, shard_bytes)) = Self::current_shard_hash_and_bytes(&shard)? else {
1232 return Ok(Vec::new());
1233 };
1234 if self.shard_is_tagged(&shard_hash).await {
1235 return Ok(Vec::new());
1236 }
1237 let tag = Self::object_tag_from_key_and_payload(b"shard", &shard_hash, shard_bytes.as_ref());
1238 Ok(vec![(shard_hash, tag)])
1239 }
1240
1241 async fn delete_shard_if_tag_matches(&self, hash: &MerkleHash, tag: &ObjectTag) -> Result<bool> {
1242 let (current_hash, current_tag) = {
1243 let shard = self.shard.read().await;
1244 let Some((current_hash, shard_bytes)) = Self::current_shard_hash_and_bytes(&shard)? else {
1245 return Err(ClientError::Other(format!("Shard not found: {}", hash.hex())));
1246 };
1247 if ¤t_hash != hash {
1248 return Err(ClientError::Other(format!("Shard not found: {}", hash.hex())));
1249 }
1250 let current_tag = Self::object_tag_from_key_and_payload(b"shard", ¤t_hash, shard_bytes.as_ref());
1251 (current_hash, current_tag)
1252 };
1253 if ¤t_tag != tag {
1254 return Ok(false);
1255 }
1256 if self.lifecycle_tag_deletion_enabled() {
1257 *self.gc_tagged_shard.write().await = Some(current_hash);
1258 } else {
1259 *self.shard.write().await = MDBInMemoryShard::default();
1260 }
1261 Ok(true)
1262 }
1263
1264 async fn verify_integrity(&self) -> Result<()> {
1265 let tagged = self.gc_tagged_xorbs.read().await;
1266 let xorbs = self.xorbs.read().await;
1267 let shard = self.shard.read().await;
1268 let shard_tagged = match Self::current_shard_hash_and_bytes(&shard)? {
1274 Some((shard_hash, _)) => self.gc_tagged_shard.read().await.is_some_and(|h| h == shard_hash),
1275 None => false,
1276 };
1277 if !shard_tagged {
1278 for file_info in shard.file_content.values() {
1279 for segment in &file_info.segments {
1280 if !xorbs.contains_key(&segment.xorb_hash) || tagged.contains(&segment.xorb_hash) {
1281 return Err(ClientError::XORBNotFound(segment.xorb_hash));
1282 }
1283 }
1284 }
1285 }
1286 Ok(())
1287 }
1288
1289 async fn verify_all_reachable(&self) -> Result<()> {
1290 self.verify_integrity().await
1291 }
1292}
1293
1294fn generate_fetch_url(hash: &MerkleHash, byte_range: &FileRange, timestamp: Instant) -> String {
1295 let timestamp_ms = timestamp.saturating_duration_since(*REFERENCE_INSTANT).as_millis() as u64;
1296 format!("{}:{}:{}:{}", hash.hex(), byte_range.start, byte_range.end, timestamp_ms)
1297}
1298
1299fn parse_fetch_url(url: &str) -> Result<(MerkleHash, FileRange, Instant)> {
1300 let mut parts = url.rsplitn(4, ':').collect::<Vec<_>>();
1301 parts.reverse();
1302
1303 if parts.len() != 4 {
1304 return Err(ClientError::InvalidArguments);
1305 }
1306
1307 let hash = MerkleHash::from_hex(parts[0]).map_err(|_| ClientError::InvalidArguments)?;
1308 let start_pos: u64 = parts[1].parse().map_err(|_| ClientError::InvalidArguments)?;
1309 let end_pos: u64 = parts[2].parse().map_err(|_| ClientError::InvalidArguments)?;
1310 let timestamp_ms: u64 = parts[3].parse().map_err(|_| ClientError::InvalidArguments)?;
1311
1312 let byte_range = FileRange::new(start_pos, end_pos);
1313 let timestamp = *REFERENCE_INSTANT + Duration::from_millis(timestamp_ms);
1314
1315 Ok((hash, byte_range, timestamp))
1316}
1317
1318fn generate_v2_fetch_url(hash: &MerkleHash, ranges: &[XorbRangeDescriptor], timestamp: Instant) -> String {
1319 xorb_utils::generate_v2_fetch_url(hash, ranges, timestamp)
1320}
1321
1322fn parse_any_fetch_url(url: &str) -> Result<(MerkleHash, Instant)> {
1324 if let Ok((hash, _, ts)) = parse_fetch_url(url) {
1325 return Ok((hash, ts));
1326 }
1327 let (hash, ts, _) = xorb_utils::parse_v2_fetch_url(url)?;
1328 Ok((hash, ts))
1329}
1330
1331#[cfg(all(test, not(target_family = "wasm")))]
1332mod tests {
1333 use xet_runtime::config::XetConfig;
1334 use xet_runtime::core::XetContext;
1335
1336 use super::super::client_testing_utils::ClientTestingUtils;
1337 use super::super::deletion_controls::DeletionControlableClient;
1338 use super::*;
1339
1340 fn test_ctx() -> XetContext {
1341 let config = XetConfig::new();
1342 XetContext::from_external(tokio::runtime::Handle::current(), config)
1343 }
1344
1345 fn new_client() -> Arc<dyn super::super::DirectAccessClient> {
1346 MemoryClient::new(test_ctx())
1347 }
1348
1349 fn new_deletion_client() -> Arc<MemoryClient> {
1350 MemoryClient::new(test_ctx())
1351 }
1352
1353 #[tokio::test]
1354 async fn test_common_client_suite() {
1355 super::super::client_unit_testing::test_client_functionality(|| async { new_client() }).await;
1356 }
1357
1358 #[tokio::test]
1359 async fn test_memory_deletion_controls_basic() {
1360 let client = new_deletion_client();
1361 let file = client.upload_random_file(&[(1, (0, 3))], 2048).await.unwrap();
1362
1363 let xorbs_and_tags = client.list_xorbs_and_tags().await.unwrap();
1364 assert!(!xorbs_and_tags.is_empty());
1365 let (xorb_hash, tag) = xorbs_and_tags[0];
1366
1367 let wrong_tag = [0xABu8; 32];
1368 assert!(!client.delete_xorb_if_tag_matches(&xorb_hash, &wrong_tag).await.unwrap());
1369 assert!(client.xorb_exists(&xorb_hash).await.unwrap());
1370
1371 assert!(client.delete_xorb_if_tag_matches(&xorb_hash, &tag).await.unwrap());
1372 assert!(!client.xorb_exists(&xorb_hash).await.unwrap());
1373
1374 client.delete_file_entry(&file.file_hash).await.unwrap();
1376 client.delete_file_entry(&file.file_hash).await.unwrap();
1377
1378 let shards_and_tags = client.list_shards_with_tags().await.unwrap();
1379 if !shards_and_tags.is_empty() {
1380 let (shard_hash, shard_tag) = shards_and_tags[0];
1381 assert!(!client.delete_shard_if_tag_matches(&shard_hash, &wrong_tag).await.unwrap());
1382 assert!(client.delete_shard_if_tag_matches(&shard_hash, &shard_tag).await.unwrap());
1383 assert!(client.list_shard_entries().await.unwrap().is_empty());
1384 }
1385 }
1386
1387 #[tokio::test(start_paused = true)]
1388 async fn test_url_expiration() {
1389 super::super::client_unit_testing::test_url_expiration_functionality(|| async { new_client() }).await;
1390 }
1391
1392 #[tokio::test(start_paused = true)]
1393 async fn test_api_delay() {
1394 super::super::client_unit_testing::test_api_delay_functionality(|| async { new_client() }).await;
1395 }
1396
1397 #[tokio::test(start_paused = true)]
1398 async fn test_global_dedup_shard_expiration() {
1399 super::super::client_unit_testing::test_global_dedup_shard_expiration_functionality(|| async { new_client() })
1400 .await;
1401 }
1402
1403 #[tokio::test]
1404 #[cfg_attr(feature = "smoke-test", ignore)]
1405 async fn test_global_dedup_shard_expiration_stress() {
1406 super::super::client_unit_testing::test_global_dedup_shard_expiration_stress(|| async { new_client() }).await;
1407 }
1408
1409 #[tokio::test]
1411 async fn test_random_xorb() {
1412 let client = MemoryClient::new(test_ctx());
1413
1414 let xorb = RandomXorb::from_seed(42, 5, 1024);
1416 let xorb_hash = xorb.xorb_hash();
1417 let returned_hash = client.insert_random_xorb(xorb.clone()).await.unwrap();
1418 assert_eq!(xorb_hash, returned_hash);
1419 assert!(client.xorb_exists(&xorb_hash).await.unwrap());
1420 assert_eq!(client.list_xorbs().await.unwrap(), vec![xorb_hash]);
1421
1422 let full_data = client.get_full_xorb(&xorb_hash).await.unwrap();
1424 assert_eq!(full_data, xorb.get_chunk_range_data(0, 5).unwrap());
1425
1426 let range_data = client.get_xorb_ranges(&xorb_hash, vec![(1, 3)]).await.unwrap();
1427 assert_eq!(range_data[0], xorb.get_chunk_range_data(1, 3).unwrap());
1428
1429 let footer = client.xorb_footer(&xorb_hash).await.unwrap();
1431 let expected_footer = xorb.get_xorb_object();
1432 assert_eq!(footer.info.num_chunks, expected_footer.info.num_chunks);
1433 assert_eq!(footer.info.xorb_hash, expected_footer.info.xorb_hash);
1434 assert_eq!(footer.info.chunk_hashes, expected_footer.info.chunk_hashes);
1435
1436 let raw_len = client.xorb_raw_length(&xorb_hash).await.unwrap();
1438 assert_eq!(raw_len, xorb.serialized_length());
1439 assert_eq!(client.get_xorb_raw_bytes(&xorb_hash, None).await.unwrap(), xorb.get_full_serialized());
1440
1441 let partial = client
1442 .get_xorb_raw_bytes(&xorb_hash, Some(FileRange::new(10, 50)))
1443 .await
1444 .unwrap();
1445 assert_eq!(partial, xorb.get_serialized_range(10, 50));
1446 }
1447
1448 #[tokio::test]
1450 async fn test_random_xorb_large() {
1451 let client = MemoryClient::new(test_ctx());
1452 let xorb = RandomXorb::from_seed(12345, 100, 4096);
1453 let xorb_hash = client.insert_random_xorb(xorb.clone()).await.unwrap();
1454
1455 let ranges = vec![(0, 10), (50, 60), (90, 100)];
1456 let results = client.get_xorb_ranges(&xorb_hash, ranges.clone()).await.unwrap();
1457
1458 for (i, (start, end)) in ranges.iter().enumerate() {
1459 assert_eq!(results[i], xorb.get_chunk_range_data(*start, *end).unwrap());
1460 }
1461 }
1462
1463 #[tokio::test]
1465 async fn test_lazy_file() {
1466 let client = MemoryClient::new(test_ctx());
1467
1468 let file = client.insert_random_lazy_file(&[(1, (0, 3))], 256).await.unwrap();
1470 assert_eq!(client.get_file_size(&file.file_hash).await.unwrap(), file.data.len() as u64);
1471 assert_eq!(client.get_file_data(&file.file_hash, None).await.unwrap(), file.data);
1472
1473 let file2 = client
1475 .insert_random_lazy_file(&[(1, (0, 2)), (2, (0, 3)), (1, (2, 4))], 512)
1476 .await
1477 .unwrap();
1478 assert_eq!(file2.terms.len(), 3);
1479 for term in &file2.terms {
1480 let xorb_data = client
1481 .get_xorb_ranges(&term.xorb_hash, vec![(term.chunk_start, term.chunk_end)])
1482 .await
1483 .unwrap();
1484 assert_eq!(xorb_data[0], term.data);
1485 }
1486 assert_eq!(client.get_file_data(&file2.file_hash, None).await.unwrap(), file2.data);
1487
1488 let (start, end) = (100u64, 500u64);
1490 let range_data = client
1491 .get_file_data(&file2.file_hash, Some(FileRange::new(start, end)))
1492 .await
1493 .unwrap();
1494 assert_eq!(range_data.as_ref(), &file2.data[start as usize..end as usize]);
1495
1496 let recon = client.get_reconstruction_v1(&file2.file_hash, None).await.unwrap().unwrap();
1498 for term in &recon.terms {
1499 let xorb_hash: MerkleHash = term.hash.into();
1500 for fetch_info in recon.fetch_info.get(&term.hash).unwrap() {
1501 let (data, _chunk_indices) = client.fetch_term_data(xorb_hash, fetch_info.clone()).await.unwrap();
1502 assert!(!data.is_empty());
1503 }
1504 }
1505 }
1506
1507 #[tokio::test]
1509 async fn test_lazy_file_deterministic() {
1510 let term_spec = &[(999, (0, 4))];
1511 let file1 = MemoryClient::new(test_ctx())
1512 .insert_random_lazy_file(term_spec, 512)
1513 .await
1514 .unwrap();
1515 let file2 = MemoryClient::new(test_ctx())
1516 .insert_random_lazy_file(term_spec, 512)
1517 .await
1518 .unwrap();
1519 assert_eq!(file1.file_hash, file2.file_hash);
1520 assert_eq!(file1.data, file2.data);
1521 }
1522
1523 #[tokio::test]
1525 async fn test_mixed_xorb_types() {
1526 let client = MemoryClient::new(test_ctx());
1527
1528 let random_xorb = RandomXorb::from_seed(111, 3, 256);
1529 let random_hash = client.insert_random_xorb(random_xorb).await.unwrap();
1530
1531 let file = client.upload_random_file(&[(222, (0, 3))], 256).await.unwrap();
1532 let materialized_hash = file.terms[0].xorb_hash;
1533
1534 assert!(client.xorb_exists(&random_hash).await.unwrap());
1535 assert!(client.xorb_exists(&materialized_hash).await.unwrap());
1536 assert_eq!(client.list_xorbs().await.unwrap().len(), 2);
1537 assert!(!client.get_full_xorb(&random_hash).await.unwrap().is_empty());
1538 assert!(!client.get_full_xorb(&materialized_hash).await.unwrap().is_empty());
1539 }
1540
1541 #[tokio::test]
1544 async fn test_lifecycle_tag_xorb_delete_hides_xorb() {
1545 let client = new_deletion_client();
1546 client.set_lifecycle_tag_deletion(true);
1547
1548 let file = client.upload_random_file(&[(1, (0, 2))], 2048).await.unwrap();
1549 let xorb_hash = file.terms[0].xorb_hash;
1550 assert!(client.xorb_exists(&xorb_hash).await.unwrap());
1551
1552 client.delete_xorb(&xorb_hash).await;
1553
1554 assert!(!client.xorb_exists(&xorb_hash).await.unwrap());
1556 assert!(client.get_full_xorb(&xorb_hash).await.is_err());
1557 assert!(!client.list_xorbs().await.unwrap().contains(&xorb_hash));
1559 }
1560
1561 #[tokio::test]
1562 async fn test_lifecycle_tag_xorb_upload_clears_tag() {
1563 let client = new_deletion_client();
1564 client.set_lifecycle_tag_deletion(true);
1565
1566 let file = client.upload_random_file(&[(1, (0, 2))], 2048).await.unwrap();
1567 let xorb_hash = file.terms[0].xorb_hash;
1568
1569 client.delete_xorb(&xorb_hash).await;
1570 assert!(!client.xorb_exists(&xorb_hash).await.unwrap());
1571
1572 let file2 = client.upload_random_file(&[(1, (0, 2))], 2048).await.unwrap();
1574 let xorb_hash2 = file2.terms[0].xorb_hash;
1575 assert_eq!(xorb_hash, xorb_hash2, "same seed should produce same xorb hash");
1576
1577 assert!(client.xorb_exists(&xorb_hash).await.unwrap(), "re-upload should clear the tag");
1578 }
1579
1580 #[tokio::test]
1581 async fn test_lifecycle_tag_shard_delete_hides_shard() {
1582 let client = new_deletion_client();
1583 client.set_lifecycle_tag_deletion(true);
1584
1585 let file = client.upload_random_file(&[(1, (0, 2))], 2048).await.unwrap();
1586
1587 client.delete_file_entry(&file.file_hash).await.unwrap();
1588 let shard_hash = client.list_shard_entries().await.unwrap().pop().unwrap();
1591 client.delete_shard_entry(&shard_hash).await.unwrap();
1592
1593 assert!(client.get_shard_bytes(&shard_hash).await.is_err());
1595 assert!(client.list_shard_entries().await.unwrap().is_empty());
1596 }
1597
1598 #[tokio::test]
1599 async fn test_lifecycle_tag_mode_off_hard_deletes() {
1600 let client = new_deletion_client();
1601 assert!(!client.lifecycle_tag_deletion_enabled());
1602
1603 let file = client.upload_random_file(&[(1, (0, 2))], 2048).await.unwrap();
1604 let xorb_hash = file.terms[0].xorb_hash;
1605
1606 client.delete_xorb(&xorb_hash).await;
1607
1608 assert!(client.xorbs.read().await.get(&xorb_hash).is_none());
1610 }
1611
1612 #[tokio::test]
1613 async fn test_lifecycle_tag_verify_integrity_flags_tagged_xorb() {
1614 let client = new_deletion_client();
1615 client.set_lifecycle_tag_deletion(true);
1616
1617 let file = client.upload_random_file(&[(1, (0, 2))], 2048).await.unwrap();
1618 let xorb_hash = file.terms[0].xorb_hash;
1619
1620 client.delete_xorb(&xorb_hash).await;
1622
1623 assert!(client.verify_integrity().await.is_err());
1624 }
1625
1626 #[tokio::test]
1627 async fn test_lifecycle_tag_verify_all_reachable_ignores_tagged() {
1628 let client = new_deletion_client();
1629 client.set_lifecycle_tag_deletion(true);
1630
1631 let file = client.upload_random_file(&[(1, (0, 2))], 2048).await.unwrap();
1632 let xorb_hash = file.terms[0].xorb_hash;
1633
1634 client.delete_file_entry(&file.file_hash).await.unwrap();
1635 let shard_hash = client.list_shard_entries().await.unwrap().pop().unwrap();
1636 client.delete_shard_entry(&shard_hash).await.unwrap();
1637 client.delete_xorb(&xorb_hash).await;
1638
1639 client
1640 .verify_all_reachable()
1641 .await
1642 .expect("tagged objects should be ignored by reachability");
1643 }
1644}