use crate::fs_store::common::{
dir_entry_is_key, get_key_from_dir_entry_path, FilesystemStoreState,
};
use lightning::util::persist::{
KVStoreSync, MigratableKVStoreSync, PageToken, PaginatedKVStoreSync, PaginatedListResponse,
};
use std::fs;
use std::path::PathBuf;
use std::time::UNIX_EPOCH;
use std::{error, fmt, io};
#[cfg(feature = "tokio")]
use core::future::Future;
#[cfg(feature = "tokio")]
use lightning::util::persist::{KVStore, PaginatedKVStore};
use std::sync::Arc;
#[derive(Debug)]
pub enum FilesystemStoreV2Error {
V1DataDetected(PathBuf),
Io(io::Error),
}
impl fmt::Display for FilesystemStoreV2Error {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::V1DataDetected(path) => write!(
f,
"Found file `{}` where FilesystemStoreV2 expects a namespace directory. \
This indicates the directory was previously used by FilesystemStore (v1). \
Please migrate your data or use a different directory.",
path.display()
),
Self::Io(err) => write!(f, "{}", err),
}
}
}
impl error::Error for FilesystemStoreV2Error {
fn source(&self) -> Option<&(dyn error::Error + 'static)> {
match self {
Self::V1DataDetected(_) => None,
Self::Io(err) => Some(err),
}
}
}
impl From<io::Error> for FilesystemStoreV2Error {
fn from(err: io::Error) -> Self {
Self::Io(err)
}
}
pub struct FilesystemStoreV2 {
inner: Arc<FilesystemStoreState>,
}
impl FilesystemStoreV2 {
pub fn new(data_dir: PathBuf) -> Result<Self, FilesystemStoreV2Error> {
if data_dir.exists() {
for entry in fs::read_dir(&data_dir)? {
let entry = entry?;
let file_type = entry.file_type()?;
if file_type.is_file() {
return Err(FilesystemStoreV2Error::V1DataDetected(entry.path()));
}
if file_type.is_dir() {
for child_entry in fs::read_dir(entry.path())? {
let child_entry = child_entry?;
if child_entry.file_type()?.is_file() {
return Err(FilesystemStoreV2Error::V1DataDetected(child_entry.path()));
}
}
}
}
}
Ok(Self { inner: Arc::new(FilesystemStoreState::new(data_dir)) })
}
pub fn get_data_dir(&self) -> PathBuf {
self.inner.get_data_dir()
}
#[cfg(any(all(feature = "tokio", test), fuzzing))]
pub fn state_size(&self) -> usize {
self.inner.state_size()
}
}
pub(crate) const PAGE_SIZE: usize = 50;
const PAGE_TOKEN_TIMESTAMP_LEN: usize = 16;
impl FilesystemStoreState {
fn list_paginated_impl(
&self, prefixed_dest: PathBuf, page_token: Option<PageToken>,
) -> Result<PaginatedListResponse, lightning::io::Error> {
if !prefixed_dest.exists() {
return Ok(PaginatedListResponse { keys: Vec::new(), next_page_token: None });
}
let mut entries: Vec<(u64, String)> = Vec::new();
for dir_entry in fs::read_dir(&prefixed_dest)? {
let dir_entry = dir_entry?;
match dir_entry_is_key(&dir_entry) {
Ok(false) => continue,
Ok(true) => {},
Err(_) => {},
}
let key =
get_key_from_dir_entry_path(&dir_entry.path(), prefixed_dest.as_path(), false)?;
let mtime_millis = dir_entry
.metadata()
.ok()
.and_then(|m| m.modified().ok())
.and_then(|t| t.duration_since(UNIX_EPOCH).ok())
.map(|d| d.as_millis() as u64)
.unwrap_or(0);
entries.push((mtime_millis, key));
}
entries.sort_by(|a, b| b.0.cmp(&a.0).then_with(|| b.1.cmp(&a.1)));
let start_idx = if let Some(token) = page_token {
let (token_mtime, token_key) = parse_page_token(token.as_str())?;
entries
.iter()
.position(|(mtime, key)| {
*mtime < token_mtime
|| (*mtime == token_mtime && key.as_str() < token_key.as_str())
})
.unwrap_or(entries.len())
} else {
0
};
let page_entries: Vec<_> =
entries.iter().skip(start_idx).take(PAGE_SIZE).cloned().collect();
let next_page_token = if start_idx + PAGE_SIZE < entries.len() {
page_entries.last().map(|(mtime, key)| PageToken::new(format_page_token(*mtime, key)))
} else {
None
};
let keys: Vec<String> = page_entries.into_iter().map(|(_, key)| key).collect();
Ok(PaginatedListResponse { keys, next_page_token })
}
}
impl KVStoreSync for FilesystemStoreV2 {
fn read(
&self, primary_namespace: &str, secondary_namespace: &str, key: &str,
) -> Result<Vec<u8>, lightning::io::Error> {
self.inner.read_impl(primary_namespace, secondary_namespace, key, true)
}
fn write(
&self, primary_namespace: &str, secondary_namespace: &str, key: &str, buf: Vec<u8>,
) -> Result<(), lightning::io::Error> {
self.inner.write_impl(primary_namespace, secondary_namespace, key, buf, true)
}
fn remove(
&self, primary_namespace: &str, secondary_namespace: &str, key: &str, lazy: bool,
) -> Result<(), lightning::io::Error> {
self.inner.remove_impl(primary_namespace, secondary_namespace, key, lazy, true)
}
fn list(
&self, primary_namespace: &str, secondary_namespace: &str,
) -> Result<Vec<String>, lightning::io::Error> {
self.inner.list_impl(primary_namespace, secondary_namespace, true)
}
}
impl PaginatedKVStoreSync for FilesystemStoreV2 {
fn list_paginated(
&self, primary_namespace: &str, secondary_namespace: &str, page_token: Option<PageToken>,
) -> Result<PaginatedListResponse, lightning::io::Error> {
let prefixed_dest = self.inner.get_checked_dest_file_path(
primary_namespace,
secondary_namespace,
None,
"list_paginated",
true,
)?;
self.inner.list_paginated_impl(prefixed_dest, page_token)
}
}
#[cfg(feature = "tokio")]
impl KVStore for FilesystemStoreV2 {
fn read(
&self, primary_namespace: &str, secondary_namespace: &str, key: &str,
) -> impl Future<Output = Result<Vec<u8>, lightning::io::Error>> + 'static + Send {
self.inner.read_async(primary_namespace, secondary_namespace, key, true)
}
fn write(
&self, primary_namespace: &str, secondary_namespace: &str, key: &str, buf: Vec<u8>,
) -> impl Future<Output = Result<(), lightning::io::Error>> + 'static + Send {
self.inner.write_async(primary_namespace, secondary_namespace, key, buf, true)
}
fn remove(
&self, primary_namespace: &str, secondary_namespace: &str, key: &str, lazy: bool,
) -> impl Future<Output = Result<(), lightning::io::Error>> + 'static + Send {
self.inner.remove_async(primary_namespace, secondary_namespace, key, lazy, true)
}
fn list(
&self, primary_namespace: &str, secondary_namespace: &str,
) -> impl Future<Output = Result<Vec<String>, lightning::io::Error>> + 'static + Send {
self.inner.list_async(primary_namespace, secondary_namespace, true)
}
}
#[cfg(feature = "tokio")]
impl PaginatedKVStore for FilesystemStoreV2 {
fn list_paginated(
&self, primary_namespace: &str, secondary_namespace: &str, page_token: Option<PageToken>,
) -> impl Future<Output = Result<PaginatedListResponse, lightning::io::Error>> + 'static + Send
{
let this = Arc::clone(&self.inner);
let path = this.get_checked_dest_file_path(
primary_namespace,
secondary_namespace,
None,
"list_paginated",
true,
);
async move {
let path = match path {
Ok(path) => path,
Err(e) => return Err(e),
};
tokio::task::spawn_blocking(move || this.list_paginated_impl(path, page_token))
.await
.unwrap_or_else(|e| {
Err(lightning::io::Error::new(lightning::io::ErrorKind::Other, e))
})
}
}
}
impl MigratableKVStoreSync for FilesystemStoreV2 {
fn list_all_keys(&self) -> Result<Vec<(String, String, String)>, lightning::io::Error> {
self.inner.list_all_keys_impl(true)
}
}
#[cfg(feature = "tokio")]
impl lightning::util::persist::MigratableKVStore for FilesystemStoreV2 {
fn list_all_keys(
&self,
) -> impl Future<Output = Result<Vec<(String, String, String)>, lightning::io::Error>> + 'static + Send
{
self.inner.list_all_keys_async(true)
}
}
pub(crate) fn format_page_token(mtime_millis: u64, key: &str) -> String {
format!("{mtime_millis:016}:{key}")
}
pub(crate) fn parse_page_token(token: &str) -> lightning::io::Result<(u64, String)> {
if token.as_bytes().get(PAGE_TOKEN_TIMESTAMP_LEN) != Some(&b':') {
return Err(lightning::io::Error::new(
lightning::io::ErrorKind::InvalidInput,
"Invalid page token format",
));
}
let mtime = token[..PAGE_TOKEN_TIMESTAMP_LEN].parse::<u64>().map_err(|_| {
lightning::io::Error::new(
lightning::io::ErrorKind::InvalidInput,
"Invalid page token timestamp",
)
})?;
let key = token[PAGE_TOKEN_TIMESTAMP_LEN + 1..].to_string();
Ok((mtime, key))
}
#[cfg(test)]
mod tests {
use super::*;
use crate::fs_store::common::EMPTY_NAMESPACE_DIR;
#[cfg(feature = "tokio")]
use crate::test_utils::do_test_data_migration_async;
use crate::test_utils::{
do_read_write_remove_list_persist, do_test_data_migration, do_test_store,
};
use std::fs::FileTimes;
use std::time::UNIX_EPOCH;
impl Drop for FilesystemStoreV2 {
fn drop(&mut self) {
match fs::remove_dir_all(&self.inner.get_data_dir()) {
Err(e) => println!("Failed to remove test persister directory: {}", e),
_ => {},
}
}
}
#[test]
fn read_write_remove_list_persist() {
let mut temp_path = std::env::temp_dir();
temp_path.push("test_read_write_remove_list_persist_v2");
let fs_store = FilesystemStoreV2::new(temp_path).unwrap();
do_read_write_remove_list_persist(&fs_store);
}
#[cfg(feature = "tokio")]
#[tokio::test]
async fn read_write_remove_list_persist_async() {
use lightning::util::persist::KVStore;
use std::sync::Arc;
let mut temp_path = std::env::temp_dir();
temp_path.push("test_read_write_remove_list_persist_async_v2");
let fs_store = Arc::new(FilesystemStoreV2::new(temp_path).unwrap());
assert_eq!(fs_store.state_size(), 0);
let async_fs_store = Arc::clone(&fs_store);
let data1 = vec![42u8; 32];
let data2 = vec![43u8; 32];
let primary = "testspace";
let secondary = "testsubspace";
let key = "testkey";
let fut1 = KVStore::write(&*async_fs_store, primary, secondary, key, data1);
assert_eq!(fs_store.state_size(), 1);
let fut2 = KVStore::remove(&*async_fs_store, primary, secondary, key, false);
assert_eq!(fs_store.state_size(), 1);
let fut3 = KVStore::write(&*async_fs_store, primary, secondary, key, data2.clone());
assert_eq!(fs_store.state_size(), 1);
fut3.await.unwrap();
assert_eq!(fs_store.state_size(), 1);
fut2.await.unwrap();
assert_eq!(fs_store.state_size(), 1);
fut1.await.unwrap();
assert_eq!(fs_store.state_size(), 0);
let listed_keys = KVStore::list(&*async_fs_store, primary, secondary).await.unwrap();
assert_eq!(listed_keys.len(), 1);
assert_eq!(listed_keys[0], key);
let read_data = KVStore::read(&*async_fs_store, primary, secondary, key).await.unwrap();
assert_eq!(data2, &*read_data);
KVStore::remove(&*async_fs_store, primary, secondary, key, false).await.unwrap();
let listed_keys = KVStore::list(&*async_fs_store, primary, secondary).await.unwrap();
assert_eq!(listed_keys.len(), 0);
}
#[cfg(feature = "tokio")]
#[tokio::test]
async fn stale_write_does_not_leak_tmp_file() {
use lightning::util::persist::KVStore;
let mut temp_path = std::env::temp_dir();
temp_path.push("test_stale_write_does_not_leak_tmp_file_v2");
let _ = fs::remove_dir_all(&temp_path);
let fs_store = FilesystemStoreV2::new(temp_path.clone()).unwrap();
let data1 = vec![1u8; 32];
let data2 = vec![2u8; 32];
let primary = "testspace";
let secondary = "testsubspace";
let key = "testkey";
let fut1 = KVStore::write(&fs_store, primary, secondary, key, data1);
let fut2 = KVStore::write(&fs_store, primary, secondary, key, data2);
fut2.await.unwrap();
fut1.await.unwrap();
let dir = temp_path.join(primary).join(secondary);
let tmp_files: Vec<_> = fs::read_dir(&dir)
.unwrap()
.filter_map(|e| e.ok())
.map(|e| e.path())
.filter(|p| p.extension().map_or(false, |ext| ext == "tmp"))
.collect();
assert!(tmp_files.is_empty(), "Found leaked tmp files: {:?}", tmp_files);
}
#[test]
fn test_data_migration() {
let mut source_temp_path = std::env::temp_dir();
source_temp_path.push("test_data_migration_source_v2");
let mut source_store = FilesystemStoreV2::new(source_temp_path).unwrap();
let mut target_temp_path = std::env::temp_dir();
target_temp_path.push("test_data_migration_target_v2");
let mut target_store = FilesystemStoreV2::new(target_temp_path).unwrap();
do_test_data_migration(&mut source_store, &mut target_store);
}
#[cfg(feature = "tokio")]
#[tokio::test]
async fn test_data_migration_async() {
let mut source_temp_path = std::env::temp_dir();
source_temp_path.push("test_data_migration_source_async_v2");
let source_store = FilesystemStoreV2::new(source_temp_path).unwrap();
let mut target_temp_path = std::env::temp_dir();
target_temp_path.push("test_data_migration_target_async_v2");
let target_store = FilesystemStoreV2::new(target_temp_path).unwrap();
do_test_data_migration_async(&source_store, &target_store).await;
}
#[test]
fn test_filesystem_store_v2() {
let store_0 = FilesystemStoreV2::new("test_filesystem_store_v2_0".into()).unwrap();
let store_1 = FilesystemStoreV2::new("test_filesystem_store_v2_1".into()).unwrap();
do_test_store(&store_0, &store_1)
}
#[test]
fn test_page_token_format() {
let mtime: u64 = 1706500000000;
let key = "test_key";
let token = format_page_token(mtime, key);
assert_eq!(token, "0001706500000000:test_key");
let parsed = parse_page_token(&token).unwrap();
assert_eq!(parsed, (mtime, key.to_string()));
assert!(parse_page_token("invalid").is_err());
assert!(parse_page_token("0001706500000000_key").is_err()); assert!(parse_page_token("0001706500000000").is_err()); assert!(parse_page_token("1706500000000:key").is_err()); }
#[test]
fn test_directory_structure() {
use lightning::util::persist::KVStoreSync;
let mut temp_path = std::env::temp_dir();
temp_path.push("test_directory_structure_v2");
let fs_store = FilesystemStoreV2::new(temp_path.clone()).unwrap();
let data = vec![42u8; 32];
KVStoreSync::write(&fs_store, "", "", "key1", data.clone()).unwrap();
assert!(temp_path.join(EMPTY_NAMESPACE_DIR).join(EMPTY_NAMESPACE_DIR).exists());
KVStoreSync::write(&fs_store, "primary", "", "key2", data.clone()).unwrap();
assert!(temp_path.join("primary").join(EMPTY_NAMESPACE_DIR).exists());
KVStoreSync::write(&fs_store, "primary", "secondary", "key3", data.clone()).unwrap();
assert!(temp_path.join("primary").join("secondary").exists());
assert_eq!(KVStoreSync::read(&fs_store, "", "", "key1").unwrap(), data);
assert_eq!(KVStoreSync::read(&fs_store, "primary", "", "key2").unwrap(), data);
assert_eq!(KVStoreSync::read(&fs_store, "primary", "secondary", "key3").unwrap(), data);
assert!(temp_path
.join(EMPTY_NAMESPACE_DIR)
.join(EMPTY_NAMESPACE_DIR)
.join("key1")
.exists());
assert!(temp_path.join("primary").join(EMPTY_NAMESPACE_DIR).join("key2").exists());
assert!(temp_path.join("primary").join("secondary").join("key3").exists());
}
#[test]
fn test_update_preserves_mtime() {
use lightning::util::persist::KVStoreSync;
let mut temp_path = std::env::temp_dir();
temp_path.push("test_update_preserves_mtime_v2");
let fs_store = FilesystemStoreV2::new(temp_path.clone()).unwrap();
let data1 = vec![42u8; 32];
let data2 = vec![43u8; 32];
KVStoreSync::write(&fs_store, "ns", "sub", "key", data1).unwrap();
let file_path = temp_path.join("ns").join("sub").join("key");
let original_mtime = fs::metadata(&file_path).unwrap().modified().unwrap();
std::thread::sleep(std::time::Duration::from_millis(50));
KVStoreSync::write(&fs_store, "ns", "sub", "key", data2.clone()).unwrap();
let updated_mtime = fs::metadata(&file_path).unwrap().modified().unwrap();
assert_eq!(original_mtime, updated_mtime);
assert_eq!(KVStoreSync::read(&fs_store, "ns", "sub", "key").unwrap(), data2);
}
#[test]
fn test_paginated_listing() {
use lightning::util::persist::{KVStoreSync, PaginatedKVStoreSync};
let mut temp_path = std::env::temp_dir();
temp_path.push("test_paginated_listing_v2");
let fs_store = FilesystemStoreV2::new(temp_path).unwrap();
let data = vec![42u8; 32];
let keys: Vec<String> = (0..5).map(|i| format!("key{}", i)).collect();
for key in &keys {
KVStoreSync::write(&fs_store, "ns", "sub", key, data.clone()).unwrap();
std::thread::sleep(std::time::Duration::from_millis(10));
}
let response = PaginatedKVStoreSync::list_paginated(&fs_store, "ns", "sub", None).unwrap();
assert_eq!(response.keys.len(), 5);
assert_eq!(response.keys[0], "key4");
assert_eq!(response.keys[4], "key0");
assert!(response.next_page_token.is_none()); }
#[test]
fn test_paginated_listing_with_pagination() {
use lightning::util::persist::{KVStoreSync, PaginatedKVStoreSync};
let mut temp_path = std::env::temp_dir();
temp_path.push("test_paginated_listing_with_pagination_v2");
let fs_store = FilesystemStoreV2::new(temp_path).unwrap();
let data = vec![42u8; 32];
let num_keys = PAGE_SIZE + 50;
for i in 0..num_keys {
let key = format!("key{:04}", i);
KVStoreSync::write(&fs_store, "ns", "sub", &key, data.clone()).unwrap();
if i % 10 == 0 {
std::thread::sleep(std::time::Duration::from_millis(1));
}
}
let response1 = PaginatedKVStoreSync::list_paginated(&fs_store, "ns", "sub", None).unwrap();
assert_eq!(response1.keys.len(), PAGE_SIZE);
assert!(response1.next_page_token.is_some());
let response2 =
PaginatedKVStoreSync::list_paginated(&fs_store, "ns", "sub", response1.next_page_token)
.unwrap();
assert_eq!(response2.keys.len(), 50);
assert!(response2.next_page_token.is_none());
let all_keys: std::collections::HashSet<_> =
response1.keys.iter().chain(response2.keys.iter()).collect();
assert_eq!(all_keys.len(), num_keys);
}
#[test]
fn test_page_token_after_deletion() {
use lightning::util::persist::{KVStoreSync, PaginatedKVStoreSync};
let mut temp_path = std::env::temp_dir();
temp_path.push("test_page_token_after_deletion_v2");
let fs_store = FilesystemStoreV2::new(temp_path).unwrap();
let data = vec![42u8; 32];
for i in 0..10 {
let key = format!("key{}", i);
KVStoreSync::write(&fs_store, "ns", "sub", &key, data.clone()).unwrap();
std::thread::sleep(std::time::Duration::from_millis(10));
}
let response1 = PaginatedKVStoreSync::list_paginated(&fs_store, "ns", "sub", None).unwrap();
assert_eq!(response1.keys.len(), 10);
KVStoreSync::remove(&fs_store, "ns", "sub", "key5", false).unwrap();
KVStoreSync::remove(&fs_store, "ns", "sub", "key3", false).unwrap();
let response2 = PaginatedKVStoreSync::list_paginated(&fs_store, "ns", "sub", None).unwrap();
assert_eq!(response2.keys.len(), 8); }
#[test]
fn test_same_mtime_sorted_by_key() {
use lightning::util::persist::PaginatedKVStoreSync;
use std::time::Duration;
let mut temp_path = std::env::temp_dir();
temp_path.push("test_same_mtime_sorted_by_key_v2");
let _ = fs::remove_dir_all(&temp_path);
let data = vec![42u8; 32];
let dir = temp_path.join("ns").join("sub");
fs::create_dir_all(&dir).unwrap();
let keys = vec!["zebra", "apple", "mango", "banana"];
let fixed_time = UNIX_EPOCH + Duration::from_secs(1706500000);
for key in &keys {
let file_path = dir.join(key);
let file = fs::File::create(&file_path).unwrap();
std::io::Write::write_all(&mut &file, &data).unwrap();
file.set_times(FileTimes::new().set_modified(fixed_time)).unwrap();
}
let fs_store = FilesystemStoreV2::new(temp_path.clone()).unwrap();
let response = PaginatedKVStoreSync::list_paginated(&fs_store, "ns", "sub", None).unwrap();
assert_eq!(response.keys.len(), 4);
assert_eq!(response.keys[0], "zebra");
assert_eq!(response.keys[1], "mango");
assert_eq!(response.keys[2], "banana");
assert_eq!(response.keys[3], "apple");
}
#[test]
fn test_paginated_listing_skips_tmp_files() {
use lightning::util::persist::{KVStoreSync, PaginatedKVStoreSync};
let mut temp_path = std::env::temp_dir();
temp_path.push("test_paginated_listing_skips_tmp_files_v2");
let fs_store = FilesystemStoreV2::new(temp_path.clone()).unwrap();
let data = vec![42u8; 32];
KVStoreSync::write(&fs_store, "ns", "sub", "key0", data.clone()).unwrap();
std::thread::sleep(std::time::Duration::from_millis(10));
KVStoreSync::write(&fs_store, "ns", "sub", "key1", data.clone()).unwrap();
let dir = temp_path.join("ns").join("sub");
fs::write(dir.join("inflight.tmp"), &data).unwrap();
fs::create_dir_all(dir.join("stray_dir")).unwrap();
let response = PaginatedKVStoreSync::list_paginated(&fs_store, "ns", "sub", None).unwrap();
assert_eq!(response.keys.len(), 2);
assert!(response.keys.contains(&"key0".to_string()));
assert!(response.keys.contains(&"key1".to_string()));
}
#[test]
fn test_rejects_v1_data_directory() {
let mut temp_path = std::env::temp_dir();
temp_path.push("test_rejects_v1_data_directory");
let _ = fs::remove_dir_all(&temp_path);
fs::create_dir_all(&temp_path).unwrap();
fs::write(temp_path.join("some_key"), b"data").unwrap();
match FilesystemStoreV2::new(temp_path.clone()) {
Err(FilesystemStoreV2Error::V1DataDetected(path)) => {
assert_eq!(path, temp_path.join("some_key"));
},
Err(err) => panic!("Expected V1DataDetected, got {:?}", err),
Ok(_) => panic!("Expected error for directory with top-level files"),
}
let _ = fs::remove_dir_all(&temp_path);
fs::create_dir_all(temp_path.join("some_namespace")).unwrap();
fs::write(temp_path.join("some_namespace").join("some_key"), b"data").unwrap();
match FilesystemStoreV2::new(temp_path.clone()) {
Err(FilesystemStoreV2Error::V1DataDetected(path)) => {
assert_eq!(path, temp_path.join("some_namespace").join("some_key"));
},
Err(err) => panic!("Expected V1DataDetected, got {:?}", err),
Ok(_) => panic!("Expected error for directory with files one level down"),
}
let _ = fs::remove_dir_all(&temp_path);
fs::create_dir_all(temp_path.join("some_secondary_namespace")).unwrap();
fs::write(temp_path.join("some_secondary_namespace").join("some_key"), b"data").unwrap();
match FilesystemStoreV2::new(temp_path.clone()) {
Err(FilesystemStoreV2Error::V1DataDetected(path)) => {
assert_eq!(path, temp_path.join("some_secondary_namespace").join("some_key"));
},
Err(err) => panic!("Expected V1DataDetected, got {:?}", err),
Ok(_) => panic!("Expected error for directory with files one level down"),
}
let _ = fs::remove_dir_all(&temp_path);
fs::create_dir_all(&temp_path).unwrap();
let result = FilesystemStoreV2::new(temp_path.clone());
assert!(result.is_ok());
fs::create_dir_all(temp_path.join("some_namespace").join("some_sub_namespace")).unwrap();
let result = FilesystemStoreV2::new(temp_path.clone());
assert!(result.is_ok());
let fs_store = result.unwrap();
KVStoreSync::write(
&fs_store,
"some_namespace",
"some_sub_namespace",
"some_key",
b"data".to_vec(),
)
.unwrap();
let result = FilesystemStoreV2::new(temp_path);
assert!(result.is_ok());
}
}