use anyhow::{Context, Result, bail};
use futures::stream::{self, StreamExt};
use opendal::{Operator, services};
use std::io::{IsTerminal, Write};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
use tokio::io::AsyncReadExt;
const UPLOAD_CONCURRENCY: usize = 8;
const CHUNK: usize = 1 << 20;
pub fn object_key(hash: &str) -> String {
format!("objects/{}/{}", &hash[..2], &hash[2..])
}
pub struct Remote {
op: Operator,
rt: tokio::runtime::Runtime,
}
pub fn open(url: &str) -> Result<Remote> {
let op = build_operator(url)?;
let rt = tokio::runtime::Builder::new_multi_thread()
.enable_all()
.build()?;
Ok(Remote { op, rt })
}
fn build_operator(url: &str) -> Result<Operator> {
let local_path = url.strip_prefix("local:").unwrap_or(url);
let is_s3 = url.starts_with("s3://");
if is_s3 {
let rest = url.trim_start_matches("s3://");
let (bucket, root) = rest.split_once('/').unwrap_or((rest, ""));
if bucket.is_empty() {
bail!("s3 remote needs a bucket: s3://<bucket>[/<root>]");
}
let mut b = services::S3::default().bucket(bucket);
if !root.is_empty() {
b = b.root(root);
}
let region = std::env::var("AWS_REGION").unwrap_or_else(|_| "auto".into());
b = b.region(®ion);
if let Ok(ep) = std::env::var("AWS_ENDPOINT_URL") {
b = b.endpoint(&ep);
}
return Ok(Operator::new(b)?.finish());
}
std::fs::create_dir_all(local_path)
.with_context(|| format!("creating remote dir {local_path}"))?;
let tmp_dir = Path::new(local_path).join(".stowe-tmp");
std::fs::create_dir_all(&tmp_dir)?;
Ok(Operator::new(
services::Fs::default()
.root(local_path)
.atomic_write_dir(&tmp_dir.to_string_lossy()),
)?
.finish())
}
impl Remote {
pub fn exists(&self, key: &str) -> Result<bool> {
self.rt.block_on(async { Ok(self.op.exists(key).await?) })
}
pub fn put_bytes(&self, key: &str, data: &[u8]) -> Result<()> {
self.rt.block_on(async {
self.op.write(key, data.to_vec()).await?;
Ok(())
})
}
pub fn get_bytes(&self, key: &str) -> Result<Vec<u8>> {
self.rt
.block_on(async { Ok(self.op.read(key).await?.to_vec()) })
}
pub fn get_file(&self, key: &str, dest: &Path) -> Result<()> {
if let Some(parent) = dest.parent() {
std::fs::create_dir_all(parent)?;
}
let bytes = self.get_bytes(key)?;
std::fs::write(dest, bytes)?;
Ok(())
}
pub fn put_files(&self, items: Vec<(String, PathBuf)>) -> Result<usize> {
let total = items.len();
let done = Arc::new(AtomicUsize::new(0));
let show = std::io::stderr().is_terminal();
self.rt.block_on(async {
let op = &self.op;
let done = &done;
let results: Vec<Result<usize>> = stream::iter(items)
.map(|(key, src)| async move {
let r = if op.exists(&key).await? {
Ok(0usize)
} else {
upload_one(op, &key, &src).await?;
Ok(1usize)
};
let n = done.fetch_add(1, Ordering::Relaxed) + 1;
if show {
eprint!("\r\x1b[Kpushing... {n}/{total}");
let _ = std::io::stderr().flush();
}
r
})
.buffer_unordered(UPLOAD_CONCURRENCY)
.collect()
.await;
if show && total > 0 {
eprint!("\r\x1b[K"); let _ = std::io::stderr().flush();
}
let mut uploaded = 0;
for r in results {
uploaded += r?;
}
Ok(uploaded)
})
}
}
async fn upload_one(op: &Operator, key: &str, src: &Path) -> Result<()> {
let mut file = tokio::fs::File::open(src)
.await
.with_context(|| format!("opening {}", src.display()))?;
let mut writer = op.writer(key).await?;
let mut buf = vec![0u8; CHUNK];
loop {
let n = file.read(&mut buf).await?;
if n == 0 {
break;
}
writer.write(buf[..n].to_vec()).await?;
}
writer.close().await?;
Ok(())
}