diff --git a/kv-service/client-rs/tests/local_cluster_e2e.rs b/kv-service/client-rs/tests/local_cluster_e2e.rs index 8ce0f80..a361147 100644 --- a/kv-service/client-rs/tests/local_cluster_e2e.rs +++ b/kv-service/client-rs/tests/local_cluster_e2e.rs @@ -176,6 +176,12 @@ redis_command_timeout_ms = 1000 [metrics] enabled = false listen = "127.0.0.1:0" +[gc] +enabled = true +interval_seconds = 1 +grace_seconds = 0 +max_tasks_per_run = 100 +task_lease_seconds = 5 [cluster] node_id = "{node}" grpc_advertise = "127.0.0.1:{port}" @@ -336,7 +342,10 @@ async fn assert_four_way_placement(cluster: &Cluster, client: &mut KvClient) { .await .expect("lookup object") .expect("object metadata must exist after put"); - assert!(lookup.descriptor.is_striped, "object must use striped placement"); + assert!( + lookup.descriptor.is_striped, + "object must use striped placement" + ); assert_eq!(lookup.descriptor.stripe_count, 4, "expected four stripes"); let locations: HashSet<(String, u32)> = lookup .placement @@ -416,7 +425,7 @@ async fn two_node_restart_recovers_shared_metadata_and_stripes() { } #[tokio::test] -async fn two_node_distributed_delete_removes_all_stripe_files() { +async fn two_node_distributed_delete_retires_then_reclaims_all_stripe_files() { let cluster = Cluster::start("distributed_delete").await; let payload = striped_payload(); let mut client = connect(&cluster, "node-a").await; @@ -429,16 +438,22 @@ async fn two_node_distributed_delete_removes_all_stripe_files() { .delete(TEST_NAMESPACE, TEST_OBJECT_KEY) .await .expect("distributed delete"), - "delete must remove the existing object", + "delete must retire the existing object", ); - for node in ["node-a", "node-b"] { - for device in 0..2 { - assert_eq!( - cluster.file_count(node, device), - 0, - "distributed delete must remove stripe from {node}/nvme{device}", - ); + + let deadline = Instant::now() + Duration::from_secs(10); + while Instant::now() < deadline { + if ["node-a", "node-b"] + .into_iter() + .all(|node| (0..2).all(|device| cluster.file_count(node, device) == 0)) + { + cluster.log("distributed_delete_reclaimed"); + return; } + tokio::time::sleep(Duration::from_millis(100)).await; } - cluster.log("distributed_delete_verified"); + panic!( + "distributed GC did not reclaim all stripe files before deadline: {}", + cluster.diagnostics() + ); } diff --git a/kv-service/configs/server-nvmeof.toml b/kv-service/configs/server-nvmeof.toml index d69a7aa..ddd2935 100644 --- a/kv-service/configs/server-nvmeof.toml +++ b/kv-service/configs/server-nvmeof.toml @@ -49,6 +49,13 @@ redis_key_prefix = "contextstore:metadata:" redis_connect_timeout_ms = 1000 redis_command_timeout_ms = 1000 +[gc] +enabled = true +interval_seconds = 300 +grace_seconds = 600 +max_tasks_per_run = 1000 +task_lease_seconds = 300 + [metrics] enabled = true listen = "0.0.0.0:9090" diff --git a/kv-service/configs/server-test.toml b/kv-service/configs/server-test.toml index 3bf499c..fe555eb 100644 --- a/kv-service/configs/server-test.toml +++ b/kv-service/configs/server-test.toml @@ -50,6 +50,13 @@ redis_key_prefix = "contextstore:metadata:" redis_connect_timeout_ms = 1000 redis_command_timeout_ms = 1000 +[gc] +enabled = true +interval_seconds = 300 +grace_seconds = 600 +max_tasks_per_run = 1000 +task_lease_seconds = 300 + [metrics] enabled = false # Set to true after building with --features metrics listen = "0.0.0.0:9090" diff --git a/kv-service/configs/server.toml b/kv-service/configs/server.toml index 3cb1122..c531e57 100644 --- a/kv-service/configs/server.toml +++ b/kv-service/configs/server.toml @@ -54,6 +54,16 @@ redis_key_prefix = "contextstore:metadata:" redis_connect_timeout_ms = 1000 redis_command_timeout_ms = 1000 +[gc] +# Retire metadata immediately and reclaim superseded physical generations after +# the grace period. TTL remains access-triggered; this worker does not scan for +# expired metadata that has not been accessed. +enabled = true +interval_seconds = 300 +grace_seconds = 600 +max_tasks_per_run = 1000 +task_lease_seconds = 300 + [metrics] # Prometheus HTTP exporter (requires --features metrics at build time). # Endpoints: GET http:///metrics, GET http:///health diff --git a/kv-service/server/src/api/service.rs b/kv-service/server/src/api/service.rs index db9ef20..94c6de7 100644 --- a/kv-service/server/src/api/service.rs +++ b/kv-service/server/src/api/service.rs @@ -67,6 +67,7 @@ struct StreamPutStats { metadata_elapsed: Duration, } +#[derive(Clone)] pub struct KVServiceImpl { ctx: Arc, write_locks: Arc>>>, @@ -76,6 +77,7 @@ pub struct KVServiceImpl { mod tests { use super::*; use crate::metadata::{ChunkLocation, StripingInfo}; + use tempfile::TempDir; fn ctx_with_nodes(nodes: Vec) -> KVServiceContext { ctx_with_local_and_nodes("coordinator", "127.0.0.1:50051", nodes) @@ -106,6 +108,15 @@ mod tests { } } + fn ctx_with_gc_storage(dir: &TempDir) -> Arc { + let mut cfg = crate::config::Config::default(); + cfg.metadata.redis_url = format!("memory://api-service-gc-{}", dir.path().display()); + cfg.storage.devices = vec![dir.path().join("nvme0")]; + cfg.storage.striping_threshold = 0; + cfg.gc.grace_seconds = 0; + Arc::new(KVServiceContext::new(cfg).unwrap()) + } + fn meta() -> BlockMeta { BlockMeta { device_id: 0, @@ -342,6 +353,44 @@ mod tests { .unwrap() .is_none()); } + + #[tokio::test] + async fn gc_worker_reclaims_only_the_retired_generation() { + let dir = TempDir::new().unwrap(); + let ctx = ctx_with_gc_storage(&dir); + let service = KVServiceImpl::new_shared(ctx.clone()); + let key = key(); + + ctx.storage + .put(&key, Bytes::from_static(b"generation-one"), meta()) + .unwrap(); + let first = ctx + .metadata + .get_block(&key.to_string_key()) + .unwrap() + .unwrap(); + ctx.storage + .put(&key, Bytes::from_static(b"generation-two"), meta()) + .unwrap(); + let current = ctx + .metadata + .get_block(&key.to_string_key()) + .unwrap() + .unwrap(); + + service.run_gc_once().await; + + assert!(!std::path::Path::new(&first.file_path).exists()); + assert!(std::path::Path::new(¤t.file_path).exists()); + assert_eq!( + ctx.metadata + .get_block(&key.to_string_key()) + .unwrap() + .unwrap() + .object_generation, + current.object_generation + ); + } } impl KVServiceImpl { @@ -469,17 +518,19 @@ impl KVServiceImpl { let metadata = self.ctx.metadata.clone(); let str_key = key.to_string_key(); let expected = meta.clone(); - let deleted = tokio::task::spawn_blocking(move || { - metadata.delete_block_if_matches(&str_key, &expected) + let not_before = chrono::Utc::now() + .timestamp() + .saturating_add(self.ctx.config.gc.grace_seconds as i64); + let retired = tokio::task::spawn_blocking(move || { + metadata.retire_block_if_matches(&str_key, &expected, not_before) }) .await .map_err(|e| Status::internal(e.to_string()))? .map_err(Status::from)?; - if !deleted { + if !retired { return Ok(false); } self.ctx.memory.invalidate(key); - self.delete_distributed_chunks(placement).await?; return Ok(true); } @@ -533,7 +584,9 @@ impl KVServiceImpl { let data_len = data.len() as u64; if is_local_node(&ctx, &node) { let device_stripe_index = local_device_stripe_index(&ctx, &key, stripe_index) - .ok_or_else(|| Status::failed_precondition("stripe assigned to a different data node"))?; + .ok_or_else(|| { + Status::failed_precondition("stripe assigned to a different data node") + })?; let key_for_write = key.clone(); let storage = ctx.storage.clone(); let generation = descriptor.object_generation; @@ -699,7 +752,9 @@ impl KVServiceImpl { .map(|loc| chunk_location_to_pb(&key, loc)) .collect::>(); if self.ctx.storage.verify_stripe_checksums() - && locations.iter().any(|location| location.checksum.is_empty()) + && locations + .iter() + .any(|location| location.checksum.is_empty()) { for chunk in rollback_chunks { let _ = Self::delete_chunk_from_placement(self.ctx.clone(), chunk).await; @@ -713,10 +768,7 @@ impl KVServiceImpl { .iter() .map(|loc| loc.storage_handle.clone()) .collect(); - let chunk_checksums = locations - .iter() - .map(|loc| loc.checksum.clone()) - .collect(); + let chunk_checksums = locations.iter().map(|loc| loc.checksum.clone()).collect(); let mut committed = prepared_meta; committed.size = total as u64; committed.file_path = String::new(); @@ -733,12 +785,17 @@ impl KVServiceImpl { let metadata = self.ctx.metadata.clone(); let str_key = key.to_string_key(); let committed_meta = committed.clone(); + let gc_not_before = chrono::Utc::now() + .timestamp() + .saturating_add(self.ctx.config.gc.grace_seconds as i64); let metadata_start = Instant::now(); let committed_result = tokio::task::spawn_blocking(move || { if if_absent { metadata.put_block_if_absent(&str_key, &committed_meta) } else { - metadata.put_block(&str_key, &committed_meta).map(|_| true) + metadata + .replace_block(&str_key, &committed_meta, gc_not_before) + .map(|_| true) } }) .await; @@ -838,7 +895,9 @@ impl KVServiceImpl { let chunk_start = Instant::now(); let data_len: u64 = segments.iter().map(|s| s.len() as u64).sum(); let device_stripe_index = local_device_stripe_index(&ctx, &key, stripe_index) - .ok_or_else(|| Status::failed_precondition("stripe assigned to a different data node"))?; + .ok_or_else(|| { + Status::failed_precondition("stripe assigned to a different data node") + })?; let key_for_write = key.clone(); let storage = ctx.storage.clone(); let generation = descriptor.object_generation; @@ -928,12 +987,8 @@ impl KVServiceImpl { .prepare_write_meta(&key, meta, declared_total as u64) .map_err(Status::from)?; let prepare_elapsed = prepare_start.elapsed(); - let descriptor = self.make_distributed_descriptor( - &key, - &prepared_meta, - stripe_count, - chunk_size as u64, - ); + let descriptor = + self.make_distributed_descriptor(&key, &prepared_meta, stripe_count, chunk_size as u64); let mut inflight: JoinSet> = JoinSet::new(); let mut locations: Vec = Vec::with_capacity(stripe_count); @@ -1085,10 +1140,7 @@ impl KVServiceImpl { .iter() .map(|loc| loc.storage_handle.clone()) .collect(); - let chunk_checksums = locations - .iter() - .map(|loc| loc.checksum.clone()) - .collect(); + let chunk_checksums = locations.iter().map(|loc| loc.checksum.clone()).collect(); let mut committed_meta = prepared_meta; committed_meta.size = declared_total as u64; committed_meta.file_path = String::new(); @@ -1105,12 +1157,17 @@ impl KVServiceImpl { let metadata = self.ctx.metadata.clone(); let str_key = key.to_string_key(); let meta_for_commit = committed_meta.clone(); + let gc_not_before = chrono::Utc::now() + .timestamp() + .saturating_add(self.ctx.config.gc.grace_seconds as i64); let metadata_start = Instant::now(); let committed = tokio::task::spawn_blocking(move || { if if_absent { metadata.put_block_if_absent(&str_key, &meta_for_commit) } else { - metadata.put_block(&str_key, &meta_for_commit).map(|_| true) + metadata + .replace_block(&str_key, &meta_for_commit, gc_not_before) + .map(|_| true) } }) .await @@ -1286,6 +1343,73 @@ impl KVServiceImpl { } Ok(()) } + + /// Process a bounded set of durable retired-generation tasks. A task is + /// checked against current metadata before physical deletion, so delayed + /// cleanup cannot remove a generation that has become current again. + pub async fn run_gc_once(&self) { + let now = chrono::Utc::now().timestamp(); + let metadata = self.ctx.metadata.clone(); + let limit = self.ctx.config.gc.max_tasks_per_run; + let lease_seconds = self.ctx.config.gc.task_lease_seconds; + let tasks = match tokio::task::spawn_blocking(move || { + metadata.take_due_gc_tasks(now, limit, lease_seconds) + }) + .await + { + Ok(Ok(tasks)) => tasks, + Ok(Err(error)) => { + tracing::warn!("GC could not claim due tasks: {}", error); + return; + } + Err(error) => { + tracing::warn!("GC task claim worker failed: {}", error); + return; + } + }; + + for mut task in tasks { + let key = match InternalKey::from_string_key(&task.key) { + Ok(key) => key, + Err(error) => { + tracing::warn!(task_id = %task.id, "GC task has invalid key: {}", error); + if let Err(error) = self.ctx.metadata.complete_gc_task(&task.id) { + tracing::warn!(task_id = %task.id, "GC could not discard invalid task: {}", error); + } + continue; + } + }; + let current = match self.ctx.metadata.get_block(&task.key) { + Ok(meta) => meta, + Err(error) => { + self.reschedule_gc_task(&mut task, now, &error.to_string()); + continue; + } + }; + if current + .as_ref() + .is_some_and(|meta| Self::meta_identity_matches(meta, &task.retired)) + { + self.reschedule_gc_task(&mut task, now, "retired generation is current"); + continue; + } + let placement = placement_from_meta(&self.ctx, &key, &task.retired); + if let Err(error) = self.delete_distributed_chunks(placement).await { + self.reschedule_gc_task(&mut task, now, error.message()); + } else if let Err(error) = self.ctx.metadata.complete_gc_task(&task.id) { + self.reschedule_gc_task(&mut task, now, &error.to_string()); + } + } + } + + fn reschedule_gc_task(&self, task: &mut crate::metadata::GcTask, now: i64, reason: &str) { + task.attempts = task.attempts.saturating_add(1); + let delay = 5_i64.saturating_mul(1_i64 << task.attempts.min(6)); + task.not_before = now.saturating_add(delay.min(300)); + if let Err(error) = self.ctx.metadata.reschedule_gc_task(task) { + tracing::warn!(task_id = %task.id, attempts = task.attempts, "GC retry could not be persisted after {}: {}", reason, error); + } + } } fn pb_key_to_internal(k: &pb::ObjectKey) -> InternalKey { @@ -1822,34 +1946,13 @@ impl pb::kv_service_server::KvService for KVServiceImpl { .key .ok_or_else(|| Status::invalid_argument("missing key"))?; let internal = pb_key_to_internal(&key); - let str_key = internal.to_string_key(); - let meta_ctx = self.ctx.clone(); - let meta = tokio::task::spawn_blocking(move || meta_ctx.metadata.get_block(&str_key)) - .await - .map_err(|e| Status::internal(e.to_string()))? - .map_err(Status::from)?; - if let Some(meta) = meta.as_ref() { - let placement = placement_from_meta(&self.ctx, &internal, meta); - if self.placement_has_remote_chunks(&placement) { - self.delete_distributed_chunks(placement).await?; - self.ctx.memory.invalidate(&internal); - let metadata = self.ctx.metadata.clone(); - let str_key = internal.to_string_key(); - let expected = meta.clone(); - tokio::task::spawn_blocking(move || { - metadata.delete_block_if_matches(&str_key, &expected) - }) - .await - .map_err(|e| Status::internal(e.to_string()))? - .map_err(Status::from)?; - return Ok(Response::new(pb::DeleteResponse { success: true })); - } - } - let ctx = self.ctx.clone(); - let ok = tokio::task::spawn_blocking(move || ctx.memory.delete(&internal)) + let storage = self.ctx.storage.clone(); + let key_for_retire = internal.clone(); + let ok = tokio::task::spawn_blocking(move || storage.retire(&key_for_retire)) .await .map_err(|e| Status::internal(e.to_string()))? .map_err(Status::from)?; + self.ctx.memory.invalidate(&internal); Ok(Response::new(pb::DeleteResponse { success: ok })) } @@ -2175,11 +2278,16 @@ impl pb::kv_service_server::KvService for KVServiceImpl { let metadata = self.ctx.metadata.clone(); let str_key = internal.to_string_key(); let meta_for_commit = committed_meta; + let gc_not_before = chrono::Utc::now() + .timestamp() + .saturating_add(self.ctx.config.gc.grace_seconds as i64); let committed = tokio::task::spawn_blocking(move || { if if_not_exists { metadata.put_block_if_absent(&str_key, &meta_for_commit) } else { - metadata.put_block(&str_key, &meta_for_commit).map(|_| true) + metadata + .replace_block(&str_key, &meta_for_commit, gc_not_before) + .map(|_| true) } }) .await @@ -2218,7 +2326,9 @@ impl pb::kv_service_server::KvService for KVServiceImpl { let stripe_index_u32 = req.stripe_index; let stripe_index = stripe_index_u32 as usize; let device_stripe_index = local_device_stripe_index(&self.ctx, &internal, stripe_index) - .ok_or_else(|| Status::failed_precondition("stripe assigned to a different data node"))?; + .ok_or_else(|| { + Status::failed_precondition("stripe assigned to a different data node") + })?; let offset = stripe_index as u64 * req.chunk_size; let local = local_node(&self.ctx); let ctx = self.ctx.clone(); @@ -2548,8 +2658,7 @@ impl pb::kv_service_server::KvService for KVServiceImpl { let start = i * SUB_CHUNK; let end = (start + SUB_CHUNK).min(seg_len); sent_bytes += (end - start) as u64; - let is_last = - next_send + 1 == stripe_count && i + 1 == n_sub; + let is_last = next_send + 1 == stripe_count && i + 1 == n_sub; let chunk = pb::DataChunk { data: seg.slice(start..end), offset: base + start as i64, @@ -2783,11 +2892,10 @@ impl pb::kv_service_server::KvService for KVServiceImpl { // a live object refuses the write; an expired one is purged then overwritten. let metadata = self.ctx.metadata.clone(); let str_key = internal.to_string_key(); - let existing = - tokio::task::spawn_blocking(move || metadata.get_block(&str_key)) - .await - .map_err(|e| Status::internal(e.to_string()))? - .map_err(Status::from)?; + let existing = tokio::task::spawn_blocking(move || metadata.get_block(&str_key)) + .await + .map_err(|e| Status::internal(e.to_string()))? + .map_err(Status::from)?; if let Some(existing) = existing { if !existing.is_expired() { let result = Ok(Response::new(pb::PutResponse { @@ -2797,7 +2905,8 @@ impl pb::kv_service_server::KvService for KVServiceImpl { self.record_request("put_stream", request_start, &result, "ok"); return result; } - self.purge_expired_object_locked(&internal, &existing).await?; + self.purge_expired_object_locked(&internal, &existing) + .await?; } } let first_is_last = first.is_last; @@ -2841,8 +2950,7 @@ impl pb::kv_service_server::KvService for KVServiceImpl { slowest_stripe_us = stats.slowest_stripe.as_micros(), metadata_us = stats.metadata_elapsed.as_micros(), total_us = t_total.as_micros(), - throughput_gib_s = - declared_total as f64 / t_total.as_secs_f64() / 1_073_741_824.0, + throughput_gib_s = declared_total as f64 / t_total.as_secs_f64() / 1_073_741_824.0, ); let result = Ok(Response::new(pb::PutResponse { success: inserted, diff --git a/kv-service/server/src/config.rs b/kv-service/server/src/config.rs index 2a185c1..096d1c9 100644 --- a/kv-service/server/src/config.rs +++ b/kv-service/server/src/config.rs @@ -27,6 +27,8 @@ pub struct Config { pub metrics: MetricsConfig, #[serde(default)] pub gds: GdsConfig, + #[serde(default)] + pub gc: GcConfig, } // ===== Cluster / Placement ===== @@ -262,6 +264,33 @@ impl Default for GdsConfig { } } +// ===== Generation Garbage Collection ===== +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct GcConfig { + /// Enable delayed reclamation of superseded object generations. + pub enabled: bool, + /// Frequency of the persistent retirement-queue worker. + pub interval_seconds: u64, + /// Minimum time a retired generation remains available on storage. + pub grace_seconds: u64, + /// Upper bound on queued tasks processed in one worker iteration. + pub max_tasks_per_run: usize, + /// Task lease duration. Expired leases are reclaimed after a worker crash. + pub task_lease_seconds: u64, +} + +impl Default for GcConfig { + fn default() -> Self { + Self { + enabled: true, + interval_seconds: 300, + grace_seconds: 600, + max_tasks_per_run: 1_000, + task_lease_seconds: 300, + } + } +} + impl Default for Config { fn default() -> Self { Self { @@ -274,6 +303,7 @@ impl Default for Config { metadata: Default::default(), metrics: Default::default(), gds: Default::default(), + gc: Default::default(), } } } @@ -350,6 +380,21 @@ impl Config { "metadata.redis_command_timeout_ms must be greater than 0".to_string(), )); } + if self.gc.enabled && self.gc.interval_seconds == 0 { + return Err(KVError::Config( + "gc.interval_seconds must be greater than 0 when GC is enabled".to_string(), + )); + } + if self.gc.enabled && self.gc.max_tasks_per_run == 0 { + return Err(KVError::Config( + "gc.max_tasks_per_run must be greater than 0 when GC is enabled".to_string(), + )); + } + if self.gc.enabled && self.gc.task_lease_seconds == 0 { + return Err(KVError::Config( + "gc.task_lease_seconds must be greater than 0 when GC is enabled".to_string(), + )); + } Ok(()) } } diff --git a/kv-service/server/src/main.rs b/kv-service/server/src/main.rs index e6562f1..aa3f797 100644 --- a/kv-service/server/src/main.rs +++ b/kv-service/server/src/main.rs @@ -11,6 +11,7 @@ use contextstore_server::{ }; use std::path::PathBuf; use std::sync::Arc; +use std::time::Duration; use tonic::transport::Server; use tracing::{info, warn}; use tracing_subscriber::EnvFilter; @@ -211,6 +212,22 @@ async fn main() -> anyhow::Result<()> { } let svc = KVServiceImpl::new_shared(ctx); + if config.gc.enabled { + let task_worker = svc.clone(); + let interval = Duration::from_secs(config.gc.interval_seconds); + tokio::spawn(async move { + let mut ticker = tokio::time::interval(interval); + loop { + ticker.tick().await; + task_worker.run_gc_once().await; + } + }); + info!( + grace_seconds = config.gc.grace_seconds, + interval_seconds = config.gc.interval_seconds, + "generation garbage collection enabled" + ); + } // Large KV payloads: a single layer can reach several MB and a batch several // hundred MB. Raise tonic's default 4MB limit to 2GiB. let kv_server = KvServiceServer::new(svc) diff --git a/kv-service/server/src/metadata.rs b/kv-service/server/src/metadata.rs index 2e5da59..dfae7ae 100644 --- a/kv-service/server/src/metadata.rs +++ b/kv-service/server/src/metadata.rs @@ -16,6 +16,9 @@ use std::collections::HashMap; const BLOCK_META_SEGMENT: &str = "block_meta"; const GENERATION_SEGMENT: &str = "generation"; +const GC_DUE_SEGMENT: &str = "gc:due"; +const GC_INFLIGHT_SEGMENT: &str = "gc:inflight"; +const GC_TASK_SEGMENT: &str = "gc:task:"; fn default_object_generation() -> u64 { 1 @@ -113,6 +116,17 @@ impl BlockMeta { } } +/// A durable request to reclaim a generation that is no longer referenced by +/// current metadata. The complete placement is retained for crash-safe cleanup. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct GcTask { + pub id: String, + pub key: String, + pub retired: BlockMeta, + pub not_before: i64, + pub attempts: u32, +} + enum MetadataBackend { Redis(RedisMetadataBackend), #[cfg(test)] @@ -161,6 +175,80 @@ impl MetadataService { } } + /// Atomically publish `meta` and enqueue the previously visible generation + /// for delayed physical reclamation. + pub fn replace_block( + &self, + key: &str, + meta: &BlockMeta, + not_before: i64, + ) -> Result> { + match &self.backend { + MetadataBackend::Redis(backend) => backend.replace_block(key, meta, not_before), + #[cfg(test)] + MetadataBackend::Memory(backend) => backend.replace_block(key, meta, not_before), + } + } + + /// Atomically make the object unreachable and enqueue its physical layout. + pub fn retire_block(&self, key: &str, not_before: i64) -> Result> { + match &self.backend { + MetadataBackend::Redis(backend) => backend.retire_block(key, not_before), + #[cfg(test)] + MetadataBackend::Memory(backend) => backend.retire_block(key, not_before), + } + } + + /// Retire only the metadata version observed by the caller. + pub fn retire_block_if_matches( + &self, + key: &str, + expected: &BlockMeta, + not_before: i64, + ) -> Result { + match &self.backend { + MetadataBackend::Redis(backend) => { + backend.retire_block_if_matches(key, expected, not_before) + } + #[cfg(test)] + MetadataBackend::Memory(backend) => { + backend.retire_block_if_matches(key, expected, not_before) + } + } + } + + /// Claim due tasks under a time-bounded lease. + pub fn take_due_gc_tasks( + &self, + now: i64, + limit: usize, + lease_seconds: u64, + ) -> Result> { + match &self.backend { + MetadataBackend::Redis(backend) => backend.take_due_gc_tasks(now, limit, lease_seconds), + #[cfg(test)] + MetadataBackend::Memory(backend) => { + backend.take_due_gc_tasks(now, limit, lease_seconds) + } + } + } + + pub fn complete_gc_task(&self, task_id: &str) -> Result<()> { + match &self.backend { + MetadataBackend::Redis(backend) => backend.complete_gc_task(task_id), + #[cfg(test)] + MetadataBackend::Memory(backend) => backend.complete_gc_task(task_id), + } + } + + pub fn reschedule_gc_task(&self, task: &GcTask) -> Result<()> { + match &self.backend { + MetadataBackend::Redis(backend) => backend.reschedule_gc_task(task), + #[cfg(test)] + MetadataBackend::Memory(backend) => backend.reschedule_gc_task(task), + } + } + pub fn get_block(&self, key: &str) -> Result> { match &self.backend { MetadataBackend::Redis(backend) => backend.get_block(key), @@ -313,6 +401,18 @@ impl RedisMetadataBackend { format!("{}{}:{}", self.prefix, GENERATION_SEGMENT, key) } + fn gc_due_key(&self) -> String { + format!("{}{}", self.prefix, GC_DUE_SEGMENT) + } + + fn gc_inflight_key(&self) -> String { + format!("{}{}", self.prefix, GC_INFLIGHT_SEGMENT) + } + + fn gc_task_prefix(&self) -> String { + format!("{}{}", self.prefix, GC_TASK_SEGMENT) + } + fn serialize(meta: &BlockMeta) -> Result> { serde_json::to_vec(meta).map_err(|e| KVError::Metadata(format!("serialize: {}", e))) } @@ -342,6 +442,201 @@ impl RedisMetadataBackend { Ok(result.is_some()) } + fn replace_block( + &self, + key: &str, + meta: &BlockMeta, + not_before: i64, + ) -> Result> { + static SCRIPT: OnceLock = OnceLock::new(); + let script = SCRIPT.get_or_init(|| { + redis::Script::new( + r#" + local old = redis.call("GET", KEYS[1]) + redis.call("SET", KEYS[1], ARGV[1]) + if old then + local id = redis.sha1hex(KEYS[1] .. ":" .. old) + local task_key = ARGV[2] .. id + local task = cjson.encode({id=id, key=ARGV[3], retired=cjson.decode(old), not_before=tonumber(ARGV[4]), attempts=0}) + redis.call("HSET", task_key, "payload", task) + redis.call("ZADD", KEYS[2], tonumber(ARGV[4]), task_key) + end + return old + "#, + ) + }); + let block_key = self.block_key(key); + let due_key = self.gc_due_key(); + let bytes = Self::serialize(meta)?; + let old = + self.run_with_reconnect("replace block metadata and enqueue GC", |connection| { + script + .key(&block_key) + .key(&due_key) + .arg(&bytes) + .arg(self.gc_task_prefix()) + .arg(key) + .arg(not_before) + .invoke::>>(connection) + })?; + old.as_deref().map(Self::deserialize).transpose() + } + + fn retire_block(&self, key: &str, not_before: i64) -> Result> { + static SCRIPT: OnceLock = OnceLock::new(); + let script = SCRIPT.get_or_init(|| { + redis::Script::new( + r#" + local old = redis.call("GET", KEYS[1]) + if old then + redis.call("DEL", KEYS[1]) + local id = redis.sha1hex(KEYS[1] .. ":" .. old) + local task_key = ARGV[1] .. id + local task = cjson.encode({id=id, key=ARGV[2], retired=cjson.decode(old), not_before=tonumber(ARGV[3]), attempts=0}) + redis.call("HSET", task_key, "payload", task) + redis.call("ZADD", KEYS[2], tonumber(ARGV[3]), task_key) + end + return old + "#, + ) + }); + let block_key = self.block_key(key); + let due_key = self.gc_due_key(); + let old = + self.run_with_reconnect("retire block metadata and enqueue GC", |connection| { + script + .key(&block_key) + .key(&due_key) + .arg(self.gc_task_prefix()) + .arg(key) + .arg(not_before) + .invoke::>>(connection) + })?; + old.as_deref().map(Self::deserialize).transpose() + } + + fn retire_block_if_matches( + &self, + key: &str, + expected: &BlockMeta, + not_before: i64, + ) -> Result { + static SCRIPT: OnceLock = OnceLock::new(); + let script = SCRIPT.get_or_init(|| { + redis::Script::new( + r#" + local current = redis.call("GET", KEYS[1]) + if not current then return 0 end + local expected = cjson.decode(ARGV[1]) + local actual = cjson.decode(current) + if actual["object_handle"] ~= expected["object_handle"] + or actual["object_generation"] ~= expected["object_generation"] + or actual["layout_version"] ~= expected["layout_version"] + or actual["content_etag"] ~= expected["content_etag"] + or actual["size"] ~= expected["size"] then + return 0 + end + redis.call("DEL", KEYS[1]) + local id = redis.sha1hex(KEYS[1] .. ":" .. current) + local task_key = ARGV[2] .. id + local task = cjson.encode({id=id, key=ARGV[3], retired=actual, not_before=tonumber(ARGV[4]), attempts=0}) + redis.call("HSET", task_key, "payload", task) + redis.call("ZADD", KEYS[2], tonumber(ARGV[4]), task_key) + return 1 + "#, + ) + }); + let expected = Self::serialize(expected)?; + let block_key = self.block_key(key); + let due_key = self.gc_due_key(); + let retired = self.run_with_reconnect("retire matching block metadata", |connection| { + script + .key(&block_key) + .key(&due_key) + .arg(&expected) + .arg(self.gc_task_prefix()) + .arg(key) + .arg(not_before) + .invoke::(connection) + })?; + Ok(retired == 1) + } + + fn take_due_gc_tasks(&self, now: i64, limit: usize, lease_seconds: u64) -> Result> { + static SCRIPT: OnceLock = OnceLock::new(); + let script = SCRIPT.get_or_init(|| redis::Script::new(r#" + local expired = redis.call("ZRANGEBYSCORE", KEYS[2], "-inf", ARGV[1]) + for _, task_key in ipairs(expired) do + if redis.call("ZREM", KEYS[2], task_key) == 1 then + redis.call("ZADD", KEYS[1], ARGV[1], task_key) + end + end + local task_keys = redis.call("ZRANGEBYSCORE", KEYS[1], "-inf", ARGV[1], "LIMIT", 0, ARGV[2]) + local payloads = {} + for _, task_key in ipairs(task_keys) do + if redis.call("ZREM", KEYS[1], task_key) == 1 then + local payload = redis.call("HGET", task_key, "payload") + if payload then + redis.call("ZADD", KEYS[2], tonumber(ARGV[1]) + tonumber(ARGV[3]), task_key) + table.insert(payloads, payload) + end + end + end + return payloads + "#)); + let due_key = self.gc_due_key(); + let inflight_key = self.gc_inflight_key(); + let payloads = self.run_with_reconnect("claim due GC tasks", |connection| { + script + .key(&due_key) + .key(&inflight_key) + .arg(now) + .arg(limit) + .arg(lease_seconds) + .invoke::>>(connection) + })?; + payloads + .into_iter() + .map(|payload| { + serde_json::from_slice(&payload) + .map_err(|e| KVError::Metadata(format!("deserialize GC task: {}", e))) + }) + .collect() + } + + fn complete_gc_task(&self, task_id: &str) -> Result<()> { + let task_key = format!("{}{}", self.gc_task_prefix(), task_id); + self.run_with_reconnect("complete GC task", |connection| { + let _: () = redis::cmd("ZREM") + .arg(self.gc_inflight_key()) + .arg(&task_key) + .query(connection)?; + redis::cmd("DEL").arg(&task_key).query::<()>(connection) + }) + } + + fn reschedule_gc_task(&self, task: &GcTask) -> Result<()> { + let task_key = format!("{}{}", self.gc_task_prefix(), task.id); + let payload = serde_json::to_vec(task) + .map_err(|e| KVError::Metadata(format!("serialize GC task: {}", e)))?; + self.run_with_reconnect("reschedule GC task", |connection| { + let _: () = redis::cmd("ZREM") + .arg(self.gc_inflight_key()) + .arg(&task_key) + .query(connection)?; + let _: () = redis::cmd("HSET") + .arg(&task_key) + .arg("payload") + .arg(&payload) + .query(connection)?; + redis::cmd("ZADD") + .arg(self.gc_due_key()) + .arg(task.not_before) + .arg(&task_key) + .query::<()>(connection) + }) + } + fn get_block(&self, key: &str) -> Result> { let redis_key = self.block_key(key); let bytes = self.run_with_reconnect("get_block", |connection| { @@ -435,6 +730,8 @@ struct MemoryMetadataBackend { struct MemoryMetadataInner { blocks: HashMap>, generations: HashMap, + gc_tasks: HashMap, + gc_leases: HashMap, fail_next_puts: usize, } @@ -482,6 +779,136 @@ impl MemoryMetadataBackend { Ok(true) } + fn task_id(key: &str, meta: &BlockMeta) -> Result { + Ok(format!( + "{}:g{}:l{}:{}", + key, meta.object_generation, meta.layout_version, meta.content_etag + )) + } + + fn enqueue_gc( + inner: &mut MemoryMetadataInner, + key: &str, + retired: BlockMeta, + not_before: i64, + ) -> Result<()> { + let id = Self::task_id(key, &retired)?; + inner.gc_tasks.insert( + id.clone(), + GcTask { + id, + key: key.to_string(), + retired, + not_before, + attempts: 0, + }, + ); + Ok(()) + } + + fn replace_block( + &self, + key: &str, + meta: &BlockMeta, + not_before: i64, + ) -> Result> { + let mut inner = self.inner.lock(); + if inner.fail_next_puts > 0 { + inner.fail_next_puts -= 1; + return Err(KVError::Metadata("injected put failure".to_string())); + } + let old = inner + .blocks + .insert(key.to_string(), Self::serialize(meta)?) + .as_deref() + .map(Self::deserialize) + .transpose()?; + if let Some(retired) = old.as_ref() { + Self::enqueue_gc(&mut inner, key, retired.clone(), not_before)?; + } + Ok(old) + } + + fn retire_block(&self, key: &str, not_before: i64) -> Result> { + let mut inner = self.inner.lock(); + let old = inner + .blocks + .remove(key) + .as_deref() + .map(Self::deserialize) + .transpose()?; + if let Some(retired) = old.as_ref() { + Self::enqueue_gc(&mut inner, key, retired.clone(), not_before)?; + } + Ok(old) + } + + fn same_version(actual: &BlockMeta, expected: &BlockMeta) -> bool { + actual.object_handle == expected.object_handle + && actual.object_generation == expected.object_generation + && actual.layout_version == expected.layout_version + && actual.content_etag == expected.content_etag + && actual.size == expected.size + } + + fn retire_block_if_matches( + &self, + key: &str, + expected: &BlockMeta, + not_before: i64, + ) -> Result { + let mut inner = self.inner.lock(); + let current = inner + .blocks + .get(key) + .map(|bytes| Self::deserialize(bytes)) + .transpose()?; + let Some(current) = current else { + return Ok(false); + }; + if !Self::same_version(¤t, expected) { + return Ok(false); + } + inner.blocks.remove(key); + Self::enqueue_gc(&mut inner, key, current, not_before)?; + Ok(true) + } + + fn take_due_gc_tasks(&self, now: i64, limit: usize, lease_seconds: u64) -> Result> { + let mut inner = self.inner.lock(); + inner.gc_leases.retain(|_, lease_until| *lease_until > now); + let ids = inner + .gc_tasks + .iter() + .filter(|(id, task)| task.not_before <= now && !inner.gc_leases.contains_key(*id)) + .map(|(id, _)| id.clone()) + .take(limit) + .collect::>(); + let tasks = ids + .iter() + .filter_map(|id| inner.gc_tasks.get(id).cloned()) + .collect(); + let lease_until = now.saturating_add(lease_seconds as i64); + for id in ids { + inner.gc_leases.insert(id, lease_until); + } + Ok(tasks) + } + + fn complete_gc_task(&self, task_id: &str) -> Result<()> { + let mut inner = self.inner.lock(); + inner.gc_tasks.remove(task_id); + inner.gc_leases.remove(task_id); + Ok(()) + } + + fn reschedule_gc_task(&self, task: &GcTask) -> Result<()> { + let mut inner = self.inner.lock(); + inner.gc_leases.remove(&task.id); + inner.gc_tasks.insert(task.id.clone(), task.clone()); + Ok(()) + } + fn get_block(&self, key: &str) -> Result> { self.inner .lock() @@ -648,6 +1075,71 @@ mod tests { assert_eq!(svc.next_generation("k1").unwrap(), 11); } + #[test] + fn replace_block_enqueues_the_superseded_generation() { + let svc = MetadataService::new(&test_config("replace-gc")).unwrap(); + let first = meta(); + let mut second = meta(); + second.object_generation = 2; + second.object_handle = "handle-2".to_string(); + second.content_etag = "etag-2".to_string(); + + svc.put_block("k1", &first).unwrap(); + let retired = svc.replace_block("k1", &second, 100).unwrap().unwrap(); + assert_eq!(retired.object_generation, first.object_generation); + assert_eq!(retired.object_handle, first.object_handle); + let current = svc.get_block("k1").unwrap().unwrap(); + assert_eq!(current.object_generation, second.object_generation); + assert_eq!(current.object_handle, second.object_handle); + + let tasks = svc.take_due_gc_tasks(100, 10, 10).unwrap(); + assert_eq!(tasks.len(), 1); + assert_eq!(tasks[0].key, "k1"); + assert_eq!(tasks[0].retired.object_generation, 1); + } + + #[test] + fn retiring_a_block_preserves_the_generation_counter() { + let svc = MetadataService::new(&test_config("retire-generation")).unwrap(); + let first = meta(); + assert_eq!(svc.next_generation("k1").unwrap(), 1); + svc.put_block("k1", &first).unwrap(); + + let retired = svc.retire_block("k1", 100).unwrap().unwrap(); + assert_eq!(retired.object_generation, first.object_generation); + assert!(svc.get_block("k1").unwrap().is_none()); + assert_eq!(svc.next_generation("k1").unwrap(), 2); + } + + #[test] + fn stale_retire_does_not_remove_a_newer_identity() { + let svc = MetadataService::new(&test_config("stale-retire")).unwrap(); + let first = meta(); + let mut replacement = first.clone(); + replacement.content_etag = "replacement-etag".to_string(); + replacement.object_handle = "replacement-handle".to_string(); + + svc.put_block("k1", &replacement).unwrap(); + assert!(!svc.retire_block_if_matches("k1", &first, 100).unwrap()); + let current = svc.get_block("k1").unwrap().unwrap(); + assert_eq!(current.object_handle, replacement.object_handle); + assert!(svc.take_due_gc_tasks(100, 10, 10).unwrap().is_empty()); + } + + #[test] + fn expired_gc_lease_can_be_reclaimed() { + let svc = MetadataService::new(&test_config("gc-lease")).unwrap(); + svc.put_block("k1", &meta()).unwrap(); + svc.retire_block("k1", 100).unwrap(); + + let first = svc.take_due_gc_tasks(100, 1, 10).unwrap(); + assert_eq!(first.len(), 1); + assert!(svc.take_due_gc_tasks(109, 1, 10).unwrap().is_empty()); + let recovered = svc.take_due_gc_tasks(110, 1, 10).unwrap(); + assert_eq!(recovered.len(), 1); + assert_eq!(recovered[0].id, first[0].id); + } + #[test] fn block_meta_expiry_uses_created_at_and_positive_ttl() { let mut meta = meta(); diff --git a/kv-service/server/src/storage_tier.rs b/kv-service/server/src/storage_tier.rs index d88a5d3..6751ea5 100644 --- a/kv-service/server/src/storage_tier.rs +++ b/kv-service/server/src/storage_tier.rs @@ -35,6 +35,7 @@ pub struct StorageTier { striping_chunk_size: u64, rdma_stream_chunk_size: usize, verify_stripe_checksums: bool, + gc_grace_seconds: u64, executor_name: String, device_labels: Vec, metrics: Option>, @@ -84,6 +85,7 @@ impl StorageTier { .max(DIRECT_IO_ALIGNMENT) & !(DIRECT_IO_ALIGNMENT - 1), verify_stripe_checksums: config.storage.verify_stripe_checksums, + gc_grace_seconds: config.gc.grace_seconds, executor_name: config.io_executor.kind.clone(), device_labels: (0..num_devices) .map(|device_id| format!("nvme{}", device_id)) @@ -129,6 +131,24 @@ impl StorageTier { .unwrap_or(0) } + /// Return the configured delay before a retired generation may be reclaimed. + pub fn gc_grace_seconds(&self) -> u64 { + self.gc_grace_seconds + } + + fn gc_not_before(&self) -> i64 { + chrono::Utc::now().timestamp() + self.gc_grace_seconds as i64 + } + + /// Atomically remove the current metadata and schedule its physical layout + /// for delayed reclamation. L1 invalidation is owned by the caller. + pub fn retire(&self, key: &ObjectKey) -> Result { + Ok(self + .metadata + .retire_block(&key.to_string_key(), self.gc_not_before())? + .is_some()) + } + pub fn device_used_bytes(&self, device_id: usize) -> u64 { self.device_used_bytes_total .get(device_id) @@ -416,6 +436,7 @@ impl StorageTier { && actual.size == expected.size } + #[cfg(test)] fn delete_files_for_meta(&self, key: &ObjectKey, meta: &BlockMeta) -> Result<()> { if let Some(stripe) = &meta.striping { let mut last_err = None; @@ -447,44 +468,6 @@ impl StorageTier { Ok(()) } - fn delete_expired_files_best_effort(&self, key: &ObjectKey, meta: &BlockMeta) { - if let Some(stripe) = &meta.striping { - for (index, path) in stripe.chunk_paths.iter().enumerate() { - let result = stripe - .chunk_devices - .get(index) - .copied() - .map(|device_id| device_id as usize) - .or_else(|| self.device_id_for_path(Path::new(path))) - .ok_or_else(|| { - KVError::InvalidArgument(format!( - "unable to identify storage device for {}", - path - )) - }) - .and_then(|device_id| self.delete_file_on_device(device_id, Path::new(path))); - if let Err(error) = result { - warn!( - path = %path, - error = %error, - "failed to delete expired striped file" - ); - } - } - } else { - let path = self.meta_path_or_route(key, meta); - if let Err(error) = - self.delete_file_on_device(self.meta_device_or_route(key, meta), &path) - { - warn!( - path = %path.display(), - error = %error, - "failed to delete expired file" - ); - } - } - } - fn delete_metadata_if_current(&self, key: &ObjectKey, meta: &BlockMeta) -> Result { self.metadata .delete_block_if_matches(&key.to_string_key(), meta) @@ -501,7 +484,9 @@ impl StorageTier { let result = if if_absent { self.metadata.put_block_if_absent(&str_key, meta) } else { - self.metadata.put_block(&str_key, meta).map(|_| true) + self.metadata + .replace_block(&str_key, meta, self.gc_not_before()) + .map(|_| true) }; match result { Ok(true) => Ok(true), @@ -539,8 +524,8 @@ impl StorageTier { if !current.is_expired() || !Self::meta_identity_matches(¤t, expected) { return Ok(false); } - self.delete_expired_files_best_effort(key, ¤t); - self.delete_metadata_if_current(key, ¤t) + self.metadata + .retire_block_if_matches(&str_key, ¤t, self.gc_not_before()) } fn purge_if_expired(&self, key: &ObjectKey, meta: &BlockMeta) -> Result { @@ -1250,17 +1235,7 @@ impl StorageTier { } pub fn delete(&self, key: &ObjectKey) -> Result { - let str_key = key.to_string_key(); - let meta = self.metadata.get_block(&str_key)?; - let existed = meta.is_some(); - - if let Some(meta) = meta { - self.delete_files_for_meta(key, &meta)?; - self.delete_metadata_if_current(key, &meta)?; - return Ok(existed); - } - self.metadata.delete_block(&str_key)?; - Ok(existed) + self.retire(key) } pub fn exists(&self, key: &ObjectKey) -> Result { @@ -3026,7 +3001,7 @@ mod tests { } #[test] - fn expired_single_object_is_purged_on_get() { + fn expired_single_object_is_retired_on_get() { let tmp = TempDir::new().unwrap(); let cfg = test_config(tmp.path()); let router = Arc::new(ShardRouter::new(&cfg).unwrap()); @@ -3047,11 +3022,22 @@ mod tests { .get_block(&key.to_string_key()) .unwrap() .is_none()); - assert!(!std::path::Path::new(&path).exists()); + assert!(std::path::Path::new(&path).exists()); + assert_eq!( + st.metadata + .take_due_gc_tasks( + chrono::Utc::now().timestamp() + cfg.gc.grace_seconds as i64, + 1, + 1 + ) + .unwrap() + .len(), + 1 + ); } #[test] - fn expired_striped_object_purges_all_chunks() { + fn expired_striped_object_is_retired_on_get() { let tmp = TempDir::new().unwrap(); let mut cfg = test_config(tmp.path()); cfg.storage.striping_threshold = 8; @@ -3074,13 +3060,22 @@ mod tests { .get_block(&key.to_string_key()) .unwrap() .is_none()); - assert!(paths - .iter() - .all(|path| !std::path::Path::new(path).exists())); + assert!(paths.iter().all(|path| std::path::Path::new(path).exists())); + assert_eq!( + st.metadata + .take_due_gc_tasks( + chrono::Utc::now().timestamp() + cfg.gc.grace_seconds as i64, + 1, + 1 + ) + .unwrap() + .len(), + 1 + ); } #[test] - fn expired_object_purges_metadata_when_chunk_cleanup_fails() { + fn expired_object_retires_metadata_without_synchronous_chunk_cleanup() { let tmp = TempDir::new().unwrap(); let cfg = test_config(tmp.path()); let router = Arc::new(ShardRouter::new(&cfg).unwrap()); @@ -3111,7 +3106,7 @@ mod tests { .get_block(&key.to_string_key()) .unwrap() .is_none()); - assert!(!valid_path.exists()); + assert!(valid_path.exists()); } #[test] @@ -3507,14 +3502,26 @@ mod tests { let (got, _) = st.get(&key).unwrap().unwrap(); assert_eq!(got.as_ref(), data.as_slice()); - // delete should clean up all chunks. + // Delete retires the metadata immediately and preserves the chunk layout + // until the background GC worker reaches its grace deadline. st.delete(&key).unwrap(); for p in &stripe.chunk_paths { assert!( - !std::path::Path::new(p).exists(), - "chunk {} still exists", + std::path::Path::new(p).exists(), + "chunk {} should remain until GC", p ); } + assert_eq!( + st.metadata + .take_due_gc_tasks( + chrono::Utc::now().timestamp() + cfg.gc.grace_seconds as i64, + 1, + 1 + ) + .unwrap() + .len(), + 1 + ); } }