// Copyright 2024 RustFS Team // // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. // You may obtain a copy of the License at // // http://www.apache.org/licenses/LICENSE-2.0 // // Unless required by applicable law or agreed to in writing, software // distributed under the License is distributed on an "AS IS" BASIS, // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. // See the License for the specific language governing permissions and // limitations under the License. use crate::error::StoreError; use rustfs_config::notify::{COMPRESS_EXT, DEFAULT_EXT}; use rustfs_config::{DEFAULT_LIMIT, DEFAULT_TARGET_STORE_COMPRESS, ENV_TARGET_STORE_COMPRESS, EnableState}; use serde::{Serialize, de::DeserializeOwned}; use snap::raw::{Decoder, Encoder}; use std::{ collections::HashMap, marker::PhantomData, path::PathBuf, sync::{ Arc, RwLock, atomic::{AtomicU64, Ordering}, }, time::{SystemTime, UNIX_EPOCH}, }; use tracing::{debug, warn}; use uuid::Uuid; const LOG_COMPONENT_TARGETS: &str = "targets"; const LOG_SUBSYSTEM_STORE: &str = "store"; const EVENT_TARGET_STORE_STATE: &str = "target_store_state"; fn resolve_queue_store_compression_from_env_value(value: Option<&str>) -> bool { value .and_then(|value| value.parse::().ok().map(|state| state.is_enabled())) .unwrap_or(DEFAULT_TARGET_STORE_COMPRESS) } fn queue_store_compression_enabled() -> bool { let value = std::env::var(ENV_TARGET_STORE_COMPRESS).ok(); resolve_queue_store_compression_from_env_value(value.as_deref()) } /// Represents a key for an entry in the store #[derive(Debug, Clone)] pub struct Key { /// The name of the key (UUID) pub name: String, /// The file extension for the entry pub extension: String, /// The number of items in the entry (for batch storage) pub item_count: usize, /// Whether the entry is compressed pub compress: bool, } impl Key { /// Converts the key to a string (filename) pub fn to_key_string(&self) -> String { let name_part = if self.item_count > 1 { format!("{}:{}", self.item_count, self.name) } else { self.name.clone() }; let mut file_name = name_part; if !self.extension.is_empty() { file_name.push_str(&self.extension); } if self.compress { file_name.push_str(COMPRESS_EXT); } file_name } } impl std::fmt::Display for Key { fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { f.write_str(&self.to_key_string()) } } /// Parses a string into a Key pub fn parse_key(s: &str) -> Key { debug!( event = EVENT_TARGET_STORE_STATE, component = LOG_COMPONENT_TARGETS, subsystem = LOG_SUBSYSTEM_STORE, action = "parse_key", key = %s, "target store state" ); let mut name = s.to_string(); let mut extension = String::new(); let mut item_count = 1; let mut compress = false; // Check for compressed suffixes if name.ends_with(COMPRESS_EXT) { compress = true; name = name[..name.len() - COMPRESS_EXT.len()].to_string(); } // Number of batch items parsed if let Some(colon_pos) = name.find(':') && let Ok(count) = name[..colon_pos].parse::() { item_count = count; name = name[colon_pos + 1..].to_string(); } // Resolve extension if let Some(dot_pos) = name.rfind('.') { extension = name[dot_pos..].to_string(); name = name[..dot_pos].to_string(); } debug!( event = EVENT_TARGET_STORE_STATE, component = LOG_COMPONENT_TARGETS, subsystem = LOG_SUBSYSTEM_STORE, action = "parse_key", result = "parsed", key_name = %name, extension = %extension, item_count, compressed = compress, "target store state" ); Key { name, extension, item_count, compress, } } pub fn ensure_store_entry_raw_readable( store: &(dyn Store + Send), key: &Key, ) -> Result where T: Send + Sync + 'static + Clone + Serialize, { match store.get_raw(key) { Ok(_) => Ok(true), Err(StoreError::NotFound) => Ok(false), Err(err) => { match store.del(key) { Ok(()) | Err(StoreError::NotFound) => {} Err(del_err) => { return Err(StoreError::Internal(format!("Failed to remove unreadable store entry {key}: {del_err}"))); } } Err(err) } } } /// Trait for a store that can store and retrieve items of type T pub trait Store: Send + Sync where T: Send + Sync + 'static + Clone + Serialize, { /// The error type for the store type Error; /// The key type for the store type Key; /// Opens the store fn open(&self) -> Result<(), Self::Error>; /// Stores a single item fn put(&self, item: Arc) -> Result; /// Stores multiple items in a single batch fn put_multiple(&self, items: Vec) -> Result; /// Stores raw bytes in a single entry. fn put_raw(&self, data: &[u8]) -> Result; /// Retrieves a single item by key fn get(&self, key: &Self::Key) -> Result; /// Retrieves multiple items by key fn get_multiple(&self, key: &Self::Key) -> Result, Self::Error>; /// Retrieves the raw bytes stored for a key. fn get_raw(&self, key: &Self::Key) -> Result, Self::Error>; /// Deletes an item by key fn del(&self, key: &Self::Key) -> Result<(), Self::Error>; /// Deletes the underlying store directory and clears all in-memory state. fn delete(&self) -> Result<(), Self::Error>; /// Lists all keys in the store fn list(&self) -> Vec; /// Returns the number of items in the store fn len(&self) -> usize; /// Returns true if the store is empty fn is_empty(&self) -> bool; /// Clones the store into a boxed trait object fn boxed_clone(&self) -> Box + Send + Sync>; } /// A store that uses the filesystem to persist events in a queue pub struct QueueStore { entry_limit: u64, directory: PathBuf, file_ext: String, compress: bool, entries: Arc>>, // key -> modtime as unix nano pending_entries: Arc, fs_guard: Arc>, _phantom: PhantomData, } impl Clone for QueueStore { fn clone(&self) -> Self { QueueStore { entry_limit: self.entry_limit, directory: self.directory.clone(), file_ext: self.file_ext.clone(), compress: self.compress, entries: Arc::clone(&self.entries), pending_entries: Arc::clone(&self.pending_entries), fs_guard: Arc::clone(&self.fs_guard), _phantom: PhantomData, } } } struct EntryReservation<'a> { pending_entries: &'a AtomicU64, } impl Drop for EntryReservation<'_> { fn drop(&mut self) { self.pending_entries.fetch_sub(1, Ordering::SeqCst); } } impl QueueStore { /// Creates a new QueueStore pub fn new(directory: impl Into, limit: u64, ext: &str) -> Self { Self::new_with_compression(directory, limit, ext, queue_store_compression_enabled()) } /// Creates a new QueueStore with an explicit compression setting. pub fn new_with_compression(directory: impl Into, limit: u64, ext: &str, compress: bool) -> Self { let file_ext = if ext.is_empty() { DEFAULT_EXT } else { ext }; let entry_limit = if limit == 0 { DEFAULT_LIMIT } else { limit }; QueueStore { directory: directory.into(), entry_limit, file_ext: file_ext.to_string(), compress, entries: Arc::new(RwLock::new(HashMap::with_capacity(entry_limit as usize))), pending_entries: Arc::new(AtomicU64::new(0)), fs_guard: Arc::new(RwLock::new(())), _phantom: PhantomData, } } /// Returns the full path for a key fn file_path(&self, key: &Key) -> PathBuf { self.directory.join(key.to_key_string()) } fn build_key(&self, item_count: usize) -> Key { Key { name: Uuid::new_v4().to_string(), extension: self.file_ext.clone(), item_count, compress: self.compress, } } /// Reads a file for the given key fn read_file(&self, key: &Key) -> Result, StoreError> { let _fs_guard = self .fs_guard .read() .map_err(|_| StoreError::Internal("Failed to acquire read lock on store filesystem".to_string()))?; let path = self.file_path(key); debug!( event = EVENT_TARGET_STORE_STATE, component = LOG_COMPONENT_TARGETS, subsystem = LOG_SUBSYSTEM_STORE, action = "read_file", key = %key, path = %path.display(), "target store state" ); let data = std::fs::read(&path).map_err(|e| { if e.kind() == std::io::ErrorKind::NotFound { StoreError::NotFound } else { StoreError::Io(e) } })?; if data.is_empty() { return Err(StoreError::NotFound); } if !key.compress { return Ok(data); } let mut decoder = Decoder::new(); decoder .decompress_vec(&data) .map_err(|e| StoreError::Compression(e.to_string())) } fn reserve_entry_slot(&self) -> Result, StoreError> { loop { let entries = self .entries .read() .map_err(|_| StoreError::Internal("Failed to acquire read lock on entries".to_string()))?; let entries_len = entries.len() as u64; let pending = self.pending_entries.load(Ordering::SeqCst); if entries_len + pending >= self.entry_limit { return Err(StoreError::LimitExceeded); } if self .pending_entries .compare_exchange(pending, pending + 1, Ordering::SeqCst, Ordering::SeqCst) .is_ok() { return Ok(EntryReservation { pending_entries: self.pending_entries.as_ref(), }); } } } /// Writes data to a file for the given key. fn write_file(&self, key: &Key, data: &[u8]) -> Result { let path = self.file_path(key); // Create directory if it doesn't exist if let Some(parent) = path.parent() { std::fs::create_dir_all(parent).map_err(StoreError::Io)?; } 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)?; } else { std::fs::write(&path, data).map_err(StoreError::Io)?; } let modified = SystemTime::now().duration_since(UNIX_EPOCH).unwrap_or_default().as_nanos() as i64; debug!( event = EVENT_TARGET_STORE_STATE, component = LOG_COMPONENT_TARGETS, subsystem = LOG_SUBSYSTEM_STORE, action = "write_file", key = %key, "target store state" ); Ok(modified) } fn insert_entry(&self, key: &Key, modified: i64) -> Result<(), StoreError> { let mut entries = self .entries .write() .map_err(|_| StoreError::Internal("Failed to acquire write lock on entries".to_string()))?; entries.insert(key.to_key_string(), modified); Ok(()) } fn remove_file_if_present(&self, key: &Key) -> Result<(), StoreError> { let path = self.file_path(key); match std::fs::remove_file(&path) { Ok(()) => Ok(()), Err(err) if err.kind() == std::io::ErrorKind::NotFound => Ok(()), Err(err) => Err(StoreError::Io(err)), } } fn write_and_index(&self, key: &Key, data: &[u8]) -> Result<(), StoreError> { let modified = self.write_file(key, data)?; if let Err(err) = self.insert_entry(key, modified) { self.remove_file_if_present(key).map_err(|cleanup_err| { StoreError::Internal(format!("Failed to index store entry {key}: {err}; cleanup failed: {cleanup_err}")) })?; return Err(err); } Ok(()) } } impl Store for QueueStore where T: Serialize + DeserializeOwned + Clone + Send + Sync + 'static, { type Error = StoreError; type Key = Key; fn open(&self) -> Result<(), Self::Error> { let _fs_guard = self .fs_guard .write() .map_err(|_| StoreError::Internal("Failed to acquire write lock on store filesystem".to_string()))?; std::fs::create_dir_all(&self.directory).map_err(StoreError::Io)?; let dir_entries = std::fs::read_dir(&self.directory).map_err(StoreError::Io)?; let mut entries_map = self .entries .write() .map_err(|_| StoreError::Internal("Failed to acquire write lock on entries".to_string()))?; self.pending_entries.store(0, Ordering::SeqCst); entries_map.clear(); 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); } } debug!( event = EVENT_TARGET_STORE_STATE, component = LOG_COMPONENT_TARGETS, subsystem = LOG_SUBSYSTEM_STORE, state = "opened", store_dir = ?self.directory, entry_count = entries_map.len(), "target store state" ); Ok(()) } fn put(&self, item: Arc) -> Result { let _fs_guard = self .fs_guard .read() .map_err(|_| StoreError::Internal("Failed to acquire read lock on store filesystem".to_string()))?; let _reservation = self.reserve_entry_slot()?; let key = self.build_key(1); let data = serde_json::to_vec(&*item).map_err(|e| StoreError::Serialization(e.to_string()))?; self.write_and_index(&key, &data)?; Ok(key) } fn put_multiple(&self, items: Vec) -> Result { if items.is_empty() { return Err(StoreError::Internal("Cannot put_multiple with empty items list".to_string())); } let _fs_guard = self .fs_guard .read() .map_err(|_| StoreError::Internal("Failed to acquire read lock on store filesystem".to_string()))?; let _reservation = self.reserve_entry_slot()?; let key = self.build_key(items.len()); let mut buffer = Vec::new(); for item in items { serde_json::to_writer(&mut buffer, &item).map_err(|e| StoreError::Serialization(e.to_string()))?; } self.write_and_index(&key, &buffer)?; Ok(key) } fn put_raw(&self, data: &[u8]) -> Result { let _fs_guard = self .fs_guard .read() .map_err(|_| StoreError::Internal("Failed to acquire read lock on store filesystem".to_string()))?; let _reservation = self.reserve_entry_slot()?; let key = self.build_key(1); self.write_and_index(&key, data)?; Ok(key) } fn get(&self, key: &Self::Key) -> Result { if key.item_count != 1 { return Err(StoreError::Internal(format!( "get() called on a batch key ({} items), use get_multiple()", key.item_count ))); } let items = self.get_multiple(key)?; items.into_iter().next().ok_or(StoreError::NotFound) } fn get_multiple(&self, key: &Self::Key) -> Result, Self::Error> { debug!( event = EVENT_TARGET_STORE_STATE, component = LOG_COMPONENT_TARGETS, subsystem = LOG_SUBSYSTEM_STORE, action = "read_batch", key = %key, "target store state" ); let data = self.get_raw(key)?; if data.is_empty() { return Err(StoreError::Deserialization("Cannot deserialize empty data".to_string())); } let mut items = Vec::with_capacity(key.item_count); // let mut deserializer = serde_json::Deserializer::from_slice(&data); // while let Ok(item) = serde::Deserialize::deserialize(&mut deserializer) { // items.push(item); // } // This deserialization logic assumes multiple JSON objects are simply concatenated in the file. // This is fragile. It's better to store a JSON array `[item1, item2, ...]` // or use a streaming deserializer that can handle multiple top-level objects if that's the format. // For now, assuming serde_json::Deserializer::from_slice can handle this if input is well-formed. let mut deserializer = serde_json::Deserializer::from_slice(&data).into_iter::(); for _ in 0..key.item_count { match deserializer.next() { Some(Ok(item)) => items.push(item), Some(Err(e)) => { return Err(StoreError::Deserialization(format!("Failed to deserialize item in batch: {e}"))); } None => { // Reached end of stream sooner than item_count if items.len() < key.item_count && !items.is_empty() { // Partial read warn!( event = EVENT_TARGET_STORE_STATE, component = LOG_COMPONENT_TARGETS, subsystem = LOG_SUBSYSTEM_STORE, action = "read_batch", key = %key, expected_items = key.item_count, actual_items = items.len(), reason = "partial_batch_read", "target store state" ); // Depending on strictness, this could be an error. } else if items.is_empty() { // No items at all, but file existed return Err(StoreError::Deserialization(format!( "No items deserialized for key {key} though file existed." ))); } break; } } } if items.is_empty() && key.item_count > 0 { return Err(StoreError::Deserialization("No items found".to_string())); } Ok(items) } fn get_raw(&self, key: &Self::Key) -> Result, Self::Error> { self.read_file(key) } fn del(&self, key: &Self::Key) -> Result<(), Self::Error> { let _fs_guard = self .fs_guard .read() .map_err(|_| StoreError::Internal("Failed to acquire read lock on store filesystem".to_string()))?; let path = self.file_path(key); match std::fs::remove_file(&path) { Ok(()) => {} Err(e) if e.kind() == std::io::ErrorKind::NotFound => { // File already gone — still clean up the entries map to avoid stale keys. warn!( event = EVENT_TARGET_STORE_STATE, component = LOG_COMPONENT_TARGETS, subsystem = LOG_SUBSYSTEM_STORE, action = "delete", key = %key, result = "file_missing", "target store state" ); } Err(e) => return Err(StoreError::Io(e)), } // Get the write lock to update the internal state let mut entries = self .entries .write() .map_err(|_| StoreError::Internal("Failed to acquire write lock on entries".to_string()))?; if entries.remove(&key.to_key_string()).is_none() { debug!( event = EVENT_TARGET_STORE_STATE, component = LOG_COMPONENT_TARGETS, subsystem = LOG_SUBSYSTEM_STORE, action = "delete", key = %key, result = "entry_missing", "target store state" ); } debug!( event = EVENT_TARGET_STORE_STATE, component = LOG_COMPONENT_TARGETS, subsystem = LOG_SUBSYSTEM_STORE, action = "delete", key = %key, result = "deleted", "target store state" ); Ok(()) } fn delete(&self) -> Result<(), Self::Error> { let _fs_guard = self .fs_guard .write() .map_err(|_| StoreError::Internal("Failed to acquire write lock on store filesystem".to_string()))?; let mut entries = self .entries .write() .map_err(|_| StoreError::Internal("Failed to acquire write lock on entries".to_string()))?; entries.clear(); self.pending_entries.store(0, Ordering::SeqCst); match std::fs::remove_dir_all(&self.directory) { Ok(()) => Ok(()), Err(err) if err.kind() == std::io::ErrorKind::NotFound => Ok(()), Err(err) => Err(StoreError::Io(err)), } } fn list(&self) -> Vec { // Get the read lock to read the internal state let entries = match self.entries.read() { Ok(entries) => entries, Err(_) => { debug!( event = EVENT_TARGET_STORE_STATE, component = LOG_COMPONENT_TARGETS, subsystem = LOG_SUBSYSTEM_STORE, action = "list", result = "lock_unavailable", "target store state" ); return Vec::new(); } }; 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 entries_vec.into_iter().map(|(k, _)| parse_key(k)).collect() } fn len(&self) -> usize { // Get the read lock to read the internal state match self.entries.read() { Ok(entries) => entries.len(), Err(_) => { debug!( event = EVENT_TARGET_STORE_STATE, component = LOG_COMPONENT_TARGETS, subsystem = LOG_SUBSYSTEM_STORE, action = "len", result = "lock_unavailable", "target store state" ); 0 } } } fn is_empty(&self) -> bool { self.len() == 0 } fn boxed_clone(&self) -> Box + Send + Sync> { Box::new(self.clone()) as Box + Send + Sync> } } #[cfg(test)] mod tests { use super::*; use std::{ sync::{Arc, Barrier}, thread, }; fn temp_store_dir(name: &str) -> PathBuf { std::env::temp_dir().join(format!("rustfs-targets-{name}-{}", Uuid::new_v4())) } #[test] fn resolve_queue_store_compression_defaults_to_true() { assert!(resolve_queue_store_compression_from_env_value(None)); } #[test] fn resolve_queue_store_compression_respects_disabled_env_value() { assert!(!resolve_queue_store_compression_from_env_value(Some("off"))); assert!(!resolve_queue_store_compression_from_env_value(Some("false"))); } #[test] fn put_uses_store_compression_setting_in_key() { let dir = temp_store_dir("put-key"); let store = QueueStore::::new_with_compression(&dir, 8, ".test", false); store.open().unwrap(); let key = store.put(Arc::new("payload".to_string())).unwrap(); assert!(!key.compress); assert!(store.file_path(&key).exists()); let _ = std::fs::remove_dir_all(dir); } #[test] fn parse_key_round_trips_batch_and_compression_suffixes() { let key = Key { name: "event-id".to_string(), extension: ".json".to_string(), item_count: 3, compress: true, }; let parsed = parse_key(&key.to_key_string()); assert_eq!(parsed.name, key.name); assert_eq!(parsed.extension, key.extension); assert_eq!(parsed.item_count, key.item_count); assert_eq!(parsed.compress, key.compress); } #[test] fn put_raw_and_get_raw_round_trip_bytes() { let dir = temp_store_dir("raw-roundtrip"); let store = QueueStore::::new_with_compression(&dir, 8, ".test", true); store.open().unwrap(); let payload = br#"{"kind":"notify","bucket":"demo","key":"alpha.txt"}"#; let key = store.put_raw(payload).unwrap(); let raw = store.get_raw(&key).unwrap(); assert_eq!(raw, payload); let _ = store.delete(); } #[test] fn delete_removes_directory_and_clears_entries() { let dir = temp_store_dir("delete-store"); let store = QueueStore::::new_with_compression(&dir, 8, ".test", false); store.open().unwrap(); let _ = store.put(Arc::new("payload".to_string())).unwrap(); store.delete().unwrap(); assert!(store.list().is_empty()); assert!(!dir.exists()); } #[test] fn put_enforces_entry_limit() { let dir = temp_store_dir("limit"); let store = QueueStore::::new_with_compression(&dir, 1, ".test", false); store.open().unwrap(); let _ = store.put(Arc::new("first".to_string())).unwrap(); let err = store.put(Arc::new("second".to_string())).unwrap_err(); assert!(matches!(err, StoreError::LimitExceeded)); let _ = store.delete(); } #[test] fn concurrent_put_raw_respects_entry_limit() { let dir = temp_store_dir("concurrent-limit"); let store = Arc::new(QueueStore::::new_with_compression(&dir, 1, ".test", true)); store.open().unwrap(); let start = Arc::new(Barrier::new(4)); let mut handles = Vec::new(); for idx in 0..4 { let store = Arc::clone(&store); let start = Arc::clone(&start); handles.push(thread::spawn(move || { let payload = vec![b'x'; 32 * 1024 + idx]; start.wait(); store.put_raw(&payload) })); } let mut successes = 0; let mut limit_errors = 0; for handle in handles { match handle.join().unwrap() { Ok(_) => successes += 1, Err(StoreError::LimitExceeded) => limit_errors += 1, Err(err) => panic!("unexpected error: {err}"), } } assert_eq!(successes, 1); assert_eq!(limit_errors, 3); assert_eq!(store.len(), 1); let _ = store.delete(); } }