Stream purchased files into durable cache and avoid duplicate concurrent payments
This commit is contained in:
@@ -212,6 +212,27 @@ pub async fn save_new_file(dir: &Path, name: &str, bytes: &[u8]) -> Result<PathB
|
||||
save_new_file_with(dir, name, bytes, write_via_userns).await
|
||||
}
|
||||
|
||||
/// Copy the already-owned file with bounded buffers, retaining the no-clobber
|
||||
/// and Files namespace rules used by small purchases.
|
||||
pub async fn save_new_file_from(dir: &Path, name: &str, mut source: fs::File) -> Result<PathBuf> {
|
||||
use tokio::io::AsyncSeekExt;
|
||||
validate_filename(name)?;
|
||||
match fs::symlink_metadata(dir).await {
|
||||
Ok(meta) => anyhow::ensure!(meta.is_dir(), "Files destination is not a directory"),
|
||||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
|
||||
Err(error) => return Err(error.into()),
|
||||
}
|
||||
source.seek(std::io::SeekFrom::Start(0)).await?;
|
||||
match write_direct_stream(dir, name, &mut source).await {
|
||||
Ok(path) => Ok(path),
|
||||
Err(error) if error.kind() == std::io::ErrorKind::PermissionDenied => {
|
||||
source.seek(std::io::SeekFrom::Start(0)).await?;
|
||||
write_via_userns_stream(dir.to_owned(), name.to_owned(), source).await
|
||||
}
|
||||
Err(error) => Err(error).context("Saving purchased file"),
|
||||
}
|
||||
}
|
||||
|
||||
fn validate_filename(name: &str) -> Result<()> {
|
||||
anyhow::ensure!(
|
||||
!name.is_empty()
|
||||
@@ -291,8 +312,15 @@ impl Drop for PendingFile {
|
||||
}
|
||||
|
||||
async fn write_direct(dir: &Path, name: &str, bytes: &[u8]) -> std::io::Result<PathBuf> {
|
||||
write_direct_stream(dir, name, &mut &bytes[..]).await
|
||||
}
|
||||
|
||||
async fn write_direct_stream<R: tokio::io::AsyncRead + Unpin>(
|
||||
dir: &Path,
|
||||
name: &str,
|
||||
source: &mut R,
|
||||
) -> std::io::Result<PathBuf> {
|
||||
use std::os::unix::fs::PermissionsExt;
|
||||
use tokio::io::AsyncWriteExt;
|
||||
fs::create_dir_all(dir).await?;
|
||||
let temp_path = dir.join(format!(".archy-saving-{}", uuid::Uuid::new_v4()));
|
||||
let mut file = fs::OpenOptions::new()
|
||||
@@ -302,7 +330,7 @@ async fn write_direct(dir: &Path, name: &str, bytes: &[u8]) -> std::io::Result<P
|
||||
.open(&temp_path)
|
||||
.await?;
|
||||
let temp = PendingFile(temp_path);
|
||||
file.write_all(bytes).await?;
|
||||
tokio::io::copy(source, &mut file).await?;
|
||||
file.set_permissions(std::fs::Permissions::from_mode(0o644))
|
||||
.await?;
|
||||
file.sync_all().await?;
|
||||
@@ -397,10 +425,77 @@ async fn write_via_userns(dir: PathBuf, name: String, bytes: Vec<u8>) -> Result<
|
||||
.context("Files namespace writer timed out")?
|
||||
}
|
||||
|
||||
async fn write_via_userns_stream(
|
||||
dir: PathBuf,
|
||||
name: String,
|
||||
mut source: fs::File,
|
||||
) -> Result<PathBuf> {
|
||||
let expected = source.metadata().await?.len();
|
||||
let mut child = tokio::process::Command::new("podman")
|
||||
.args(["unshare", "sh", "-c", WRITE_VIA_USERNS, "sh"])
|
||||
.arg(&dir)
|
||||
.arg(&name)
|
||||
.arg(expected.to_string())
|
||||
.kill_on_drop(true)
|
||||
.stdin(std::process::Stdio::piped())
|
||||
.stdout(std::process::Stdio::piped())
|
||||
.stderr(std::process::Stdio::null())
|
||||
.spawn()
|
||||
.context("Starting Files namespace writer")?;
|
||||
let mut stdin = child.stdin.take().context("Files writer stdin missing")?;
|
||||
tokio::time::timeout(std::time::Duration::from_secs(900), async {
|
||||
let count = tokio::io::copy(&mut source, &mut stdin).await?;
|
||||
drop(stdin);
|
||||
let output = child.wait_with_output().await?;
|
||||
anyhow::ensure!(
|
||||
count == expected && output.status.success(),
|
||||
"Files copy failed; purchased cache is retained"
|
||||
);
|
||||
let chosen = String::from_utf8(output.stdout).context("Invalid Files response")?;
|
||||
validate_filename(&chosen)?;
|
||||
anyhow::ensure!(
|
||||
(1..=100).any(|n| numbered_name(&name, n) == chosen),
|
||||
"Unexpected Files destination"
|
||||
);
|
||||
Ok(dir.join(chosen))
|
||||
})
|
||||
.await
|
||||
.context("Files namespace writer timed out")?
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
#[tokio::test]
|
||||
async fn streamed_purchase_preserves_existing_file_and_exact_large_copy() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
fs::write(dir.path().join("film.mp4"), b"keep")
|
||||
.await
|
||||
.unwrap();
|
||||
let source = dir.path().join("source");
|
||||
let bytes = vec![17; 2 * 1024 * 1024];
|
||||
fs::write(&source, &bytes).await.unwrap();
|
||||
let path = save_new_file_from(
|
||||
dir.path(),
|
||||
"film.mp4",
|
||||
fs::File::open(&source).await.unwrap(),
|
||||
)
|
||||
.await
|
||||
.unwrap();
|
||||
assert_eq!(path.file_name().unwrap(), "film (2).mp4");
|
||||
assert_eq!(fs::read(path).await.unwrap(), bytes);
|
||||
assert_eq!(
|
||||
fs::read(dir.path().join("film.mp4")).await.unwrap(),
|
||||
b"keep"
|
||||
);
|
||||
assert!(!std::fs::read_dir(dir.path()).unwrap().any(|entry| entry
|
||||
.unwrap()
|
||||
.file_name()
|
||||
.to_string_lossy()
|
||||
.starts_with(".archy-saving")));
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn cloud_credentials_use_unique_record_and_never_default_password() {
|
||||
let dir = tempfile::tempdir().unwrap();
|
||||
|
||||
Reference in New Issue
Block a user