1use std::time::Duration;
4
5use async_trait::async_trait;
6use jiff::Timestamp;
7use serde::{Deserialize, Serialize};
8
9use crate::Result;
10
11#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
13#[serde(rename_all = "snake_case")]
14pub enum HealthProbe {
15 Liveness,
17 Readiness,
19}
20
21impl HealthProbe {
22 pub const fn path(self) -> &'static str {
24 match self {
25 Self::Liveness => "/health",
26 Self::Readiness => "/health/ready",
27 }
28 }
29
30 pub const fn as_str(self) -> &'static str {
32 match self {
33 Self::Liveness => "liveness",
34 Self::Readiness => "readiness",
35 }
36 }
37}
38
39#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
41pub struct HealthReport {
42 pub probe: HealthProbe,
44 pub endpoint: String,
46 pub path: String,
48 pub status_code: u16,
50 pub healthy: bool,
52 pub latency_ms: u64,
54 pub status: Option<String>,
56 pub service: Option<String>,
58 pub server_version: Option<String>,
60}
61
62#[async_trait]
64pub trait HealthApi: Send + Sync {
65 async fn check_health(&self, probe: HealthProbe, timeout: Duration) -> Result<HealthReport>;
67}
68
69#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
71#[serde(rename_all = "snake_case")]
72pub enum UsageSource {
73 ServerSnapshot,
75 ClientScan,
77}
78
79impl UsageSource {
80 pub const fn as_str(self) -> &'static str {
82 match self {
83 Self::ServerSnapshot => "server_snapshot",
84 Self::ClientScan => "client_scan",
85 }
86 }
87}
88
89#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
91#[serde(rename_all = "snake_case")]
92pub enum UsageScope {
93 Cluster,
95 Bucket,
97 Prefix,
99}
100
101impl UsageScope {
102 pub const fn as_str(self) -> &'static str {
104 match self {
105 Self::Cluster => "cluster",
106 Self::Bucket => "bucket",
107 Self::Prefix => "prefix",
108 }
109 }
110}
111
112#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
114pub struct UsageBucket {
115 pub name: String,
117 pub total_bytes: u64,
119 pub object_count: u64,
121 pub version_count: Option<u64>,
123 pub delete_marker_count: Option<u64>,
125 pub incomplete_upload_count: Option<u64>,
127 pub incomplete_upload_bytes: Option<u64>,
129}
130
131#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
133pub struct UsageFailure {
134 pub bucket: String,
136 pub message: String,
138}
139
140#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
142pub struct UsageReport {
143 pub source: UsageSource,
145 pub scope: UsageScope,
147 pub path: Option<String>,
149 pub snapshot_at: Option<Timestamp>,
151 pub total_bytes: u64,
153 pub object_count: u64,
155 pub version_count: Option<u64>,
157 pub delete_marker_count: Option<u64>,
159 pub incomplete_upload_count: Option<u64>,
161 pub incomplete_upload_bytes: Option<u64>,
163 pub buckets: Vec<UsageBucket>,
165 pub partial: bool,
167 pub failures: Vec<UsageFailure>,
169}
170
171impl UsageReport {
172 pub fn empty(source: UsageSource, scope: UsageScope, path: Option<String>) -> Self {
174 Self {
175 source,
176 scope,
177 path,
178 snapshot_at: None,
179 total_bytes: 0,
180 object_count: 0,
181 version_count: None,
182 delete_marker_count: None,
183 incomplete_upload_count: None,
184 incomplete_upload_bytes: None,
185 buckets: Vec::new(),
186 partial: false,
187 failures: Vec::new(),
188 }
189 }
190
191 pub fn push_bucket(&mut self, bucket: UsageBucket) {
193 self.buckets.push(bucket);
194 }
195
196 pub fn push_failure(&mut self, failure: UsageFailure) {
198 self.partial = true;
199 self.failures.push(failure);
200 }
201
202 pub fn finish(&mut self) {
204 let versions_were_requested = self.version_count.is_some();
205 let delete_markers_were_requested = self.delete_marker_count.is_some();
206 let incomplete_uploads_were_requested = self.incomplete_upload_count.is_some();
207 let incomplete_bytes_were_requested = self.incomplete_upload_bytes.is_some();
208 self.buckets
209 .sort_by(|left, right| left.name.cmp(&right.name));
210 self.failures
211 .sort_by(|left, right| left.bucket.cmp(&right.bucket));
212 self.total_bytes = self.buckets.iter().fold(0_u64, |total, bucket| {
213 total.saturating_add(bucket.total_bytes)
214 });
215 self.object_count = self.buckets.iter().fold(0_u64, |total, bucket| {
216 total.saturating_add(bucket.object_count)
217 });
218 self.version_count = sum_optional(self.buckets.iter().map(|bucket| bucket.version_count))
219 .or_else(|| versions_were_requested.then_some(0));
220 self.delete_marker_count =
221 sum_optional(self.buckets.iter().map(|bucket| bucket.delete_marker_count))
222 .or_else(|| delete_markers_were_requested.then_some(0));
223 self.incomplete_upload_count = sum_optional(
224 self.buckets
225 .iter()
226 .map(|bucket| bucket.incomplete_upload_count),
227 )
228 .or_else(|| incomplete_uploads_were_requested.then_some(0));
229 self.incomplete_upload_bytes = sum_optional(
230 self.buckets
231 .iter()
232 .map(|bucket| bucket.incomplete_upload_bytes),
233 )
234 .or_else(|| incomplete_bytes_were_requested.then_some(0));
235 }
236}
237
238fn sum_optional(values: impl Iterator<Item = Option<u64>>) -> Option<u64> {
239 let mut values = values;
240 let first = values.next()??;
241 values.try_fold(first, |total, value| Some(total.saturating_add(value?)))
242}
243
244#[derive(Debug, Clone, PartialEq, Eq)]
246pub struct UsageScanRequest {
247 pub bucket: Option<String>,
249 pub prefix: Option<String>,
251 pub include_versions: bool,
253 pub include_incomplete_uploads: bool,
255}
256
257impl UsageScanRequest {
258 pub const fn scope(&self) -> UsageScope {
260 if self.prefix.is_some() {
261 UsageScope::Prefix
262 } else if self.bucket.is_some() {
263 UsageScope::Bucket
264 } else {
265 UsageScope::Cluster
266 }
267 }
268
269 pub const fn requires_client_scan(&self) -> bool {
271 self.prefix.is_some() || self.include_incomplete_uploads
272 }
273
274 pub fn path(&self, alias: &str) -> Option<String> {
276 let bucket = self.bucket.as_deref()?;
277 Some(match self.prefix.as_deref() {
278 Some(prefix) => format!("{alias}/{bucket}/{prefix}"),
279 None => format!("{alias}/{bucket}"),
280 })
281 }
282}
283
284#[async_trait]
286pub trait UsageSnapshotApi: Send + Sync {
287 async fn usage_snapshot(&self) -> Result<UsageReport>;
289}
290
291#[async_trait]
293pub trait UsageScanApi: Send + Sync {
294 async fn scan_usage(&self, request: &UsageScanRequest) -> Result<UsageReport>;
296}
297
298#[cfg(test)]
299mod tests {
300 use super::*;
301
302 #[test]
303 fn optional_totals_remain_unknown_when_any_bucket_is_unknown() {
304 assert_eq!(sum_optional(std::iter::empty()), None);
305 assert_eq!(sum_optional([Some(1), None].into_iter()), None);
306 assert_eq!(sum_optional([Some(1), Some(2)].into_iter()), Some(3));
307 }
308}