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