use std::io::Cursor;
use std::io::Result as IoResult;
use std::sync::Arc;
use std::sync::Mutex;
use std::sync::atomic::AtomicU64;
use std::sync::atomic::Ordering;
use std::time::Duration;
use qubit_fs::FileSystem;
use qubit_fs::FsError;
use qubit_fs::FsResult;
use qubit_fs::Path;
use qubit_fs::copy::CopyOptions;
use qubit_fs::directory::CreateDirectoryOptions;
use qubit_fs::directory::CreateDirectoryOutcome;
use qubit_fs::directory::DeleteOptions;
use qubit_fs::directory::DeleteOutcome;
use qubit_fs::directory::ListOptions;
use qubit_fs::directory::ListScope;
use qubit_fs::error::FsErrorKind;
use qubit_fs::error::FsOperation;
use qubit_fs::metadata::AchievedAtomicity;
use qubit_fs::metadata::DirEntry;
use qubit_fs::metadata::FileKind;
use qubit_fs::metadata::FileMetadata;
use qubit_fs::metadata::FileSystemCapabilities;
use qubit_fs::metadata::FileSystemCapability;
use qubit_fs::metadata::FileSystemCapabilitySupport;
use qubit_fs::metadata::FileSystemId;
use qubit_fs::metadata::FileSystemInfo;
use qubit_fs::metadata::FileSystemLimit;
use qubit_fs::metadata::FileSystemLimits;
use qubit_fs::metadata::OpenedFileInfo;
use qubit_fs::metadata::PublicationMethod;
use qubit_fs::metadata::SymlinkPolicy;
use qubit_fs::metadata::WriteOutcome;
use qubit_fs::path::PathConstraints;
use qubit_fs::path::PathSemantics;
use qubit_fs::read::ReadOptions;
use qubit_fs::rename::RenameFailureState;
use qubit_fs::rename::RenameOutcome;
use qubit_fs::spi::CreateDirectoryRequest;
use qubit_fs::spi::CreateTempDirectoryRequest;
use qubit_fs::spi::CreateTempFileRequest;
use qubit_fs::spi::DeleteDirectoryRequest;
use qubit_fs::spi::DeleteFileRequest;
use qubit_fs::spi::DirectoryStreamSpi;
use qubit_fs::spi::FileSystemSpi;
use qubit_fs::spi::FileWriterSpi;
use qubit_fs::spi::ListRequest;
use qubit_fs::spi::OpenReaderRequest;
use qubit_fs::spi::OpenWriterRequest;
use qubit_fs::spi::OpenedDirectoryStream;
use qubit_fs::spi::OpenedReader;
use qubit_fs::spi::OpenedTempDirectory;
use qubit_fs::spi::OpenedTempFile;
use qubit_fs::spi::OpenedWriter;
use qubit_fs::spi::PersistRequest;
use qubit_fs::spi::ProviderOperation;
use qubit_fs::spi::ProviderOperations;
use qubit_fs::spi::ProviderProperties;
use qubit_fs::spi::RenameRequest;
use qubit_fs::spi::SpiPersistFailure;
use qubit_fs::spi::SpiRenameFailure;
use qubit_fs::spi::SpiWriteFailure;
use qubit_fs::spi::StatRequest;
use qubit_fs::spi::StatResponse;
use qubit_fs::spi::TempResourceSpi;
use qubit_fs::temp::PersistFailureState;
use qubit_fs::temp::PersistOutcome;
use qubit_fs::temp::TempOptions;
use qubit_fs::write::WriteAbortOutcome;
use qubit_fs::write::WriteFailureState;
use qubit_fs::write::WriteOptions;
use qubit_io::Input;
use qubit_io::Output;
static WRITER_FLUSH_DELAY_MS: AtomicU64 = AtomicU64::new(0);
static WRITER_COMMIT_DELAY_MS: AtomicU64 = AtomicU64::new(0);
static WRITER_DELAY_LOCK: Mutex<()> = Mutex::new(());
pub(crate) fn writer_delay_guard() -> std::sync::MutexGuard<'static, ()> {
WRITER_DELAY_LOCK.lock().expect("writer delay lock should succeed")
}
pub(crate) fn set_writer_delays(flush: Duration, commit: Duration) {
WRITER_FLUSH_DELAY_MS.store(flush.as_millis() as u64, Ordering::Relaxed);
WRITER_COMMIT_DELAY_MS.store(commit.as_millis() as u64, Ordering::Relaxed);
}
pub(crate) struct BehaviorSpi {
pub(crate) fail_commit: bool,
pub(crate) fail_write: bool,
pub(crate) commit_failure: Option<WriteFailureState>,
pub(crate) abort_failure: Option<FsErrorKind>,
pub(crate) limits: FileSystemLimits,
pub(crate) entries: Mutex<Vec<DirEntry>>,
pub(crate) cleanup_calls: Arc<Mutex<usize>>,
pub(crate) persist_calls: Arc<Mutex<usize>>,
pub(crate) commit_calls: Arc<Mutex<usize>>,
pub(crate) abort_calls: Arc<Mutex<usize>>,
pub(crate) temp_path: Path,
pub(crate) temp_failure: Option<PersistFailureState>,
pub(crate) temp_keep_error: Option<FsErrorKind>,
pub(crate) temp_cleanup_error: Option<FsErrorKind>,
pub(crate) directory_persist_non_atomic: bool,
pub(crate) provider_open_error: bool,
}
pub(crate) fn filesystem(
fail_commit: bool,
entries: Vec<DirEntry>,
) -> (FileSystem, Arc<Mutex<usize>>, Arc<Mutex<usize>>) {
let cleanup_calls = Arc::new(Mutex::new(0));
let persist_calls = Arc::new(Mutex::new(0));
let commit_calls = Arc::new(Mutex::new(0));
let abort_calls = Arc::new(Mutex::new(0));
let spi = BehaviorSpi {
fail_commit,
fail_write: false,
commit_failure: None,
abort_failure: None,
limits: FileSystemLimits::unknown(),
entries: Mutex::new(entries),
cleanup_calls: Arc::clone(&cleanup_calls),
persist_calls: Arc::clone(&persist_calls),
commit_calls: Arc::clone(&commit_calls),
abort_calls: Arc::clone(&abort_calls),
temp_path: Path::parse("/temporary").expect("test path should parse"),
temp_failure: None,
temp_keep_error: None,
temp_cleanup_error: None,
directory_persist_non_atomic: false,
provider_open_error: false,
};
(
FileSystem::from_spi(spi).expect("facade should construct"),
cleanup_calls,
persist_calls,
)
}
pub(crate) fn limited_write_filesystem(maximum: u64) -> FileSystem {
FileSystem::from_spi(BehaviorSpi {
fail_commit: false,
fail_write: false,
commit_failure: None,
abort_failure: None,
limits: FileSystemLimits::unknown().with_max_write_bytes(FileSystemLimit::Maximum(maximum)),
entries: Mutex::new(Vec::new()),
cleanup_calls: Arc::new(Mutex::new(0)),
persist_calls: Arc::new(Mutex::new(0)),
commit_calls: Arc::new(Mutex::new(0)),
abort_calls: Arc::new(Mutex::new(0)),
temp_path: Path::parse("/temporary").expect("test path should parse"),
temp_failure: None,
temp_keep_error: None,
temp_cleanup_error: None,
directory_persist_non_atomic: false,
provider_open_error: false,
})
.expect("facade should construct")
}
pub(crate) fn stream_failure_filesystem() -> FileSystem {
FileSystem::from_spi(BehaviorSpi {
fail_commit: false,
fail_write: true,
commit_failure: None,
abort_failure: None,
limits: FileSystemLimits::unknown(),
entries: Mutex::new(Vec::new()),
cleanup_calls: Arc::new(Mutex::new(0)),
persist_calls: Arc::new(Mutex::new(0)),
commit_calls: Arc::new(Mutex::new(0)),
abort_calls: Arc::new(Mutex::new(0)),
temp_path: Path::parse("/temporary").expect("test path should parse"),
temp_failure: None,
temp_keep_error: None,
temp_cleanup_error: None,
directory_persist_non_atomic: false,
provider_open_error: false,
})
.expect("facade should construct")
}
pub(crate) fn provider_open_failure_filesystem() -> FileSystem {
FileSystem::from_spi(BehaviorSpi {
fail_commit: false,
fail_write: false,
commit_failure: None,
abort_failure: None,
limits: FileSystemLimits::unknown(),
entries: Mutex::new(Vec::new()),
cleanup_calls: Arc::new(Mutex::new(0)),
persist_calls: Arc::new(Mutex::new(0)),
commit_calls: Arc::new(Mutex::new(0)),
abort_calls: Arc::new(Mutex::new(0)),
temp_path: Path::parse("/temporary").expect("test path should parse"),
temp_failure: None,
temp_keep_error: None,
temp_cleanup_error: None,
directory_persist_non_atomic: false,
provider_open_error: true,
})
.expect("facade should construct")
}
pub(crate) fn writer_lifecycle_filesystem(
commit_failure: Option<WriteFailureState>,
abort_failure: Option<FsErrorKind>,
) -> FileSystem {
writer_lifecycle_filesystem_with_counts(commit_failure, abort_failure).0
}
pub(crate) fn writer_lifecycle_filesystem_with_counts(
commit_failure: Option<WriteFailureState>,
abort_failure: Option<FsErrorKind>,
) -> (FileSystem, Arc<Mutex<usize>>, Arc<Mutex<usize>>) {
let commit_calls = Arc::new(Mutex::new(0));
let abort_calls = Arc::new(Mutex::new(0));
let filesystem = FileSystem::from_spi(BehaviorSpi {
fail_commit: false,
fail_write: false,
commit_failure,
abort_failure,
limits: FileSystemLimits::unknown(),
entries: Mutex::new(Vec::new()),
cleanup_calls: Arc::new(Mutex::new(0)),
persist_calls: Arc::new(Mutex::new(0)),
commit_calls: Arc::clone(&commit_calls),
abort_calls: Arc::clone(&abort_calls),
temp_path: Path::parse("/temporary").expect("test path should parse"),
temp_failure: None,
temp_keep_error: None,
temp_cleanup_error: None,
directory_persist_non_atomic: false,
provider_open_error: false,
})
.expect("facade should construct");
(filesystem, commit_calls, abort_calls)
}
pub(crate) fn invalid_temp_path_filesystem() -> FileSystem {
FileSystem::from_spi(BehaviorSpi {
fail_commit: false,
fail_write: false,
commit_failure: None,
abort_failure: None,
limits: FileSystemLimits::unknown(),
entries: Mutex::new(Vec::new()),
cleanup_calls: Arc::new(Mutex::new(0)),
persist_calls: Arc::new(Mutex::new(0)),
commit_calls: Arc::new(Mutex::new(0)),
abort_calls: Arc::new(Mutex::new(0)),
temp_path: Path::parse("relative").expect("test path should parse"),
temp_failure: None,
temp_keep_error: None,
temp_cleanup_error: None,
directory_persist_non_atomic: false,
provider_open_error: false,
})
.expect("facade should construct")
}
pub(crate) fn wrong_temp_kind_filesystem() -> FileSystem {
FileSystem::from_spi(BehaviorSpi {
fail_commit: false,
fail_write: false,
commit_failure: None,
abort_failure: None,
limits: FileSystemLimits::unknown(),
entries: Mutex::new(Vec::new()),
cleanup_calls: Arc::new(Mutex::new(0)),
persist_calls: Arc::new(Mutex::new(0)),
commit_calls: Arc::new(Mutex::new(0)),
abort_calls: Arc::new(Mutex::new(0)),
temp_path: Path::parse("/wrong-kind").expect("test path should parse"),
temp_failure: None,
temp_keep_error: None,
temp_cleanup_error: None,
directory_persist_non_atomic: false,
provider_open_error: false,
})
.expect("facade should construct")
}
fn invalid_temp_cleanup_filesystem() -> FileSystem {
FileSystem::from_spi(BehaviorSpi {
fail_commit: false,
fail_write: false,
commit_failure: None,
abort_failure: None,
limits: FileSystemLimits::unknown(),
entries: Mutex::new(Vec::new()),
cleanup_calls: Arc::new(Mutex::new(0)),
persist_calls: Arc::new(Mutex::new(0)),
commit_calls: Arc::new(Mutex::new(0)),
abort_calls: Arc::new(Mutex::new(0)),
temp_path: Path::parse("/foreign").expect("test path should parse"),
temp_failure: None,
temp_keep_error: None,
temp_cleanup_error: Some(FsErrorKind::Io),
directory_persist_non_atomic: false,
provider_open_error: false,
})
.expect("facade should construct")
}
pub(crate) fn temp_failure_filesystem(
state: PersistFailureState,
) -> (FileSystem, Arc<Mutex<usize>>, Arc<Mutex<usize>>) {
temp_failure_with_cleanup_error_filesystem(state, None)
}
pub(crate) fn temp_failure_with_cleanup_error_filesystem(
state: PersistFailureState,
cleanup_error: Option<FsErrorKind>,
) -> (FileSystem, Arc<Mutex<usize>>, Arc<Mutex<usize>>) {
let cleanup_calls = Arc::new(Mutex::new(0));
let persist_calls = Arc::new(Mutex::new(0));
let filesystem = FileSystem::from_spi(BehaviorSpi {
fail_commit: false,
fail_write: false,
commit_failure: None,
abort_failure: None,
limits: FileSystemLimits::unknown(),
entries: Mutex::new(Vec::new()),
cleanup_calls: Arc::clone(&cleanup_calls),
persist_calls: Arc::clone(&persist_calls),
commit_calls: Arc::new(Mutex::new(0)),
abort_calls: Arc::new(Mutex::new(0)),
temp_path: Path::parse("/temporary").expect("test path should parse"),
temp_failure: Some(state),
temp_keep_error: None,
temp_cleanup_error: cleanup_error,
directory_persist_non_atomic: false,
provider_open_error: false,
})
.expect("facade should construct");
(filesystem, cleanup_calls, persist_calls)
}
pub(crate) fn temp_lifecycle_error_filesystem(
keep_error: Option<FsErrorKind>,
cleanup_error: Option<FsErrorKind>,
) -> (FileSystem, Arc<Mutex<usize>>) {
let cleanup_calls = Arc::new(Mutex::new(0));
let filesystem = FileSystem::from_spi(BehaviorSpi {
fail_commit: false,
fail_write: false,
commit_failure: None,
abort_failure: None,
limits: FileSystemLimits::unknown(),
entries: Mutex::new(Vec::new()),
cleanup_calls: Arc::clone(&cleanup_calls),
persist_calls: Arc::new(Mutex::new(0)),
commit_calls: Arc::new(Mutex::new(0)),
abort_calls: Arc::new(Mutex::new(0)),
temp_path: Path::parse("/temporary").expect("test path should parse"),
temp_failure: None,
temp_keep_error: keep_error,
temp_cleanup_error: cleanup_error,
directory_persist_non_atomic: false,
provider_open_error: false,
})
.expect("facade should construct");
(filesystem, cleanup_calls)
}
pub(crate) fn non_atomic_temp_directory_filesystem() -> FileSystem {
FileSystem::from_spi(BehaviorSpi {
fail_commit: false,
fail_write: false,
commit_failure: None,
abort_failure: None,
limits: FileSystemLimits::unknown(),
entries: Mutex::new(Vec::new()),
cleanup_calls: Arc::new(Mutex::new(0)),
persist_calls: Arc::new(Mutex::new(0)),
commit_calls: Arc::new(Mutex::new(0)),
abort_calls: Arc::new(Mutex::new(0)),
temp_path: Path::parse("/temporary").expect("test path should parse"),
temp_failure: None,
temp_keep_error: None,
temp_cleanup_error: None,
directory_persist_non_atomic: true,
provider_open_error: false,
})
.expect("facade should construct")
}
#[test]
fn test_handle_support_constructs_file_system() {
let (file_system, _, _) = filesystem(false, Vec::new());
assert_eq!("handles-test", file_system.properties().info().id().as_str());
}
#[test]
fn test_handle_support_dispatches_directory_and_delete_operations() {
let (file_system, _, _) = filesystem(false, Vec::new());
let path = Path::parse("/target").expect("test path should parse");
assert!(
!file_system
.create_directory(&path, CreateDirectoryOptions::default())
.expect("directory creation should succeed")
.already_existed()
);
assert!(
!file_system
.delete_file(&path, DeleteOptions::default())
.expect("file deletion should succeed")
.already_missing()
);
assert!(
!file_system
.delete_directory(&path, DeleteOptions::default())
.expect("directory deletion should succeed")
.already_missing()
);
}
#[test]
fn test_handle_support_dispatches_successful_facade_operations() {
let (file_system, _, _) = filesystem(false, Vec::new());
let source = Path::parse("/source").expect("test path should parse");
let target = Path::parse("/target").expect("test path should parse");
assert!(file_system.exists(&source).expect("stat should succeed"));
let mut directory = file_system
.list(&ListScope::Path((source).clone()), ListOptions::default())
.expect("list should succeed");
assert!(directory.next_entry().expect("stream should succeed").is_none());
let mut reader = file_system
.open_reader(&source, ReadOptions::default())
.expect("reader should open");
let mut bytes = [0_u8; 5];
assert_eq!(
5,
Input::read(&mut reader, &mut bytes).expect("reader should transfer bytes")
);
let mut writer = file_system
.open_writer(&target, WriteOptions::default())
.expect("writer should open");
Output::write_fully(&mut writer, b"bytes").expect("writer should accept bytes");
writer.commit().expect("writer should commit");
let mut temporary_file = file_system
.create_temp_file(TempOptions::default())
.expect("temporary file should open");
temporary_file.keep().expect("temporary file should be kept");
let mut temporary_directory = file_system
.create_temp_directory(TempOptions::default())
.expect("temporary directory should open");
temporary_directory.keep().expect("temporary directory should be kept");
}
#[test]
fn test_handle_support_uses_default_spi_copy_decline() {
let (file_system, _, _) = filesystem(false, Vec::new());
assert_eq!(
FileSystemCapabilitySupport::Guaranteed,
file_system
.properties()
.capabilities()
.support(FileSystemCapability::Copy),
"the fixture must reach the default TryCopy implementation",
);
let outcome = file_system
.copy(
&Path::parse("/source").expect("test path should parse"),
&Path::parse("/target").expect("test path should parse"),
CopyOptions::default(),
)
.expect("fallback should use reader and writer capabilities");
assert!(outcome.used_fallback());
}
#[test]
fn test_handle_support_enriches_open_and_temp_provider_failures() {
let file_system = provider_open_failure_filesystem();
let path = Path::parse("/target").expect("test path should parse");
let list = file_system
.list(&ListScope::Path(path.clone()), ListOptions::default())
.expect_err("list failure");
let writer = file_system
.open_writer(&path, WriteOptions::default())
.expect_err("writer failure");
let reader = file_system
.open_reader(&path, ReadOptions::default())
.expect_err("reader failure");
let temp_file = file_system
.create_temp_file(TempOptions::default())
.expect_err("temp failure");
let temp_directory = file_system
.create_temp_directory(TempOptions::default())
.expect_err("temp failure");
for error in [
&list,
writer.error(),
&reader,
temp_file.error(),
temp_directory.error(),
] {
assert_eq!(FsErrorKind::UnsupportedOperation, error.kind());
assert_eq!(Some("handles-test"), error.provider());
}
assert!(writer.recovery().is_none());
assert!(temp_file.recovery().is_none());
assert!(temp_directory.recovery().is_none());
}
#[test]
fn test_handle_support_keeps_pathless_temp_errors_pathless() {
let file_system = provider_open_failure_filesystem();
let file_error = file_system
.create_temp_file(TempOptions::default())
.expect_err("temporary-file provider failure should propagate");
let directory_error = file_system
.create_temp_directory(TempOptions::default())
.expect_err("temporary-directory provider failure should propagate");
for error in [file_error, directory_error] {
assert_eq!(FsOperation::CreateTemp, error.error().operation());
assert_eq!(None, error.error().path());
assert_eq!(Some("handles-test"), error.error().provider());
}
}
#[test]
fn test_handle_support_validates_temp_parent_before_provider_call() {
let file_system = provider_open_failure_filesystem();
let parent = Path::parse("relative").expect("relative test path should parse");
let file_error = file_system
.create_temp_file(TempOptions::default().with_parent(Some(parent.clone())))
.expect_err("invalid temporary-file parent must fail in the facade");
assert_eq!(FsErrorKind::InvalidPath, file_error.error().kind());
assert_eq!(FsOperation::CreateTemp, file_error.error().operation());
assert_eq!(Some(&parent), file_error.error().path());
let directory_error = file_system
.create_temp_directory(TempOptions::default().with_parent(Some(parent.clone())))
.expect_err("invalid temporary-directory parent must fail in the facade");
assert_eq!(FsErrorKind::InvalidPath, directory_error.error().kind());
assert_eq!(FsOperation::CreateTemp, directory_error.error().operation());
assert_eq!(Some(&parent), directory_error.error().path());
}
#[test]
fn test_handle_support_rejects_invalid_temp_identities_with_cleanup_failure() {
let file_system = invalid_temp_cleanup_filesystem();
for mut error in [
file_system
.create_temp_file(TempOptions::default())
.expect_err("foreign temporary file identity must be rejected"),
file_system
.create_temp_directory(TempOptions::default())
.expect_err("foreign temporary directory identity must be rejected"),
] {
assert_eq!(FsErrorKind::ProviderContractViolation, error.error().kind());
assert_eq!(FsOperation::ValidateProviderOutcome, error.error().operation());
let cleanup = error
.recovery_mut()
.expect("isolated session")
.cleanup()
.expect_err("cleanup injection");
assert!(cleanup.to_string().contains("injected temporary cleanup failure"));
assert_eq!(FsErrorKind::ProviderContractViolation, error.error().kind());
}
}
#[test]
fn test_handle_support_rejects_invalid_temp_paths_after_cleanup() {
for error in [
invalid_temp_path_filesystem()
.create_temp_file(TempOptions::default())
.expect_err("relative temporary file path must be rejected"),
invalid_temp_path_filesystem()
.create_temp_directory(TempOptions::default())
.expect_err("relative temporary directory path must be rejected"),
] {
assert_eq!(FsErrorKind::ProviderContractViolation, error.error().kind());
}
}
#[test]
fn test_handle_support_rejects_wrong_temp_kind() {
let error = wrong_temp_kind_filesystem()
.create_temp_file(TempOptions::default())
.expect_err("temporary-file kind must be validated");
assert_eq!(FsErrorKind::ProviderContractViolation, error.error().kind());
}
impl BehaviorSpi {
fn unsupported() -> FsError {
FsError::new(
FsErrorKind::UnsupportedOperation,
FsOperation::Other,
"unused test operation",
)
}
fn info(&self, kind: FileKind) -> OpenedFileInfo {
OpenedFileInfo::new(
FileSystemId::new(if self.temp_path.as_str() == "/foreign" {
"foreign"
} else {
"handles-test"
})
.expect("valid test id"),
self.temp_path.clone(),
)
.with_metadata(FileMetadata::new(kind))
}
}
impl FileSystemSpi for BehaviorSpi {
fn properties(&self) -> ProviderProperties {
ProviderProperties::new(
FileSystemInfo::new(
FileSystemId::new("handles-test").expect("valid test id"),
"handles-test",
PathSemantics::Hierarchical,
),
ProviderOperations::new()
.with(ProviderOperation::Stat)
.with(ProviderOperation::List)
.with(ProviderOperation::OpenReader)
.with(ProviderOperation::OpenWriter)
.with(ProviderOperation::CreateDirectory)
.with(ProviderOperation::DeleteFile)
.with(ProviderOperation::DeleteDirectory)
.with(ProviderOperation::TryCopy)
.with(ProviderOperation::Rename)
.with(ProviderOperation::CreateTempFile)
.with(ProviderOperation::CreateTempDirectory),
FileSystemCapabilities::new()
.with_guaranteed(FileSystemCapability::List)
.with_guaranteed(FileSystemCapability::Copy)
.with_guaranteed(FileSystemCapability::Read)
.with_guaranteed(FileSystemCapability::Write)
.with_guaranteed(FileSystemCapability::DurableWrite)
.with_guaranteed(FileSystemCapability::CreateDirectory)
.with_guaranteed(FileSystemCapability::Delete)
.with_guaranteed(FileSystemCapability::AtomicReplace)
.with_guaranteed(FileSystemCapability::TempFile)
.with_guaranteed(FileSystemCapability::TempDirectory)
.with_guaranteed(FileSystemCapability::AtomicTempPersist),
self.limits,
PathConstraints::absolute(),
SymlinkPolicy::Reject,
)
.expect("valid test properties")
}
fn stat(&self, request: StatRequest<'_>) -> FsResult<StatResponse> {
let _ = request.options();
Ok(StatResponse::new(
request.path().clone(),
FileMetadata::new(FileKind::File),
))
}
fn list(&self, request: ListRequest<'_>) -> FsResult<OpenedDirectoryStream> {
let _ = request.scope().path().expect("path-scoped test request");
let _ = request.options();
if self.provider_open_error {
return Err(Self::unsupported());
}
Ok(OpenedDirectoryStream::new(Box::new(Entries(std::mem::take(
&mut *self.entries.lock().expect("entries lock should succeed"),
)))))
}
fn open_reader(&self, request: OpenReaderRequest<'_>) -> FsResult<OpenedReader> {
if self.provider_open_error {
return Err(Self::unsupported());
}
let _ = request.options();
Ok(OpenedReader::new(
OpenedFileInfo::new(
FileSystemId::new("handles-test").expect("valid test id"),
request.path().clone(),
),
Box::new(Cursor::new(b"bytes".to_vec())),
))
}
fn open_writer(&self, request: OpenWriterRequest<'_>) -> FsResult<OpenedWriter> {
let _ = request.options();
if self.provider_open_error {
return Err(Self::unsupported());
}
Ok(OpenedWriter::new(
OpenedFileInfo::new(
FileSystemId::new("handles-test").expect("valid test id"),
request.path().clone(),
),
Box::new(Writer {
fail_commit: self.fail_commit,
fail_write: self.fail_write,
commit_failure: self.commit_failure,
abort_failure: self.abort_failure,
non_atomic_commit: self.directory_persist_non_atomic,
commit_calls: Arc::clone(&self.commit_calls),
abort_calls: Arc::clone(&self.abort_calls),
}),
))
}
fn create_directory(&self, request: CreateDirectoryRequest<'_>) -> FsResult<CreateDirectoryOutcome> {
let _ = request.path();
let _ = request.options();
Ok(CreateDirectoryOutcome::new(false))
}
fn delete_file(&self, request: DeleteFileRequest<'_>) -> FsResult<DeleteOutcome> {
let _ = request.path();
let _ = request.options();
Ok(DeleteOutcome::new(false))
}
fn delete_directory(&self, request: DeleteDirectoryRequest<'_>) -> FsResult<DeleteOutcome> {
let _ = request.path();
let _ = request.options();
Ok(DeleteOutcome::new(false))
}
fn rename(&self, request: RenameRequest<'_>) -> Result<RenameOutcome, SpiRenameFailure> {
let _ = request.source();
let _ = request.target();
let _ = request.options();
Err(SpiRenameFailure::new(
Self::unsupported(),
RenameFailureState::Unchanged,
))
}
fn create_temp_file(&self, request: CreateTempFileRequest) -> FsResult<OpenedTempFile> {
let _ = request.options();
if self.provider_open_error {
return Err(Self::unsupported());
}
Ok(OpenedTempFile::new(
self.info(if self.temp_path.as_str() == "/wrong-kind" {
FileKind::Directory
} else {
FileKind::File
}),
Box::new(Temp {
cleanup_calls: Arc::clone(&self.cleanup_calls),
persist_calls: Arc::clone(&self.persist_calls),
non_atomic: true,
failure: self.temp_failure,
keep_error: self.temp_keep_error,
cleanup_error: self.temp_cleanup_error,
}),
))
}
fn create_temp_directory(&self, request: CreateTempDirectoryRequest) -> FsResult<OpenedTempDirectory> {
let _ = request.options();
if self.provider_open_error {
return Err(Self::unsupported());
}
Ok(OpenedTempDirectory::new(
self.info(FileKind::Directory),
Box::new(Temp {
cleanup_calls: Arc::clone(&self.cleanup_calls),
persist_calls: Arc::clone(&self.persist_calls),
non_atomic: self.directory_persist_non_atomic,
failure: self.temp_failure,
keep_error: self.temp_keep_error,
cleanup_error: self.temp_cleanup_error,
}),
))
}
}
struct Writer {
fail_commit: bool,
fail_write: bool,
commit_failure: Option<WriteFailureState>,
abort_failure: Option<FsErrorKind>,
non_atomic_commit: bool,
commit_calls: Arc<Mutex<usize>>,
abort_calls: Arc<Mutex<usize>>,
}
impl Output for Writer {
type Item = u8;
unsafe fn write_unchecked(&mut self, _: &[u8], _: usize, count: usize) -> IoResult<usize> {
if self.fail_write {
Err(std::io::Error::new(
std::io::ErrorKind::InvalidData,
"stream secret=top-secret",
))
} else {
Ok(count)
}
}
fn flush(&mut self) -> IoResult<()> {
let delay = WRITER_FLUSH_DELAY_MS.load(Ordering::Relaxed);
if delay != 0 {
std::thread::sleep(Duration::from_millis(delay));
}
Ok(())
}
}
impl FileWriterSpi for Writer {
fn commit(&mut self) -> Result<WriteOutcome, SpiWriteFailure> {
let delay = WRITER_COMMIT_DELAY_MS.load(Ordering::Relaxed);
if delay != 0 {
std::thread::sleep(Duration::from_millis(delay));
}
*self.commit_calls.lock().expect("commit counter lock should succeed") += 1;
if let Some(state) = self.commit_failure {
return Err(SpiWriteFailure::new(
FsError::new(FsErrorKind::Io, FsOperation::CommitWriter, "injected commit failure"),
state,
));
}
if self.fail_commit {
Err(SpiWriteFailure::new(
FsError::new(FsErrorKind::Io, FsOperation::CommitWriter, "injected commit failure"),
WriteFailureState::RetryableNotPublished,
))
} else {
Ok(WriteOutcome::new(
if self.non_atomic_commit {
AchievedAtomicity::NonAtomic
} else {
AchievedAtomicity::Atomic
},
PublicationMethod::Direct,
))
}
}
fn abort(&mut self) -> FsResult<WriteAbortOutcome> {
*self.abort_calls.lock().expect("abort counter lock should succeed") += 1;
match self.abort_failure {
Some(kind) => Err(FsError::new(kind, FsOperation::AbortWriter, "injected abort failure")),
None => Ok(match self.commit_failure {
Some(WriteFailureState::Published) => WriteAbortOutcome::Published,
Some(WriteFailureState::Indeterminate) => WriteAbortOutcome::Indeterminate,
Some(WriteFailureState::RetryableNotPublished) | Some(WriteFailureState::NotPublished) | None => {
WriteAbortOutcome::NotPublished
}
}),
}
}
}
struct Entries(Vec<DirEntry>);
impl DirectoryStreamSpi for Entries {
fn next_entry(&mut self) -> FsResult<Option<DirEntry>> {
Ok(self.0.pop())
}
}
struct Temp {
cleanup_calls: Arc<Mutex<usize>>,
persist_calls: Arc<Mutex<usize>>,
non_atomic: bool,
failure: Option<PersistFailureState>,
keep_error: Option<FsErrorKind>,
cleanup_error: Option<FsErrorKind>,
}
impl TempResourceSpi for Temp {
fn persist(&mut self, request: PersistRequest<'_>) -> Result<PersistOutcome, SpiPersistFailure> {
*self.persist_calls.lock().expect("persist lock should succeed") += 1;
if let Some(state) = self.failure {
return Err(SpiPersistFailure::new(
FsError::new(
FsErrorKind::Io,
FsOperation::PersistTemp,
"injected temporary persist failure",
),
state,
));
}
let target = if request.target().as_str() == "/wrong-persist-target" {
Path::parse("/reported-persist-target").expect("generated path should parse")
} else {
request.target().clone()
};
Ok(PersistOutcome::new(
target,
if self.non_atomic {
AchievedAtomicity::NonAtomic
} else {
AchievedAtomicity::Atomic
},
PublicationMethod::Direct,
))
}
fn keep(&mut self) -> Result<PersistOutcome, SpiPersistFailure> {
if let Some(state) = self.failure {
return Err(SpiPersistFailure::new(
FsError::new(FsErrorKind::Io, FsOperation::KeepTemp, "injected keep failure")
.with_target(Path::parse("/kept-resource").expect("generated target")),
state,
));
}
self.keep_error.map_or_else(
|| {
Ok(PersistOutcome::new(
Path::parse("/kept-resource").expect("generated path should parse"),
AchievedAtomicity::Atomic,
PublicationMethod::Direct,
))
},
|kind| {
Err(SpiPersistFailure::new(
FsError::new(kind, FsOperation::KeepTemp, "injected temporary keep failure"),
if kind == FsErrorKind::Indeterminate {
PersistFailureState::Indeterminate
} else {
PersistFailureState::NotPublished
},
))
},
)
}
fn cleanup(&mut self) -> FsResult<()> {
*self.cleanup_calls.lock().expect("cleanup lock should succeed") += 1;
self.cleanup_error.map_or(Ok(()), |kind| {
Err(FsError::new(
kind,
FsOperation::CleanupTemp,
"injected temporary cleanup failure",
))
})
}
}