Files
rustfs/crates/targets/src/store.rs
T

839 lines
28 KiB
Rust

// 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::<EnableState>().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::<usize>()
{
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<T>(
store: &(dyn Store<T, Error = StoreError, Key = Key> + Send),
key: &Key,
) -> Result<bool, StoreError>
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<T>: 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<T>) -> Result<Self::Key, Self::Error>;
/// Stores multiple items in a single batch
fn put_multiple(&self, items: Vec<T>) -> Result<Self::Key, Self::Error>;
/// Stores raw bytes in a single entry.
fn put_raw(&self, data: &[u8]) -> Result<Self::Key, Self::Error>;
/// Retrieves a single item by key
fn get(&self, key: &Self::Key) -> Result<T, Self::Error>;
/// Retrieves multiple items by key
fn get_multiple(&self, key: &Self::Key) -> Result<Vec<T>, Self::Error>;
/// Retrieves the raw bytes stored for a key.
fn get_raw(&self, key: &Self::Key) -> Result<Vec<u8>, 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<Self::Key>;
/// 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<dyn Store<T, Error = Self::Error, Key = Self::Key> + Send + Sync>;
}
/// A store that uses the filesystem to persist events in a queue
pub struct QueueStore<T> {
entry_limit: u64,
directory: PathBuf,
file_ext: String,
compress: bool,
entries: Arc<RwLock<HashMap<String, i64>>>, // key -> modtime as unix nano
pending_entries: Arc<AtomicU64>,
fs_guard: Arc<RwLock<()>>,
_phantom: PhantomData<T>,
}
impl<T> Clone for QueueStore<T> {
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<T: Serialize + DeserializeOwned + Send + Sync> QueueStore<T> {
/// Creates a new QueueStore
pub fn new(directory: impl Into<PathBuf>, 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<PathBuf>, 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<Vec<u8>, 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<EntryReservation<'_>, 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<i64, StoreError> {
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<T> Store<T> for QueueStore<T>
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<T>) -> Result<Self::Key, Self::Error> {
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<T>) -> Result<Self::Key, Self::Error> {
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<Self::Key, Self::Error> {
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<T, Self::Error> {
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<Vec<T>, 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::<T>();
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<Vec<u8>, 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<Self::Key> {
// 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<dyn Store<T, Error = Self::Error, Key = Self::Key> + Send + Sync> {
Box::new(self.clone()) as Box<dyn Store<T, Error = Self::Error, Key = Self::Key> + 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::<String>::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::<String>::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::<String>::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::<String>::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::<String>::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();
}
}