diff --git a/crates/targets/Cargo.toml b/crates/targets/Cargo.toml index 570ca0b09..470362719 100644 --- a/crates/targets/Cargo.toml +++ b/crates/targets/Cargo.toml @@ -42,7 +42,7 @@ tokio-postgres-rustls = { workspace = true } tracing = { workspace = true } url = { workspace = true } urlencoding = { workspace = true } -uuid = { workspace = true, features = ["v4", "serde"] } +uuid = { workspace = true, features = ["v4", "v7", "serde"] } sysinfo = { workspace = true, features = ["multithread"] } rustfs-kafka-async = { workspace = true } mysql_async = { workspace = true } diff --git a/crates/targets/src/runtime/adapter.rs b/crates/targets/src/runtime/adapter.rs index 5b3649b8c..9e7485634 100644 --- a/crates/targets/src/runtime/adapter.rs +++ b/crates/targets/src/runtime/adapter.rs @@ -142,8 +142,15 @@ where replay_workers: &mut ReplayWorkerManager, activation: RuntimeActivation, ) -> Result<(), TargetError> { + // Stop (and join) the old replay workers before installing the new set so + // no two workers ever drain the same store concurrently, then close the + // old targets. A close failure during reload is logged but does not abort + // the reload — the new configuration must still take effect. self.stop_replay_workers(replay_workers).await; - runtime.clear_and_close().await; + let close_errors = runtime.clear_and_close().await; + if !close_errors.is_empty() { + tracing::warn!(failed_targets = close_errors.len(), "Some targets failed to close during runtime reload"); + } for target in activation.targets { runtime.add_arc(target); @@ -174,8 +181,19 @@ where runtime: &mut TargetRuntimeManager, replay_workers: &mut ReplayWorkerManager, ) -> Result<(), TargetError> { + // On explicit shutdown, propagate any close/flush failures instead of + // swallowing them: the runtime is still fully torn down, but the caller + // learns that a target could not be flushed/closed cleanly. self.stop_replay_workers(replay_workers).await; - runtime.clear_and_close().await; + let close_errors = runtime.clear_and_close().await; + if !close_errors.is_empty() { + let detail = close_errors + .into_iter() + .map(|(target_id, err)| format!("{target_id}: {err}")) + .collect::>() + .join("; "); + return Err(TargetError::Storage(format!("Failed to close {detail}"))); + } Ok(()) } } diff --git a/crates/targets/src/runtime/mod.rs b/crates/targets/src/runtime/mod.rs index 24d1e9061..4db68527f 100644 --- a/crates/targets/src/runtime/mod.rs +++ b/crates/targets/src/runtime/mod.rs @@ -31,14 +31,36 @@ use std::sync::Arc; use std::{collections::HashMap, fmt::Debug}; use std::{future::Future, pin::Pin, time::Duration}; use tokio::sync::{Semaphore, mpsc}; +use tokio::task::JoinHandle; /// Shared target trait object used by the runtime manager. pub type SharedTarget = Arc + Send + Sync>; type ReplayHook = Arc) -> Pin + Send>> + Send + Sync>; -#[derive(Debug, Default)] +/// Upper bound on how long [`ReplayWorkerManager::stop_all`] waits for a single +/// replay worker to observe its cancel signal and exit before it is forcibly +/// aborted. Workers observe cancellation promptly (including during retry +/// backoff), so this only guards against a wedged task. +const STOP_JOIN_TIMEOUT: Duration = Duration::from_secs(5); + +/// Tracks a running replay worker: its cancel channel and, when the worker was +/// spawned in-process, the [`JoinHandle`] used to await its exit on shutdown. +struct ReplayWorkerHandle { + cancel_tx: mpsc::Sender<()>, + join: Option>, +} + +#[derive(Default)] pub struct ReplayWorkerManager { - cancellers: HashMap>, + cancellers: HashMap, +} + +impl Debug for ReplayWorkerManager { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("ReplayWorkerManager") + .field("worker_count", &self.cancellers.len()) + .finish() + } } impl ReplayWorkerManager { @@ -48,8 +70,27 @@ impl ReplayWorkerManager { } } + /// Registers a cancel channel without a join handle. + /// + /// Used where the worker's lifetime is managed elsewhere (or in tests). Such + /// workers are signalled on `stop_all` but not awaited. Prefer + /// [`Self::insert_with_handle`] for in-process workers so shutdown can join + /// them and avoid orphaned tasks. pub fn insert(&mut self, target_id: String, cancel_tx: mpsc::Sender<()>) { - self.cancellers.insert(target_id, cancel_tx); + self.cancellers + .insert(target_id, ReplayWorkerHandle { cancel_tx, join: None }); + } + + /// Registers a cancel channel together with the worker's join handle so + /// `stop_all` can await the worker's exit (bounded by [`STOP_JOIN_TIMEOUT`]). + pub fn insert_with_handle(&mut self, target_id: String, cancel_tx: mpsc::Sender<()>, join: JoinHandle<()>) { + self.cancellers.insert( + target_id, + ReplayWorkerHandle { + cancel_tx, + join: Some(join), + }, + ); } pub fn len(&self) -> usize { @@ -67,10 +108,39 @@ impl ReplayWorkerManager { } } + /// Stops every replay worker: it first signals cancellation to all of them, + /// then awaits each worker's exit (bounded by [`STOP_JOIN_TIMEOUT`], after + /// which the task is aborted). Signalling before joining lets all workers + /// wind down concurrently, and joining guarantees no worker keeps draining + /// the shared store after this returns — preventing duplicate delivery and + /// orphaned tasks across reloads and shutdown. pub async fn stop_all(&mut self, log_prefix: &str) { - for (target_id, cancel_tx) in self.cancellers.drain() { + let mut handles: Vec<(String, ReplayWorkerHandle)> = self.cancellers.drain().collect(); + + // Phase 1: signal cancellation to all workers. + for (target_id, handle) in &handles { tracing::info!(target_id = %target_id, "{log_prefix}"); - let _ = cancel_tx.send(()).await; + let _ = handle.cancel_tx.send(()).await; + } + + // Phase 2: await each worker's exit, forcibly aborting any that overrun. + for (target_id, handle) in handles.drain(..) { + let Some(mut join) = handle.join else { + continue; + }; + match tokio::time::timeout(STOP_JOIN_TIMEOUT, &mut join).await { + Ok(Ok(())) => {} + Ok(Err(err)) => { + tracing::warn!(target_id = %target_id, error = %err, "Replay worker terminated abnormally"); + } + Err(_) => { + join.abort(); + tracing::warn!( + target_id = %target_id, + "Timed out awaiting replay worker exit; task aborted" + ); + } + } } } } @@ -232,12 +302,24 @@ where self.remove_and_close(&target_id.to_string()).await } - pub async fn clear_and_close(&mut self) { + /// Closes and removes every target, returning the id/error of each target + /// whose `close()` failed. Previously these flush/close errors were logged + /// and dropped, so a shutdown that failed to flush a target reported success. + /// Callers can now surface them (e.g. fail an explicit shutdown) while still + /// tearing down the rest of the runtime. + pub async fn clear_and_close(&mut self) -> Vec<(String, TargetError)> { let target_ids: Vec = self.targets.keys().cloned().collect(); + let mut errors = Vec::new(); for target_id in target_ids { - let _ = self.remove_and_close(&target_id).await; + if let Some(target) = self.targets.remove(&target_id) + && let Err(err) = target.close().await + { + tracing::error!(target_id = %target_id, error = %err, "Failed to close target during shutdown"); + errors.push((target_id, err)); + } } self.targets.clear(); + errors } pub fn target_ids(&self) -> Vec { @@ -317,11 +399,14 @@ pub async fn init_target_and_optionally_start_replay( target: Box + Send + Sync>, on_replay_start: F, start_replay: G, -) -> Option<(SharedTarget, Option>)> +) -> Option<(SharedTarget, Option<(mpsc::Sender<()>, JoinHandle<()>)>)> where E: PluginEvent, F: FnOnce(&str, bool), - G: FnOnce(Box + Send>, SharedTarget) -> mpsc::Sender<()>, + G: FnOnce( + Box + Send>, + SharedTarget, + ) -> (mpsc::Sender<()>, JoinHandle<()>), { let target_id = target.id().to_string(); let has_store = target.store().is_some(); @@ -350,6 +435,8 @@ where Some((shared, cancel)) } +type ActivatedTarget = (SharedTarget, Option<(mpsc::Sender<()>, JoinHandle<()>)>); + pub async fn activate_targets_with_replay( targets: Vec + Send + Sync>>, mut activate_one: F, @@ -357,16 +444,16 @@ pub async fn activate_targets_with_replay( where E: PluginEvent, F: FnMut(Box + Send + Sync>) -> Fut, - Fut: Future, Option>)>>, + Fut: Future>>, { let mut replay_workers = ReplayWorkerManager::new(); let mut shared_targets = Vec::new(); for target in targets { - if let Some((shared_target, cancel_tx)) = activate_one(target).await { + if let Some((shared_target, replay)) = activate_one(target).await { let target_id = shared_target.id().to_string(); - if let Some(cancel_tx) = cancel_tx { - replay_workers.insert(target_id, cancel_tx); + if let Some((cancel_tx, join)) = replay { + replay_workers.insert_with_handle(target_id, cancel_tx, join); } shared_targets.push(shared_target); } @@ -385,17 +472,35 @@ pub fn start_replay_worker( semaphore: Option>, batch_timeout: Duration, idle_sleep: Duration, -) -> mpsc::Sender<()> +) -> (mpsc::Sender<()>, JoinHandle<()>) where E: PluginEvent, { let (cancel_tx, cancel_rx) = mpsc::channel(1); - tokio::spawn(async move { + let join = tokio::spawn(async move { stream_replay_worker(&mut *store, target, cancel_rx, hook, semaphore, batch_timeout, idle_sleep).await; }); - cancel_tx + (cancel_tx, join) +} + +/// Number of readable entries accumulated before a replay batch is flushed under +/// a single semaphore permit. The previous `!batch_keys.is_empty()` flush +/// condition was always true, so this effectively defaulted to 1 (a permit per +/// entry) and made `batch_timeout` dead code. +const REPLAY_BATCH_SIZE: usize = 16; + +/// Sleeps for `dur` unless a cancel signal arrives first. Returns `true` if +/// cancellation was observed. Used so idle waits, inter-scan pauses, and retry +/// backoff all react promptly to shutdown instead of blocking for the full +/// duration. +async fn sleep_or_cancelled(dur: Duration, cancel_rx: &mut mpsc::Receiver<()>) -> bool { + tokio::select! { + biased; + _ = cancel_rx.recv() => true, + _ = tokio::time::sleep(dur) => false, + } } async fn stream_replay_worker( @@ -409,10 +514,7 @@ async fn stream_replay_worker( ) where E: PluginEvent, { - const MAX_RETRIES: usize = 5; - const BASE_RETRY_DELAY: Duration = Duration::from_secs(2); - - let mut batch_keys = Vec::with_capacity(1); + let mut batch_keys = Vec::with_capacity(REPLAY_BATCH_SIZE); let mut last_flush = tokio::time::Instant::now(); loop { @@ -423,17 +525,21 @@ async fn stream_replay_worker( let keys = store.list(); if keys.is_empty() { if !batch_keys.is_empty() && last_flush.elapsed() >= batch_timeout { - process_replay_batch(&mut batch_keys, target.clone(), &hook, semaphore.clone()).await; + if process_replay_batch(&mut batch_keys, target.clone(), &hook, semaphore.clone(), &mut cancel_rx).await { + return; + } last_flush = tokio::time::Instant::now(); } - tokio::time::sleep(idle_sleep).await; + if sleep_or_cancelled(idle_sleep, &mut cancel_rx).await { + return; + } continue; } for key in keys { if cancel_rx.try_recv().is_ok() { if !batch_keys.is_empty() { - process_replay_batch(&mut batch_keys, target.clone(), &hook, semaphore.clone()).await; + process_replay_batch(&mut batch_keys, target.clone(), &hook, semaphore.clone(), &mut cancel_rx).await; } return; } @@ -452,101 +558,136 @@ async fn stream_replay_worker( } } + // Skip keys already pending in the current batch: an un-flushed + // partial batch carries across scans, and `store.list()` keeps + // returning not-yet-delivered keys, so without this guard the same + // key would be enqueued repeatedly. + if batch_keys + .iter() + .any(|pending: &Key| pending.to_key_string() == key.to_key_string()) + { + continue; + } + batch_keys.push(key); - if !batch_keys.is_empty() || last_flush.elapsed() >= batch_timeout { - process_replay_batch(&mut batch_keys, target.clone(), &hook, semaphore.clone()).await; + // Flush once a full batch has accumulated or the batch has aged past + // batch_timeout — real size/time-based batching, not once-per-entry. + if batch_keys.len() >= REPLAY_BATCH_SIZE || last_flush.elapsed() >= batch_timeout { + if process_replay_batch(&mut batch_keys, target.clone(), &hook, semaphore.clone(), &mut cancel_rx).await { + return; + } last_flush = tokio::time::Instant::now(); } } - tokio::time::sleep(Duration::from_millis(100)).await; - } - - async fn process_replay_batch( - batch_keys: &mut Vec, - target: SharedTarget, - hook: &ReplayHook, - semaphore: Option>, - ) where - E: PluginEvent, - { - if batch_keys.is_empty() { + if sleep_or_cancelled(Duration::from_millis(100), &mut cancel_rx).await { return; } + } +} - let _permit = match semaphore { - Some(ref semaphore) => match semaphore.clone().acquire_owned().await { - Ok(permit) => Some(permit), - Err(err) => { - tracing::error!(error = %err, "Failed to acquire replay semaphore permit"); - return; +/// Delivers a batch of queued entries under a single semaphore permit. +/// +/// Returns `true` if a cancel signal was observed while processing (e.g. during +/// retry backoff), so the caller can stop promptly instead of continuing to +/// drain a store that a replacement worker may already own. +async fn process_replay_batch( + batch_keys: &mut Vec, + target: SharedTarget, + hook: &ReplayHook, + semaphore: Option>, + cancel_rx: &mut mpsc::Receiver<()>, +) -> bool +where + E: PluginEvent, +{ + const MAX_RETRIES: usize = 5; + const BASE_RETRY_DELAY: Duration = Duration::from_secs(2); + + if batch_keys.is_empty() { + return false; + } + + let _permit = match semaphore { + Some(ref semaphore) => match semaphore.clone().acquire_owned().await { + Ok(permit) => Some(permit), + Err(err) => { + tracing::error!(error = %err, "Failed to acquire replay semaphore permit"); + batch_keys.clear(); + return false; + } + }, + None => None, + }; + + let mut cancelled = false; + 'keys: for key in batch_keys.iter() { + let mut retry_count = 0usize; + let mut success = false; + + while retry_count < MAX_RETRIES && !success { + match target.send_from_store(key.clone()).await { + Ok(_) => { + hook(ReplayEvent::Delivered { + key: key.clone(), + target: target.clone(), + }) + .await; + success = true; } - }, - None => None, - }; + Err(err) => match err { + TargetError::NotConnected | TargetError::Timeout(_) => { + retry_count += 1; + hook(ReplayEvent::RetryableError { + error: err, + key: key.clone(), + retry_count, + target: target.clone(), + }) + .await; - for key in batch_keys.iter() { - let mut retry_count = 0usize; - let mut success = false; - - while retry_count < MAX_RETRIES && !success { - match target.send_from_store(key.clone()).await { - Ok(_) => { - hook(ReplayEvent::Delivered { + let jitter = Duration::from_millis(key.to_string().len() as u64 % 500); + let backoff = 1u32 << retry_count as u32; + // Observe cancellation during backoff so shutdown/reload is + // not blocked for the full (potentially many-second) delay. + if sleep_or_cancelled(BASE_RETRY_DELAY * backoff + jitter, cancel_rx).await { + cancelled = true; + break 'keys; + } + } + TargetError::Dropped(reason) => { + hook(ReplayEvent::Dropped { + key: key.clone(), + reason, + target: target.clone(), + }) + .await; + break; + } + other => { + hook(ReplayEvent::PermanentFailure { + error: other, key: key.clone(), target: target.clone(), }) .await; - success = true; + break; } - Err(err) => match err { - TargetError::NotConnected | TargetError::Timeout(_) => { - retry_count += 1; - hook(ReplayEvent::RetryableError { - error: err, - key: key.clone(), - retry_count, - target: target.clone(), - }) - .await; - - let jitter = Duration::from_millis(key.to_string().len() as u64 % 500); - let backoff = 1u32 << retry_count as u32; - tokio::time::sleep(BASE_RETRY_DELAY * backoff + jitter).await; - } - TargetError::Dropped(reason) => { - hook(ReplayEvent::Dropped { - key: key.clone(), - reason, - target: target.clone(), - }) - .await; - break; - } - other => { - hook(ReplayEvent::PermanentFailure { - error: other, - key: key.clone(), - target: target.clone(), - }) - .await; - break; - } - }, - } - } - - if retry_count >= MAX_RETRIES && !success { - hook(ReplayEvent::RetryExhausted { - key: key.clone(), - target: target.clone(), - }) - .await; + }, } } - batch_keys.clear(); + if retry_count >= MAX_RETRIES && !success { + hook(ReplayEvent::RetryExhausted { + key: key.clone(), + target: target.clone(), + }) + .await; + } } + + batch_keys.clear(); + cancelled } #[cfg(test)] @@ -641,4 +782,57 @@ mod tests { assert_eq!(snapshots[0].target_id, "primary:webhook"); assert_eq!(snapshots[0].target_type, "webhook"); } + + #[tokio::test] + async fn sleep_or_cancelled_returns_immediately_on_cancel() { + let (cancel_tx, mut cancel_rx) = tokio::sync::mpsc::channel::<()>(1); + cancel_tx.send(()).await.unwrap(); + + // A pending cancel signal must short-circuit a long sleep. + let start = std::time::Instant::now(); + let cancelled = super::sleep_or_cancelled(std::time::Duration::from_secs(30), &mut cancel_rx).await; + assert!(cancelled); + assert!( + start.elapsed() < std::time::Duration::from_secs(5), + "cancel should not wait for the full sleep" + ); + } + + #[tokio::test] + async fn sleep_or_cancelled_returns_false_when_not_cancelled() { + let (_cancel_tx, mut cancel_rx) = tokio::sync::mpsc::channel::<()>(1); + let cancelled = super::sleep_or_cancelled(std::time::Duration::from_millis(10), &mut cancel_rx).await; + assert!(!cancelled); + } + + #[tokio::test] + async fn stop_all_joins_and_awaits_worker_exit() { + use super::ReplayWorkerManager; + use std::sync::atomic::AtomicBool; + + let mut manager = ReplayWorkerManager::new(); + let exited = Arc::new(AtomicBool::new(false)); + + let (cancel_tx, mut cancel_rx) = tokio::sync::mpsc::channel::<()>(1); + let exited_task = Arc::clone(&exited); + let join = tokio::spawn(async move { + // Run until cancelled, then record clean exit. + loop { + if super::sleep_or_cancelled(std::time::Duration::from_millis(50), &mut cancel_rx).await { + break; + } + } + exited_task.store(true, Ordering::SeqCst); + }); + + manager.insert_with_handle("primary:webhook".to_string(), cancel_tx, join); + assert_eq!(manager.len(), 1); + + // stop_all must signal AND await the worker: once it returns, the worker + // has actually exited (no orphaned task). + manager.stop_all("stopping test worker").await; + + assert!(manager.is_empty()); + assert!(exited.load(Ordering::SeqCst), "stop_all must await the worker to completion"); + } } diff --git a/crates/targets/src/store.rs b/crates/targets/src/store.rs index 5fd58f48e..e1e42353c 100644 --- a/crates/targets/src/store.rs +++ b/crates/targets/src/store.rs @@ -19,8 +19,10 @@ use serde::{Serialize, de::DeserializeOwned}; use snap::raw::{Decoder, Encoder}; use std::{ collections::HashMap, + fs::File, + io::Write, marker::PhantomData, - path::PathBuf, + path::{Path, PathBuf}, sync::{ Arc, RwLock, atomic::{AtomicU64, Ordering}, @@ -34,6 +36,31 @@ const LOG_COMPONENT_TARGETS: &str = "targets"; const LOG_SUBSYSTEM_STORE: &str = "store"; const EVENT_TARGET_STORE_STATE: &str = "target_store_state"; +/// Suffix used for the temporary file of an in-progress atomic write. A crash +/// between `File::create` and `rename` leaves one of these behind; `open()` +/// removes them so they are never mistaken for committed queue entries. +const TMP_SUFFIX: &str = ".tmp"; + +/// Upper bound applied to the initial `HashMap`/`Vec` capacities derived from +/// untrusted inputs (`entry_limit`, batch `item_count`). Growth is still lazy, +/// so a huge configured limit or a malicious/corrupt filename can no longer +/// trigger a giant up-front allocation or a capacity-overflow panic. +const MAX_PREALLOC_CAPACITY: usize = 4096; + +/// Returns true if `file_name` looks like a committed queue entry for `file_ext`. +/// +/// A committed entry is `<...>` optionally followed by [`COMPRESS_EXT`] +/// (e.g. `.event` or `3:.event.snappy`). Any other file in the queue +/// directory (foreign files, leftover temp files) is ignored so it can never be +/// indexed and replayed as if it were a real event. +fn is_queue_file_name(file_name: &str, file_ext: &str) -> bool { + if file_ext.is_empty() { + return false; + } + let base = file_name.strip_suffix(COMPRESS_EXT).unwrap_or(file_name); + base.ends_with(file_ext) +} + fn resolve_queue_store_compression_from_env_value(value: Option<&str>) -> bool { value .and_then(|value| value.parse::().ok().map(|state| state.is_enabled())) @@ -267,7 +294,7 @@ impl QueueStore { entry_limit, file_ext: file_ext.to_string(), compress, - entries: Arc::new(RwLock::new(HashMap::with_capacity(entry_limit as usize))), + entries: Arc::new(RwLock::new(HashMap::with_capacity((entry_limit as usize).min(MAX_PREALLOC_CAPACITY)))), pending_entries: Arc::new(AtomicU64::new(0)), fs_guard: Arc::new(RwLock::new(())), _phantom: PhantomData, @@ -281,13 +308,51 @@ impl QueueStore { fn build_key(&self, item_count: usize) -> Key { Key { - name: Uuid::new_v4().to_string(), + // UUIDv7 is time-ordered: sorting entries by name reproduces FIFO + // enqueue order deterministically and, crucially, identically after + // a restart — the ordering is intrinsic to the persisted filename + // rather than derived from coarse, clock-dependent file mtimes. + name: Uuid::now_v7().to_string(), extension: self.file_ext.clone(), item_count, compress: self.compress, } } + /// Best-effort `fsync` of the queue directory so a freshly `rename`d entry + /// (and its containing directory entry) survives a power loss. Failure is + /// logged at debug and not propagated: the data file itself is already + /// durably `fsync`ed, and not every filesystem/platform supports directory + /// fsync. + fn fsync_dir(dir: &Path) { + match File::open(dir) { + Ok(dir_file) => { + if let Err(err) = dir_file.sync_all() { + debug!( + event = EVENT_TARGET_STORE_STATE, + component = LOG_COMPONENT_TARGETS, + subsystem = LOG_SUBSYSTEM_STORE, + action = "fsync_dir", + dir = %dir.display(), + error = %err, + "target store state" + ); + } + } + Err(err) => { + debug!( + event = EVENT_TARGET_STORE_STATE, + component = LOG_COMPONENT_TARGETS, + subsystem = LOG_SUBSYSTEM_STORE, + action = "fsync_dir_open", + dir = %dir.display(), + error = %err, + "target store state" + ); + } + } + } + /// Reads a file for the given key fn read_file(&self, key: &Key) -> Result, StoreError> { let _fs_guard = self @@ -351,7 +416,15 @@ impl QueueStore { } } - /// Writes data to a file for the given key. + /// Durably and atomically writes data to the file for the given key. + /// + /// The write is crash-safe: bytes are written to a per-key temporary file, + /// `fsync`ed (`sync_all`), and only then `rename`d onto the final path. A + /// same-directory `rename` is atomic, so a reader (including `open()` after a + /// restart) observes either the complete previous file or the complete new + /// one — never a half-written payload. On a crash mid-write the temp file is + /// left behind and cleaned up by `open()`, so an acknowledged event is never + /// lost and no ghost/truncated entry is ever indexed. fn write_file(&self, key: &Key, data: &[u8]) -> Result { let path = self.file_path(key); // Create directory if it doesn't exist @@ -359,15 +432,46 @@ impl QueueStore { std::fs::create_dir_all(parent).map_err(StoreError::Io)?; } - if key.compress { + let payload: std::borrow::Cow<'_, [u8]> = if key.compress { let mut encoder = Encoder::new(); let compressed = encoder .compress_vec(data) .map_err(|e| StoreError::Compression(e.to_string()))?; - std::fs::write(&path, &compressed).map_err(StoreError::Io)?; + std::borrow::Cow::Owned(compressed) } else { - std::fs::write(&path, data).map_err(StoreError::Io)?; + std::borrow::Cow::Borrowed(data) + }; + + let tmp_path = { + let mut file_name = key.to_key_string(); + file_name.push_str(TMP_SUFFIX); + self.directory.join(file_name) + }; + + // Write + fsync the full payload into the temp file, then atomically + // rename it into place. Any error path removes the temp file so a failed + // write cannot leave residue that would otherwise be cleaned only on the + // next open(). + let write_result = (|| -> Result<(), StoreError> { + let mut file = File::create(&tmp_path).map_err(StoreError::Io)?; + file.write_all(&payload).map_err(StoreError::Io)?; + file.sync_all().map_err(StoreError::Io)?; + Ok(()) + })(); + + if let Err(err) = write_result { + let _ = std::fs::remove_file(&tmp_path); + return Err(err); } + + if let Err(err) = std::fs::rename(&tmp_path, &path) { + let _ = std::fs::remove_file(&tmp_path); + return Err(StoreError::Io(err)); + } + + // Best-effort: persist the new directory entry created by rename. + Self::fsync_dir(&self.directory); + let modified = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_nanos() as i64; debug!( event = EVENT_TARGET_STORE_STATE, @@ -434,13 +538,36 @@ where for entry in dir_entries { let entry = entry.map_err(StoreError::Io)?; let metadata = entry.metadata().map_err(StoreError::Io)?; - if metadata.is_file() { - let modified = metadata.modified().map_err(StoreError::Io)?; - let unix_nano = modified.duration_since(UNIX_EPOCH).unwrap_or_default().as_nanos() as i64; - - let file_name = entry.file_name().to_string_lossy().to_string(); - entries_map.insert(file_name, unix_nano); + if !metadata.is_file() { + continue; } + + let file_name = entry.file_name().to_string_lossy().to_string(); + + // Remove leftover temp files from an interrupted atomic write; they + // are never valid committed entries. + if file_name.ends_with(TMP_SUFFIX) { + let _ = std::fs::remove_file(entry.path()); + continue; + } + + // Ignore foreign files that do not match our queue file extension so + // externally-dropped files cannot pollute the queue. + if !is_queue_file_name(&file_name, &self.file_ext) { + continue; + } + + // Drop zero-byte files (truncated/empty writes from older code paths + // or crashes) instead of indexing a ghost entry that read_file would + // report as NotFound forever. + if metadata.len() == 0 { + let _ = std::fs::remove_file(entry.path()); + continue; + } + + let modified = metadata.modified().map_err(StoreError::Io)?; + let unix_nano = modified.duration_since(UNIX_EPOCH).unwrap_or_default().as_nanos() as i64; + entries_map.insert(file_name, unix_nano); } debug!( @@ -525,7 +652,10 @@ where if data.is_empty() { return Err(StoreError::Deserialization("Cannot deserialize empty data".to_string())); } - let mut items = Vec::with_capacity(key.item_count); + // `item_count` is parsed from the (untrusted) filename; clamp the + // up-front allocation so a corrupt/malicious name cannot drive a huge + // reservation. The loop below still reads exactly `item_count` items. + let mut items = Vec::with_capacity(key.item_count.min(MAX_PREALLOC_CAPACITY)); // let mut deserializer = serde_json::Deserializer::from_slice(&data); // while let Ok(item) = serde::Deserialize::deserialize(&mut deserializer) { @@ -669,11 +799,18 @@ where } }; - let mut entries_vec: Vec<_> = entries.iter().collect(); - // Sort by modtime (value in HashMap) to process oldest first - entries_vec.sort_by(|a, b| a.1.cmp(b.1)); // Oldest first + // Order by the entry name (UUIDv7), which is time-ordered by + // construction. This yields a stable FIFO order that is identical in a + // single run and after a restart, because it derives purely from the + // persisted filename rather than from coarse, clock-dependent file + // mtimes. The full filename is used as a deterministic tie-breaker. + let mut entries_vec: Vec<(String, String)> = entries + .keys() + .map(|file_name| (parse_key(file_name).name, file_name.clone())) + .collect(); + entries_vec.sort_by(|a, b| a.0.cmp(&b.0).then_with(|| a.1.cmp(&b.1))); - entries_vec.into_iter().map(|(k, _)| parse_key(k)).collect() + entries_vec.into_iter().map(|(_, file_name)| parse_key(&file_name)).collect() } fn len(&self) -> usize { @@ -869,4 +1006,126 @@ mod tests { let _ = store.delete(); } + + #[test] + fn open_cleans_leftover_tmp_and_zero_byte_files() { + let dir = temp_store_dir("open-cleanup"); + let store = QueueStore::::new_with_compression(&dir, 8, ".test", false); + store.open().unwrap(); + + // A committed entry that must survive open(). + let key = store.put(Arc::new("payload".to_string())).unwrap(); + let good_name = key.to_key_string(); + + // Simulate an interrupted atomic write: a temp file that was fsynced but + // never renamed into place. It must be removed and never indexed. + let tmp_path = dir.join(format!("{good_name}{TMP_SUFFIX}")); + std::fs::write(&tmp_path, b"half-written").unwrap(); + + // Simulate a zero-byte ghost file with a valid queue extension. + let zero_name = format!("{}.test", Uuid::now_v7()); + let zero_path = dir.join(&zero_name); + std::fs::write(&zero_path, b"").unwrap(); + + store.open().unwrap(); + + // The queue is clean: only the committed entry remains indexed, and both + // the temp and zero-byte residues are gone from disk. + assert_eq!(store.len(), 1); + let listed: Vec = store.list().iter().map(|k| k.to_key_string()).collect(); + assert_eq!(listed, vec![good_name]); + assert!(!tmp_path.exists(), "leftover temp file should be removed"); + assert!(!zero_path.exists(), "zero-byte file should be removed"); + + let _ = store.delete(); + } + + #[test] + fn open_ignores_foreign_extension_files() { + let dir = temp_store_dir("open-foreign"); + let store = QueueStore::::new_with_compression(&dir, 8, ".test", false); + store.open().unwrap(); + let key = store.put(Arc::new("payload".to_string())).unwrap(); + + // A non-empty file with an unrelated extension must not be indexed and + // must be left untouched (we only clean up our own residue). + let foreign_path = dir.join("intruder.txt"); + std::fs::write(&foreign_path, b"not ours").unwrap(); + + store.open().unwrap(); + + assert_eq!(store.len(), 1); + let listed: Vec = store.list().iter().map(|k| k.to_key_string()).collect(); + assert_eq!(listed, vec![key.to_key_string()]); + assert!(foreign_path.exists(), "foreign file should be left in place, just not indexed"); + + let _ = store.delete(); + } + + #[test] + fn list_order_is_stable_across_reopen() { + let dir = temp_store_dir("list-order"); + let store = QueueStore::::new_with_compression(&dir, 32, ".test", false); + store.open().unwrap(); + + for idx in 0..8 { + store.put(Arc::new(format!("event-{idx}"))).unwrap(); + } + + let order_before: Vec = store.list().iter().map(|k| k.to_key_string()).collect(); + assert_eq!(order_before.len(), 8); + + // Re-open from disk (simulating a restart) and assert the FIFO order is + // reproduced exactly — it is derived from the persisted UUIDv7 names, not + // from coarse file mtimes. + let reopened = QueueStore::::new_with_compression(&dir, 32, ".test", false); + reopened.open().unwrap(); + let order_after: Vec = reopened.list().iter().map(|k| k.to_key_string()).collect(); + + assert_eq!(order_before, order_after); + + let _ = store.delete(); + } + + #[test] + fn new_with_huge_limit_does_not_panic() { + // Previously `HashMap::with_capacity(entry_limit as usize)` would attempt + // a giant allocation / capacity overflow for an absurd configured limit. + let dir = temp_store_dir("huge-limit"); + let store = QueueStore::::new_with_compression(&dir, u64::MAX, ".test", false); + store.open().unwrap(); + assert!(store.is_empty()); + let _ = store.delete(); + } + + #[test] + fn get_multiple_does_not_overallocate_on_huge_item_count() { + let dir = temp_store_dir("huge-item-count"); + let store = QueueStore::::new_with_compression(&dir, 8, ".test", false); + store.open().unwrap(); + std::fs::create_dir_all(&dir).unwrap(); + + // Craft a batch file whose name claims a colossal item_count (far beyond + // MAX_PREALLOC_CAPACITY) but whose body holds only two items — the shape a + // corrupt/malicious filename would take. The read must fail cleanly with a + // Deserialization error and must not attempt a giant pre-allocation + // (Vec::with_capacity is clamped to MAX_PREALLOC_CAPACITY). + let claimed_count = 1_000_000usize; + let file_name = format!("{claimed_count}:{}.test", Uuid::now_v7()); + let mut body = Vec::new(); + body.extend_from_slice(&serde_json::to_vec(&"aa".to_string()).unwrap()); + body.extend_from_slice(&serde_json::to_vec(&"bb".to_string()).unwrap()); + std::fs::write(dir.join(&file_name), &body).unwrap(); + + let key = parse_key(&file_name); + assert_eq!(key.item_count, claimed_count); + + let err = store.get_multiple(&key).unwrap_err(); + assert!( + matches!(err, StoreError::Deserialization(_)), + "expected Deserialization error, got {err:?}" + ); + + let _ = store.delete(); + } } diff --git a/crates/targets/src/target/mod.rs b/crates/targets/src/target/mod.rs index 60975b877..e34bd7125 100644 --- a/crates/targets/src/target/mod.rs +++ b/crates/targets/src/target/mod.rs @@ -123,7 +123,15 @@ where let raw = match store.get_raw(&key) { Ok(raw) => raw, - Err(StoreError::NotFound) => return Ok(()), + Err(StoreError::NotFound) => { + // The backing file is missing or empty (a zero-byte file reads as + // NotFound). Left in the index it would be "replayed" forever and + // permanently occupy a queue slot, eventually rejecting new events + // with LimitExceeded. Purge the stale index entry (and any residual + // file) before returning. + delete_stored_payload(store, &key)?; + return Ok(()); + } Err(err) => return Err(TargetError::Storage(format!("Failed to read queued payload from store: {err}"))), }; @@ -278,10 +286,22 @@ impl QueuedPayload { return Err(TargetError::Serialization("Queued payload metadata length exceeds input".to_string())); } - let meta = serde_json::from_slice(&raw[meta_start..meta_end]) + let meta: QueuedPayloadMeta = serde_json::from_slice(&raw[meta_start..meta_end]) .map_err(|err| TargetError::Serialization(format!("Failed to deserialize queued payload metadata: {err}")))?; let body = raw[meta_end..].to_vec(); + // Reject torn/truncated writes: the body length recorded at encode time + // must match the bytes actually present. Without this, a partially + // written file (e.g. a crash mid-write) would decode into a silently + // truncated payload and be delivered as if complete. + if body.len() != meta.payload_len { + return Err(TargetError::Serialization(format!( + "Queued payload body length mismatch: header declares {} bytes but {} were present", + meta.payload_len, + body.len() + ))); + } + Ok(Self { meta, body }) } } @@ -382,17 +402,56 @@ impl std::fmt::Display for TargetType { } } +/// Stable, deterministic 64-bit FNV-1a hash used only to disambiguate queue +/// directory names. It must stay identical across restarts and releases so a +/// target keeps resolving to the same on-disk queue directory, hence a fixed +/// inline implementation rather than `DefaultHasher` (whose algorithm is not +/// contractually stable). +fn fnv1a_hash(bytes: &[u8]) -> u64 { + const FNV_OFFSET: u64 = 0xcbf2_9ce4_8422_2325; + const FNV_PRIME: u64 = 0x0000_0100_0000_01b3; + let mut hash = FNV_OFFSET; + for &byte in bytes { + hash ^= byte as u64; + hash = hash.wrapping_mul(FNV_PRIME); + } + hash +} + +/// Maps a target-id component to a filesystem-safe queue directory name. +/// +/// Path-unsafe characters are replaced with `_`. Because that replacement is +/// lossy, two distinct ids (e.g. `a/b` and `a_b`) could otherwise collapse to +/// the same directory and interleave their persisted events. To prevent that, +/// whenever any character had to be replaced we append a short hash of the +/// original component, guaranteeing distinct ids map to distinct directories. +/// +/// Ids that are already path-safe are returned unchanged, preserving the +/// on-disk directory layout for existing deployments (no queue migration). pub(crate) fn sanitize_queue_dir_component(component: &str) -> String { let mut sanitized = String::with_capacity(component.len()); + let mut lossy = false; for ch in component.chars() { if ch.is_ascii_alphanumeric() || matches!(ch, '-' | '_' | '.') { sanitized.push(ch); } else { sanitized.push('_'); + lossy = true; } } - if sanitized.is_empty() { "_".to_string() } else { sanitized } + if sanitized.is_empty() { + // An entirely non-safe id would otherwise all collapse to "_"; key it by + // the original bytes so distinct ids stay distinct. + return format!("_{:016x}", fnv1a_hash(component.as_bytes())); + } + + if lossy { + // Disambiguate the lossy replacement so different originals cannot alias. + return format!("{sanitized}-{:016x}", fnv1a_hash(component.as_bytes())); + } + + sanitized } pub(crate) fn queue_store_subdir_name(target_type: &str, target_id: &str) -> String { @@ -686,14 +745,15 @@ mod tests { #[test] fn queued_payload_round_trips_meta_and_body() { + let body = br#"{"ok":true}"#.to_vec(); let meta = QueuedPayloadMeta::new( EventName::ObjectCreatedPut, "bucket-a".to_string(), "folder/object.txt".to_string(), "application/json", - 12, + body.len(), ); - let payload = QueuedPayload::new(meta.clone(), br#"{"ok":true}"#.to_vec()); + let payload = QueuedPayload::new(meta.clone(), body); let encoded = payload.encode().unwrap(); let decoded = QueuedPayload::decode(&encoded).unwrap(); @@ -863,12 +923,137 @@ mod tests { #[test] fn sanitize_queue_dir_component_replaces_non_path_safe_characters() { let sanitized = sanitize_queue_dir_component("tenant:alpha/beta\\gamma?*"); - assert_eq!(sanitized, "tenant_alpha_beta_gamma__"); + // The readable, path-safe prefix is preserved, followed by a disambiguating + // hash suffix because the replacement was lossy. + assert!( + sanitized.starts_with("tenant_alpha_beta_gamma__-"), + "unexpected sanitized value: {sanitized}" + ); + // Deterministic across calls (must be stable across restarts). + assert_eq!(sanitized, sanitize_queue_dir_component("tenant:alpha/beta\\gamma?*")); + } + + #[test] + fn sanitize_queue_dir_component_preserves_path_safe_ids() { + // Path-safe ids are returned unchanged so existing on-disk queue + // directories keep resolving (no migration on upgrade). + assert_eq!(sanitize_queue_dir_component("plain-id_1.2"), "plain-id_1.2"); + } + + #[test] + fn sanitize_queue_dir_component_disambiguates_colliding_ids() { + // Two distinct ids that used to collapse onto the same directory must now + // map to different directories. + let a = sanitize_queue_dir_component("a/b"); + let b = sanitize_queue_dir_component("a_b"); + assert_ne!(a, b, "distinct ids must not share a queue directory"); } #[test] fn queue_store_subdir_name_sanitizes_target_id() { let dir = queue_store_subdir_name("redis", "tenant:alpha"); - assert_eq!(dir, "rustfs-redis-tenant_alpha"); + assert!(dir.starts_with("rustfs-redis-tenant_alpha-"), "unexpected subdir: {dir}"); + } + + #[derive(Clone)] + struct StoreBackedTarget { + id: TargetID, + store: QueueStore, + } + + #[async_trait] + impl Target for StoreBackedTarget { + fn id(&self) -> TargetID { + self.id.clone() + } + + async fn is_active(&self) -> Result { + Ok(true) + } + + async fn save(&self, _event: Arc>) -> Result<(), TargetError> { + Ok(()) + } + + async fn send_raw_from_store(&self, _key: Key, _body: Vec, _meta: QueuedPayloadMeta) -> Result<(), TargetError> { + Ok(()) + } + + async fn close(&self) -> Result<(), TargetError> { + Ok(()) + } + + fn store(&self) -> Option<&(dyn Store + Send + Sync)> { + Some(&self.store) + } + + fn clone_dyn(&self) -> Box + Send + Sync> { + Box::new(self.clone()) + } + + fn is_enabled(&self) -> bool { + true + } + } + + #[tokio::test] + async fn send_from_store_purges_missing_or_empty_entry() { + let dir = std::env::temp_dir().join(format!("rustfs-send-from-store-{}", Uuid::new_v4())); + let store = QueueStore::::new_with_compression(&dir, 8, ".event", false); + store.open().unwrap(); + + // Enqueue a valid payload, then truncate its backing file to zero bytes to + // simulate a torn write: read_file now reports NotFound while the index + // still counts the entry. + let meta = QueuedPayloadMeta::new( + EventName::ObjectCreatedPut, + "bucket-a".to_string(), + "obj.txt".to_string(), + "application/json", + 7, + ); + let encoded = QueuedPayload::new(meta, br#"{"x":1}"#.to_vec()).encode().unwrap(); + let key = store.put_raw(&encoded).unwrap(); + assert_eq!(store.len(), 1); + + let event_file = std::fs::read_dir(&dir) + .unwrap() + .filter_map(|e| e.ok()) + .map(|e| e.path()) + .find(|p| p.is_file()) + .expect("event file should exist"); + std::fs::write(&event_file, b"").unwrap(); + + let target = StoreBackedTarget { + id: TargetID::new("primary".to_string(), "webhook".to_string()), + store: store.clone(), + }; + + // A NotFound/empty entry must be purged (index + file) rather than + // silently skipped and replayed forever. + target.send_from_store(key).await.unwrap(); + assert_eq!(store.len(), 0, "stale entry must be removed from the index"); + + let _ = store.delete(); + } + + #[test] + fn queued_payload_decode_rejects_body_length_mismatch() { + let meta = QueuedPayloadMeta::new( + EventName::ObjectCreatedPut, + "bucket-a".to_string(), + "obj.txt".to_string(), + "application/json", + 11, + ); + let payload = QueuedPayload::new(meta, br#"{"ok":true}"#.to_vec()); + let mut encoded = payload.encode().unwrap(); + + // Drop the final body byte, simulating a torn/truncated write. The header + // still declares the original payload_len, so decode must reject it rather + // than hand back a silently truncated body. + encoded.pop(); + let err = QueuedPayload::decode(&encoded).unwrap_err(); + assert!(err.to_string().contains("body length mismatch"), "unexpected error: {err}"); } }