Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 5 additions & 1 deletion xet_client/src/cas_client/simulation/deletion_controls.rs
Original file line number Diff line number Diff line change
Expand Up @@ -36,7 +36,11 @@ pub trait DeletionControlableClient: Send + Sync {

/// Removes all global-dedup table entries contributed by the given shard.
/// Called by GC Stage 4 before replacing or discarding a shard.
async fn remove_shard_dedup_entries(&self, shard_hash: &MerkleHash) -> Result<()>;
///
/// Returns the chunk hashes that were removed from the global-dedup table,
/// mirroring the production CAS `gc_delete_shard_dedup` endpoint so callers
/// can audit the reclaimed dedup keys.
async fn remove_shard_dedup_entries(&self, shard_hash: &MerkleHash) -> Result<Vec<MerkleHash>>;

/// Deletes a XORB by hash.
async fn delete_xorb(&self, hash: &MerkleHash);
Expand Down
14 changes: 12 additions & 2 deletions xet_client/src/cas_client/simulation/deletion_unit_testing.rs
Original file line number Diff line number Diff line change
Expand Up @@ -290,21 +290,30 @@ async fn test_remove_shard_dedup_entries_removes_correct_entries<
assert!(dedup.is_some(), "Chunk from file B should have a dedup entry");
}

client.remove_shard_dedup_entries(&shard_a).await.unwrap();
let removed_chunks: HashSet<MerkleHash> =
client.remove_shard_dedup_entries(&shard_a).await.unwrap().into_iter().collect();

for t in &file_a.terms {
let dedup = client
.query_for_global_dedup_shard("default", &t.chunk_hashes[0])
.await
.unwrap();
assert!(dedup.is_none(), "Dedup entries for shard A's chunks should be removed");
assert!(
removed_chunks.contains(&t.chunk_hashes[0]),
"Returned chunk list should report shard A's deregistered chunk"
);
}
for t in &file_b.terms {
let dedup = client
.query_for_global_dedup_shard("default", &t.chunk_hashes[0])
.await
.unwrap();
assert!(dedup.is_some(), "Dedup entries for shard B's chunks should be preserved");
assert!(
!removed_chunks.contains(&t.chunk_hashes[0]),
"Returned chunk list must not include shard B's chunks"
);
}
}

Expand All @@ -317,7 +326,8 @@ async fn test_remove_shard_dedup_entries_noop_on_unknown_hash<
let file = client.upload_random_file(&[(1, (0, 2))], 2048).await.unwrap();

let bogus_hash = MerkleHash::from([0xFFu8; 32]);
client.remove_shard_dedup_entries(&bogus_hash).await.unwrap();
let removed = client.remove_shard_dedup_entries(&bogus_hash).await.unwrap();
assert!(removed.is_empty(), "Removing dedup entries for an unknown shard should report no removed chunks");

for t in &file.terms {
let dedup = client
Expand Down
8 changes: 5 additions & 3 deletions xet_client/src/cas_client/simulation/local_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -884,8 +884,9 @@ impl super::DeletionControlableClient for LocalClient {
Ok(())
}

async fn remove_shard_dedup_entries(&self, shard_hash: &MerkleHash) -> Result<()> {
async fn remove_shard_dedup_entries(&self, shard_hash: &MerkleHash) -> Result<Vec<MerkleHash>> {
let shard_redb = RedbHash::from(*shard_hash);
let mut removed_chunks = Vec::new();
for _ in 0..4 {
let to_delete: Vec<RedbHash> = {
let read_txn = self.db().begin_read().map_err(map_redb_db_error)?;
Expand All @@ -900,7 +901,7 @@ impl super::DeletionControlableClient for LocalClient {
};

if to_delete.is_empty() {
return Ok(());
return Ok(removed_chunks);
}

let write_txn = self.db().begin_write().map_err(map_redb_db_error)?;
Expand All @@ -911,6 +912,7 @@ impl super::DeletionControlableClient for LocalClient {
}
}
write_txn.commit().map_err(map_redb_db_error)?;
removed_chunks.extend(to_delete.into_iter().map(MerkleHash::from));
}

let still_present = {
Expand All @@ -930,7 +932,7 @@ impl super::DeletionControlableClient for LocalClient {
)));
}

Ok(())
Ok(removed_chunks)
}

async fn delete_xorb(&self, hash: &MerkleHash) {
Expand Down
6 changes: 5 additions & 1 deletion xet_client/src/cas_client/simulation/local_server/server.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1431,9 +1431,13 @@ mod tests {
.is_some()
);

DeletionControlableClient::remove_shard_dedup_entries(&sc, &shard_hash)
let removed = DeletionControlableClient::remove_shard_dedup_entries(&sc, &shard_hash)
.await
.unwrap();
assert!(
removed.contains(&first_chunk),
"removed chunk list should round-trip the shard's deregistered chunk over HTTP"
);
assert!(
Client::query_for_global_dedup_shard(&sc, "default", &first_chunk)
.await
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -532,7 +532,7 @@ impl DeletionControlableClient for SimulationControlClient {
}

/// Removes all global-dedup table entries for a shard via `/simulation/shards/{hash}/dedup_entries`.
async fn remove_shard_dedup_entries(&self, shard_hash: &MerkleHash) -> Result<()> {
async fn remove_shard_dedup_entries(&self, shard_hash: &MerkleHash) -> Result<Vec<MerkleHash>> {
let hex = HexMerkleHash::from(*shard_hash);
let url = self.sim_url(&format!("/shards/{hex}/dedup_entries"));
let resp = self
Expand All @@ -541,8 +541,9 @@ impl DeletionControlableClient for SimulationControlClient {
.send()
.await
.map_err(|e| ClientError::Other(e.to_string()))?;
Self::check_status(resp).await?;
Ok(())
let resp = Self::check_status(resp).await?;
let removed: Vec<HexMerkleHash> = resp.json().await.map_err(|e| ClientError::Other(e.to_string()))?;
Ok(removed.into_iter().map(MerkleHash::from).collect())
}

async fn delete_xorb(&self, hash: &MerkleHash) {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -327,7 +327,10 @@ async fn remove_shard_dedup_entries(
return not_implemented();
};
match dc.remove_shard_dedup_entries(&hash).await {
Ok(()) => StatusCode::NO_CONTENT.into_response(),
Ok(removed) => {
let response: Vec<HexMerkleHash> = removed.into_iter().map(HexMerkleHash::from).collect();
Json(response).into_response()
},
Err(e) => error_to_response(e),
}
}
Expand Down
14 changes: 9 additions & 5 deletions xet_client/src/cas_client/simulation/memory_client.rs
Original file line number Diff line number Diff line change
Expand Up @@ -1042,18 +1042,22 @@ impl super::DeletionControlableClient for MemoryClient {
Ok(())
}

async fn remove_shard_dedup_entries(&self, shard_hash: &MerkleHash) -> Result<()> {
async fn remove_shard_dedup_entries(&self, shard_hash: &MerkleHash) -> Result<Vec<MerkleHash>> {
let shard = self.shard.read().await;
let Some((current_hash, _)) = Self::current_shard_hash_and_bytes(&shard)? else {
return Ok(());
return Ok(Vec::new());
};
if &current_hash != shard_hash {
return Ok(());
return Ok(Vec::new());
}
drop(shard);

self.global_dedup.write().await.clear();
Ok(())
// The in-memory client holds a single shard, so every global-dedup entry
// belongs to it; return all keys as the removed set before clearing.
let mut dedup = self.global_dedup.write().await;
let removed_chunks = dedup.keys().copied().collect();
dedup.clear();
Ok(removed_chunks)
}

async fn delete_xorb(&self, hash: &MerkleHash) {
Expand Down
Loading