use anyhow::{bail, Result};
use opendal::{services, Operator};
pub fn operator_for_uri(uri: &str) -> Result<(Operator, String)> {
let (scheme, rest) = uri
.split_once("://")
.ok_or_else(|| anyhow::anyhow!("not a URI: {uri}"))?;
match scheme {
"mem" => {
let op = Operator::new(services::Memory::default())?.finish();
Ok((op, rest.to_string()))
}
"s3" => {
let (bucket, key) = rest.split_once('/').unwrap_or((rest, ""));
let region = std::env::var("AWS_DEFAULT_REGION")
.unwrap_or_else(|_| "us-east-1".into());
let builder = services::S3::default().bucket(bucket).region(®ion);
let op = Operator::new(builder)?.finish();
Ok((op, key.to_string()))
}
"gcs" => {
let (bucket, key) = rest.split_once('/').unwrap_or((rest, ""));
let builder = services::Gcs::default().bucket(bucket);
let op = Operator::new(builder)?.finish();
Ok((op, key.to_string()))
}
"azblob" => {
let (container, blob) = rest.split_once('/').unwrap_or((rest, ""));
let account = std::env::var("AZURE_STORAGE_ACCOUNT")
.unwrap_or_else(|_| "devstoreaccount1".into());
let builder = services::Azblob::default()
.container(container)
.account_name(&account);
let op = Operator::new(builder)?.finish();
Ok((op, blob.to_string()))
}
"azdls" => {
let (filesystem, path) = rest.split_once('/').unwrap_or((rest, ""));
let account = std::env::var("AZURE_STORAGE_ACCOUNT")
.unwrap_or_else(|_| "devstoreaccount1".into());
let endpoint = std::env::var("AZDLS_ENDPOINT").unwrap_or_else(|_| {
format!("https://{account}.dfs.core.windows.net")
});
let builder = services::Azdls::default()
.filesystem(filesystem)
.endpoint(&endpoint)
.account_name(&account);
let op = Operator::new(builder)?.finish();
Ok((op, path.to_string()))
}
"azfile" => {
let (share, path) = rest.split_once('/').unwrap_or((rest, ""));
let account = std::env::var("AZURE_STORAGE_ACCOUNT")
.unwrap_or_else(|_| "devstoreaccount1".into());
let endpoint = std::env::var("AZFILE_ENDPOINT").unwrap_or_else(|_| {
format!("https://{account}.file.core.windows.net")
});
let builder = services::Azfile::default()
.share_name(share)
.endpoint(&endpoint)
.account_name(&account);
let op = Operator::new(builder)?.finish();
Ok((op, path.to_string()))
}
"b2" => {
let (bucket, key) = rest.split_once('/').unwrap_or((rest, ""));
let key_id = std::env::var("BACKBLAZE_APPLICATION_KEY_ID").unwrap_or_default();
let app_key = std::env::var("BACKBLAZE_APPLICATION_KEY").unwrap_or_default();
let builder = services::B2::default()
.bucket(bucket)
.application_key_id(&key_id)
.application_key(&app_key);
let op = Operator::new(builder)?.finish();
Ok((op, key.to_string()))
}
"cos" => {
let (bucket, key) = rest.split_once('/').unwrap_or((rest, ""));
let region = std::env::var("TENCENTCLOUD_REGION")
.unwrap_or_else(|_| "ap-guangzhou".into());
let endpoint = std::env::var("COS_ENDPOINT").unwrap_or_else(|_| {
format!("https://{bucket}.cos.{region}.myqcloud.com")
});
let secret_id = std::env::var("TENCENTCLOUD_SECRET_ID").unwrap_or_default();
let secret_key = std::env::var("TENCENTCLOUD_SECRET_KEY").unwrap_or_default();
let builder = services::Cos::default()
.bucket(bucket)
.endpoint(&endpoint)
.secret_id(&secret_id)
.secret_key(&secret_key);
let op = Operator::new(builder)?.finish();
Ok((op, key.to_string()))
}
"obs" => {
let (bucket, key) = rest.split_once('/').unwrap_or((rest, ""));
let region = std::env::var("HUAWEI_REGION")
.unwrap_or_else(|_| "cn-north-4".into());
let endpoint = std::env::var("OBS_ENDPOINT").unwrap_or_else(|_| {
format!("https://obs.{region}.myhuaweicloud.com")
});
let access_key = std::env::var("HUAWEI_ACCESS_KEY_ID").unwrap_or_default();
let secret_key = std::env::var("HUAWEI_SECRET_ACCESS_KEY").unwrap_or_default();
let builder = services::Obs::default()
.bucket(bucket)
.endpoint(&endpoint)
.access_key_id(&access_key)
.secret_access_key(&secret_key);
let op = Operator::new(builder)?.finish();
Ok((op, key.to_string()))
}
"oss" => {
let (bucket, key) = rest.split_once('/').unwrap_or((rest, ""));
let region = std::env::var("ALIBABA_CLOUD_REGION")
.unwrap_or_else(|_| "cn-hangzhou".into());
let endpoint = std::env::var("OSS_ENDPOINT").unwrap_or_else(|_| {
format!("https://oss-{region}.aliyuncs.com")
});
let access_key = std::env::var("ALIBABA_CLOUD_ACCESS_KEY_ID").unwrap_or_default();
let access_secret = std::env::var("ALIBABA_CLOUD_ACCESS_KEY_SECRET").unwrap_or_default();
let builder = services::Oss::default()
.bucket(bucket)
.endpoint(&endpoint)
.access_key_id(&access_key)
.access_key_secret(&access_secret);
let op = Operator::new(builder)?.finish();
Ok((op, key.to_string()))
}
"swift" => {
let (container, path) = rest.split_once('/').unwrap_or((rest, ""));
let endpoint = std::env::var("SWIFT_ENDPOINT")
.unwrap_or_else(|_| "https://object.example.com".into());
let token = std::env::var("SWIFT_TOKEN").unwrap_or_default();
let builder = services::Swift::default()
.endpoint(&endpoint)
.container(container)
.token(&token);
let op = Operator::new(builder)?.finish();
Ok((op, path.to_string()))
}
"upyun" => {
let (bucket, key) = rest.split_once('/').unwrap_or((rest, ""));
let operator_name = std::env::var("UPYUN_OPERATOR").unwrap_or_default();
let password = std::env::var("UPYUN_PASSWORD").unwrap_or_default();
let builder = services::Upyun::default()
.bucket(bucket)
.operator(&operator_name)
.password(&password);
let op = Operator::new(builder)?.finish();
Ok((op, key.to_string()))
}
"onedrive" => {
let token = std::env::var("ONEDRIVE_ACCESS_TOKEN").unwrap_or_default();
let builder = services::Onedrive::default()
.root("/")
.access_token(&token);
let op = Operator::new(builder)?.finish();
Ok((op, rest.to_string()))
}
"dropbox" => {
let token = std::env::var("DROPBOX_ACCESS_TOKEN").unwrap_or_default();
let builder = services::Dropbox::default()
.root("/")
.access_token(&token);
let op = Operator::new(builder)?.finish();
Ok((op, rest.to_string()))
}
"aliyun-drive" => {
let token = std::env::var("ALIYUN_DRIVE_ACCESS_TOKEN").unwrap_or_default();
let builder = services::AliyunDrive::default()
.root("/")
.access_token(&token);
let op = Operator::new(builder)?.finish();
Ok((op, rest.to_string()))
}
"yandex-disk" => {
let token = std::env::var("YANDEX_DISK_ACCESS_TOKEN").unwrap_or_default();
let builder = services::YandexDisk::default()
.root("/")
.access_token(&token);
let op = Operator::new(builder)?.finish();
Ok((op, rest.to_string()))
}
"pcloud" => {
let endpoint = std::env::var("PCLOUD_ENDPOINT")
.unwrap_or_else(|_| "https://api.pcloud.com".into());
let username = std::env::var("PCLOUD_USERNAME").unwrap_or_default();
let password = std::env::var("PCLOUD_PASSWORD").unwrap_or_default();
let builder = services::Pcloud::default()
.root("/")
.endpoint(&endpoint)
.username(&username)
.password(&password);
let op = Operator::new(builder)?.finish();
Ok((op, rest.to_string()))
}
"koofr" => {
let endpoint = std::env::var("KOOFR_ENDPOINT")
.unwrap_or_else(|_| "https://app.koofr.net".into());
let email = std::env::var("KOOFR_EMAIL").unwrap_or_default();
let password = std::env::var("KOOFR_PASSWORD").unwrap_or_default();
let builder = services::Koofr::default()
.root("/")
.endpoint(&endpoint)
.email(&email)
.password(&password);
let op = Operator::new(builder)?.finish();
Ok((op, rest.to_string()))
}
"seafile" => {
let (server, rest_path) = rest.split_once('/').unwrap_or((rest, ""));
let (repo, path) = rest_path.split_once('/').unwrap_or((rest_path, ""));
let endpoint = format!("https://{server}");
let username = std::env::var("SEAFILE_USERNAME").unwrap_or_default();
let password = std::env::var("SEAFILE_PASSWORD").unwrap_or_default();
let repo_name = if repo.is_empty() {
std::env::var("SEAFILE_REPO").unwrap_or_else(|_| "My Library".into())
} else {
repo.to_string()
};
let builder = services::Seafile::default()
.endpoint(&endpoint)
.username(&username)
.password(&password)
.repo_name(&repo_name);
let op = Operator::new(builder)?.finish();
Ok((op, path.to_string()))
}
"github" => {
let mut parts = rest.splitn(3, '/');
let owner = parts.next().unwrap_or("");
let repo = parts.next().unwrap_or("");
let path = parts.next().unwrap_or("").to_string();
let token = std::env::var("GITHUB_TOKEN").unwrap_or_default();
let builder = services::Github::default()
.token(&token)
.owner(owner)
.repo(repo);
let op = Operator::new(builder)?.finish();
Ok((op, path))
}
"huggingface" => {
let (repo_id, path) = rest.split_once('/').map(|(a, b)| {
let full = format!("{a}/{b}");
if let Some(idx) = full.find('/') {
let second = full[idx+1..].find('/');
if let Some(second_idx) = second {
let split_at = idx + 1 + second_idx;
(full[..split_at].to_string(), full[split_at+1..].to_string())
} else {
(full, String::new())
}
} else {
(full, String::new())
}
}).unwrap_or((rest.to_string(), String::new()));
let token = std::env::var("HUGGINGFACE_TOKEN").unwrap_or_default();
let builder = services::Huggingface::default()
.repo_id(&repo_id)
.token(&token);
let op = Operator::new(builder)?.finish();
Ok((op, path))
}
"vercel-blob" => {
let token = std::env::var("BLOB_READ_WRITE_TOKEN").unwrap_or_default();
let builder = services::VercelBlob::default().token(&token);
let op = Operator::new(builder)?.finish();
Ok((op, rest.to_string()))
}
"vercel-artifacts" => {
let token = std::env::var("VERCEL_ARTIFACTS_TOKEN").unwrap_or_default();
let builder = services::VercelArtifacts::default().access_token(&token);
let op = Operator::new(builder)?.finish();
Ok((op, rest.to_string()))
}
"ghac" => {
let version = std::env::var("GHAC_VERSION").unwrap_or_else(|_| "v1".into());
let builder = services::Ghac::default().version(&version);
let op = Operator::new(builder)?.finish();
Ok((op, rest.to_string()))
}
"dbfs" => {
let endpoint = std::env::var("DATABRICKS_HOST")
.unwrap_or_else(|_| "https://adb-example.azuredatabricks.net".into());
let token = std::env::var("DATABRICKS_TOKEN").unwrap_or_default();
let builder = services::Dbfs::default()
.root("/")
.endpoint(&endpoint)
.token(&token);
let op = Operator::new(builder)?.finish();
Ok((op, rest.to_string()))
}
"alluxio" => {
let (hostport, path) = rest.split_once('/').unwrap_or((rest, ""));
let endpoint = format!("http://{hostport}");
let builder = services::Alluxio::default()
.root("/")
.endpoint(&endpoint);
let op = Operator::new(builder)?.finish();
Ok((op, path.to_string()))
}
"webhdfs" => {
let (hostport, path) = rest.split_once('/').unwrap_or((rest, ""));
let endpoint = format!("http://{hostport}");
let user = std::env::var("WEBHDFS_USER").unwrap_or_default();
let mut builder = services::Webhdfs::default()
.root("/")
.endpoint(&endpoint);
if !user.is_empty() {
builder = builder.user_name(&user);
}
let op = Operator::new(builder)?.finish();
Ok((op, path.to_string()))
}
"hdfs" => {
let (hostport, path) = rest.split_once('/').unwrap_or((rest, ""));
let name_node = format!("hdfs://{hostport}");
let builder = services::HdfsNative::default()
.name_node(&name_node)
.root("/");
let op = Operator::new(builder)?.finish();
Ok((op, path.to_string()))
}
"lakefs" => {
let mut parts = rest.splitn(3, '/');
let repo = parts.next().unwrap_or("").to_string();
let branch = parts.next().unwrap_or("main").to_string();
let path = parts.next().unwrap_or("").to_string();
let endpoint = std::env::var("LAKEFS_ENDPOINT")
.unwrap_or_else(|_| "http://localhost:8000".into());
let username = std::env::var("LAKEFS_ACCESS_KEY_ID").unwrap_or_default();
let password = std::env::var("LAKEFS_SECRET_ACCESS_KEY").unwrap_or_default();
let builder = services::Lakefs::default()
.endpoint(&endpoint)
.username(&username)
.password(&password)
.repository(&repo)
.branch(&branch);
let op = Operator::new(builder)?.finish();
Ok((op, path))
}
"ipfs" => {
let gateway = std::env::var("IPFS_GATEWAY")
.unwrap_or_else(|_| "http://127.0.0.1:8080".into());
let builder = services::Ipfs::default()
.root("/")
.endpoint(&gateway);
let op = Operator::new(builder)?.finish();
Ok((op, rest.to_string()))
}
"ipmfs" => {
let endpoint = std::env::var("IPFS_ENDPOINT")
.unwrap_or_else(|_| "http://127.0.0.1:5001".into());
let builder = services::Ipmfs::default().endpoint(&endpoint);
let op = Operator::new(builder)?.finish();
Ok((op, rest.to_string()))
}
#[cfg(feature = "rocksdb-storage")]
"rocksdb" => {
let (db_path, key) = rest.rsplit_once('/').unwrap_or((rest, ""));
let builder = services::Rocksdb::default().datadir(db_path);
let op = Operator::new(builder)?.finish();
Ok((op, key.to_string()))
}
"rediss" => {
let (hostport, key) = rest.split_once('/').unwrap_or((rest, ""));
let redis_url = format!("rediss://{hostport}");
let builder = services::Redis::default().endpoint(&redis_url);
let op = Operator::new(builder)?.finish();
Ok((op, key.to_string()))
}
"redis" => {
let (conn, path) = rest.split_once('/').unwrap_or((rest, ""));
let redis_url = format!("redis://{conn}");
let builder = services::Redis::default().endpoint(&redis_url);
let op = Operator::new(builder)?.finish();
Ok((op, path.to_string()))
}
"memcached" => {
let (hostport, path) = rest.split_once('/').unwrap_or((rest, ""));
let endpoint = format!("tcp://{hostport}");
let builder = services::Memcached::default().endpoint(&endpoint);
let op = Operator::new(builder)?.finish();
Ok((op, path.to_string()))
}
"etcd" => {
let (hostport, path) = rest.split_once('/').unwrap_or((rest, ""));
let endpoint = format!("http://{hostport}");
let builder = services::Etcd::default().endpoints(&endpoint);
let op = Operator::new(builder)?.finish();
Ok((op, path.to_string()))
}
"tikv" => {
let (hostport, path) = rest.split_once('/').unwrap_or((rest, ""));
let builder = services::Tikv::default().endpoints(vec![hostport.to_string()]);
let op = Operator::new(builder)?.finish();
Ok((op, path.to_string()))
}
"mongodb" => {
let conn_str = format!("mongodb://{rest}");
let mut parts = rest.splitn(3, '/');
let _ = parts.next(); let database = parts.next().unwrap_or("blazehash");
let rest_path = parts.next().unwrap_or("");
let (collection, path) = rest_path.split_once('/').unwrap_or((rest_path, ""));
let builder = services::Mongodb::default()
.connection_string(&conn_str)
.database(database)
.collection(collection);
let op = Operator::new(builder)?.finish();
Ok((op, path.to_string()))
}
"gridfs" => {
let conn_str = format!("mongodb://{rest}");
let mut parts = rest.splitn(3, '/');
let _ = parts.next(); let database = parts.next().unwrap_or("blazehash");
let rest_path = parts.next().unwrap_or("");
let (bucket, path) = rest_path.split_once('/').unwrap_or((rest_path, ""));
let builder = services::Gridfs::default()
.connection_string(&conn_str)
.database(database)
.bucket(bucket);
let op = Operator::new(builder)?.finish();
Ok((op, path.to_string()))
}
"mysql" => {
let conn_str = format!("mysql://{rest}");
let path = rest
.split_once('/')
.and_then(|(_, after_host)| after_host.split_once('/'))
.map(|(_, key)| key.to_string())
.unwrap_or_default();
let builder = services::Mysql::default().connection_string(&conn_str);
let op = Operator::new(builder)?.finish();
Ok((op, path))
}
"postgresql" => {
let conn_str = format!("postgresql://{rest}");
let mut parts = rest.splitn(3, '/');
let _host = parts.next().unwrap_or("");
let _db = parts.next().unwrap_or("");
let path = parts.next().unwrap_or("").to_string();
let builder = services::Postgresql::default()
.connection_string(&conn_str);
let op = Operator::new(builder)?.finish();
Ok((op, path))
}
"sqlite" => {
let (db_path, key) = rest.rsplit_once('/').unwrap_or((rest, ""));
let conn = format!("sqlite://{db_path}");
let builder = services::Sqlite::default().connection_string(&conn);
let op = Operator::new(builder)?.finish();
Ok((op, key.to_string()))
}
"cloudflare-kv" => {
let (namespace, key) = rest.split_once('/').unwrap_or((rest, ""));
let account_id = std::env::var("CLOUDFLARE_ACCOUNT_ID").unwrap_or_default();
let token = std::env::var("CLOUDFLARE_API_TOKEN").unwrap_or_default();
let builder = services::CloudflareKv::default()
.account_id(&account_id)
.api_token(&token)
.namespace_id(namespace);
let op = Operator::new(builder)?.finish();
Ok((op, key.to_string()))
}
"d1" => {
let (db_id, key) = rest.split_once('/').unwrap_or((rest, ""));
let account_id = std::env::var("CLOUDFLARE_ACCOUNT_ID").unwrap_or_default();
let token = std::env::var("CLOUDFLARE_API_TOKEN").unwrap_or_default();
let builder = services::D1::default()
.account_id(&account_id)
.token(&token)
.database_id(db_id);
let op = Operator::new(builder)?.finish();
Ok((op, key.to_string()))
}
"webdav" => {
let host = rest.split('/').next().unwrap_or("");
let endpoint = format!("https://{host}");
let path = rest.split_once('/').map(|(_, p)| p).unwrap_or("");
let builder = services::Webdav::default().endpoint(&endpoint);
let op = Operator::new(builder)?.finish();
Ok((op, path.to_string()))
}
"http" | "https" => {
let host = rest.split('/').next().unwrap_or("");
let endpoint = format!("{scheme}://{host}");
let path = rest.split_once('/').map(|(_, p)| p).unwrap_or("");
let builder = services::Http::default().endpoint(&endpoint);
let op = Operator::new(builder)?.finish();
Ok((op, path.to_string()))
}
"sftp" => {
anyhow::bail!(
"sftp:// URIs must be fetched with \
`crate::remote::sftp::fetch_sftp_bytes(uri)` \
(ssh2/libssh2, cross-platform). \
operator_for_uri() does not support sftp://."
)
}
"ftp" | "ftps" => {
anyhow::bail!(
"ftp:// and ftps:// URIs must be fetched with \
crate::remote::ftp::fetch_ftp_bytes(), not operator_for_uri()"
)
}
"monoiofs" => {
let full = format!("/{rest}");
let (dir, file) = full.rsplit_once('/').unwrap_or(("/", &full));
#[cfg(target_os = "linux")]
{
let builder = services::Monoiofs::default().root(dir);
let op = Operator::new(builder)?.finish();
Ok((op, file.to_string()))
}
#[cfg(not(target_os = "linux"))]
{
let _ = (dir, file);
bail!("monoiofs:// is only supported on Linux (io_uring required)")
}
}
"compfs" => {
let full = format!("/{rest}"); let (dir, file) = full.rsplit_once('/').unwrap_or(("/", &full));
let builder = services::Compfs::default().root(dir);
let op = Operator::new(builder)?.finish();
Ok((op, file.to_string()))
}
"file" => {
let (dir, file) = rest.rsplit_once('/').unwrap_or(("/", rest));
let builder = services::Fs::default().root(dir);
let op = Operator::new(builder)?.finish();
Ok((op, file.to_string()))
}
other => bail!("unsupported URI scheme: {other}://"),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn hdfs_uri_is_not_unsupported() {
let result = operator_for_uri("hdfs://namenode:9000/data/evidence.zip");
assert!(
result.is_ok() || result
.as_ref()
.unwrap_err()
.to_string()
.contains("unsupported URI scheme: hdfs://")
.then(|| false)
.unwrap_or(true),
"hdfs:// should not return unsupported-scheme error"
);
}
#[test]
fn hdfs_uri_path_extraction() {
let (_, path) = operator_for_uri("hdfs://namenode:9000/data/evidence.zip")
.expect("hdfs:// should be supported");
assert_eq!(path, "data/evidence.zip");
}
#[test]
fn hdfs_is_recognised_as_remote() {
assert!(
crate::remote::is_remote_uri("hdfs://namenode:9000/path"),
"hdfs:// should be recognised as a remote URI"
);
}
#[test]
fn sftp_uri_is_supported() {
let result = operator_for_uri("sftp://user@host/path/file.zip");
assert!(
result.is_ok() || !result.unwrap_err().to_string().contains("unsupported"),
"sftp:// should be supported"
);
}
#[test]
fn webhdfs_uri_is_supported() {
let result = operator_for_uri("webhdfs://namenode:50070/path/file");
assert!(
result.is_ok() || !result.unwrap_err().to_string().contains("unsupported"),
"webhdfs:// should be supported"
);
}
}