use std::time::Duration;
use reqwest::blocking::Client;
use reqwest::StatusCode;
use crate::{Error, ObjectStore};
pub struct HttpReadOnlyObjectStore {
origin: String,
client: Client,
}
impl HttpReadOnlyObjectStore {
pub fn new(origin: impl Into<String>) -> Result<Self, Error> {
let origin = origin.into().trim_end_matches('/').to_string();
if origin.is_empty() {
return Err(Error::Backend("bundle origin must not be empty".into()));
}
let client = Client::builder()
.timeout(Duration::from_secs(60))
.connect_timeout(Duration::from_secs(10))
.build()
.map_err(|e| Error::Backend(format!("building http client: {e}")))?;
Ok(Self { origin, client })
}
pub fn origin(&self) -> &str {
&self.origin
}
fn url(&self, key: &str) -> String {
format!("{}/{}", self.origin, key.trim_start_matches('/'))
}
fn read_only(op: &str) -> Error {
Error::Backend(format!(
"{op} is not supported by the read-only bundle origin — publishing \
goes through the credentialed R2 store on the publisher, never a node"
))
}
}
impl ObjectStore for HttpReadOnlyObjectStore {
fn locate(&self, key: &str) -> String {
self.url(key)
}
fn get(&self, key: &str) -> Result<Option<Vec<u8>>, Error> {
let url = self.url(key);
let resp = self
.client
.get(&url)
.send()
.map_err(|e| Error::Io(format!("GET {url}: {e}")))?;
match resp.status() {
StatusCode::NOT_FOUND | StatusCode::FORBIDDEN => {
Ok(None)
}
s if s.is_success() => {
let bytes = resp
.bytes()
.map_err(|e| Error::Io(format!("reading body of {url}: {e}")))?;
Ok(Some(bytes.to_vec()))
}
s => Err(Error::Backend(format!("GET {url}: unexpected status {s}"))),
}
}
fn head(&self, key: &str) -> Result<bool, Error> {
let url = self.url(key);
let resp = self
.client
.head(&url)
.send()
.map_err(|e| Error::Io(format!("HEAD {url}: {e}")))?;
match resp.status() {
StatusCode::NOT_FOUND | StatusCode::FORBIDDEN => Ok(false),
s if s.is_success() => Ok(true),
s => Err(Error::Backend(format!("HEAD {url}: unexpected status {s}"))),
}
}
fn put(&self, _key: &str, _data: Vec<u8>) -> Result<(), Error> {
Err(Self::read_only("put"))
}
fn delete(&self, _key: &str) -> Result<(), Error> {
Err(Self::read_only("delete"))
}
fn list_prefix(&self, _prefix: &str) -> Result<Vec<String>, Error> {
Err(Self::read_only("list_prefix"))
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn origin_trailing_slash_is_trimmed() {
let s = HttpReadOnlyObjectStore::new("https://cdn.yah.dev/").unwrap();
assert_eq!(s.origin(), "https://cdn.yah.dev");
assert_eq!(s.url("blobs/abc"), "https://cdn.yah.dev/blobs/abc");
}
#[test]
fn key_leading_slash_does_not_double_up() {
let s = HttpReadOnlyObjectStore::new("https://cdn.yah.dev").unwrap();
assert_eq!(s.url("/blobs/abc"), "https://cdn.yah.dev/blobs/abc");
}
#[test]
fn empty_origin_is_rejected() {
assert!(HttpReadOnlyObjectStore::new("").is_err());
assert!(HttpReadOnlyObjectStore::new("///").is_err());
}
#[test]
fn mutating_ops_are_refused() {
let s = HttpReadOnlyObjectStore::new("https://cdn.yah.dev").unwrap();
assert!(s.put("blobs/x", vec![1]).is_err());
assert!(s.delete("blobs/x").is_err());
assert!(s.list_prefix("blobs/").is_err());
assert!(s.etag("blobs/x").is_err());
}
}