mod s3;
use std::path::Path;
use std::time::Duration;
use axum::body::Bytes;
use futures_util::Stream;
use crate::error::Error;
pub(crate) use s3::S3Keys;
pub(crate) struct Listing {
pub(crate) entries: Vec<Entry>,
pub(crate) complete: bool,
}
pub(crate) struct Entry {
pub(crate) key: String,
pub(crate) modified: Option<time::OffsetDateTime>,
pub(crate) size: u64,
}
impl Entry {
pub(crate) fn age(&self) -> Option<Duration> {
Duration::try_from(time::OffsetDateTime::now_utc() - self.modified?).ok()
}
}
pub struct Presigned {
pub href: String,
pub headers: Vec<(String, String)>,
}
#[derive(Clone)]
pub enum Keyspace {
S3(S3Keys),
}
impl Keyspace {
pub(crate) async fn reachable(&self) -> Result<(), Error> {
match self {
Self::S3(keys) => keys.reachable().await,
}
}
pub(crate) fn signed_download(&self, key: &str) -> Option<String> {
match self {
Self::S3(keys) => Some(keys.signed_download(key)),
}
}
pub(crate) fn signed_upload(&self, key: &str, digest: &str) -> Option<Presigned> {
match self {
Self::S3(keys) => Some(keys.signed_upload(key, digest)),
}
}
pub(crate) async fn get_range(
&self,
key: &str,
start: u64,
length: u64,
) -> Result<impl Stream<Item = Result<Bytes, reqwest::Error>> + use<>, Error> {
let response = match self {
Self::S3(keys) => keys.get_range(key, start, length).await?,
};
Ok(response.bytes_stream())
}
pub(crate) async fn head(&self, key: &str) -> Result<u64, Error> {
match self {
Self::S3(keys) => keys.head(key).await,
}
}
pub(crate) async fn put(&self, key: &str, body: Vec<u8>) -> Result<(), Error> {
let length = body.len() as u64;
match self {
Self::S3(keys) => keys.put(key, reqwest::Body::from(body), length).await,
}
}
pub(crate) async fn put_file(&self, key: &str, staged: &Path) -> Result<(), Error> {
match self {
Self::S3(keys) => keys.put_file(key, staged).await,
}
}
pub(crate) async fn copy(&self, from: &str, to: &str) -> Result<(), Error> {
match self {
Self::S3(keys) => keys.copy(from, to).await,
}
}
pub(crate) async fn put_if_absent(&self, key: &str, body: Vec<u8>) -> Result<bool, Error> {
match self {
Self::S3(keys) => keys.put_if_absent(key, body).await,
}
}
pub(crate) async fn get_bytes(&self, key: &str) -> Result<Option<Vec<u8>>, Error> {
match self {
Self::S3(keys) => keys.get_bytes(key).await,
}
}
pub(crate) async fn delete(&self, key: &str) -> Result<bool, Error> {
match self {
Self::S3(keys) => keys.delete(key).await,
}
}
pub(crate) async fn entries(&self, prefix: &str) -> Result<Vec<Entry>, Error> {
match self {
Self::S3(keys) => keys.entries(prefix).await,
}
}
pub(crate) async fn keys(&self, prefix: &str) -> Result<Vec<String>, Error> {
Ok(self
.entries(prefix)
.await?
.into_iter()
.map(|entry| entry.key)
.collect())
}
pub(crate) async fn listing(&self, prefix: &str) -> Listing {
match self.entries(prefix).await {
Ok(entries) => Listing {
entries,
complete: true,
},
Err(error) => {
tracing::warn!(%error, prefix, "the listing could not be finished");
Listing {
entries: Vec::new(),
complete: false,
}
}
}
}
}
pub(crate) async fn read_retrying(
request: reqwest::RequestBuilder,
) -> Result<reqwest::Response, Error> {
let retry = request.try_clone();
match request.send().await {
Ok(response) => Ok(response),
Err(_) => match retry {
Some(retry) => retry.send().await.map_err(|_| unreachable_store()),
None => Err(unreachable_store()),
},
}
}
pub(crate) fn unreachable_store() -> Error {
Error::Storage(std::io::Error::other("the object store is unreachable"))
}
pub(crate) async fn expect_success(
response: reqwest::Response,
what: &str,
) -> Result<reqwest::Response, Error> {
let status = response.status();
if status.is_success() {
return Ok(response);
}
let detail = response.text().await.unwrap_or_default();
Err(Error::Storage(std::io::Error::other(format!(
"the object store refused a {what} with {status}: {}",
detail.trim()
))))
}
pub(crate) fn content_length(response: &reqwest::Response) -> Result<u64, Error> {
response
.headers()
.get(reqwest::header::CONTENT_LENGTH)
.and_then(|value| value.to_str().ok())
.and_then(|value| value.parse().ok())
.ok_or_else(|| {
Error::Storage(std::io::Error::other(
"the object store gave no object size",
))
})
}