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