use std::path::Path;
use anyhow::{Context, Result, anyhow, bail};
use reqwest::multipart::{Form, Part};
use reqwest::{Body, Client, StatusCode, Url};
use serde::Deserialize;
use tokio_util::io::ReaderStream;
pub struct Bucket {
pub route: &'static str,
pub store: &'static str,
pub empty_msg: &'static str,
pub limit_note: &'static str,
}
#[derive(Debug, Deserialize)]
pub struct Published {
pub size: u64,
pub digest: Option<String>,
}
#[derive(Debug, Deserialize)]
struct Row {
name: String,
version: String,
size: u64,
digest: Option<String>,
modified: Option<String>,
}
pub fn object_url(base: &str, route: &str, name: &str, version: &str) -> Result<Url> {
for (label, value) in [("name", name), ("version", version)] {
if value == "." || value == ".." {
bail!(
"{label} must not be '.' or '..' (it would be normalised out of the request path)"
);
}
}
let mut url = Url::parse(base).with_context(|| format!("invalid backend URL {base:?}"))?;
url.path_segments_mut()
.map_err(|_| anyhow!("backend URL {base:?} cannot carry a path"))?
.pop_if_empty()
.extend(route.split('/'))
.push(name)
.push(version);
Ok(url)
}
fn collection_url(base: &str, route: &str) -> Result<Url> {
let mut url = Url::parse(base).with_context(|| format!("invalid backend URL {base:?}"))?;
url.path_segments_mut()
.map_err(|_| anyhow!("backend URL {base:?} cannot carry a path"))?
.pop_if_empty()
.extend(route.split('/'));
Ok(url)
}
async fn rejected(op: &str, resp: reqwest::Response, limit_note: Option<&str>) -> anyhow::Error {
let status = resp.status();
let body = resp.text().await.unwrap_or_default();
let mut msg = format!("{op} failed: {status} — {body}");
match status {
StatusCode::UNAUTHORIZED | StatusCode::FORBIDDEN => msg.push_str(
"\nhint: this command goes through the backend API — set KANADE_AUTH_TOKEN \
(`kanade login`) to an account with the operator role",
),
StatusCode::PAYLOAD_TOO_LARGE => {
if let Some(note) = limit_note {
msg.push_str(&format!("\nhint: {note}"));
}
}
_ => {}
}
anyhow!(msg)
}
pub async fn upload(
client: &Client,
base: &str,
bucket: &Bucket,
name: &str,
version: &str,
file: &Path,
) -> Result<Published> {
let url = object_url(base, bucket.route, name, version)?;
let handle = tokio::fs::File::open(file)
.await
.with_context(|| format!("open {file:?}"))?;
let len = handle
.metadata()
.await
.with_context(|| format!("stat {file:?}"))?
.len();
let part = Part::stream_with_length(Body::wrap_stream(ReaderStream::new(handle)), len)
.file_name(
file.file_name()
.map(|n| n.to_string_lossy().into_owned())
.unwrap_or_else(|| "file".into()),
)
.mime_str("application/octet-stream")
.context("build multipart part")?;
let resp = client
.post(url.clone())
.multipart(Form::new().part("file", part))
.send()
.await
.with_context(|| {
format!(
"POST {url} (if the backend closed the connection mid-upload it rejected the \
body early — check the size limit and KANADE_AUTH_TOKEN)"
)
})?;
if !resp.status().is_success() {
return Err(rejected("publish", resp, Some(bucket.limit_note)).await);
}
resp.json()
.await
.context("parse publish response from server")
}
pub async fn list(client: &Client, base: &str, bucket: &Bucket) -> Result<()> {
let url = collection_url(base, bucket.route)?;
let resp = client
.get(url.clone())
.send()
.await
.with_context(|| format!("GET {url}"))?;
if !resp.status().is_success() {
return Err(rejected("list", resp, None).await);
}
let rows: Vec<Row> = resp
.json()
.await
.context("parse list response from server")?;
for line in render_rows(rows, bucket.empty_msg) {
println!("{line}");
}
Ok(())
}
fn render_rows(rows: Vec<Row>, empty_msg: &str) -> Vec<String> {
let mut rows: Vec<(String, Row)> = rows
.into_iter()
.map(|r| (format!("{}/{}", r.name, r.version), r))
.collect();
rows.sort_by(|a, b| a.0.cmp(&b.0));
if rows.is_empty() {
return vec![empty_msg.to_string()];
}
rows.into_iter()
.map(|(key, row)| {
let dgst = row.digest.as_deref().unwrap_or("—");
let modt = row.modified.as_deref().unwrap_or("—");
format!("{key}\t{}\t{modt}\t{dgst}", row.size)
})
.collect()
}
pub async fn delete(
client: &Client,
base: &str,
bucket: &Bucket,
name: &str,
version: &str,
) -> Result<()> {
let url = object_url(base, bucket.route, name, version)?;
let key = format!("{name}/{version}");
let resp = client
.delete(url.clone())
.send()
.await
.with_context(|| format!("DELETE {url}"))?;
match resp.status() {
s if s.is_success() => println!("deleted: {key}"),
StatusCode::NOT_FOUND => println!("not present: {key} (idempotent no-op)"),
_ => return Err(rejected("delete", resp, None).await),
}
Ok(())
}
pub fn print_published(bucket: &Bucket, key: &str, published: &Published) {
println!("published: {key}");
println!(" object_store : {}/{key}", bucket.store);
println!(" size : {} bytes", published.size);
if let Some(d) = published.digest.as_deref() {
println!(" digest : {d}");
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::test_http::{fake_backend, seen};
fn row(name: &str, version: &str, digest: Option<&str>) -> Row {
Row {
name: name.into(),
version: version.into(),
size: 7,
digest: digest.map(String::from),
modified: None,
}
}
#[test]
fn list_is_sorted_by_key_not_by_backend_order() {
let lines = render_rows(
vec![
row("b", "1", Some("SHA-256=x")),
row("a", "2", None),
row("a", "10", None),
],
"(none)",
);
assert_eq!(
lines,
vec![
"a/10\t7\t—\t—".to_string(),
"a/2\t7\t—\t—".to_string(),
"b/1\t7\t—\tSHA-256=x".to_string(),
]
);
}
#[test]
fn empty_list_prints_the_bucket_message() {
assert_eq!(
render_rows(vec![], "(no app packages)"),
["(no app packages)"]
);
}
#[test]
fn url_encodes_each_segment_and_keeps_a_base_path_prefix() {
let u = object_url("http://h:1/prefix", "api/app-packages", "a b?#%x", "1.0").unwrap();
assert_eq!(
u.as_str(),
"http://h:1/prefix/api/app-packages/a%20b%3F%23%25x/1.0"
);
assert_eq!(u.query(), None);
assert_eq!(u.fragment(), None);
}
#[test]
fn dot_segments_are_rejected_rather_than_normalised() {
for (n, v) in [(".", "1"), ("..", "1"), ("a", "."), ("a", "..")] {
assert!(object_url("http://h", "api/app-packages", n, v).is_err());
}
assert!(object_url("http://h", "api/app-packages", "a..", "...").is_ok());
}
#[tokio::test]
async fn encoded_name_reaches_the_backend_as_one_segment() {
let (base, log) = fake_backend(vec![(204, "")]).await;
let bucket = Bucket {
route: "api/app-packages",
store: "S",
empty_msg: "",
limit_note: "",
};
let client = crate::http_client::authed_client().unwrap();
delete(&client, &base, &bucket, "we b?x", "v#1")
.await
.unwrap();
let got = seen(&log);
assert_eq!(got[0].target, "/api/app-packages/we%20b%3Fx/v%231");
}
}