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::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
42/// Stored XORB data: the serialized data and the deserialized XorbObject (header/footer).
43struct MaterializedXorb {
44    serialized_data: Bytes,
45    xorb_object: XorbObject,
46}
47
48/// Storage for a XORB - either fully materialized or generated on-the-fly.
49/// Each variant carries a monotonic `generation` counter that is bumped on
50/// every insert, so the tag changes even when content is identical (matching
51/// production ETag semantics).
52enum XorbStorage {
53    Materialized { entry: MaterializedXorb, generation: u64 },
54    Random { xorb: RandomXorb, generation: u64 },
55}
56
57/// In-memory client for testing purposes. Stores all data in memory using hash tables.
58pub struct MemoryClient {
59    /// XORBs stored by hash
60    xorbs: RwLock<MerkleHashMap<XorbStorage>>,
61    /// In-memory shard for file reconstruction info
62    shard: RwLock<MDBInMemoryShard>,
63    /// Global dedup lookup: chunk_hash -> shard bytes
64    global_dedup: RwLock<MerkleHashMap<Bytes>>,
65    /// Upload concurrency controller
66    upload_concurrency_controller: Arc<AdaptiveConcurrencyController>,
67    /// Monotonic counter for xorb upload generations (tag freshness).
68    xorb_generation: AtomicU64,
69    /// URL expiration in milliseconds
70    url_expiration_ms: AtomicU64,
71    /// Global dedup shard expiration in seconds (0 = disabled).
72    global_dedup_expiration_secs: AtomicU64,
73    /// API delay range in milliseconds as (min_ms, max_ms). (0, 0) means disabled.
74    random_ms_delay_window: (AtomicU64, AtomicU64),
75    /// Max ranges per XorbMultiRangeFetch entry. usize::MAX means no splitting.
76    max_ranges_per_fetch: AtomicUsize,
77    /// HTTP status code to return when V2 is disabled (0 = enabled).
78    v2_disabled_status: AtomicU16,
79    /// When true, deletes move the object's identity into the tagged sets
80    /// below instead of dropping its data, mirroring production's S3
81    /// lifecycle-tag deletion. Reads of a tagged object
82    /// then fail as if it were gone, while a subsequent upload of the same
83    /// hash clears the tag (matching S3 PutObject overwriting a tagged
84    /// object). Off by default; opt in via [`Self::set_lifecycle_tag_deletion`].
85    lifecycle_tag_deletion: AtomicBool,
86    /// XORB hashes currently tagged for lifecycle deletion. Data is retained
87    /// in `xorbs` so a re-upload can clear the tag.
88    gc_tagged_xorbs: RwLock<HashSet<MerkleHash>>,
89    /// Shard hash currently tagged for lifecycle deletion. Shard data is
90    /// retained in `shard` so a re-upload can clear the tag.
91    gc_tagged_shard: RwLock<Option<MerkleHash>>,
92}
93
94impl MemoryClient {
95    /// Create a new in-memory client.
96    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    /// Toggle lifecycle-tag deletion mode (see [`Self::lifecycle_tag_deletion`]).
115    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    /// Inserts a RandomXorb into the client's storage and registers it in the shard
132    /// for XORB block lookups. Returns the xorb hash.
133    ///
134    /// This allows storing XORBs that generate their data on-the-fly,
135    /// enabling testing with massive files without storing all data in memory.
136    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        // Create XORB block entries for each chunk
145        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    /// Upload a random file using lazy generation with RandomXorb.
180    ///
181    /// This uses `RandomXorb` to generate XORB data on-the-fly, allowing testing
182    /// with massive files without storing all data in memory.
183    ///
184    /// Each term is defined as `(xorb_seed, (chunk_start, chunk_end))` where:
185    /// - `xorb_seed` determines the random data for that XORB
186    /// - `chunk_start` and `chunk_end` define the range of chunks to include
187    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        // Collect max chunk count needed for each xorb seed
197        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        // Create and register RandomXorbs
204        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        // Build file info from terms
212        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        // Add file reconstruction info to shard
245        {
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        // Treat the shard as gone only when the *current* shard hash matches
490        // the tagged one — delete_file_entry mutates the shard hash, so a
491        // stale `gc_tagged_shard` must not block files of the new shard.
492        #[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        // Refuse reads while the *current* shard hash matches the tagged
506        // shard hash — a stale tag must not block files of a new shard
507        // after delete_file_entry rewrites the shard.
508        #[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        // Check if URL has expired
611        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        // Validate byte range matches url_range
618        // Note: url_byte_range is FileRange (exclusive end), url_range is HttpRange (inclusive end)
619        // We convert url_range to FileRange for comparison
620        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        // A lifecycle-tagged xorb is "gone" from the namespace; mirror the
626        // get_full_xorb behavior so reconstruction paths cannot resurrect
627        // tagged data.
628        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        // Snapshot the tagged set so the sync closure can refuse tagged xorbs
688        // without an async lock acquire — a lifecycle-tagged xorb must appear
689        // "gone" to reconstruction, matching get_full_xorb.
690        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    /// V1 reconstruction: returns per-range presigned URLs.
708    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    /// V2 reconstruction: returns per-xorb multi-range fetch descriptors.
750    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        // Parse the shard using the streaming parser (handles shards without footer)
854        let mut reader = Cursor::new(&shard_data);
855        let minimal_shard = MDBMinimalShard::from_reader(&mut reader, true, true)?;
856
857        // Merge file info into our shard
858        {
859            let mut shard_lg = self.shard.write().await;
860
861            // Add file info from the views
862            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            // Add XORB info from the views
868            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        // Update global dedup lookup using the minimal shard's method
875        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        // A re-upload supersedes any prior lifecycle-tagged shard: clear the
885        // tag so the shard is readable again. Mirrors S3 PutObject
886        // overwriting a tagged object.
887        *self.gc_tagged_shard.write().await = None;
888
889        // No NDJSON stream in the memory sim; synthesize transfer + terminal Result so
890        // SessionShardInterface progress counters complete like production V1.
891        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        // Always overwrite: even if the xorb already exists, we must store it
912        // with a fresh generation so its tag changes, matching production ETag
913        // semantics and ensuring delete_xorb_if_tag_matches is safe under
914        // concurrent uploads.
915
916        info!("Storing XORB {hash:?} in memory");
917
918        // Reconstruct XorbObject from chunk data if no footer, or deserialize if footer present
919        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        // A re-upload of the same xorb hash supersedes any prior
957        // lifecycle-tagged copy: clear the tag so the xorb is readable again.
958        // Mirrors S3 PutObject overwriting a tagged object.
959        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        // Check if URL has expired
1018        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        // A lifecycle-tagged xorb is "gone" from the namespace; refuse to
1025        // serve its bytes via the term-data path, matching get_full_xorb.
1026        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        // Extract each byte range from the serialized data and deserialize
1034        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        // Snapshot tagged xorbs so a lifecycle-tagged xorb appears "gone" and
1093        // is surfaced as XORBNotFound, matching get_xorb_ranges.
1094        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 &current_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 &current_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 &current_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 &current_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 &current_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", &current_hash, shard_bytes.as_ref());
1260            (current_hash, current_tag)
1261        };
1262        if &current_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        // Files living in a lifecycle-tagged shard are "gone" from the
1278        // namespace (list APIs hide them too), so do not walk their
1279        // segments — `delete_shard_entry` retains `file_content` for
1280        // re-upload semantics, but those entries must not be treated as
1281        // active for integrity verification.
1282        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
1331/// Parse either a V1 or V2 fetch URL, returning (hash, timestamp).
1332fn 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        // file deletion is idempotent for parity with the disk-backed behavior.
1384        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    /// Comprehensive test for RandomXorb insertion and data access.
1419    #[tokio::test]
1420    async fn test_random_xorb() {
1421        let client = MemoryClient::new(test_ctx());
1422
1423        // Basic insertion and existence
1424        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        // Full and partial data retrieval
1432        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        // Footer/XorbObject correctness
1439        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        // Raw serialized bytes
1446        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    /// Test RandomXorb with large chunk count and scattered range access.
1458    #[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    /// Comprehensive test for lazy file insertion with on-the-fly xorb generation.
1473    #[tokio::test]
1474    async fn test_lazy_file() {
1475        let client = MemoryClient::new(test_ctx());
1476
1477        // Single-term file
1478        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        // Multi-term file with reused xorb
1483        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        // Range access
1498        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        // Reconstruction workflow
1506        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    /// Same term_spec on two clients should produce identical file hashes and data.
1517    #[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    /// Verify materialized and random xorbs coexist correctly.
1533    #[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    // ── Lifecycle-tag deletion mode tests ───────────────────────────────
1551
1552    #[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        // Reads fail — tagged xorb is "gone" from the namespace.
1564        assert!(!client.xorb_exists(&xorb_hash).await.unwrap());
1565        assert!(client.get_full_xorb(&xorb_hash).await.is_err());
1566        // list_xorbs excludes the tagged xorb.
1567        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        // Re-upload the same xorb (same content seed) — tag is cleared.
1582        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        // After removing the file entry, the shard's content (and thus hash)
1598        // changes. Re-query the shard hash before GC deletes it.
1599        let shard_hash = client.list_shard_entries().await.unwrap().pop().unwrap();
1600        client.delete_shard_entry(&shard_hash).await.unwrap();
1601
1602        // Reads fail — tagged shard is "gone" from the namespace.
1603        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        // Hard-deleted: data is gone.
1618        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        // Tag the xorb WITHOUT removing the file entry (simulates a GC bug).
1630        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}