use std::time::Duration;
use async_trait::async_trait;
use jiff::Timestamp;
use serde::{Deserialize, Serialize};
use crate::Result;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum HealthProbe {
Liveness,
Readiness,
}
impl HealthProbe {
pub const fn path(self) -> &'static str {
match self {
Self::Liveness => "/health",
Self::Readiness => "/health/ready",
}
}
pub const fn as_str(self) -> &'static str {
match self {
Self::Liveness => "liveness",
Self::Readiness => "readiness",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct HealthReport {
pub probe: HealthProbe,
pub endpoint: String,
pub path: String,
pub status_code: u16,
pub healthy: bool,
pub latency_ms: u64,
pub status: Option<String>,
pub service: Option<String>,
pub server_version: Option<String>,
}
#[async_trait]
pub trait HealthApi: Send + Sync {
async fn check_health(&self, probe: HealthProbe, timeout: Duration) -> Result<HealthReport>;
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum UsageSource {
ServerSnapshot,
ClientScan,
}
impl UsageSource {
pub const fn as_str(self) -> &'static str {
match self {
Self::ServerSnapshot => "server_snapshot",
Self::ClientScan => "client_scan",
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum UsageScope {
Cluster,
Bucket,
Prefix,
}
impl UsageScope {
pub const fn as_str(self) -> &'static str {
match self {
Self::Cluster => "cluster",
Self::Bucket => "bucket",
Self::Prefix => "prefix",
}
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct UsageBucket {
pub name: String,
pub total_bytes: u64,
pub object_count: u64,
pub version_count: Option<u64>,
pub delete_marker_count: Option<u64>,
pub incomplete_upload_count: Option<u64>,
pub incomplete_upload_bytes: Option<u64>,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct UsageFailure {
pub bucket: String,
pub message: String,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct UsageReport {
pub source: UsageSource,
pub scope: UsageScope,
pub path: Option<String>,
pub snapshot_at: Option<Timestamp>,
pub total_bytes: u64,
pub object_count: u64,
pub version_count: Option<u64>,
pub delete_marker_count: Option<u64>,
pub incomplete_upload_count: Option<u64>,
pub incomplete_upload_bytes: Option<u64>,
pub buckets: Vec<UsageBucket>,
pub partial: bool,
pub failures: Vec<UsageFailure>,
}
impl UsageReport {
pub fn empty(source: UsageSource, scope: UsageScope, path: Option<String>) -> Self {
Self {
source,
scope,
path,
snapshot_at: None,
total_bytes: 0,
object_count: 0,
version_count: None,
delete_marker_count: None,
incomplete_upload_count: None,
incomplete_upload_bytes: None,
buckets: Vec::new(),
partial: false,
failures: Vec::new(),
}
}
pub fn push_bucket(&mut self, bucket: UsageBucket) {
self.buckets.push(bucket);
}
pub fn push_failure(&mut self, failure: UsageFailure) {
self.partial = true;
self.failures.push(failure);
}
pub fn finish(&mut self) {
let versions_were_requested = self.version_count.is_some();
let delete_markers_were_requested = self.delete_marker_count.is_some();
let incomplete_uploads_were_requested = self.incomplete_upload_count.is_some();
let incomplete_bytes_were_requested = self.incomplete_upload_bytes.is_some();
self.buckets
.sort_by(|left, right| left.name.cmp(&right.name));
self.failures
.sort_by(|left, right| left.bucket.cmp(&right.bucket));
self.total_bytes = self.buckets.iter().fold(0_u64, |total, bucket| {
total.saturating_add(bucket.total_bytes)
});
self.object_count = self.buckets.iter().fold(0_u64, |total, bucket| {
total.saturating_add(bucket.object_count)
});
self.version_count = sum_optional(self.buckets.iter().map(|bucket| bucket.version_count))
.or_else(|| versions_were_requested.then_some(0));
self.delete_marker_count =
sum_optional(self.buckets.iter().map(|bucket| bucket.delete_marker_count))
.or_else(|| delete_markers_were_requested.then_some(0));
self.incomplete_upload_count = sum_optional(
self.buckets
.iter()
.map(|bucket| bucket.incomplete_upload_count),
)
.or_else(|| incomplete_uploads_were_requested.then_some(0));
self.incomplete_upload_bytes = sum_optional(
self.buckets
.iter()
.map(|bucket| bucket.incomplete_upload_bytes),
)
.or_else(|| incomplete_bytes_were_requested.then_some(0));
}
}
fn sum_optional(values: impl Iterator<Item = Option<u64>>) -> Option<u64> {
let mut values = values;
let first = values.next()??;
values.try_fold(first, |total, value| Some(total.saturating_add(value?)))
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct UsageScanRequest {
pub bucket: Option<String>,
pub prefix: Option<String>,
pub include_versions: bool,
pub include_incomplete_uploads: bool,
}
impl UsageScanRequest {
pub const fn scope(&self) -> UsageScope {
if self.prefix.is_some() {
UsageScope::Prefix
} else if self.bucket.is_some() {
UsageScope::Bucket
} else {
UsageScope::Cluster
}
}
pub const fn requires_client_scan(&self) -> bool {
self.prefix.is_some() || self.include_incomplete_uploads
}
pub fn path(&self, alias: &str) -> Option<String> {
let bucket = self.bucket.as_deref()?;
Some(match self.prefix.as_deref() {
Some(prefix) => format!("{alias}/{bucket}/{prefix}"),
None => format!("{alias}/{bucket}"),
})
}
}
#[async_trait]
pub trait UsageSnapshotApi: Send + Sync {
async fn usage_snapshot(&self) -> Result<UsageReport>;
}
#[async_trait]
pub trait UsageScanApi: Send + Sync {
async fn scan_usage(&self, request: &UsageScanRequest) -> Result<UsageReport>;
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn optional_totals_remain_unknown_when_any_bucket_is_unknown() {
assert_eq!(sum_optional(std::iter::empty()), None);
assert_eq!(sum_optional([Some(1), None].into_iter()), None);
assert_eq!(sum_optional([Some(1), Some(2)].into_iter()), Some(3));
}
}