pub(super) mod attributes;
mod journal;
mod native;
mod objects;
mod publication;
pub(super) mod supervision;
#[cfg(test)]
mod tests;
use super::{MacFiles, policy, protocol};
use cueward_core::files::mutation::record_error;
pub use cueward_core::files::mutation::{
CopyVerification, Creation, MutationAction, MutationPlatform, MutationReceipt, MutationRequest,
MutationStage, MutationStatus, Publication, Staging,
};
use cueward_core::files::*;
use serde::{Deserialize, Serialize};
use std::path::{Path, PathBuf};
#[derive(Debug, Serialize, Deserialize)]
pub struct MutationWorkerRequest {
pub operation_id: String,
}
pub fn run(
executable: &Path,
request: &MutationRequest,
timeout_ms: u64,
) -> Result<MutationReceipt, FileError> {
if !(1..=30000).contains(&timeout_ms) {
return Err(FileError::new(
FileErrorCode::InvalidOptions,
"timeout-ms must be 1..30000",
));
}
let payload = serde_json::to_vec(request)
.map_err(|e| FileError::new(FileErrorCode::InvalidOptions, e.to_string()))?;
if payload.len() > protocol::MAX_REQUEST_BYTES {
return Err(FileError::new(
FileErrorCode::InvalidOptions,
"request exceeds 16 KiB",
));
}
let prepared = journal::create(request)?;
let input = MutationWorkerRequest {
operation_id: prepared.operation_id.clone(),
};
supervision::announce(&prepared.operation_id, "receipt")?;
let result = supervision::run(executable, "files-mutation-worker", &input, timeout_ms)
.and_then(|receipt| {
protocol::verify_receipt(
&prepared,
receipt,
&journal::load(&prepared.operation_id)?,
"worker result differs from its durable operation receipt",
)
});
match result {
Ok(receipt) => Ok(receipt),
Err(mut error) => {
error.message.push_str("; destination may exist or be partial; inspect the receipt and paths, do not automatically retry");
let mut receipt = journal::load(&prepared.operation_id).unwrap_or(prepared);
receipt.status = MutationStatus::Uncertain;
receipt.completion_verified = false;
receipt.error = Some(error);
journal::save(&receipt)?;
Ok(receipt)
}
}
}
pub fn read_receipt(operation_id: &str) -> Result<MutationReceipt, FileError> {
journal::load(operation_id)
}
pub fn execute_worker(input: &MutationWorkerRequest) -> Result<MutationReceipt, FileError> {
journal::claim(&input.operation_id)?;
let mut receipt = journal::load(&input.operation_id)?;
if receipt.stage != MutationStage::Preflight
|| receipt.mutation_attempted
|| receipt.error.is_some()
{
return Err(FileError::new(
FileErrorCode::InvalidOptions,
"operation has already progressed; replay is not supported",
));
}
let result = policy::NoMaterialization::enter().and_then(|_policy| {
cueward_core::files::mutation::execute(&MacFiles, &mut receipt, &mut journal::save)
});
if let Err(error) = result {
record_error(&mut receipt, error);
}
journal::save(&receipt)?;
Ok(receipt)
}
pub fn read_supervised_request() -> Result<MutationWorkerRequest, FileError> {
supervision::read_request()
}