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