Skip to main content

xet_client/cas_client/simulation/
memory_client.rs

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
41/// Stored XORB data: the serialized data and the deserialized XorbObject (header/footer).
42struct MaterializedXorb {
43    serialized_data: Bytes,
44    xorb_object: XorbObject,
45}
46
47/// Storage for a XORB - either fully materialized or generated on-the-fly.
48/// Each variant carries a monotonic `generation` counter that is bumped on
49/// every insert, so the tag changes even when content is identical (matching
50/// production ETag semantics).
51enum XorbStorage {
52    Materialized { entry: MaterializedXorb, generation: u64 },
53    Random { xorb: RandomXorb, generation: u64 },
54}
55
56/// In-memory client for testing purposes. Stores all data in memory using hash tables.
57pub struct MemoryClient {
58    /// XORBs stored by hash
59    xorbs: RwLock<MerkleHashMap<XorbStorage>>,
60    /// In-memory shard for file reconstruction info
61    shard: RwLock<MDBInMemoryShard>,
62    /// Global dedup lookup: chunk_hash -> shard bytes
63    global_dedup: RwLock<MerkleHashMap<Bytes>>,
64    /// Upload concurrency controller
65    upload_concurrency_controller: Arc<AdaptiveConcurrencyController>,
66    /// Monotonic counter for xorb upload generations (tag freshness).
67    xorb_generation: AtomicU64,
68    /// URL expiration in milliseconds
69    url_expiration_ms: AtomicU64,
70    /// Global dedup shard expiration in seconds (0 = disabled).
71    global_dedup_expiration_secs: AtomicU64,
72    /// API delay range in milliseconds as (min_ms, max_ms). (0, 0) means disabled.
73    random_ms_delay_window: (AtomicU64, AtomicU64),
74    /// Max ranges per XorbMultiRangeFetch entry. usize::MAX means no splitting.
75    max_ranges_per_fetch: AtomicUsize,
76    /// HTTP status code to return when V2 is disabled (0 = enabled).
77    v2_disabled_status: AtomicU16,
78    /// When true, deletes move the object's identity into the tagged sets
79    /// below instead of dropping its data, mirroring production's S3
80    /// lifecycle-tag deletion. Reads of a tagged object
81    /// then fail as if it were gone, while a subsequent upload of the same
82    /// hash clears the tag (matching S3 PutObject overwriting a tagged
83    /// object). Off by default; opt in via [`Self::set_lifecycle_tag_deletion`].
84    lifecycle_tag_deletion: AtomicBool,
85    /// XORB hashes currently tagged for lifecycle deletion. Data is retained
86    /// in `xorbs` so a re-upload can clear the tag.
87    gc_tagged_xorbs: RwLock<HashSet<MerkleHash>>,
88    /// Shard hash currently tagged for lifecycle deletion. Shard data is
89    /// retained in `shard` so a re-upload can clear the tag.
90    gc_tagged_shard: RwLock<Option<MerkleHash>>,
91}
92
93impl MemoryClient {
94    /// Create a new in-memory client.
95    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    /// Toggle lifecycle-tag deletion mode (see [`Self::lifecycle_tag_deletion`]).
114    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    /// Inserts a RandomXorb into the client's storage and registers it in the shard
131    /// for XORB block lookups. Returns the xorb hash.
132    ///
133    /// This allows storing XORBs that generate their data on-the-fly,
134    /// enabling testing with massive files without storing all data in memory.
135    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        // Create XORB block entries for each chunk
144        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    /// Upload a random file using lazy generation with RandomXorb.
179    ///
180    /// This uses `RandomXorb` to generate XORB data on-the-fly, allowing testing
181    /// with massive files without storing all data in memory.
182    ///
183    /// Each term is defined as `(xorb_seed, (chunk_start, chunk_end))` where:
184    /// - `xorb_seed` determines the random data for that XORB
185    /// - `chunk_start` and `chunk_end` define the range of chunks to include
186    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        // Collect max chunk count needed for each xorb seed
196        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        // Create and register RandomXorbs
203        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        // Build file info from terms
211        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        // Add file reconstruction info to shard
244        {
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        // Treat the shard as gone only when the *current* shard hash matches
489        // the tagged one — delete_file_entry mutates the shard hash, so a
490        // stale `gc_tagged_shard` must not block files of the new shard.
491        #[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        // Refuse reads while the *current* shard hash matches the tagged
505        // shard hash — a stale tag must not block files of a new shard
506        // after delete_file_entry rewrites the shard.
507        #[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        // Check if URL has expired
610        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        // Validate byte range matches url_range
617        // Note: url_byte_range is FileRange (exclusive end), url_range is HttpRange (inclusive end)
618        // We convert url_range to FileRange for comparison
619        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        // A lifecycle-tagged xorb is "gone" from the namespace; mirror the
625        // get_full_xorb behavior so reconstruction paths cannot resurrect
626        // tagged data.
627        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        // Snapshot the tagged set so the sync closure can refuse tagged xorbs
687        // without an async lock acquire — a lifecycle-tagged xorb must appear
688        // "gone" to reconstruction, matching get_full_xorb.
689        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    /// V1 reconstruction: returns per-range presigned URLs.
707    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    /// V2 reconstruction: returns per-xorb multi-range fetch descriptors.
749    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        // Parse the shard using the streaming parser (handles shards without footer)
852        let mut reader = Cursor::new(&shard_data);
853        let minimal_shard = MDBMinimalShard::from_reader(&mut reader, true, true)?;
854
855        // Merge file info into our shard
856        {
857            let mut shard_lg = self.shard.write().await;
858
859            // Add file info from the views
860            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            // Add XORB info from the views
866            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        // Update global dedup lookup using the minimal shard's method
873        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        // A re-upload supersedes any prior lifecycle-tagged shard: clear the
883        // tag so the shard is readable again. Mirrors S3 PutObject
884        // overwriting a tagged object.
885        *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        // Always overwrite: even if the xorb already exists, we must store it
903        // with a fresh generation so its tag changes, matching production ETag
904        // semantics and ensuring delete_xorb_if_tag_matches is safe under
905        // concurrent uploads.
906
907        info!("Storing XORB {hash:?} in memory");
908
909        // Reconstruct XorbObject from chunk data if no footer, or deserialize if footer present
910        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        // A re-upload of the same xorb hash supersedes any prior
948        // lifecycle-tagged copy: clear the tag so the xorb is readable again.
949        // Mirrors S3 PutObject overwriting a tagged object.
950        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        // Check if URL has expired
1009        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        // A lifecycle-tagged xorb is "gone" from the namespace; refuse to
1016        // serve its bytes via the term-data path, matching get_full_xorb.
1017        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        // Extract each byte range from the serialized data and deserialize
1025        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        // Snapshot tagged xorbs so a lifecycle-tagged xorb appears "gone" and
1084        // is surfaced as XORBNotFound, matching get_xorb_ranges.
1085        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 &current_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 &current_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 &current_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 &current_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 &current_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", &current_hash, shard_bytes.as_ref());
1251            (current_hash, current_tag)
1252        };
1253        if &current_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        // Files living in a lifecycle-tagged shard are "gone" from the
1269        // namespace (list APIs hide them too), so do not walk their
1270        // segments — `delete_shard_entry` retains `file_content` for
1271        // re-upload semantics, but those entries must not be treated as
1272        // active for integrity verification.
1273        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
1322/// Parse either a V1 or V2 fetch URL, returning (hash, timestamp).
1323fn 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        // file deletion is idempotent for parity with the disk-backed behavior.
1375        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    /// Comprehensive test for RandomXorb insertion and data access.
1410    #[tokio::test]
1411    async fn test_random_xorb() {
1412        let client = MemoryClient::new(test_ctx());
1413
1414        // Basic insertion and existence
1415        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        // Full and partial data retrieval
1423        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        // Footer/XorbObject correctness
1430        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        // Raw serialized bytes
1437        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    /// Test RandomXorb with large chunk count and scattered range access.
1449    #[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    /// Comprehensive test for lazy file insertion with on-the-fly xorb generation.
1464    #[tokio::test]
1465    async fn test_lazy_file() {
1466        let client = MemoryClient::new(test_ctx());
1467
1468        // Single-term file
1469        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        // Multi-term file with reused xorb
1474        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        // Range access
1489        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        // Reconstruction workflow
1497        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    /// Same term_spec on two clients should produce identical file hashes and data.
1508    #[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    /// Verify materialized and random xorbs coexist correctly.
1524    #[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    // ── Lifecycle-tag deletion mode tests ───────────────────────────────
1542
1543    #[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        // Reads fail — tagged xorb is "gone" from the namespace.
1555        assert!(!client.xorb_exists(&xorb_hash).await.unwrap());
1556        assert!(client.get_full_xorb(&xorb_hash).await.is_err());
1557        // list_xorbs excludes the tagged xorb.
1558        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        // Re-upload the same xorb (same content seed) — tag is cleared.
1573        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        // After removing the file entry, the shard's content (and thus hash)
1589        // changes. Re-query the shard hash before GC deletes it.
1590        let shard_hash = client.list_shard_entries().await.unwrap().pop().unwrap();
1591        client.delete_shard_entry(&shard_hash).await.unwrap();
1592
1593        // Reads fail — tagged shard is "gone" from the namespace.
1594        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        // Hard-deleted: data is gone.
1609        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        // Tag the xorb WITHOUT removing the file entry (simulates a GC bug).
1621        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}