Skip to main content

xet_client/cas_client/simulation/
local_client.rs

1use std::collections::HashMap;
2use std::fs::{File, metadata};
3use std::io::{BufReader, Cursor, Read, Seek, SeekFrom, Write};
4use std::mem::size_of;
5use std::ops::Range;
6use std::path::{Path, PathBuf};
7use std::sync::atomic::{AtomicBool, AtomicU16, AtomicU64, AtomicUsize, Ordering};
8use std::sync::{Arc, Mutex, Weak};
9
10use anyhow::anyhow;
11use async_trait::async_trait;
12use bytes::Bytes;
13use rand::RngExt;
14use redb::{ReadableDatabase, ReadableTable, TableDefinition};
15use tempfile::TempDir;
16use tokio::time::{Duration, Instant};
17use tracing::{error, info, warn};
18use xet_core_structures::merklehash::{MerkleHash, compute_data_hash};
19use xet_core_structures::metadata_shard::file_structs::{FileDataSequenceHeader, MDBFileInfo, MDBFileInfoView};
20use xet_core_structures::metadata_shard::shard_file_reconstructor::FileReconstructor;
21use xet_core_structures::metadata_shard::shard_format::MDB_FILE_INFO_ENTRY_SIZE;
22use xet_core_structures::metadata_shard::shard_in_memory::MDBInMemoryShard;
23use xet_core_structures::metadata_shard::streaming_shard::MDBMinimalShard;
24use xet_core_structures::metadata_shard::utils::{parse_shard_filename, shard_file_name};
25use xet_core_structures::metadata_shard::xorb_structs::MDBXorbInfo;
26use xet_core_structures::metadata_shard::{MDBShardFile, MDBShardFileHeader, ShardFileManager};
27use xet_core_structures::serialization_utils::read_u32;
28use xet_core_structures::xorb_object::{SerializedXorbObject, XorbObject};
29use xet_runtime::core::XetContext;
30#[cfg(feature = "fd-track")]
31use xet_runtime::fd_diagnostics::{report_fd_count, track_fd_scope};
32use xet_runtime::file_utils::SafeFileCreator;
33
34use super::deletion_controls::ObjectTag;
35use super::direct_access_client::DirectAccessClient;
36use super::xorb_utils::{self, REFERENCE_INSTANT, duration_to_expiration_secs_ceil};
37use crate::cas_client::Client;
38use crate::cas_client::adaptive_concurrency::AdaptiveConcurrencyController;
39use crate::cas_client::chunk_window_builder::build_file_chunk_hashes_response;
40use crate::cas_client::progress_tracked_streams::ProgressCallback;
41use crate::cas_types::{
42    BatchQueryReconstructionResponse, FileChunkHashesResponse, FileRange, HexMerkleHash, HttpRange,
43    QueryReconstructionResponse, QueryReconstructionResponseV2, XorbMultiRangeFetch, XorbRangeDescriptor,
44    XorbReconstructionFetchInfo,
45};
46use crate::error::{ClientError, Result};
47
48/// Newtype wrapper for MerkleHash to implement redb Key/Value traits.
49/// MerkleHash is DataHash([u64; 4]) = 32 bytes, stored as fixed-width little-endian.
50#[derive(Debug, Clone, Copy, PartialEq, Eq)]
51struct RedbHash(MerkleHash);
52
53impl From<MerkleHash> for RedbHash {
54    fn from(h: MerkleHash) -> Self {
55        RedbHash(h)
56    }
57}
58
59impl From<RedbHash> for MerkleHash {
60    fn from(h: RedbHash) -> Self {
61        h.0
62    }
63}
64
65impl redb::Value for RedbHash {
66    type SelfType<'a> = RedbHash;
67    type AsBytes<'a> = [u8; 32];
68
69    fn fixed_width() -> Option<usize> {
70        Some(32)
71    }
72
73    fn from_bytes<'a>(data: &'a [u8]) -> Self::SelfType<'a>
74    where
75        Self: 'a,
76    {
77        let mut hash = MerkleHash::default();
78        let u64s: &mut [u64; 4] = &mut hash;
79        for (i, chunk) in data.chunks_exact(8).enumerate() {
80            u64s[i] = u64::from_le_bytes(chunk.try_into().unwrap());
81        }
82        RedbHash(hash)
83    }
84
85    fn as_bytes<'a, 'b: 'a>(value: &'a Self::SelfType<'b>) -> Self::AsBytes<'a>
86    where
87        Self: 'a + 'b,
88    {
89        let mut bytes = [0u8; 32];
90        let u64s: &[u64; 4] = &value.0;
91        for (i, &val) in u64s.iter().enumerate() {
92            bytes[i * 8..(i + 1) * 8].copy_from_slice(&val.to_le_bytes());
93        }
94        bytes
95    }
96
97    fn type_name() -> redb::TypeName {
98        redb::TypeName::new("MerkleHash")
99    }
100}
101
102impl redb::Key for RedbHash {
103    fn compare(data1: &[u8], data2: &[u8]) -> std::cmp::Ordering {
104        data1.cmp(data2)
105    }
106}
107
108const GLOBAL_DEDUP_TABLE: TableDefinition<RedbHash, RedbHash> = TableDefinition::new("global_dedup");
109
110/// Maps each active file hash to the shard that owns it.  Absence means deleted.
111const FILE_TO_SHARD_TABLE: TableDefinition<RedbHash, FileShardRef> = TableDefinition::new("file_to_shard");
112
113/// Points a file hash at the shard that canonically owns it, along with the
114/// byte offset and length of the file entry within the shard.  This enables
115/// direct-seek reads without parsing the entire shard.
116/// Stored as 48 bytes: 32 (shard_hash) + 8 (offset) + 8 (length).
117#[derive(Debug, Clone, Copy, PartialEq, Eq)]
118struct FileShardRef {
119    shard_hash: MerkleHash,
120    offset: u64,
121    length: u64,
122}
123
124impl redb::Value for FileShardRef {
125    type SelfType<'a> = FileShardRef;
126    type AsBytes<'a> = [u8; 48];
127
128    fn fixed_width() -> Option<usize> {
129        Some(48)
130    }
131
132    fn from_bytes<'a>(data: &'a [u8]) -> Self::SelfType<'a>
133    where
134        Self: 'a,
135    {
136        let mut hash = MerkleHash::default();
137        let u64s: &mut [u64; 4] = &mut hash;
138        for (i, chunk) in data[..32].chunks_exact(8).enumerate() {
139            u64s[i] = u64::from_le_bytes(chunk.try_into().unwrap());
140        }
141        let offset = u64::from_le_bytes(data[32..40].try_into().unwrap());
142        let length = u64::from_le_bytes(data[40..48].try_into().unwrap());
143        FileShardRef {
144            shard_hash: hash,
145            offset,
146            length,
147        }
148    }
149
150    fn as_bytes<'a, 'b: 'a>(value: &'a Self::SelfType<'b>) -> Self::AsBytes<'a>
151    where
152        Self: 'a + 'b,
153    {
154        let mut bytes = [0u8; 48];
155        let u64s: &[u64; 4] = &value.shard_hash;
156        for (i, &val) in u64s.iter().enumerate() {
157            bytes[i * 8..(i + 1) * 8].copy_from_slice(&val.to_le_bytes());
158        }
159        bytes[32..40].copy_from_slice(&value.offset.to_le_bytes());
160        bytes[40..48].copy_from_slice(&value.length.to_le_bytes());
161        bytes
162    }
163
164    fn type_name() -> redb::TypeName {
165        redb::TypeName::new("FileShardRef")
166    }
167}
168
169/// Process-global cache of open redb databases, keyed by canonicalized path.
170/// Stores [`Weak`] pointers so the cache never keeps a database alive on its
171/// own; only [`Arc`]s held by [`LocalClient`] instances do.
172///
173/// Both [`get_or_open_db`] and [`LocalClient::drop`] hold this mutex while
174/// creating/destroying the inner [`redb::Database`], which eliminates the race
175/// window where a database file lock could still be held between a cache miss
176/// and a `Database::create` call.
177static DB_CACHE: std::sync::LazyLock<Mutex<HashMap<PathBuf, Weak<redb::Database>>>> =
178    std::sync::LazyLock::new(|| Mutex::new(HashMap::new()));
179
180/// Opens or returns a shared [`Arc<redb::Database>`] for `db_path`.
181///
182/// Must be called (and the returned `Arc` kept alive) while no other thread is
183/// dropping the last `Arc` for the same path -- which is guaranteed because
184/// both this function and `LocalClient::drop` serialize on `DB_CACHE`.
185fn get_or_open_db(db_path: &Path) -> std::result::Result<Arc<redb::Database>, redb::DatabaseError> {
186    #[cfg(feature = "fd-track")]
187    let _fd_scope = track_fd_scope(format!("LocalClient::get_or_open_db({})", db_path.display()));
188
189    let mut map = DB_CACHE.lock().unwrap();
190
191    if let Some(weak) = map.get(db_path)
192        && let Some(db) = weak.upgrade()
193    {
194        tracing::trace!(target: "xet_client::local_cas_redb", path = %db_path.display(), "DB_CACHE hit");
195        #[cfg(feature = "fd-track")]
196        report_fd_count("LocalClient::get_or_open_db cache hit");
197        return Ok(db);
198    }
199
200    // Purge dead entries while we hold the lock.
201    map.retain(|_, weak| weak.strong_count() > 0);
202
203    tracing::trace!(target: "xet_client::local_cas_redb", path = %db_path.display(), "DB_CACHE miss");
204
205    let db = Arc::new(redb::Database::create(db_path)?);
206    map.insert(db_path.to_owned(), Arc::downgrade(&db));
207    #[cfg(feature = "fd-track")]
208    report_fd_count("LocalClient::get_or_open_db opened new DB");
209    Ok(db)
210}
211
212/// Scans the file-info section of a serialized shard and returns
213/// `(file_hash, byte_offset, byte_length)` for every file entry.
214/// The offset/length pair identifies the contiguous blob (header + data entries +
215/// verification + metadata_ext) that `MDBFileInfoView::new()` can parse directly.
216fn file_entry_byte_ranges(shard_bytes: &[u8]) -> std::result::Result<Vec<(MerkleHash, u64, u64)>, ClientError> {
217    let mut cursor = Cursor::new(shard_bytes);
218    let _ = MDBShardFileHeader::deserialize(&mut cursor)?;
219
220    let mut entries = Vec::new();
221    loop {
222        let start = cursor.position();
223        let header = FileDataSequenceHeader::deserialize(&mut cursor)?;
224        if header.is_bookend() {
225            break;
226        }
227
228        let n = header.num_entries as usize;
229        let mut n_data = n;
230        if header.contains_verification() {
231            n_data += n;
232        }
233        if header.contains_metadata_ext() {
234            n_data += 1;
235        }
236
237        cursor.set_position(cursor.position() + (n_data * MDB_FILE_INFO_ENTRY_SIZE) as u64);
238        entries.push((header.file_hash, start, cursor.position() - start));
239    }
240    Ok(entries)
241}
242
243pub struct LocalClient {
244    /// Wrapped in `Option` so `Drop` can `take()` it under the `DB_CACHE` lock,
245    /// ensuring the redb file lock is released before another caller can open
246    /// the same path. Always `Some` during normal operation.
247    db: Option<Arc<redb::Database>>,
248    db_path: PathBuf,
249    shard_manager: Arc<ShardFileManager>,
250    xorb_dir: PathBuf,
251    shard_dir: PathBuf,
252    upload_concurrency_controller: Arc<AdaptiveConcurrencyController>,
253    url_expiration_ms: AtomicU64,
254    /// Global dedup shard expiration in seconds (0 = disabled).
255    global_dedup_expiration_secs: AtomicU64,
256    /// API delay range in milliseconds as (min_ms, max_ms). (0, 0) means disabled.
257    random_ms_delay_window: (AtomicU64, AtomicU64),
258    /// Max ranges per XorbMultiRangeFetch entry. usize::MAX means no splitting.
259    max_ranges_per_fetch: AtomicUsize,
260    /// HTTP status code to return when V2 is disabled (0 = enabled).
261    v2_disabled_status: AtomicU16,
262    /// When true, deletes rename the on-disk object to `<canonical>.gctag`
263    /// instead of removing it, mirroring production's S3 lifecycle-tag
264    /// deletion. Reads of the canonical path then fail as
265    /// if the object were gone, while a subsequent upload of the same hash
266    /// clears the `.gctag` file (matching S3 PutObject overwriting a tagged
267    /// object). Off by default; opt in via [`Self::set_lifecycle_tag_deletion`].
268    lifecycle_tag_deletion: AtomicBool,
269    _tmp_dir: Option<TempDir>,
270}
271
272impl LocalClient {
273    /// Create a local client hosted in a temporary directory for testing.
274    /// This is an async function to allow use with current-thread tokio runtime.
275    pub async fn temporary(ctx: XetContext) -> Result<Arc<Self>> {
276        let tmp_dir = TempDir::new().unwrap();
277        let path = tmp_dir.path().to_owned();
278        let s = Self::new_internal(ctx, path, Some(tmp_dir)).await?;
279        Ok(Arc::new(s))
280    }
281
282    /// Create a local client hosted in a directory.  Effectively, this directory
283    /// is the CAS endpoint and persists across instances of LocalClient.
284    pub async fn new(ctx: XetContext, path: impl AsRef<Path>) -> Result<Arc<Self>> {
285        let path = path.as_ref().to_owned();
286        Ok(Arc::new(Self::new_internal(ctx, path, None).await?))
287    }
288
289    async fn new_internal(ctx: XetContext, path: impl AsRef<Path>, tmp_dir: Option<TempDir>) -> Result<Self> {
290        let base_dir = std::path::absolute(path)?;
291        if !base_dir.exists() {
292            std::fs::create_dir_all(&base_dir)?;
293        }
294        // `std::path::absolute` does not resolve symlinks; without canonicalizing, two path strings
295        // for the same directory (e.g. via symlinks or /var vs /private/var on macOS) would open
296        // multiple `redb::Database` handles and duplicate file descriptors for one CAS root.
297        let base_dir = std::fs::canonicalize(&base_dir).unwrap_or(base_dir);
298        #[cfg(feature = "fd-track")]
299        let _fd_scope = track_fd_scope(format!("LocalClient::new_internal({})", base_dir.display()));
300        #[cfg(feature = "fd-track")]
301        report_fd_count("LocalClient::new_internal start");
302
303        let shard_dir = base_dir.join("shards");
304        if !shard_dir.exists() {
305            std::fs::create_dir_all(&shard_dir)?;
306        }
307
308        let xorb_dir = base_dir.join("xorbs");
309        if !xorb_dir.exists() {
310            std::fs::create_dir_all(&xorb_dir)?;
311        }
312
313        let db_path = base_dir.join("global_dedup_lookup.redb");
314        let db =
315            get_or_open_db(&db_path).map_err(|e| ClientError::Other(format!("Error opening redb database: {e}")))?;
316        #[cfg(feature = "fd-track")]
317        report_fd_count("LocalClient::new_internal after DB open");
318
319        // Ensure tables exist by opening a write transaction.
320        {
321            let write_txn = db.begin_write().map_err(map_redb_db_error)?;
322            let _ = write_txn.open_table(GLOBAL_DEDUP_TABLE).map_err(map_redb_db_error)?;
323            let _ = write_txn.open_table(FILE_TO_SHARD_TABLE).map_err(map_redb_db_error)?;
324            write_txn.commit().map_err(map_redb_db_error)?;
325        }
326
327        // Open / set up the shard lookup
328        let shard_manager = ShardFileManager::new_in_session_directory(&ctx, shard_dir.clone(), true).await?;
329        #[cfg(feature = "fd-track")]
330        report_fd_count("LocalClient::new_internal after shard manager init");
331
332        Ok(Self {
333            db: Some(db),
334            db_path,
335            shard_manager,
336            xorb_dir,
337            shard_dir,
338            upload_concurrency_controller: AdaptiveConcurrencyController::new_upload(ctx, "local_uploads"),
339            url_expiration_ms: AtomicU64::new(u64::MAX),
340            global_dedup_expiration_secs: AtomicU64::new(0),
341            random_ms_delay_window: (AtomicU64::new(0), AtomicU64::new(0)),
342            max_ranges_per_fetch: AtomicUsize::new(usize::MAX),
343            v2_disabled_status: AtomicU16::new(0),
344            lifecycle_tag_deletion: AtomicBool::new(false),
345            _tmp_dir: tmp_dir,
346        })
347    }
348
349    fn db(&self) -> &redb::Database {
350        self.db.as_deref().expect("db used after close")
351    }
352
353    /// Internal function to get the path for a given hash entry
354    fn get_path_for_entry(&self, hash: &MerkleHash) -> PathBuf {
355        self.xorb_dir.join(format!("default.{hash:?}"))
356    }
357
358    /// Toggle lifecycle-tag deletion mode (see [`Self::lifecycle_tag_deletion`]).
359    pub fn set_lifecycle_tag_deletion(&self, on: bool) {
360        self.lifecycle_tag_deletion.store(on, Ordering::Relaxed);
361    }
362
363    fn lifecycle_tag_deletion_enabled(&self) -> bool {
364        self.lifecycle_tag_deletion.load(Ordering::Relaxed)
365    }
366
367    /// Path used to park a tagged-for-deletion xorb: `<canonical>.gctag`.
368    /// Bytes are retained on disk until an upload of the same hash clears
369    /// the tag (matching S3 PutObject overwriting a tagged object).
370    fn gctag_xorb_path(&self, hash: &MerkleHash) -> PathBuf {
371        let canonical = self.get_path_for_entry(hash);
372        let mut name = canonical.into_os_string();
373        name.push(".gctag");
374        PathBuf::from(name)
375    }
376
377    /// Path used to park a tagged-for-deletion shard: `<hex>.mdb.gctag`.
378    fn gctag_shard_path(&self, hash: &MerkleHash) -> PathBuf {
379        let canonical = self.shard_dir.join(shard_file_name(hash));
380        let mut name = canonical.into_os_string();
381        name.push(".gctag");
382        PathBuf::from(name)
383    }
384
385    #[cfg(test)]
386    fn is_file_deleted(&self, file_hash: &MerkleHash) -> bool {
387        let Ok(read_txn) = self.db().begin_read() else {
388            return true;
389        };
390        let Ok(table) = read_txn.open_table(FILE_TO_SHARD_TABLE) else {
391            return true;
392        };
393        table.get(&RedbHash::from(*file_hash)).ok().flatten().is_none()
394    }
395
396    /// Returns all shard files in the shard directory as (shard_hash, path) pairs.
397    fn shard_file_paths(&self) -> Result<Vec<(MerkleHash, PathBuf)>> {
398        let mut result = Vec::new();
399        for entry in std::fs::read_dir(&self.shard_dir).map_err(ClientError::internal)? {
400            let entry = entry.map_err(ClientError::internal)?;
401            let path = entry.path();
402            if path.file_name().and_then(|n| n.to_str()).is_some_and(|n| n.ends_with(".gctag")) {
403                continue;
404            }
405            if let Some(hash) = parse_shard_filename(&path)
406                && path.is_file()
407            {
408                result.push((hash, path));
409            }
410        }
411        Ok(result)
412    }
413
414    /// Finds the path for a shard file by hash.
415    fn shard_path_for_hash(&self, hash: &MerkleHash) -> Result<PathBuf> {
416        let path = self.shard_dir.join(shard_file_name(hash));
417        if path.exists() {
418            Ok(path)
419        } else {
420            Err(ClientError::Other(format!("Shard file not found for hash {}", hash.hex())))
421        }
422    }
423
424    /// Builds an `ObjectTag` from file metadata at the given path.
425    ///
426    /// We hash multiple metadata fields to increase entropy and reduce false
427    /// matches during rapid rewrite/delete races.
428    fn object_tag_from_path(path: &Path) -> Result<ObjectTag> {
429        let meta = std::fs::metadata(path).map_err(ClientError::internal)?;
430        let modified = meta.modified().map_err(ClientError::internal)?;
431        let modified_nanos = modified.duration_since(std::time::UNIX_EPOCH).unwrap_or_default().as_nanos();
432        let created_nanos = meta
433            .created()
434            .ok()
435            .and_then(|ts| ts.duration_since(std::time::UNIX_EPOCH).ok())
436            .map_or(0u128, |d| d.as_nanos());
437
438        let mut entropy = Vec::with_capacity(16 + 16 + 8 + 1);
439        entropy.extend_from_slice(&modified_nanos.to_le_bytes());
440        entropy.extend_from_slice(&created_nanos.to_le_bytes());
441        entropy.extend_from_slice(&meta.len().to_le_bytes());
442        entropy.push(u8::from(meta.permissions().readonly()));
443
444        Ok(compute_data_hash(&entropy).into())
445    }
446
447    /// Restores `tmp_path` back to `original_path`.  Tries `hard_link` first
448    /// (fails with EEXIST if a concurrent upload recreated the path — safe).
449    /// Falls back to `rename` if hard links aren't supported.  Only removes
450    /// `tmp_path` after confirming the original is in place.
451    fn restore_from_tmp(tmp_path: &Path, original_path: &Path) {
452        if std::fs::hard_link(tmp_path, original_path).is_ok() {
453            let _ = std::fs::remove_file(tmp_path);
454        } else if original_path.exists() {
455            // Original path was recreated by a concurrent upload; discard stale copy.
456            let _ = std::fs::remove_file(tmp_path);
457        } else {
458            // hard_link failed (e.g. unsupported fs) and no concurrent upload —
459            // fall back to rename which always works on the same filesystem.
460            let _ = std::fs::rename(tmp_path, original_path);
461        }
462    }
463
464    /// Clears the readonly permission on a file so it can be deleted on Windows.
465    #[cfg(windows)]
466    fn clear_readonly(path: &Path) {
467        if let Ok(metadata) = std::fs::metadata(path) {
468            let mut permissions = metadata.permissions();
469            #[allow(clippy::permissions_set_readonly_false)]
470            permissions.set_readonly(false);
471            let _ = std::fs::set_permissions(path, permissions);
472        }
473    }
474
475    /// Loads all shard data from disk into an in-memory shard.
476    #[cfg(test)]
477    fn load_all_shard_data(&self) -> Result<MDBInMemoryShard> {
478        let mut in_memory = MDBInMemoryShard::default();
479        for (_, path) in self.shard_file_paths()? {
480            let shard_bytes = std::fs::read(&path)?;
481            let minimal_shard = MDBMinimalShard::from_reader(&mut Cursor::new(&shard_bytes), true, true)?;
482
483            for i in 0..minimal_shard.num_files() {
484                in_memory.add_file_reconstruction_info(MDBFileInfo::from(minimal_shard.file(i).unwrap()))?;
485            }
486            for i in 0..minimal_shard.num_xorb() {
487                in_memory.add_xorb_block(MDBXorbInfo::from(minimal_shard.xorb(i).unwrap()))?;
488            }
489        }
490        Ok(in_memory)
491    }
492
493    /// Clears shard files from disk, writes the given in-memory shard, and
494    /// registers the new shard file with the shard manager.
495    #[cfg(test)]
496    async fn write_shard_data_and_register(&self, in_memory: &MDBInMemoryShard) -> Result<()> {
497        for (_, path) in self.shard_file_paths()? {
498            std::fs::remove_file(&path)?;
499        }
500
501        if !in_memory.is_empty() {
502            let shard_path = in_memory.write_to_directory(&self.shard_dir, None)?;
503            let shard = MDBShardFile::load_from_file(&shard_path, self.shard_manager.shard_file_cache())?;
504            let shard_hash = shard.shard_hash;
505            self.shard_manager.register_shards(&[shard]).await?;
506
507            // Update FILE_TO_SHARD_TABLE with byte-accurate offsets.
508            let shard_bytes = std::fs::read(&shard_path)?;
509            let file_ranges = file_entry_byte_ranges(&shard_bytes)?;
510            let write_txn = self.db().begin_write().map_err(map_redb_db_error)?;
511            {
512                let mut file_table = write_txn.open_table(FILE_TO_SHARD_TABLE).map_err(map_redb_db_error)?;
513                for (file_hash, offset, length) in &file_ranges {
514                    file_table
515                        .insert(
516                            &RedbHash::from(*file_hash),
517                            &FileShardRef {
518                                shard_hash,
519                                offset: *offset,
520                                length: *length,
521                            },
522                        )
523                        .map_err(map_redb_db_error)?;
524                }
525            }
526            write_txn.commit().map_err(map_redb_db_error)?;
527        }
528
529        Ok(())
530    }
531}
532
533impl Drop for LocalClient {
534    fn drop(&mut self) {
535        #[cfg(feature = "fd-track")]
536        let _fd_scope = track_fd_scope(format!("LocalClient::drop({})", self.xorb_dir.display()));
537        #[cfg(feature = "fd-track")]
538        report_fd_count("LocalClient::drop start");
539
540        // Drop the database handle while holding the cache lock. This
541        // serializes with `get_or_open_db`, ensuring the redb file lock is
542        // fully released before any other caller can attempt to reopen the
543        // same path.
544        if let Ok(mut map) = DB_CACHE.lock() {
545            let db = self.db.take();
546            if db.as_ref().is_some_and(|d| Arc::strong_count(d) == 1) {
547                map.remove(&self.db_path);
548            }
549            drop(db);
550        }
551
552        #[cfg(feature = "fd-track")]
553        report_fd_count("LocalClient::drop end");
554    }
555}
556
557#[async_trait]
558impl DirectAccessClient for LocalClient {
559    fn set_fetch_term_url_expiration(&self, expiration: Duration) {
560        self.url_expiration_ms.store(expiration.as_millis() as u64, Ordering::Relaxed);
561    }
562
563    fn set_global_dedup_shard_expiration(&self, expiration: Option<Duration>) {
564        self.global_dedup_expiration_secs
565            .store(duration_to_expiration_secs_ceil(expiration), Ordering::Relaxed);
566    }
567
568    fn set_max_ranges_per_fetch(&self, max_ranges: usize) {
569        self.max_ranges_per_fetch.store(max_ranges, Ordering::Relaxed);
570    }
571
572    fn disable_v2_reconstruction(&self, status_code: u16) {
573        self.v2_disabled_status.store(status_code, Ordering::Relaxed);
574    }
575
576    fn v2_disabled_status_code(&self) -> u16 {
577        self.v2_disabled_status.load(Ordering::Relaxed)
578    }
579
580    async fn get_reconstruction_v1(
581        &self,
582        file_id: &MerkleHash,
583        bytes_range: Option<FileRange>,
584    ) -> Result<Option<QueryReconstructionResponse>> {
585        LocalClient::get_reconstruction_v1(self, file_id, bytes_range).await
586    }
587
588    async fn get_reconstruction_v2(
589        &self,
590        file_id: &MerkleHash,
591        bytes_range: Option<FileRange>,
592    ) -> Result<Option<QueryReconstructionResponseV2>> {
593        LocalClient::get_reconstruction_v2(self, file_id, bytes_range).await
594    }
595
596    fn set_api_delay_range(&self, delay_range: Option<Range<Duration>>) {
597        match delay_range {
598            Some(range) => {
599                self.random_ms_delay_window
600                    .0
601                    .store(range.start.as_millis() as u64, Ordering::Relaxed);
602                self.random_ms_delay_window
603                    .1
604                    .store(range.end.as_millis() as u64, Ordering::Relaxed);
605            },
606            None => {
607                self.random_ms_delay_window.0.store(0, Ordering::Relaxed);
608                self.random_ms_delay_window.1.store(0, Ordering::Relaxed);
609            },
610        }
611    }
612
613    async fn apply_api_delay(&self) {
614        let min_ms = self.random_ms_delay_window.0.load(Ordering::Relaxed);
615        let max_ms = self.random_ms_delay_window.1.load(Ordering::Relaxed);
616
617        if min_ms == 0 && max_ms == 0 {
618            return;
619        }
620
621        let delay_ms = if min_ms == max_ms {
622            min_ms
623        } else {
624            rand::rng().random_range(min_ms..max_ms)
625        };
626
627        tokio::time::sleep(Duration::from_millis(delay_ms)).await;
628    }
629
630    async fn list_xorbs(&self) -> Result<Vec<MerkleHash>> {
631        let mut ret = Vec::new();
632        self.xorb_dir
633            .read_dir()
634            .map_err(ClientError::internal)?
635            .filter_map(|x| x.ok())
636            .filter_map(|x| x.file_name().into_string().ok())
637            .for_each(|x| {
638                if x.ends_with(".gctag") {
639                    return;
640                }
641                if let Some(pos) = x.rfind('.') {
642                    let hash = &x[(pos + 1)..];
643                    if let Ok(hash) = MerkleHash::from_hex(hash) {
644                        ret.push(hash);
645                    }
646                }
647            });
648        Ok(ret)
649    }
650
651    async fn get_full_xorb(&self, hash: &MerkleHash) -> Result<Bytes> {
652        let file_path = self.get_path_for_entry(hash);
653        let file = File::open(&file_path).map_err(|_| {
654            error!("Unable to find file in local CAS {:?}", file_path);
655            ClientError::XORBNotFound(*hash)
656        })?;
657
658        let mut reader = BufReader::new(file);
659        let xorb_obj = XorbObject::deserialize(&mut reader)?;
660        let result = xorb_obj.get_all_bytes(&mut reader)?;
661        Ok(Bytes::from(result))
662    }
663
664    async fn get_xorb_ranges(&self, hash: &MerkleHash, chunk_ranges: Vec<(u32, u32)>) -> Result<Vec<Bytes>> {
665        if chunk_ranges.is_empty() {
666            return Ok(vec![Bytes::new()]);
667        }
668
669        let file_path = self.get_path_for_entry(hash);
670        let file = File::open(&file_path).map_err(|_| {
671            error!("Unable to find file in local CAS {:?}", file_path);
672            ClientError::XORBNotFound(*hash)
673        })?;
674
675        let mut reader = BufReader::new(file);
676        let xorb_obj = XorbObject::deserialize(&mut reader)?;
677
678        let mut ret: Vec<Bytes> = Vec::new();
679        for r in chunk_ranges {
680            if r.0 >= r.1 {
681                ret.push(Bytes::new());
682                continue;
683            }
684
685            let data = xorb_obj.get_bytes_by_chunk_range(&mut reader, r.0, r.1)?;
686            ret.push(Bytes::from(data));
687        }
688        Ok(ret)
689    }
690
691    async fn xorb_length(&self, hash: &MerkleHash) -> Result<u32> {
692        let file_path = self.get_path_for_entry(hash);
693        match File::open(file_path) {
694            Ok(file) => {
695                let mut reader = BufReader::new(file);
696                let xorb_obj = XorbObject::deserialize(&mut reader)?;
697                let length = xorb_obj.get_all_bytes(&mut reader)?.len();
698                Ok(length as u32)
699            },
700            Err(_) => Err(ClientError::XORBNotFound(*hash)),
701        }
702    }
703
704    async fn xorb_exists(&self, hash: &MerkleHash) -> Result<bool> {
705        let file_path = self.get_path_for_entry(hash);
706
707        let Ok(md) = metadata(&file_path) else {
708            return Ok(false);
709        };
710
711        if !md.is_file() {
712            return Err(ClientError::InternalError(anyhow!(
713                "Attempting to write to {file_path:?}, but it is not a file"
714            )));
715        }
716
717        let Ok(file) = File::open(&file_path) else {
718            return Err(ClientError::XORBNotFound(*hash));
719        };
720
721        let mut reader = BufReader::new(file);
722        XorbObject::deserialize(&mut reader)?;
723        Ok(true)
724    }
725
726    async fn xorb_footer(&self, hash: &MerkleHash) -> Result<XorbObject> {
727        let file_path = self.get_path_for_entry(hash);
728        let mut file = File::open(&file_path).map_err(|_| {
729            error!("Unable to find xorb in local CAS {:?}", file_path);
730            ClientError::XORBNotFound(*hash)
731        })?;
732
733        file.seek(SeekFrom::End(-(size_of::<u32>() as i64)))?;
734        let info_length = read_u32(&mut file)?;
735
736        file.seek(SeekFrom::End(-(info_length as i64)))?;
737
738        let mut reader = BufReader::new(file);
739        let xorb_obj = XorbObject::deserialize(&mut reader)?;
740        Ok(xorb_obj)
741    }
742
743    async fn get_file_size(&self, hash: &MerkleHash) -> Result<u64> {
744        let Some((file_info, _)) = self.get_file_info_from_table(hash)? else {
745            return Err(ClientError::FileNotFound(*hash));
746        };
747        Ok(file_info.file_size())
748    }
749
750    async fn get_file_data(&self, hash: &MerkleHash, byte_range: Option<FileRange>) -> Result<Bytes> {
751        let Some((file_info, _)) = self.get_file_info_from_table(hash)? else {
752            return Err(ClientError::FileNotFound(*hash));
753        };
754
755        let mut file_vec = Vec::new();
756        for entry in &file_info.segments {
757            let entry_bytes = self
758                .get_xorb_ranges(&entry.xorb_hash, vec![(entry.chunk_index_start, entry.chunk_index_end)])
759                .await?
760                .pop()
761                .unwrap();
762            file_vec.extend_from_slice(&entry_bytes);
763        }
764
765        let file_size = file_vec.len();
766
767        let start = byte_range.as_ref().map(|range| range.start as usize).unwrap_or(0);
768
769        if byte_range.is_some() && start >= file_size {
770            return Err(ClientError::InvalidRange);
771        }
772
773        let end = byte_range
774            .as_ref()
775            .map(|range| range.end as usize)
776            .unwrap_or(file_size)
777            .min(file_size);
778
779        Ok(Bytes::from(file_vec[start..end].to_vec()))
780    }
781
782    async fn get_xorb_raw_bytes(&self, hash: &MerkleHash, byte_range: Option<FileRange>) -> Result<Bytes> {
783        let file_path = self.get_path_for_entry(hash);
784        let data = std::fs::read(&file_path).map_err(|_| ClientError::XORBNotFound(*hash))?;
785
786        let start = byte_range.as_ref().map(|r| r.start as usize).unwrap_or(0);
787        let end = byte_range
788            .as_ref()
789            .map(|r| r.end as usize)
790            .unwrap_or(data.len())
791            .min(data.len());
792
793        if start >= data.len() {
794            return Err(ClientError::InvalidRange);
795        }
796
797        Ok(Bytes::from(data[start..end].to_vec()))
798    }
799
800    async fn xorb_raw_length(&self, hash: &MerkleHash) -> Result<u64> {
801        let file_path = self.get_path_for_entry(hash);
802        let metadata = std::fs::metadata(&file_path).map_err(|_| ClientError::XORBNotFound(*hash))?;
803        Ok(metadata.len())
804    }
805
806    async fn fetch_term_data(
807        &self,
808        hash: MerkleHash,
809        fetch_term: XorbReconstructionFetchInfo,
810    ) -> Result<(Bytes, Vec<u32>)> {
811        self.apply_api_delay().await;
812        let (file_path, url_byte_range, url_timestamp) = parse_fetch_url(&fetch_term.url)?;
813
814        // Check if URL has expired
815        let expiration_ms = self.url_expiration_ms.load(Ordering::Relaxed);
816        let elapsed_ms = Instant::now().saturating_duration_since(url_timestamp).as_millis() as u64;
817        if elapsed_ms > expiration_ms {
818            return Err(ClientError::PresignedUrlExpirationError);
819        }
820
821        // Validate byte range matches url_range
822        // Note: url_byte_range is FileRange (exclusive end), url_range is HttpRange (inclusive end)
823        // We convert url_range to FileRange for comparison
824        let fetch_byte_range = FileRange::from(fetch_term.url_range);
825        if url_byte_range.start != fetch_byte_range.start || url_byte_range.end != fetch_byte_range.end {
826            return Err(ClientError::InvalidArguments);
827        }
828        let file = File::open(&file_path).map_err(|_| {
829            error!("Unable to find xorb in local CAS {:?}", file_path);
830            ClientError::XORBNotFound(hash)
831        })?;
832
833        let mut reader = BufReader::new(file);
834        let xorb_obj = XorbObject::deserialize(&mut reader)?;
835
836        let data = xorb_obj.get_bytes_by_chunk_range(&mut reader, fetch_term.range.start, fetch_term.range.end)?;
837
838        let chunk_byte_indices = {
839            let mut indices = Vec::new();
840            let mut cumulative = 0u32;
841            // Start with 0, matching the format from deserialize_chunks_from_stream
842            indices.push(0);
843            // ChunkRange is exclusive-end, so we iterate from start to end (exclusive)
844            for chunk_idx in fetch_term.range.start..fetch_term.range.end {
845                let chunk_len = xorb_obj
846                    .uncompressed_chunk_length(chunk_idx)
847                    .map_err(|e| ClientError::Other(format!("Failed to get chunk length: {e}")))?;
848                cumulative += chunk_len;
849                indices.push(cumulative);
850            }
851            indices
852        };
853
854        Ok((data.into(), chunk_byte_indices))
855    }
856}
857
858impl LocalClient {
859    /// Removes all FILE_TO_SHARD_TABLE entries whose shard_hash equals `shard_hash`.
860    fn remove_file_entries_for_shard(&self, shard_hash: &MerkleHash) -> Result<()> {
861        let to_remove: Vec<RedbHash> = {
862            let read_txn = self.db().begin_read().map_err(map_redb_db_error)?;
863            let table = read_txn.open_table(FILE_TO_SHARD_TABLE).map_err(map_redb_db_error)?;
864            table
865                .iter()
866                .map_err(map_redb_db_error)?
867                .filter_map(|e| e.ok())
868                .filter(|(_, v)| v.value().shard_hash == *shard_hash)
869                .map(|(k, _)| k.value())
870                .collect()
871        };
872        if !to_remove.is_empty() {
873            let write_txn = self.db().begin_write().map_err(map_redb_db_error)?;
874            {
875                let mut table = write_txn.open_table(FILE_TO_SHARD_TABLE).map_err(map_redb_db_error)?;
876                for key in &to_remove {
877                    table.remove(key).map_err(map_redb_db_error)?;
878                }
879            }
880            write_txn.commit().map_err(map_redb_db_error)?;
881        }
882        Ok(())
883    }
884}
885
886#[async_trait]
887impl super::DeletionControlableClient for LocalClient {
888    async fn list_shard_entries(&self) -> Result<Vec<MerkleHash>> {
889        Ok(self.shard_file_paths()?.into_iter().map(|(h, _)| h).collect())
890    }
891
892    async fn get_shard_bytes(&self, hash: &MerkleHash) -> Result<Bytes> {
893        let path = self.shard_path_for_hash(hash)?;
894        let data = std::fs::read(&path)?;
895        Ok(Bytes::from(data))
896    }
897
898    async fn delete_shard_entry(&self, hash: &MerkleHash) -> Result<()> {
899        let path = self.shard_path_for_hash(hash)?;
900        self.remove_file_entries_for_shard(hash)?;
901        if self.lifecycle_tag_deletion_enabled() {
902            let gctag = self.gctag_shard_path(hash);
903            std::fs::rename(&path, &gctag)?;
904        } else {
905            std::fs::remove_file(&path)?;
906        }
907        Ok(())
908    }
909
910    async fn list_file_shard_entries(&self) -> Result<Vec<(MerkleHash, MerkleHash)>> {
911        let read_txn = self.db().begin_read().map_err(map_redb_db_error)?;
912        let table = read_txn.open_table(FILE_TO_SHARD_TABLE).map_err(map_redb_db_error)?;
913        let mut entries = Vec::new();
914        for entry in table.iter().map_err(map_redb_db_error)? {
915            let (key, value) = entry.map_err(map_redb_db_error)?;
916            let file_hash: MerkleHash = key.value().into();
917            let shard_ref: FileShardRef = value.value();
918            entries.push((file_hash, shard_ref.shard_hash));
919        }
920        Ok(entries)
921    }
922
923    async fn delete_file_entry(&self, file_hash: &MerkleHash) -> Result<()> {
924        let write_txn = self.db().begin_write().map_err(map_redb_db_error)?;
925        {
926            let mut table = write_txn.open_table(FILE_TO_SHARD_TABLE).map_err(map_redb_db_error)?;
927            table.remove(&RedbHash::from(*file_hash)).map_err(map_redb_db_error)?;
928        }
929        write_txn.commit().map_err(map_redb_db_error)?;
930        Ok(())
931    }
932
933    async fn remove_shard_dedup_entries(&self, shard_hash: &MerkleHash) -> Result<()> {
934        let shard_redb = RedbHash::from(*shard_hash);
935        for _ in 0..4 {
936            let to_delete: Vec<RedbHash> = {
937                let read_txn = self.db().begin_read().map_err(map_redb_db_error)?;
938                let table = read_txn.open_table(GLOBAL_DEDUP_TABLE).map_err(map_redb_db_error)?;
939                table
940                    .iter()
941                    .map_err(map_redb_db_error)?
942                    .filter_map(|entry| entry.ok())
943                    .filter(|(_, v)| v.value() == shard_redb)
944                    .map(|(k, _)| k.value())
945                    .collect()
946            };
947
948            if to_delete.is_empty() {
949                return Ok(());
950            }
951
952            let write_txn = self.db().begin_write().map_err(map_redb_db_error)?;
953            {
954                let mut table = write_txn.open_table(GLOBAL_DEDUP_TABLE).map_err(map_redb_db_error)?;
955                for chunk_hash in &to_delete {
956                    table.remove(chunk_hash).map_err(map_redb_db_error)?;
957                }
958            }
959            write_txn.commit().map_err(map_redb_db_error)?;
960        }
961
962        let still_present = {
963            let read_txn = self.db().begin_read().map_err(map_redb_db_error)?;
964            let table = read_txn.open_table(GLOBAL_DEDUP_TABLE).map_err(map_redb_db_error)?;
965            table
966                .iter()
967                .map_err(map_redb_db_error)?
968                .filter_map(|entry| entry.ok())
969                .any(|(_, v)| v.value() == shard_redb)
970        };
971
972        if still_present {
973            return Err(ClientError::Other(format!(
974                "Unable to fully remove dedup entries for shard {} due to concurrent updates",
975                shard_hash.hex()
976            )));
977        }
978
979        Ok(())
980    }
981
982    async fn delete_xorb(&self, hash: &MerkleHash) {
983        let file_path = self.get_path_for_entry(hash);
984
985        #[cfg(windows)]
986        Self::clear_readonly(&file_path);
987
988        if self.lifecycle_tag_deletion_enabled() {
989            let gctag = self.gctag_xorb_path(hash);
990            #[cfg(windows)]
991            Self::clear_readonly(&gctag);
992            let _ = std::fs::rename(&file_path, &gctag);
993        } else {
994            let _ = std::fs::remove_file(file_path);
995        }
996    }
997
998    async fn list_xorbs_and_tags(&self) -> Result<Vec<(MerkleHash, ObjectTag)>> {
999        let mut ret = Vec::new();
1000        for entry in self.xorb_dir.read_dir().map_err(ClientError::internal)? {
1001            let entry = entry.map_err(ClientError::internal)?;
1002            let path = entry.path();
1003            let Some(name) = entry.file_name().into_string().ok() else {
1004                continue;
1005            };
1006            if let Some(pos) = name.rfind('.') {
1007                let hex = &name[(pos + 1)..];
1008                if let Ok(hash) = MerkleHash::from_hex(hex) {
1009                    let tag = Self::object_tag_from_path(&path)?;
1010                    ret.push((hash, tag));
1011                }
1012            }
1013        }
1014        Ok(ret)
1015    }
1016
1017    async fn delete_xorb_if_tag_matches(&self, hash: &MerkleHash, tag: &ObjectTag) -> Result<bool> {
1018        let file_path = self.get_path_for_entry(hash);
1019
1020        // Atomically move the file out of the namespace before checking the
1021        // tag.  This closes the TOCTOU window with concurrent upload_xorb
1022        // (which always rewrites via SafeFileCreator atomic rename).
1023        let tmp_path = file_path.with_extension(format!("gc_del_{:x}", rand::random::<u64>()));
1024        if std::fs::rename(&file_path, &tmp_path).is_err() {
1025            return Err(ClientError::XORBNotFound(*hash));
1026        }
1027
1028        let current_tag = match Self::object_tag_from_path(&tmp_path) {
1029            Ok(t) => t,
1030            Err(e) => {
1031                Self::restore_from_tmp(&tmp_path, &file_path);
1032                return Err(e);
1033            },
1034        };
1035
1036        if &current_tag != tag {
1037            Self::restore_from_tmp(&tmp_path, &file_path);
1038            return Ok(false);
1039        }
1040
1041        #[cfg(windows)]
1042        Self::clear_readonly(&tmp_path);
1043
1044        if self.lifecycle_tag_deletion_enabled() {
1045            let gctag = self.gctag_xorb_path(hash);
1046            #[cfg(windows)]
1047            Self::clear_readonly(&gctag);
1048            std::fs::rename(&tmp_path, &gctag)?;
1049        } else {
1050            std::fs::remove_file(&tmp_path)?;
1051        }
1052        Ok(true)
1053    }
1054
1055    async fn list_shards_with_tags(&self) -> Result<Vec<(MerkleHash, ObjectTag)>> {
1056        let mut ret = Vec::new();
1057        for (hash, path) in self.shard_file_paths()? {
1058            let tag = Self::object_tag_from_path(&path)?;
1059            ret.push((hash, tag));
1060        }
1061        Ok(ret)
1062    }
1063
1064    async fn delete_shard_if_tag_matches(&self, hash: &MerkleHash, tag: &ObjectTag) -> Result<bool> {
1065        let path = self.shard_path_for_hash(hash)?;
1066
1067        let tmp_path = path.with_extension(format!("gc_del_{:x}", rand::random::<u64>()));
1068        if std::fs::rename(&path, &tmp_path).is_err() {
1069            return Err(ClientError::Other(format!("Shard not found: {}", hash.hex())));
1070        }
1071
1072        let current_tag = match Self::object_tag_from_path(&tmp_path) {
1073            Ok(t) => t,
1074            Err(e) => {
1075                Self::restore_from_tmp(&tmp_path, &path);
1076                return Err(e);
1077            },
1078        };
1079
1080        if &current_tag != tag {
1081            Self::restore_from_tmp(&tmp_path, &path);
1082            return Ok(false);
1083        }
1084
1085        if let Err(e) = self.remove_file_entries_for_shard(hash) {
1086            Self::restore_from_tmp(&tmp_path, &path);
1087            return Err(e);
1088        }
1089        if self.lifecycle_tag_deletion_enabled() {
1090            let gctag = self.gctag_shard_path(hash);
1091            std::fs::rename(&tmp_path, &gctag)?;
1092        } else {
1093            std::fs::remove_file(&tmp_path)?;
1094        }
1095        Ok(true)
1096    }
1097
1098    /// Verifies referential integrity of all shards on disk:
1099    /// 1. For each XORB entry listed in any shard, the corresponding XORB file must exist on disk.
1100    /// 2. For each file entry in any shard, every referenced XORB must exist on disk. (Global-dedup allows file entries
1101    ///    to reference XORBs described in a different shard, so we do a cross-shard check against disk rather than a
1102    ///    within-shard check.)
1103    async fn verify_all_reachable(&self) -> Result<()> {
1104        let shard_files = self.shard_file_paths()?;
1105
1106        // Build a map of file_hash -> shard_hash from the authoritative table.
1107        // A file entry in a shard is only considered active if the table maps that
1108        // file hash to that specific shard, preventing stale entries from resurrecting.
1109        let file_to_shard: HashMap<MerkleHash, MerkleHash> = {
1110            let read_txn = self.db().begin_read().map_err(map_redb_db_error)?;
1111            let table = read_txn.open_table(FILE_TO_SHARD_TABLE).map_err(map_redb_db_error)?;
1112            let mut map = HashMap::new();
1113            for entry in table.iter().map_err(map_redb_db_error)? {
1114                let (k, v) = entry.map_err(map_redb_db_error)?;
1115                let fh: MerkleHash = k.value().into();
1116                let sr: FileShardRef = v.value();
1117                map.insert(fh, sr.shard_hash);
1118            }
1119            map
1120        };
1121
1122        // Collect xorbs claimed by shard xorb-entries, xorbs referenced by active file
1123        // entries (cross-shard dedup), and which shards have at least one active file.
1124        let mut xorbs_in_shard_entries: std::collections::HashSet<MerkleHash> = std::collections::HashSet::new();
1125        let mut xorbs_in_active_file_entries: std::collections::HashSet<MerkleHash> = std::collections::HashSet::new();
1126        let mut shards_with_active_files: std::collections::HashSet<MerkleHash> = std::collections::HashSet::new();
1127        let mut shard_xorbs: std::collections::HashMap<MerkleHash, Vec<MerkleHash>> = std::collections::HashMap::new();
1128
1129        for (shard_hash, path) in &shard_files {
1130            let shard_bytes = std::fs::read(path)?;
1131            let minimal_shard = MDBMinimalShard::from_reader(&mut Cursor::new(&shard_bytes), true, true)?;
1132
1133            for i in 0..minimal_shard.num_xorb() {
1134                let xorb_hash = minimal_shard.xorb(i).unwrap().xorb_hash();
1135                xorbs_in_shard_entries.insert(xorb_hash);
1136                shard_xorbs.entry(*shard_hash).or_default().push(xorb_hash);
1137            }
1138
1139            let mut has_active_file = false;
1140            for i in 0..minimal_shard.num_files() {
1141                let file_view = minimal_shard.file(i).unwrap();
1142                let fh = file_view.file_hash();
1143                if file_to_shard.get(&fh) == Some(shard_hash) {
1144                    has_active_file = true;
1145                    for seg_idx in 0..file_view.num_entries() {
1146                        xorbs_in_active_file_entries.insert(file_view.entry(seg_idx).xorb_hash);
1147                    }
1148                }
1149            }
1150            if has_active_file {
1151                shards_with_active_files.insert(*shard_hash);
1152            }
1153        }
1154
1155        let mut errors: Vec<String> = Vec::new();
1156
1157        // Check 1: every on-disk shard must be "reachable".
1158        // A shard is reachable if it has at least one active (non-deleted) file entry, OR if it is
1159        // a compact shard (no file entries) that holds xorb-entries for xorbs referenced by some
1160        // active file (cross-shard dedup case — the compact shard is needed for dedup lookups).
1161        // A shard that satisfies neither condition is truly orphaned and GC should have deleted it.
1162        for (shard_hash, _) in &shard_files {
1163            if !shards_with_active_files.contains(shard_hash) {
1164                let has_file_referenced_xorb = shard_xorbs
1165                    .get(shard_hash)
1166                    .is_some_and(|xorbs| xorbs.iter().any(|x| xorbs_in_active_file_entries.contains(x)));
1167                if !has_file_referenced_xorb {
1168                    errors.push(format!(
1169                        "Reachability error: shard {} has no active file entries and no \
1170                         xorbs referenced by any active file (GC should have deleted it)",
1171                        shard_hash.hex()
1172                    ));
1173                }
1174            }
1175        }
1176
1177        // Check 2: every on-disk xorb must be reachable — either indexed by a shard's
1178        // xorb entries, or referenced directly by an active file's file entries (cross-shard
1179        // dedup). Xorbs that satisfy neither are orphaned and GC should have deleted them.
1180        for xorb_hash in self.list_xorbs().await? {
1181            if !xorbs_in_shard_entries.contains(&xorb_hash) && !xorbs_in_active_file_entries.contains(&xorb_hash) {
1182                errors.push(format!(
1183                    "Reachability error: xorb {} exists on disk but is not referenced by \
1184                     any shard xorb entry or active file entry (GC should have deleted it)",
1185                    xorb_hash.hex()
1186                ));
1187            }
1188        }
1189
1190        if errors.is_empty() {
1191            Ok(())
1192        } else {
1193            Err(ClientError::Other(errors.join("\n")))
1194        }
1195    }
1196
1197    async fn verify_integrity(&self) -> Result<()> {
1198        // Snapshot the dedup table before reading the directory.  This prevents
1199        // a TOCTOU race: upload_shard writes the file then commits the dedup
1200        // entry, so any entry visible in this MVCC snapshot is guaranteed to
1201        // have its shard file already on disk when we read the directory below.
1202        let read_txn = self.db().begin_read().map_err(map_redb_db_error)?;
1203
1204        let shard_files = self.shard_file_paths()?;
1205
1206        // Pass 1: collect all XORB hashes listed across all shards and build
1207        // a global chunk-count index.  Also verify each listed XORB file exists.
1208        let mut global_xorb_chunk_counts: HashMap<MerkleHash, usize> = HashMap::new();
1209        for (shard_hash, path) in &shard_files {
1210            let shard_bytes = std::fs::read(path)?;
1211            let minimal_shard = MDBMinimalShard::from_reader(&mut Cursor::new(&shard_bytes), true, true)?;
1212
1213            for i in 0..minimal_shard.num_xorb() {
1214                let xorb_view = minimal_shard.xorb(i).unwrap();
1215                let xorb_hash = xorb_view.xorb_hash();
1216
1217                let xorb_path = self.get_path_for_entry(&xorb_hash);
1218                if !xorb_path.exists() {
1219                    return Err(ClientError::Other(format!(
1220                        "Integrity error: shard {} references non-existent XORB {}",
1221                        shard_hash.hex(),
1222                        xorb_hash.hex()
1223                    )));
1224                }
1225
1226                global_xorb_chunk_counts.entry(xorb_hash).or_insert(xorb_view.num_entries());
1227            }
1228        }
1229
1230        // Pass 2: validate every file entry registered in FILE_TO_SHARD_TABLE.
1231        // Uses offset/length for direct-seek reads — no full-shard parsing needed.
1232        // Reuses `read_txn` from the top so all passes see a consistent MVCC snapshot.
1233        let file_table = read_txn.open_table(FILE_TO_SHARD_TABLE).map_err(map_redb_db_error)?;
1234
1235        for entry in file_table.iter().map_err(map_redb_db_error)? {
1236            let (key, value) = entry.map_err(map_redb_db_error)?;
1237            let file_hash: MerkleHash = key.value().into();
1238            let shard_ref: FileShardRef = value.value();
1239
1240            let shard_path = self.shard_dir.join(shard_file_name(&shard_ref.shard_hash));
1241            if !shard_path.exists() {
1242                return Err(ClientError::Other(format!(
1243                    "Integrity error: FILE_TO_SHARD_TABLE maps file {} to shard {} which does not exist on disk",
1244                    file_hash.hex(),
1245                    shard_ref.shard_hash.hex()
1246                )));
1247            }
1248
1249            let mut shard_file = File::open(&shard_path)?;
1250            shard_file.seek(SeekFrom::Start(shard_ref.offset))?;
1251            let mut buf = vec![0u8; shard_ref.length as usize];
1252            shard_file.read_exact(&mut buf)?;
1253
1254            let file_view = MDBFileInfoView::new(Bytes::from(buf)).map_err(|e| {
1255                ClientError::Other(format!(
1256                    "Integrity error: cannot parse file entry for {} in shard {} at offset {}: {}",
1257                    file_hash.hex(),
1258                    shard_ref.shard_hash.hex(),
1259                    shard_ref.offset,
1260                    e
1261                ))
1262            })?;
1263
1264            if file_view.file_hash() != file_hash {
1265                return Err(ClientError::Other(format!(
1266                    "Integrity error: FILE_TO_SHARD_TABLE maps file {} to shard {} offset {} but found file {} there",
1267                    file_hash.hex(),
1268                    shard_ref.shard_hash.hex(),
1269                    shard_ref.offset,
1270                    file_view.file_hash().hex()
1271                )));
1272            }
1273
1274            for seg_idx in 0..file_view.num_entries() {
1275                let segment = file_view.entry(seg_idx);
1276                let xorb_path = self.get_path_for_entry(&segment.xorb_hash);
1277
1278                if let Some(&chunk_count) = global_xorb_chunk_counts.get(&segment.xorb_hash) {
1279                    if segment.chunk_index_end as usize > chunk_count {
1280                        return Err(ClientError::Other(format!(
1281                            "Integrity error: file {} references chunk range {}..{} \
1282                             but XORB block {} only has {} chunks",
1283                            file_hash.hex(),
1284                            segment.chunk_index_start,
1285                            segment.chunk_index_end,
1286                            segment.xorb_hash.hex(),
1287                            chunk_count
1288                        )));
1289                    }
1290                } else if xorb_path.exists() {
1291                    // XORB not in any shard index but file exists — dedup reference, OK.
1292                } else {
1293                    return Err(ClientError::Other(format!(
1294                        "Integrity error: file {} in shard {} references XORB {} \
1295                         that has no shard index entry and no XORB file on disk",
1296                        file_hash.hex(),
1297                        shard_ref.shard_hash.hex(),
1298                        segment.xorb_hash.hex()
1299                    )));
1300                }
1301            }
1302        }
1303
1304        // Pass 3: verify that all shards referenced by the global dedup chunk table
1305        // are present on disk.  Uses the read transaction snapshotted at the top
1306        // of this function to avoid TOCTOU races with concurrent uploads.
1307        let shard_hashes_on_disk: std::collections::HashSet<MerkleHash> = shard_files.iter().map(|(h, _)| *h).collect();
1308
1309        let dedup_table = read_txn.open_table(GLOBAL_DEDUP_TABLE).map_err(map_redb_db_error)?;
1310        for entry in dedup_table.iter().map_err(map_redb_db_error)? {
1311            let (chunk_key, shard_val) = entry.map_err(map_redb_db_error)?;
1312            let shard_hash: MerkleHash = shard_val.value().into();
1313            if !shard_hashes_on_disk.contains(&shard_hash) {
1314                let chunk_hash: MerkleHash = chunk_key.value().into();
1315                return Err(ClientError::Other(format!(
1316                    "Integrity error: global dedup table maps chunk {} to shard {} \
1317                     which does not exist on disk",
1318                    chunk_hash.hex(),
1319                    shard_hash.hex()
1320                )));
1321            }
1322        }
1323
1324        Ok(())
1325    }
1326}
1327
1328impl LocalClient {
1329    /// Looks up a file hash in FILE_TO_SHARD_TABLE and reads its reconstruction info
1330    /// via a direct-seek into the canonical shard on disk.  Returns `None` if the file
1331    /// is not registered (i.e. deleted or never uploaded).
1332    fn get_file_info_from_table(&self, file_hash: &MerkleHash) -> Result<Option<(MDBFileInfo, MerkleHash)>> {
1333        let read_txn = self.db().begin_read().map_err(map_redb_db_error)?;
1334        let table = read_txn.open_table(FILE_TO_SHARD_TABLE).map_err(map_redb_db_error)?;
1335        let Some(entry) = table.get(&RedbHash::from(*file_hash)).map_err(map_redb_db_error)? else {
1336            return Ok(None);
1337        };
1338        let shard_ref: FileShardRef = entry.value();
1339        let shard_path = self.shard_dir.join(shard_file_name(&shard_ref.shard_hash));
1340
1341        let mut file = File::open(&shard_path)?;
1342        file.seek(SeekFrom::Start(shard_ref.offset))?;
1343        let mut buf = vec![0u8; shard_ref.length as usize];
1344        file.read_exact(&mut buf)?;
1345
1346        let file_view = MDBFileInfoView::new(Bytes::from(buf))?;
1347        Ok(Some((MDBFileInfo::from(&file_view), shard_ref.shard_hash)))
1348    }
1349
1350    async fn compute_reconstruction_ranges(
1351        &self,
1352        file_id: &MerkleHash,
1353        bytes_range: Option<FileRange>,
1354    ) -> Result<xorb_utils::ReconstructionRangesResult> {
1355        let Some((file_info, _)) = self.get_file_info_from_table(file_id)? else {
1356            return Ok(None);
1357        };
1358
1359        xorb_utils::compute_reconstruction_ranges(&file_info, bytes_range, &mut |hash| self.xorb_footer_sync(hash))
1360    }
1361
1362    fn xorb_footer_sync(&self, hash: &MerkleHash) -> Result<XorbObject> {
1363        let file_path = self.get_path_for_entry(hash);
1364        let mut file = File::open(&file_path).map_err(|_| {
1365            error!("Unable to find file in local CAS {:?}", file_path);
1366            ClientError::XORBNotFound(*hash)
1367        })?;
1368        XorbObject::deserialize(&mut file).map_err(Into::into)
1369    }
1370
1371    /// V1 reconstruction: returns per-range presigned URLs.
1372    pub async fn get_reconstruction_v1(
1373        &self,
1374        file_id: &MerkleHash,
1375        bytes_range: Option<FileRange>,
1376    ) -> Result<Option<QueryReconstructionResponse>> {
1377        self.apply_api_delay().await;
1378
1379        let result = self.compute_reconstruction_ranges(file_id, bytes_range).await?;
1380        let Some((offset_into_first_range, terms, merged_ranges)) = result else {
1381            return Ok(None);
1382        };
1383
1384        if terms.is_empty() {
1385            return Ok(Some(QueryReconstructionResponse {
1386                offset_into_first_range,
1387                terms,
1388                fetch_info: HashMap::new(),
1389            }));
1390        }
1391
1392        let timestamp = Instant::now();
1393        let mut fetch_info: HashMap<HexMerkleHash, Vec<XorbReconstructionFetchInfo>> = HashMap::new();
1394        for (hash, ranges) in merged_ranges {
1395            let file_path = self.get_path_for_entry(&hash);
1396            let entries = ranges
1397                .into_iter()
1398                .map(|r| XorbReconstructionFetchInfo {
1399                    range: r.chunk_range,
1400                    url: generate_fetch_url(&file_path, &r.byte_range, timestamp),
1401                    url_range: HttpRange::from(r.byte_range),
1402                })
1403                .collect();
1404            fetch_info.insert(hash.into(), entries);
1405        }
1406
1407        Ok(Some(QueryReconstructionResponse {
1408            offset_into_first_range,
1409            terms,
1410            fetch_info,
1411        }))
1412    }
1413
1414    /// V2 reconstruction: returns per-xorb multi-range fetch descriptors.
1415    pub async fn get_reconstruction_v2(
1416        &self,
1417        file_id: &MerkleHash,
1418        bytes_range: Option<FileRange>,
1419    ) -> Result<Option<QueryReconstructionResponseV2>> {
1420        self.apply_api_delay().await;
1421
1422        let result = self.compute_reconstruction_ranges(file_id, bytes_range).await?;
1423        let Some((offset_into_first_range, terms, merged_ranges)) = result else {
1424            return Ok(None);
1425        };
1426
1427        if terms.is_empty() {
1428            return Ok(Some(QueryReconstructionResponseV2 {
1429                offset_into_first_range,
1430                terms,
1431                xorbs: HashMap::new(),
1432            }));
1433        }
1434
1435        let timestamp = Instant::now();
1436        let max_ranges = self.max_ranges_per_fetch.load(Ordering::Relaxed);
1437
1438        let mut xorbs: HashMap<HexMerkleHash, Vec<XorbMultiRangeFetch>> = HashMap::new();
1439        for (hash, ranges) in merged_ranges {
1440            let mut fetch_entries = Vec::new();
1441
1442            for chunk in ranges.chunks(max_ranges) {
1443                let range_descriptors: Vec<XorbRangeDescriptor> = chunk
1444                    .iter()
1445                    .map(|r| XorbRangeDescriptor {
1446                        chunks: r.chunk_range,
1447                        bytes: HttpRange::from(r.byte_range),
1448                    })
1449                    .collect();
1450
1451                let url = generate_v2_fetch_url(&hash, &range_descriptors, timestamp);
1452                fetch_entries.push(XorbMultiRangeFetch {
1453                    url,
1454                    ranges: range_descriptors,
1455                });
1456            }
1457
1458            xorbs.insert(hash.into(), fetch_entries);
1459        }
1460
1461        Ok(Some(QueryReconstructionResponseV2 {
1462            offset_into_first_range,
1463            terms,
1464            xorbs,
1465        }))
1466    }
1467}
1468
1469#[async_trait]
1470impl Client for LocalClient {
1471    async fn get_file_reconstruction_info(
1472        &self,
1473        file_hash: &MerkleHash,
1474    ) -> Result<Option<(MDBFileInfo, Option<MerkleHash>)>> {
1475        self.apply_api_delay().await;
1476        Ok(self.get_file_info_from_table(file_hash)?.map(|(info, sh)| (info, Some(sh))))
1477    }
1478
1479    async fn query_for_global_dedup_shard(&self, _prefix: &str, chunk_hash: &MerkleHash) -> Result<Option<Bytes>> {
1480        self.apply_api_delay().await;
1481        let read_txn = self.db().begin_read().map_err(map_redb_db_error)?;
1482        let table = read_txn.open_table(GLOBAL_DEDUP_TABLE).map_err(map_redb_db_error)?;
1483
1484        if let Some(shard) = table.get(&RedbHash::from(*chunk_hash)).map_err(map_redb_db_error)? {
1485            let shard_hash: MerkleHash = shard.value().into();
1486            let filename = self.shard_dir.join(shard_file_name(&shard_hash));
1487
1488            let expiration_secs = self.global_dedup_expiration_secs.load(Ordering::Relaxed);
1489            if expiration_secs == 0 {
1490                return Ok(Some(std::fs::read(filename)?.into()));
1491            }
1492
1493            let expiry = std::time::SystemTime::now() + Duration::from_secs(expiration_secs);
1494            let shard_bytes = std::fs::read(filename)?;
1495
1496            let mut reader = Cursor::new(&shard_bytes);
1497            let minimal_shard = MDBMinimalShard::from_reader(&mut reader, true, true)?;
1498
1499            let mut out = Vec::new();
1500            minimal_shard.serialize_xorb_subset_with_expiry(&mut out, Some(expiry), |_| true)?;
1501            Ok(Some(out.into()))
1502        } else {
1503            Ok(None)
1504        }
1505    }
1506
1507    async fn acquire_upload_permit(&self) -> Result<super::super::adaptive_concurrency::ConnectionPermit> {
1508        self.apply_api_delay().await;
1509        self.upload_concurrency_controller.acquire_connection_permit().await
1510    }
1511
1512    async fn upload_shard(
1513        &self,
1514        shard_data: Bytes,
1515        _permit: super::super::adaptive_concurrency::ConnectionPermit,
1516    ) -> Result<bool> {
1517        self.apply_api_delay().await;
1518
1519        // Parse the shard using the streaming parser (handles shards without footer)
1520        let mut reader = Cursor::new(&shard_data);
1521        let minimal_shard = MDBMinimalShard::from_reader(&mut reader, true, true)?;
1522
1523        // Rebuild a full in-memory shard to rebuild the new shard.  Quick and convenient.
1524        let mut in_memory_shard = MDBInMemoryShard::default();
1525
1526        // Add file info from the views
1527        for i in 0..minimal_shard.num_files() {
1528            let file_view = minimal_shard.file(i).unwrap();
1529            in_memory_shard.add_file_reconstruction_info(MDBFileInfo::from(file_view))?;
1530        }
1531
1532        // Add XORB info from the views
1533        for i in 0..minimal_shard.num_xorb() {
1534            let xorb_view = minimal_shard.xorb(i).unwrap();
1535            in_memory_shard.add_xorb_block(MDBXorbInfo::from(xorb_view))?;
1536        }
1537
1538        // Write the rebuilt shard to disk (creates proper lookup tables)
1539        let shard_path = in_memory_shard.write_to_directory(&self.shard_dir, None)?;
1540        let shard = MDBShardFile::load_from_file(&shard_path, self.shard_manager.shard_file_cache())?;
1541        let shard_hash = shard.shard_hash;
1542
1543        self.shard_manager.register_shards(&[shard]).await?;
1544
1545        // Get global dedup chunks from the minimal shard
1546        let chunk_hashes = minimal_shard.global_dedup_eligible_chunks();
1547
1548        // Compute byte ranges for each file entry in the written shard
1549        let written_shard_bytes = std::fs::read(&shard_path)?;
1550        let file_ranges = file_entry_byte_ranges(&written_shard_bytes)?;
1551
1552        let shard_hash_redb = RedbHash::from(shard_hash);
1553        let write_txn = self.db().begin_write().map_err(map_redb_db_error)?;
1554        {
1555            let mut dedup_table = write_txn.open_table(GLOBAL_DEDUP_TABLE).map_err(map_redb_db_error)?;
1556            for chunk in chunk_hashes {
1557                dedup_table
1558                    .insert(&RedbHash::from(chunk), &shard_hash_redb)
1559                    .map_err(map_redb_db_error)?;
1560            }
1561
1562            let mut file_table = write_txn.open_table(FILE_TO_SHARD_TABLE).map_err(map_redb_db_error)?;
1563            for (file_hash, offset, length) in &file_ranges {
1564                file_table
1565                    .insert(
1566                        &RedbHash::from(*file_hash),
1567                        &FileShardRef {
1568                            shard_hash,
1569                            offset: *offset,
1570                            length: *length,
1571                        },
1572                    )
1573                    .map_err(map_redb_db_error)?;
1574            }
1575        }
1576        write_txn.commit().map_err(map_redb_db_error)?;
1577
1578        // A re-upload of the same shard hash supersedes any prior
1579        // lifecycle-tagged copy: clear the `.gctag` file so the shard is
1580        // readable again. Mirrors S3 PutObject overwriting a tagged object.
1581        let _ = std::fs::remove_file(self.gctag_shard_path(&shard_hash));
1582
1583        Ok(true)
1584    }
1585
1586    async fn upload_xorb(
1587        &self,
1588        _prefix: &str,
1589        serialized_xorb_object: SerializedXorbObject,
1590        progress_callback: Option<ProgressCallback>,
1591        _permit: super::super::adaptive_concurrency::ConnectionPermit,
1592    ) -> Result<u64> {
1593        self.apply_api_delay().await;
1594        let hash = serialized_xorb_object.hash;
1595        let footer_start = serialized_xorb_object.footer_start;
1596        let serialized_data = serialized_xorb_object.serialized_data;
1597
1598        // Always rewrite: even if the xorb already exists, the file must be
1599        // re-created so its filesystem metadata (mtime/ctime) changes, producing
1600        // a new tag for delete_xorb_if_tag_matches.  SafeFileCreator uses
1601        // temp-file + atomic rename, so concurrent readers are safe.
1602
1603        // Reconstruct footer if not present
1604        let data_to_write = if footer_start.is_some() {
1605            serialized_data
1606        } else {
1607            let mut data_with_footer = Vec::new();
1608            let (_, computed_hash) = xet_core_structures::xorb_object::reconstruct_xorb_with_footer(
1609                &mut data_with_footer,
1610                &serialized_data,
1611            )?;
1612            if computed_hash != hash {
1613                return Err(ClientError::Other(format!(
1614                    "XORB hash mismatch: expected {}, got {}",
1615                    hash.hex(),
1616                    computed_hash.hex(),
1617                )));
1618            }
1619            data_with_footer
1620        };
1621
1622        let file_path = self.get_path_for_entry(&hash);
1623        info!("Writing XORB {hash:?} to local path {file_path:?}");
1624
1625        let total = data_to_write.len() as u64;
1626        let mut file = SafeFileCreator::new(&file_path)?;
1627
1628        for i in 0..10 {
1629            let start = (i * data_to_write.len()) / 10;
1630            let end = ((i + 1) * data_to_write.len()) / 10;
1631            let chunk_len = end - start;
1632
1633            file.write_all(&data_to_write[start..end])?;
1634
1635            if let Some(ref cb) = progress_callback {
1636                let completed = end as u64;
1637                let delta = chunk_len as u64;
1638                cb(delta, completed, total);
1639            }
1640        }
1641
1642        let bytes_written = data_to_write.len();
1643        file.close()?;
1644
1645        #[cfg(unix)]
1646        if let Ok(metadata) = metadata(&file_path) {
1647            let mut permissions = metadata.permissions();
1648            permissions.set_readonly(true);
1649            let _ = std::fs::set_permissions(&file_path, permissions);
1650        }
1651
1652        // A re-upload of the same xorb hash supersedes any prior
1653        // lifecycle-tagged copy: clear the `.gctag` file so the xorb is
1654        // readable again. Mirrors S3 PutObject overwriting a tagged object.
1655        let _ = std::fs::remove_file(self.gctag_xorb_path(&hash));
1656
1657        info!("{file_path:?} successfully written with {bytes_written} bytes.");
1658
1659        Ok(bytes_written as u64)
1660    }
1661
1662    async fn get_reconstruction(
1663        &self,
1664        file_id: &MerkleHash,
1665        bytes_range: Option<FileRange>,
1666    ) -> Result<Option<QueryReconstructionResponseV2>> {
1667        self.get_reconstruction_v2(file_id, bytes_range).await
1668    }
1669
1670    async fn batch_get_reconstruction(&self, file_ids: &[MerkleHash]) -> Result<BatchQueryReconstructionResponse> {
1671        self.apply_api_delay().await;
1672        let mut files = HashMap::new();
1673        let mut fetch_info_map: HashMap<HexMerkleHash, Vec<XorbReconstructionFetchInfo>> = HashMap::new();
1674
1675        for file_id in file_ids {
1676            if let Some(response) = self.get_reconstruction_v1(file_id, None).await? {
1677                let hex_hash: HexMerkleHash = (*file_id).into();
1678                files.insert(hex_hash, response.terms);
1679
1680                for (hash, fetch_infos) in response.fetch_info {
1681                    fetch_info_map.entry(hash).or_default().extend(fetch_infos);
1682                }
1683            }
1684        }
1685
1686        Ok(BatchQueryReconstructionResponse {
1687            files,
1688            fetch_info: fetch_info_map,
1689        })
1690    }
1691
1692    async fn acquire_download_permit(&self) -> Result<super::super::adaptive_concurrency::ConnectionPermit> {
1693        self.apply_api_delay().await;
1694        self.upload_concurrency_controller.acquire_connection_permit().await
1695    }
1696
1697    async fn get_file_term_data(
1698        &self,
1699        url_info: Box<dyn super::super::interface::URLProvider>,
1700        _download_permit: super::super::adaptive_concurrency::ConnectionPermit,
1701        progress_callback: Option<ProgressCallback>,
1702        uncompressed_size_if_known: Option<usize>,
1703    ) -> Result<(Bytes, Vec<u32>)> {
1704        // Retry loop: try to fetch, and if URL expired, refresh and retry once.
1705        for attempt in 0..2 {
1706            self.apply_api_delay().await;
1707            let (url, http_ranges) = url_info.retrieve_url().await?;
1708
1709            let (file_path, url_timestamp) = if let Ok((path, _, ts)) = parse_fetch_url(&url) {
1710                (path, ts)
1711            } else {
1712                let (hash, ts, _) = xorb_utils::parse_v2_fetch_url(&url)?;
1713                (self.get_path_for_entry(&hash), ts)
1714            };
1715
1716            // Check if URL has expired
1717            let expiration_ms = self.url_expiration_ms.load(Ordering::Relaxed);
1718            let elapsed_ms = Instant::now().saturating_duration_since(url_timestamp).as_millis() as u64;
1719            if elapsed_ms > expiration_ms {
1720                if attempt == 0 {
1721                    // First attempt failed due to expiration - refresh URL and retry.
1722                    url_info.refresh_url().await?;
1723                    continue;
1724                }
1725                return Err(ClientError::PresignedUrlExpirationError);
1726            }
1727
1728            // Read each byte range from the serialized file and deserialize the chunks.
1729            let mut file = File::open(&file_path).map_err(|_| ClientError::XORBNotFound(MerkleHash::default()))?;
1730
1731            let mut all_decompressed = Vec::new();
1732            let mut all_chunk_indices = Vec::<u32>::new();
1733            let mut total_transfer = 0u64;
1734
1735            for http_range in &http_ranges {
1736                let len = http_range.length() as usize;
1737                total_transfer += http_range.length();
1738
1739                file.seek(SeekFrom::Start(http_range.start))?;
1740                let mut data = vec![0u8; len];
1741                std::io::Read::read_exact(&mut file, &mut data)?;
1742
1743                let (decompressed, chunk_indices) =
1744                    xet_core_structures::xorb_object::deserialize_chunks(&mut Cursor::new(&data))?;
1745
1746                xet_core_structures::xorb_object::append_chunk_segment(
1747                    &mut all_decompressed,
1748                    &mut all_chunk_indices,
1749                    &decompressed,
1750                    &chunk_indices,
1751                );
1752            }
1753
1754            if let Some(expected) = uncompressed_size_if_known {
1755                debug_assert_eq!(
1756                    all_decompressed.len(),
1757                    expected,
1758                    "get_file_term_data: expected {} bytes, got {}",
1759                    expected,
1760                    all_decompressed.len()
1761                );
1762            }
1763
1764            if let Some(ref cb) = progress_callback {
1765                cb(total_transfer, total_transfer, total_transfer);
1766            }
1767            return Ok((Bytes::from(all_decompressed), all_chunk_indices));
1768        }
1769
1770        // Should not reach here, but return error if we do.
1771        Err(ClientError::PresignedUrlExpirationError)
1772    }
1773
1774    async fn get_file_chunk_hashes(
1775        &self,
1776        file_id: &MerkleHash,
1777        dirty_ranges: Vec<FileRange>,
1778    ) -> Result<FileChunkHashesResponse> {
1779        self.apply_api_delay().await;
1780
1781        let Some((file_info, _)) = self.shard_manager.get_file_reconstruction_info(file_id).await? else {
1782            return Err(ClientError::FileNotFound(*file_id));
1783        };
1784
1785        let mut chunks: Vec<(MerkleHash, u64)> = Vec::new();
1786        for segment in &file_info.segments {
1787            let xorb_obj = self.xorb_footer(&segment.xorb_hash).await?;
1788            chunks.extend(
1789                xorb_obj
1790                    .chunk_hash_sizes(segment.chunk_index_start, segment.chunk_index_end)
1791                    .map_err(|err| ClientError::Other(format!("chunk_hash_sizes error: {err}")))?,
1792            );
1793        }
1794
1795        build_file_chunk_hashes_response(&file_info, dirty_ranges, chunks)
1796    }
1797}
1798
1799fn map_redb_db_error(e: impl std::fmt::Debug) -> ClientError {
1800    let msg = format!("Global shard dedup database error: {e:?}");
1801    warn!("{msg}");
1802    ClientError::Other(msg)
1803}
1804
1805fn generate_fetch_url(file_path: &Path, byte_range: &FileRange, timestamp: Instant) -> String {
1806    let timestamp_ms = timestamp.saturating_duration_since(*REFERENCE_INSTANT).as_millis() as u64;
1807    format!("{}:{}:{}:{}", file_path.display(), byte_range.start, byte_range.end, timestamp_ms)
1808}
1809
1810fn parse_fetch_url(url: &str) -> Result<(PathBuf, FileRange, Instant)> {
1811    let mut parts = url.rsplitn(4, ':').collect::<Vec<_>>();
1812    parts.reverse();
1813
1814    if parts.len() != 4 {
1815        return Err(ClientError::InvalidArguments);
1816    }
1817
1818    let file_path_str = parts[0];
1819    let start_pos: u64 = parts[1].parse().map_err(|_| ClientError::InvalidArguments)?;
1820    let end_pos: u64 = parts[2].parse().map_err(|_| ClientError::InvalidArguments)?;
1821    let timestamp_ms: u64 = parts[3].parse().map_err(|_| ClientError::InvalidArguments)?;
1822
1823    let file_path: PathBuf = file_path_str.into();
1824    let byte_range = FileRange::new(start_pos, end_pos);
1825    let timestamp = *REFERENCE_INSTANT + Duration::from_millis(timestamp_ms);
1826
1827    Ok((file_path, byte_range, timestamp))
1828}
1829
1830fn generate_v2_fetch_url(hash: &MerkleHash, ranges: &[XorbRangeDescriptor], timestamp: Instant) -> String {
1831    xorb_utils::generate_v2_fetch_url(hash, ranges, timestamp)
1832}
1833#[cfg(test)]
1834mod tests {
1835    use xet_core_structures::xorb_object::CompressionScheme;
1836    use xet_core_structures::xorb_object::xorb_format_test_utils::{
1837        ChunkSize, build_and_verify_xorb_object, build_raw_xorb,
1838    };
1839    use xet_runtime::config::XetConfig;
1840    use xet_runtime::core::XetContext;
1841
1842    use super::*;
1843
1844    fn test_context() -> XetContext {
1845        let config = XetConfig::new();
1846        XetContext::from_external(tokio::runtime::Handle::current(), config)
1847    }
1848    use crate::cas_client::simulation::DeletionControlableClient;
1849    use crate::cas_client::simulation::client_testing_utils::ClientTestingUtils;
1850    use crate::cas_types::{ChunkRange, XorbReconstructionFetchInfo};
1851
1852    /// Runs the common TestingClient trait test suite for LocalClient.
1853    #[tokio::test]
1854    async fn test_common_client_suite() {
1855        crate::cas_client::simulation::client_unit_testing::test_client_functionality(|| async {
1856            LocalClient::temporary(test_context()).await.unwrap()
1857                as std::sync::Arc<dyn crate::cas_client::simulation::DirectAccessClient>
1858        })
1859        .await;
1860    }
1861
1862    /// Two different path strings for the same directory (symlink) must share one `redb::Database`
1863    /// in `DB_CACHE`. Without `canonicalize` in `new_internal`, `redb` returns
1864    /// `Database already open. Cannot acquire lock.` (duplicate opens for the same file).
1865    /// That can also surface as EMFILE under parallel tests (extra FDs per duplicate handle).
1866    #[cfg(unix)]
1867    #[tokio::test]
1868    async fn db_cache_unifies_symlink_equivalent_paths() {
1869        let tmp = tempfile::tempdir().unwrap();
1870        let real = tmp.path().join("real");
1871        std::fs::create_dir_all(&real).unwrap();
1872        let link = tmp.path().join("link");
1873        std::os::unix::fs::symlink(&real, &link).unwrap();
1874
1875        let ctx = test_context();
1876        let c1 = LocalClient::new(ctx.clone(), &link).await.unwrap();
1877        let c2 = LocalClient::new(ctx, &real).await.unwrap();
1878        assert!(Arc::ptr_eq(c1.db.as_ref().unwrap(), c2.db.as_ref().unwrap()));
1879    }
1880
1881    #[tokio::test]
1882    async fn test_download_fetch_term_data_validation() {
1883        // Setup: Create a client and upload a xorb
1884        let xorb = build_raw_xorb(3, ChunkSize::Fixed(2048));
1885        let xorb_obj = build_and_verify_xorb_object(xorb, CompressionScheme::Auto);
1886        let hash = xorb_obj.hash;
1887
1888        let client = LocalClient::temporary(test_context()).await.unwrap();
1889        let permit = client.acquire_upload_permit().await.unwrap();
1890        client.upload_xorb("default", xorb_obj, None, permit).await.unwrap();
1891
1892        // Get the actual byte offsets for a chunk range
1893        let file_path = client.get_path_for_entry(&hash);
1894        let file = File::open(&file_path).unwrap();
1895        let mut reader = BufReader::new(file);
1896        let xorb_obj = XorbObject::deserialize(&mut reader).unwrap();
1897        let (fetch_byte_start, fetch_byte_end) = xorb_obj.get_byte_offset(0, 1).unwrap();
1898
1899        let timestamp = Instant::now();
1900        let byte_range = FileRange::new(fetch_byte_start as u64, fetch_byte_end as u64);
1901        let valid_url = generate_fetch_url(&file_path, &byte_range, timestamp);
1902        // HttpRange uses inclusive end, FileRange uses exclusive end
1903        let valid_url_range = HttpRange::from(byte_range);
1904
1905        // Test 1: Valid URL and fetch_term should succeed
1906        let valid_fetch_term = XorbReconstructionFetchInfo {
1907            range: ChunkRange::new(0, 1),
1908            url: valid_url.clone(),
1909            url_range: valid_url_range,
1910        };
1911        let result = client.fetch_term_data(hash, valid_fetch_term).await;
1912        assert!(result.is_ok(), "Valid fetch_term should succeed");
1913
1914        // Test 2: Invalid URL format - too few parts (3 instead of 4)
1915        let too_few_parts = "filename:123:456";
1916        let invalid_fetch_term = XorbReconstructionFetchInfo {
1917            range: ChunkRange::new(0, 1),
1918            url: too_few_parts.to_string(),
1919            url_range: valid_url_range,
1920        };
1921        let result = client.fetch_term_data(hash, invalid_fetch_term).await;
1922        assert!(result.is_err(), "URL with too few parts should fail");
1923        assert!(matches!(result.unwrap_err(), ClientError::InvalidArguments));
1924
1925        // Test 3: Invalid start_pos - doesn't match url_range.start
1926        let wrong_byte_range = FileRange::new(fetch_byte_start as u64 + 1, fetch_byte_end as u64);
1927        let wrong_start_pos = generate_fetch_url(&file_path, &wrong_byte_range, timestamp);
1928        let invalid_fetch_term = XorbReconstructionFetchInfo {
1929            range: ChunkRange::new(0, 1),
1930            url: wrong_start_pos,
1931            url_range: valid_url_range,
1932        };
1933        let result = client.fetch_term_data(hash, invalid_fetch_term).await;
1934        assert!(result.is_err(), "Wrong start_pos should fail");
1935        assert!(matches!(result.unwrap_err(), ClientError::InvalidArguments));
1936
1937        // Test 4: Invalid end_pos - doesn't match url_range.end
1938        let wrong_byte_range = FileRange::new(fetch_byte_start as u64, fetch_byte_end as u64 + 1);
1939        let wrong_end_pos = generate_fetch_url(&file_path, &wrong_byte_range, timestamp);
1940        let invalid_fetch_term = XorbReconstructionFetchInfo {
1941            range: ChunkRange::new(0, 1),
1942            url: wrong_end_pos,
1943            url_range: valid_url_range,
1944        };
1945        let result = client.fetch_term_data(hash, invalid_fetch_term).await;
1946        assert!(result.is_err(), "Wrong end_pos should fail");
1947        assert!(matches!(result.unwrap_err(), ClientError::InvalidArguments));
1948
1949        // Test 5: Invalid start_pos - non-numeric
1950        let timestamp_ms = timestamp.saturating_duration_since(*REFERENCE_INSTANT).as_millis() as u64;
1951        let non_numeric_start = format!("{}:not_a_number:{}:{}", file_path.display(), fetch_byte_end, timestamp_ms);
1952        let invalid_fetch_term = XorbReconstructionFetchInfo {
1953            range: ChunkRange::new(0, 1),
1954            url: non_numeric_start,
1955            url_range: valid_url_range,
1956        };
1957        let result = client.fetch_term_data(hash, invalid_fetch_term).await;
1958        assert!(result.is_err(), "Non-numeric start_pos should fail");
1959        assert!(matches!(result.unwrap_err(), ClientError::InvalidArguments));
1960
1961        // Test 6: Invalid end_pos - non-numeric
1962        let non_numeric_end = format!("{}:{}:not_a_number:{}", file_path.display(), fetch_byte_start, timestamp_ms);
1963        let invalid_fetch_term = XorbReconstructionFetchInfo {
1964            range: ChunkRange::new(0, 1),
1965            url: non_numeric_end,
1966            url_range: valid_url_range,
1967        };
1968        let result = client.fetch_term_data(hash, invalid_fetch_term).await;
1969        assert!(result.is_err(), "Non-numeric end_pos should fail");
1970        assert!(matches!(result.unwrap_err(), ClientError::InvalidArguments));
1971
1972        // Test 7: Empty URL
1973        let invalid_fetch_term = XorbReconstructionFetchInfo {
1974            range: ChunkRange::new(0, 1),
1975            url: String::new(),
1976            url_range: valid_url_range,
1977        };
1978        let result = client.fetch_term_data(hash, invalid_fetch_term).await;
1979        assert!(result.is_err(), "Empty URL should fail");
1980        assert!(matches!(result.unwrap_err(), ClientError::InvalidArguments));
1981
1982        // Test 8: Invalid timestamp - non-numeric
1983        let non_numeric_timestamp =
1984            format!("{}:{}:{}:not_a_number", file_path.display(), fetch_byte_start, fetch_byte_end);
1985        let invalid_fetch_term = XorbReconstructionFetchInfo {
1986            range: ChunkRange::new(0, 1),
1987            url: non_numeric_timestamp,
1988            url_range: valid_url_range,
1989        };
1990        let result = client.fetch_term_data(hash, invalid_fetch_term).await;
1991        assert!(result.is_err(), "Non-numeric timestamp should fail");
1992        assert!(matches!(result.unwrap_err(), ClientError::InvalidArguments));
1993
1994        // Test 9: Non-existent file path
1995        let non_existent_path = PathBuf::from("/nonexistent/path/file.xorb");
1996        let non_existent_url = generate_fetch_url(&non_existent_path, &byte_range, timestamp);
1997        let invalid_fetch_term = XorbReconstructionFetchInfo {
1998            range: ChunkRange::new(0, 1),
1999            url: non_existent_url,
2000            url_range: valid_url_range,
2001        };
2002        let result = client.fetch_term_data(hash, invalid_fetch_term).await;
2003        assert!(result.is_err(), "Non-existent file should fail");
2004    }
2005
2006    #[tokio::test(start_paused = true)]
2007    async fn test_url_expiration() {
2008        super::super::client_unit_testing::test_url_expiration_functionality(|| async {
2009            LocalClient::temporary(test_context()).await.unwrap()
2010                as std::sync::Arc<dyn crate::cas_client::simulation::DirectAccessClient>
2011        })
2012        .await;
2013    }
2014
2015    #[tokio::test(start_paused = true)]
2016    async fn test_api_delay() {
2017        super::super::client_unit_testing::test_api_delay_functionality(|| async {
2018            LocalClient::temporary(test_context()).await.unwrap()
2019                as std::sync::Arc<dyn crate::cas_client::simulation::DirectAccessClient>
2020        })
2021        .await;
2022    }
2023
2024    #[tokio::test(start_paused = true)]
2025    async fn test_global_dedup_shard_expiration() {
2026        super::super::client_unit_testing::test_global_dedup_shard_expiration_functionality(|| async {
2027            LocalClient::temporary(test_context()).await.unwrap()
2028                as std::sync::Arc<dyn crate::cas_client::simulation::DirectAccessClient>
2029        })
2030        .await;
2031    }
2032
2033    #[tokio::test]
2034    #[cfg_attr(feature = "smoke-test", ignore)]
2035    async fn test_global_dedup_shard_expiration_stress() {
2036        super::super::client_unit_testing::test_global_dedup_shard_expiration_stress(|| async {
2037            LocalClient::temporary(test_context()).await.unwrap()
2038                as std::sync::Arc<dyn crate::cas_client::simulation::DirectAccessClient>
2039        })
2040        .await;
2041    }
2042
2043    #[tokio::test]
2044    async fn test_deletion_suite() {
2045        super::super::deletion_unit_testing::test_deletion_functionality(|| async {
2046            LocalClient::temporary(test_context()).await.unwrap()
2047        })
2048        .await;
2049    }
2050
2051    #[tokio::test]
2052    async fn test_verify_integrity_detects_missing_cas_block_reference() {
2053        let client = LocalClient::temporary(test_context()).await.unwrap();
2054        client.upload_random_file(&[(3, (0, 3)), (4, (0, 2))], 2048).await.unwrap();
2055        client.verify_integrity().await.unwrap();
2056
2057        let mut in_memory = client.load_all_shard_data().unwrap();
2058        let removed_hash = *in_memory.xorb_content.keys().next().unwrap();
2059        in_memory.xorb_content.remove(&removed_hash);
2060        client.delete_xorb(&removed_hash).await;
2061        client.write_shard_data_and_register(&in_memory).await.unwrap();
2062
2063        assert!(client.verify_integrity().await.is_err());
2064    }
2065
2066    #[tokio::test]
2067    async fn test_verify_integrity_detects_invalid_chunk_range() {
2068        let client = LocalClient::temporary(test_context()).await.unwrap();
2069        client.upload_random_file(&[(5, (0, 3))], 2048).await.unwrap();
2070        client.verify_integrity().await.unwrap();
2071
2072        let mut in_memory = client.load_all_shard_data().unwrap();
2073        let file_info = in_memory.file_content.values_mut().next().unwrap();
2074        let segment = file_info.segments.first_mut().unwrap();
2075        let xorb_entry_count = in_memory.xorb_content.get(&segment.xorb_hash).unwrap().metadata.num_entries;
2076        segment.chunk_index_end = xorb_entry_count + 1;
2077        client.write_shard_data_and_register(&in_memory).await.unwrap();
2078
2079        assert!(client.verify_integrity().await.is_err());
2080    }
2081
2082    /// Verifies that delete_file_entry does not rewrite shard files (shard hashes remain stable).
2083    #[tokio::test]
2084    async fn test_delete_file_entry_does_not_rewrite_shards() {
2085        let client = LocalClient::temporary(test_context()).await.unwrap();
2086        client.upload_random_file(&[(1, (0, 3))], 2048).await.unwrap();
2087
2088        let shard_hashes_before: Vec<_> = client.shard_file_paths().unwrap().into_iter().map(|(h, _)| h).collect();
2089        assert!(!shard_hashes_before.is_empty());
2090
2091        client
2092            .delete_file_entry(&client.list_file_shard_entries().await.unwrap()[0].0)
2093            .await
2094            .unwrap();
2095
2096        let shard_hashes_after: Vec<_> = client.shard_file_paths().unwrap().into_iter().map(|(h, _)| h).collect();
2097        assert_eq!(shard_hashes_before, shard_hashes_after, "Shard file hashes must not change after delete");
2098    }
2099
2100    /// Verifies that file deletion (removal from FILE_TO_SHARD_TABLE) persists across restarts.
2101    #[tokio::test]
2102    async fn test_deletion_status_persists_across_restart() {
2103        let tmp_dir = TempDir::new().unwrap();
2104        let path = tmp_dir.path().to_owned();
2105
2106        let file_hash;
2107        {
2108            let client = LocalClient::new(test_context(), &path).await.unwrap();
2109            let file = client.upload_random_file(&[(1, (0, 3)), (2, (0, 2))], 2048).await.unwrap();
2110            file_hash = file.file_hash;
2111            assert!(!client.list_file_shard_entries().await.unwrap().is_empty());
2112
2113            client.delete_file_entry(&file_hash).await.unwrap();
2114            assert!(client.list_file_shard_entries().await.unwrap().is_empty());
2115        }
2116
2117        {
2118            let client = LocalClient::new(test_context(), &path).await.unwrap();
2119            assert!(
2120                client.is_file_deleted(&file_hash),
2121                "Entry should be absent from FILE_TO_SHARD_TABLE after restart"
2122            );
2123            assert!(
2124                client.list_file_shard_entries().await.unwrap().is_empty(),
2125                "Deleted files should remain hidden after restart"
2126            );
2127        }
2128    }
2129
2130    /// Tests cross-shard dedup integrity: a file entry referencing an XORB that is not indexed
2131    /// in any shard but exists on disk should pass verify_integrity (dedup case).
2132    #[tokio::test]
2133    async fn test_verify_integrity_cross_shard_dedup_ok() {
2134        let client = LocalClient::temporary(test_context()).await.unwrap();
2135        client.upload_random_file(&[(1, (0, 3))], 2048).await.unwrap();
2136        client.verify_integrity().await.unwrap();
2137
2138        // Clear dedup entries for old shards before rewriting, so pass 3 doesn't
2139        // flag stale references to the about-to-be-replaced shard files.
2140        for h in client.list_shard_entries().await.unwrap() {
2141            client.remove_shard_dedup_entries(&h).await.unwrap();
2142        }
2143
2144        let mut in_memory = client.load_all_shard_data().unwrap();
2145        in_memory.xorb_content.clear();
2146        client.write_shard_data_and_register(&in_memory).await.unwrap();
2147
2148        client
2149            .verify_integrity()
2150            .await
2151            .expect("Integrity should pass: XORB files exist on disk even though no shard indexes them");
2152    }
2153
2154    /// Tests that verify_integrity ignores deleted files (absent from FILE_TO_SHARD_TABLE),
2155    /// so missing XORBs for deleted files do not cause false integrity failures.
2156    #[tokio::test]
2157    async fn test_verify_integrity_skips_deleted_files() {
2158        let client = LocalClient::temporary(test_context()).await.unwrap();
2159        let deleted_file = client.upload_random_file(&[(1, (0, 3))], 2048).await.unwrap();
2160        let live_file = client.upload_random_file(&[(2, (0, 2))], 2048).await.unwrap();
2161        client.verify_integrity().await.unwrap();
2162
2163        client.delete_file_entry(&deleted_file.file_hash).await.unwrap();
2164
2165        for t in &deleted_file.terms {
2166            client.delete_xorb(&t.xorb_hash).await;
2167        }
2168
2169        // Clear dedup entries for old shards before rewriting, so pass 3 doesn't
2170        // flag stale references to the about-to-be-replaced shard files.
2171        for h in client.list_shard_entries().await.unwrap() {
2172            client.remove_shard_dedup_entries(&h).await.unwrap();
2173        }
2174
2175        // Remove deleted-file entries from shard metadata too, so pass 1 doesn't fail
2176        // on deliberately removed XORB files and the deleted file isn't re-registered.
2177        let mut in_memory = client.load_all_shard_data().unwrap();
2178        in_memory.file_content.remove(&deleted_file.file_hash);
2179        for t in &deleted_file.terms {
2180            in_memory.xorb_content.remove(&t.xorb_hash);
2181        }
2182        client.write_shard_data_and_register(&in_memory).await.unwrap();
2183
2184        client
2185            .verify_integrity()
2186            .await
2187            .expect("Integrity should pass: missing XORBs are only referenced by a deleted file");
2188
2189        // Sanity check: the surviving file remains readable.
2190        let live_data = client.get_file_data(&live_file.file_hash, None).await.unwrap();
2191        assert_eq!(live_data, live_file.data);
2192    }
2193
2194    /// Tests that verify_integrity catches stale global dedup table entries pointing
2195    /// to shard files that have been removed.
2196    #[tokio::test]
2197    async fn test_verify_integrity_detects_stale_dedup_shard_reference() {
2198        let client = LocalClient::temporary(test_context()).await.unwrap();
2199        let file = client.upload_random_file(&[(10, (0, 3))], 2048).await.unwrap();
2200        client.verify_integrity().await.unwrap();
2201
2202        // Confirm dedup entries exist for the file's chunks.
2203        let has_dedup = client
2204            .query_for_global_dedup_shard("default", &file.terms[0].chunk_hashes[0])
2205            .await
2206            .unwrap()
2207            .is_some();
2208        assert!(has_dedup, "Dedup entry should exist after upload");
2209
2210        // Remove file entry so Pass 2 doesn't fail on the missing shard.
2211        client.delete_file_entry(&file.file_hash).await.unwrap();
2212
2213        // Delete the shard file without clearing its dedup entries.
2214        let shard_hashes = client.list_shard_entries().await.unwrap();
2215        assert!(!shard_hashes.is_empty());
2216        for h in &shard_hashes {
2217            let path = client.shard_path_for_hash(h).unwrap();
2218            std::fs::remove_file(&path).unwrap();
2219        }
2220
2221        let result = client.verify_integrity().await;
2222        assert!(result.is_err(), "verify_integrity should fail when dedup table references a missing shard");
2223        let err_msg = format!("{:?}", result.unwrap_err());
2224        assert!(err_msg.contains("global dedup table"), "Error should mention global dedup table");
2225    }
2226
2227    /// Exercises the root-cause scenario: re-uploading the same file hash in a new shard
2228    /// after deleting the original file and its xorbs must not resurrect stale entries.
2229    #[tokio::test]
2230    async fn test_reupload_same_file_hash_does_not_resurrect_stale_entries() {
2231        let client = LocalClient::temporary(test_context()).await.unwrap();
2232
2233        // 1. Upload file F in shard S1 referencing xorb X.
2234        let file = client.upload_random_file(&[(1, (0, 3))], 2048).await.unwrap();
2235        let file_hash = file.file_hash;
2236        let xorb_hash = file.terms[0].xorb_hash;
2237        client.verify_integrity().await.unwrap();
2238
2239        // 2. Delete file F (removes from FILE_TO_SHARD_TABLE).
2240        client.delete_file_entry(&file_hash).await.unwrap();
2241        assert!(client.is_file_deleted(&file_hash));
2242
2243        // 3. Delete xorb X from disk.
2244        client.delete_xorb(&xorb_hash).await;
2245        assert!(!client.get_path_for_entry(&xorb_hash).exists());
2246
2247        // 4. Clean up old shard's xorb metadata so Pass 1 doesn't fail on the deliberately removed xorb file.  Also
2248        //    clean up dedup entries.
2249        for h in client.list_shard_entries().await.unwrap() {
2250            client.remove_shard_dedup_entries(&h).await.unwrap();
2251        }
2252        let mut in_memory = client.load_all_shard_data().unwrap();
2253        in_memory.file_content.remove(&file_hash);
2254        for t in &file.terms {
2255            in_memory.xorb_content.remove(&t.xorb_hash);
2256        }
2257        client.write_shard_data_and_register(&in_memory).await.unwrap();
2258
2259        // 5. Upload a new, different file in a new shard S2.
2260        let file2 = client.upload_random_file(&[(2, (0, 2))], 2048).await.unwrap();
2261        let file2_hash = file2.file_hash;
2262        assert!(!client.is_file_deleted(&file2_hash));
2263
2264        // 6. verify_integrity should pass — S1's stale file entry for F is not checked because FILE_TO_SHARD_TABLE no
2265        //    longer maps F to S1.
2266        client
2267            .verify_integrity()
2268            .await
2269            .expect("Integrity should pass: the old shard's stale file entry with dangling xorb refs is not consulted");
2270    }
2271
2272    /// Tests that list_xorbs_and_tags tags change after file re-creation with a timestamp delay.
2273    #[tokio::test]
2274    async fn test_list_xorbs_and_tags_timestamp_changes() {
2275        let client = LocalClient::temporary(test_context()).await.unwrap();
2276
2277        let file1 = client.upload_random_file(&[(1, (0, 2))], 2048).await.unwrap();
2278        let xorb_hash = file1.terms[0].xorb_hash;
2279
2280        let tags1 = client.list_xorbs_and_tags().await.unwrap();
2281        let (_, tag1) = tags1.iter().find(|(h, _)| *h == xorb_hash).unwrap();
2282
2283        // Delete and wait 1 second so the filesystem timestamp advances.
2284        client.delete_xorb(&xorb_hash).await;
2285        std::thread::sleep(Duration::from_secs(1));
2286
2287        // Re-upload a file that creates a new xorb with the same hash seed.
2288        let file2 = client.upload_random_file(&[(1, (0, 2))], 2048).await.unwrap();
2289        let xorb_hash2 = file2.terms[0].xorb_hash;
2290
2291        let tags2 = client.list_xorbs_and_tags().await.unwrap();
2292        let (_, tag2) = tags2.iter().find(|(h, _)| *h == xorb_hash2).unwrap();
2293
2294        assert_ne!(tag1, tag2, "Tags should differ after re-creation with timestamp delay");
2295    }
2296
2297    // ── Lifecycle-tag deletion mode tests ───────────────────────────────
2298
2299    #[tokio::test]
2300    async fn test_lifecycle_tag_xorb_delete_renames_to_gctag() {
2301        let client = LocalClient::temporary(test_context()).await.unwrap();
2302        client.set_lifecycle_tag_deletion(true);
2303
2304        let file = client.upload_random_file(&[(1, (0, 2))], 2048).await.unwrap();
2305        let xorb_hash = file.terms[0].xorb_hash;
2306        assert!(client.xorb_exists(&xorb_hash).await.unwrap());
2307
2308        client.delete_xorb(&xorb_hash).await;
2309
2310        // Canonical path is gone — reads fail.
2311        assert!(!client.xorb_exists(&xorb_hash).await.unwrap());
2312        assert!(client.get_full_xorb(&xorb_hash).await.is_err());
2313        // .gctag file retains the bytes.
2314        let gctag = client.gctag_xorb_path(&xorb_hash);
2315        assert!(gctag.exists(), "gctag file should exist after tag-delete");
2316        // list_xorbs excludes the tagged xorb.
2317        let listed = client.list_xorbs().await.unwrap();
2318        assert!(!listed.contains(&xorb_hash));
2319    }
2320
2321    #[tokio::test]
2322    async fn test_lifecycle_tag_xorb_upload_clears_tag() {
2323        let client = LocalClient::temporary(test_context()).await.unwrap();
2324        client.set_lifecycle_tag_deletion(true);
2325
2326        let file = client.upload_random_file(&[(1, (0, 2))], 2048).await.unwrap();
2327        let xorb_hash = file.terms[0].xorb_hash;
2328
2329        client.delete_xorb(&xorb_hash).await;
2330        let gctag = client.gctag_xorb_path(&xorb_hash);
2331        assert!(gctag.exists());
2332
2333        // Re-upload the same xorb (same content seed) — tag is cleared.
2334        let file2 = client.upload_random_file(&[(1, (0, 2))], 2048).await.unwrap();
2335        let xorb_hash2 = file2.terms[0].xorb_hash;
2336        assert_eq!(xorb_hash, xorb_hash2, "same seed should produce same xorb hash");
2337
2338        assert!(!gctag.exists(), "gctag file should be removed after re-upload");
2339        assert!(client.xorb_exists(&xorb_hash).await.unwrap());
2340    }
2341
2342    #[tokio::test]
2343    async fn test_lifecycle_tag_shard_delete_renames_to_gctag() {
2344        let client = LocalClient::temporary(test_context()).await.unwrap();
2345        client.set_lifecycle_tag_deletion(true);
2346
2347        let file = client.upload_random_file(&[(1, (0, 2))], 2048).await.unwrap();
2348        let shard_hash = client.list_shard_entries().await.unwrap().pop().unwrap();
2349
2350        // Remove the file entry first (as the harness does on user delete).
2351        client.delete_file_entry(&file.file_hash).await.unwrap();
2352        // Now GC deletes the shard.
2353        client.delete_shard_entry(&shard_hash).await.unwrap();
2354
2355        // Canonical shard path is gone.
2356        assert!(client.get_shard_bytes(&shard_hash).await.is_err());
2357        // .gctag file retains the bytes.
2358        let gctag = client.gctag_shard_path(&shard_hash);
2359        assert!(gctag.exists(), "gctag shard file should exist after tag-delete");
2360        // list_shard_entries excludes the tagged shard.
2361        let listed = client.list_shard_entries().await.unwrap();
2362        assert!(!listed.contains(&shard_hash));
2363    }
2364
2365    #[tokio::test]
2366    async fn test_lifecycle_tag_mode_off_hard_deletes() {
2367        let client = LocalClient::temporary(test_context()).await.unwrap();
2368        // Mode is off by default.
2369        assert!(!client.lifecycle_tag_deletion_enabled());
2370
2371        let file = client.upload_random_file(&[(1, (0, 2))], 2048).await.unwrap();
2372        let xorb_hash = file.terms[0].xorb_hash;
2373
2374        client.delete_xorb(&xorb_hash).await;
2375
2376        let gctag = client.gctag_xorb_path(&xorb_hash);
2377        assert!(!gctag.exists(), "gctag file should NOT exist when mode is off");
2378        let canonical = client.get_path_for_entry(&xorb_hash);
2379        assert!(!canonical.exists(), "canonical file should be hard-deleted");
2380    }
2381
2382    #[tokio::test]
2383    async fn test_lifecycle_tag_verify_integrity_flags_tagged_xorb() {
2384        let client = LocalClient::temporary(test_context()).await.unwrap();
2385        client.set_lifecycle_tag_deletion(true);
2386
2387        // Upload a file, then tag its xorb WITHOUT removing the file entry
2388        // (simulates a GC bug: tagging an xorb still referenced by an active file).
2389        let file = client.upload_random_file(&[(1, (0, 2))], 2048).await.unwrap();
2390        let xorb_hash = file.terms[0].xorb_hash;
2391
2392        client.delete_xorb(&xorb_hash).await;
2393
2394        // verify_integrity should fail: the file entry references a tagged xorb
2395        // whose canonical path no longer exists.
2396        assert!(client.verify_integrity().await.is_err());
2397    }
2398
2399    #[tokio::test]
2400    async fn test_lifecycle_tag_verify_all_reachable_ignores_tagged() {
2401        let client = LocalClient::temporary(test_context()).await.unwrap();
2402        client.set_lifecycle_tag_deletion(true);
2403
2404        let file = client.upload_random_file(&[(1, (0, 2))], 2048).await.unwrap();
2405        let xorb_hash = file.terms[0].xorb_hash;
2406        let shard_hash = client.list_shard_entries().await.unwrap().pop().unwrap();
2407
2408        // Properly delete: remove file entry, then tag shard + xorb.
2409        client.delete_file_entry(&file.file_hash).await.unwrap();
2410        client.delete_shard_entry(&shard_hash).await.unwrap();
2411        client.delete_xorb(&xorb_hash).await;
2412
2413        // verify_all_reachable should pass: tagged objects are treated as
2414        // already collected, not as orphaned on-disk data.
2415        client
2416            .verify_all_reachable()
2417            .await
2418            .expect("tagged objects should be ignored by reachability");
2419    }
2420}