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#[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
111const FILE_TO_SHARD_TABLE: TableDefinition<RedbHash, FileShardRef> = TableDefinition::new("file_to_shard");
113
114#[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
170static DB_CACHE: std::sync::LazyLock<Mutex<HashMap<PathBuf, Weak<redb::Database>>>> =
179 std::sync::LazyLock::new(|| Mutex::new(HashMap::new()));
180
181fn 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 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
213fn 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 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_expiration_secs: AtomicU64,
257 random_ms_delay_window: (AtomicU64, AtomicU64),
259 max_ranges_per_fetch: AtomicUsize,
261 v2_disabled_status: AtomicU16,
263 lifecycle_tag_deletion: AtomicBool,
270 _tmp_dir: Option<TempDir>,
271}
272
273impl LocalClient {
274 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 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 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 {
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 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 fn get_path_for_entry(&self, hash: &MerkleHash) -> PathBuf {
356 self.xorb_dir.join(format!("default.{hash:?}"))
357 }
358
359 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 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 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 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 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 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 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 let _ = std::fs::remove_file(tmp_path);
458 } else {
459 let _ = std::fs::rename(tmp_path, original_path);
462 }
463 }
464
465 #[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 #[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 #[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 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 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 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 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 indices.push(0);
844 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 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 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 ¤t_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 ¤t_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 async fn verify_all_reachable(&self) -> Result<()> {
1105 let shard_files = self.shard_file_paths()?;
1106
1107 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 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 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 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 let read_txn = self.db().begin_read().map_err(map_redb_db_error)?;
1204
1205 let shard_files = self.shard_file_paths()?;
1206
1207 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 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 } 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 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 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 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 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 let mut reader = Cursor::new(&shard_data);
1523 let minimal_shard = MDBMinimalShard::from_reader(&mut reader, true, true)?;
1524
1525 let mut in_memory_shard = MDBInMemoryShard::default();
1527
1528 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 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 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 let chunk_hashes = minimal_shard.global_dedup_eligible_chunks();
1549
1550 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 let _ = std::fs::remove_file(self.gctag_shard_path(&shard_hash));
1584
1585 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 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 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 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 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 url_info.refresh_url().await?;
1732 continue;
1733 }
1734 return Err(ClientError::PresignedUrlExpirationError);
1735 }
1736
1737 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 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 #[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 #[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 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 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 let valid_url_range = HttpRange::from(byte_range);
1913
1914 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 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 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 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 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 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 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 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 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 #[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 #[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 #[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 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 #[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 for h in client.list_shard_entries().await.unwrap() {
2181 client.remove_shard_dedup_entries(&h).await.unwrap();
2182 }
2183
2184 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 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 #[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 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 client.delete_file_entry(&file.file_hash).await.unwrap();
2221
2222 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 #[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 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 client.delete_file_entry(&file_hash).await.unwrap();
2250 assert!(client.is_file_deleted(&file_hash));
2251
2252 client.delete_xorb(&xorb_hash).await;
2254 assert!(!client.get_path_for_entry(&xorb_hash).exists());
2255
2256 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 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 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 #[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 client.delete_xorb(&xorb_hash).await;
2294 std::thread::sleep(Duration::from_secs(1));
2295
2296 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 #[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 assert!(!client.xorb_exists(&xorb_hash).await.unwrap());
2321 assert!(client.get_full_xorb(&xorb_hash).await.is_err());
2322 let gctag = client.gctag_xorb_path(&xorb_hash);
2324 assert!(gctag.exists(), "gctag file should exist after tag-delete");
2325 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 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 client.delete_file_entry(&file.file_hash).await.unwrap();
2361 client.delete_shard_entry(&shard_hash).await.unwrap();
2363
2364 assert!(client.get_shard_bytes(&shard_hash).await.is_err());
2366 let gctag = client.gctag_shard_path(&shard_hash);
2368 assert!(gctag.exists(), "gctag shard file should exist after tag-delete");
2369 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 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 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 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 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 client
2425 .verify_all_reachable()
2426 .await
2427 .expect("tagged objects should be ignored by reachability");
2428 }
2429}