use std::sync::Arc;
use axum::Extension;
use axum::Json;
use axum::extract::{Query, State};
use axum::http::StatusCode;
use futures::StreamExt;
use serde::{Deserialize, Serialize};
use time::OffsetDateTime;
use crate::server::AppState;
use crate::server::credentials::MaybeCreds;
use crate::server::error::storage_error_response;
use crate::server::metrics::OperationMetrics;
use crate::storage::{ObjectInfo, StorageBackend, StorageError};
const PER_PREFIX_CAP: usize = 10_000;
const MAX_PREFIXES: usize = 2_000;
const LIST_CONCURRENCY: usize = 32;
const MINUTE_GRANULARITY_THRESHOLD_SECS: i64 = 600;
#[derive(Deserialize)]
pub struct BrowseParams {
pub bucket: Option<String>,
pub prefix: Option<String>,
pub service: Option<String>,
pub from: i64,
pub to: i64,
}
#[derive(Serialize)]
pub struct BrowseResponse {
pub objects: Vec<ObjectInfo>,
pub truncated: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Granularity {
Hour,
#[cfg_attr(not(test), allow(dead_code))]
TenMinute,
Minute,
}
type BrowseOk = (Extension<OperationMetrics>, Json<BrowseResponse>);
pub async fn browse(
State(state): State<AppState>,
creds: MaybeCreds,
Query(params): Query<BrowseParams>,
) -> Result<BrowseOk, (StatusCode, String)> {
let service = normalize_service(params.service.as_deref())?;
let backend = state.resolve(creds).await?;
let bucket = params
.bucket
.or(state.default_bucket.clone())
.ok_or((StatusCode::BAD_REQUEST, "bucket is required".to_string()))?;
if params.to < params.from {
return Err((
StatusCode::BAD_REQUEST,
"`to` must be greater than or equal to `from`".to_string(),
));
}
let key_prefix = params
.prefix
.as_deref()
.map(str::trim)
.filter(|s| !s.is_empty());
let base = resolve_base(state.default_prefix.as_deref(), key_prefix);
let window = params.to - params.from;
if !state.time_partitioned_source {
return browse_local(backend, &bucket, &base, service).await;
}
browse_s3(
backend,
&bucket,
&base,
params.from,
params.to,
window,
service,
)
.await
}
fn normalize_service(service: Option<&str>) -> Result<Option<&str>, (StatusCode, String)> {
let Some(service) = service.map(str::trim).filter(|s| !s.is_empty()) else {
return Ok(None);
};
if service.contains('/') || service.chars().any(char::is_control) {
return Err((
StatusCode::BAD_REQUEST,
"`service` must be a single key path segment".to_string(),
));
}
Ok(Some(service))
}
async fn browse_local(
backend: Arc<dyn StorageBackend>,
bucket: &str,
base: &str,
service: Option<&str>,
) -> Result<BrowseOk, (StatusCode, String)> {
let page = backend
.list_objects(bucket, base, PER_PREFIX_CAP)
.await
.map_err(storage_error_response)?;
let objects: Vec<ObjectInfo> = page
.objects
.into_iter()
.filter(|o| crate::ingest::aggregate::is_trace_segment(&o.key))
.filter(|o| service.is_none_or(|wanted| key_service(&o.key) == Some(wanted)))
.collect();
let op = OperationMetrics::browse(objects.len(), 0, page.truncated, false);
Ok((
Extension(op),
Json(BrowseResponse {
objects,
truncated: page.truncated,
}),
))
}
async fn browse_s3(
backend: Arc<dyn StorageBackend>,
bucket: &str,
base: &str,
from: i64,
to: i64,
window: i64,
service: Option<&str>,
) -> Result<BrowseOk, (StatusCode, String)> {
let gran = if service.is_some() || window < MINUTE_GRANULARITY_THRESHOLD_SECS {
Granularity::Minute
} else {
Granularity::Hour
};
let (prefixes, range_truncated) = match service {
Some(service) => service_time_prefixes(base, from, to, service),
None => time_prefixes(base, from, to, gran),
};
tracing::debug!(
bucket = %bucket,
prefixes = prefixes.len(),
granularity = ?gran,
"browse fan-out"
);
let results: Vec<Result<crate::storage::ListPage, StorageError>> =
futures::stream::iter(prefixes.clone())
.map(|p| {
let backend = backend.clone();
let bucket = bucket.to_string();
async move { backend.list_objects(&bucket, &p, PER_PREFIX_CAP).await }
})
.buffered(LIST_CONCURRENCY)
.collect()
.await;
let mut objects = Vec::new();
let mut truncated = range_truncated;
let mut overflow_prefixes = Vec::new();
for (i, result) in results.into_iter().enumerate() {
let page = result.map_err(storage_error_response)?;
if page.truncated && gran == Granularity::Hour {
overflow_prefixes.push(prefixes[i].clone());
} else {
truncated |= page.truncated;
objects.extend(page.objects);
}
}
let refined = !overflow_prefixes.is_empty();
if !overflow_prefixes.is_empty() {
let refined: Vec<String> = overflow_prefixes
.iter()
.flat_map(|p| (0..6).map(move |d| format!("{p}{d}")))
.collect();
tracing::debug!(
refined_prefixes = refined.len(),
overflowed_hours = overflow_prefixes.len(),
"browse refining overflowed hours at 10-minute granularity"
);
let refined_results: Vec<Result<crate::storage::ListPage, StorageError>> =
futures::stream::iter(refined)
.map(|p| {
let backend = backend.clone();
let bucket = bucket.to_string();
async move { backend.list_objects(&bucket, &p, PER_PREFIX_CAP).await }
})
.buffer_unordered(LIST_CONCURRENCY)
.collect()
.await;
for result in refined_results {
let page = result.map_err(storage_error_response)?;
truncated |= page.truncated;
objects.extend(page.objects);
}
}
let op = OperationMetrics::browse(objects.len(), prefixes.len(), truncated, refined);
Ok((Extension(op), Json(BrowseResponse { objects, truncated })))
}
pub(super) fn key_service(key: &str) -> Option<&str> {
let parts: Vec<&str> = key.split('/').collect();
let date = parts.iter().position(|part| is_date_segment(part))?;
let hhmm = *parts.get(date + 1)?;
if hhmm.len() != 4 || !hhmm.bytes().all(|b| b.is_ascii_digit()) {
return None;
}
parts.get(date + 2).copied().filter(|s| !s.is_empty())
}
pub(super) fn key_host(key: &str) -> Option<&str> {
let parts: Vec<&str> = key.split('/').collect();
let date = parts.iter().position(|part| is_date_segment(part))?;
let hhmm = *parts.get(date + 1)?;
if hhmm.len() != 4 || !hhmm.bytes().all(|b| b.is_ascii_digit()) {
return None;
}
parts.get(date + 3).copied().filter(|s| !s.is_empty())
}
fn is_date_segment(segment: &str) -> bool {
let b = segment.as_bytes();
b.len() == 10
&& b[4] == b'-'
&& b[7] == b'-'
&& b[..4].iter().all(u8::is_ascii_digit)
&& b[5..7].iter().all(u8::is_ascii_digit)
&& b[8..].iter().all(u8::is_ascii_digit)
}
fn service_time_prefixes(base: &str, from: i64, to: i64, service: &str) -> (Vec<String>, bool) {
let (mut prefixes, truncated) = minute_time_prefixes(base, from, to);
for prefix in &mut prefixes {
prefix.push('/');
prefix.push_str(service);
prefix.push('/');
}
(prefixes, truncated)
}
pub(super) fn resolve_base(default_prefix: Option<&str>, key_prefix: Option<&str>) -> String {
match (default_prefix, key_prefix) {
(Some(pfx), Some(kp)) => {
let default = pfx.trim_matches('/');
let requested = kp.trim_matches('/');
if default.is_empty()
|| requested == default
|| requested.starts_with(&format!("{default}/"))
{
requested.to_string()
} else {
format!("{default}/{requested}")
}
}
(Some(pfx), None) => pfx.to_string(),
(None, Some(kp)) => kp.to_string(),
(None, None) => String::new(),
}
}
pub(super) fn minute_time_prefixes(base: &str, from: i64, to: i64) -> (Vec<String>, bool) {
time_prefixes(base, from, to, Granularity::Minute)
}
fn time_prefixes(base: &str, from: i64, to: i64, gran: Granularity) -> (Vec<String>, bool) {
let step = match gran {
Granularity::Hour => 3600,
Granularity::TenMinute => 600,
Granularity::Minute => 60,
};
let start = from - from.rem_euclid(step);
let mut prefixes = Vec::new();
let mut truncated = false;
let mut t = start;
while t <= to {
if prefixes.len() >= MAX_PREFIXES {
truncated = true;
break;
}
if let Some(p) = bucket_prefix(t, gran) {
prefixes.push(join_prefix(base, &p));
}
t += step;
}
(prefixes, truncated)
}
fn bucket_prefix(epoch: i64, gran: Granularity) -> Option<String> {
let dt = OffsetDateTime::from_unix_timestamp(epoch).ok()?;
let date = format!(
"{:04}-{:02}-{:02}",
dt.year(),
u8::from(dt.month()),
dt.day()
);
let time = match gran {
Granularity::Hour => format!("{:02}", dt.hour()),
Granularity::TenMinute => format!("{:02}{}", dt.hour(), dt.minute() / 10),
Granularity::Minute => format!("{:02}{:02}", dt.hour(), dt.minute()),
};
Some(format!("{date}/{time}"))
}
fn join_prefix(base: &str, tail: &str) -> String {
if base.is_empty() {
tail.to_string()
} else {
format!("{}/{}", base.trim_end_matches('/'), tail)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn ten_minute_prefixes_cover_range() {
let from = OffsetDateTime::from_unix_timestamp(1_781_032_200).unwrap();
assert_eq!(from.hour(), 19);
assert_eq!(from.minute(), 10);
let to = from.unix_timestamp() + 25 * 60;
let (prefixes, truncated) =
time_prefixes("traces", from.unix_timestamp(), to, Granularity::TenMinute);
assert!(!truncated);
assert_eq!(
prefixes,
vec![
"traces/2026-06-09/191",
"traces/2026-06-09/192",
"traces/2026-06-09/193",
]
);
}
#[test]
fn hour_prefixes_cover_range() {
let from = OffsetDateTime::from_unix_timestamp(1_781_032_200).unwrap(); let to = from.unix_timestamp() + 3 * 3600; let (prefixes, truncated) =
time_prefixes("traces", from.unix_timestamp(), to, Granularity::Hour);
assert!(!truncated);
assert_eq!(
prefixes,
vec![
"traces/2026-06-09/19",
"traces/2026-06-09/20",
"traces/2026-06-09/21",
"traces/2026-06-09/22",
]
);
}
#[test]
fn unaligned_start_aligns_down() {
let base = OffsetDateTime::from_unix_timestamp(1_781_032_200).unwrap(); let from = base.unix_timestamp() + 7 * 60; let (prefixes, _) = time_prefixes("", from, from + 60, Granularity::TenMinute);
assert_eq!(prefixes, vec!["2026-06-09/191"]);
}
#[test]
fn minute_prefixes_are_four_char() {
let from = OffsetDateTime::from_unix_timestamp(1_781_032_200).unwrap(); let to = from.unix_timestamp() + 2 * 60; let (prefixes, _) = time_prefixes("p", from.unix_timestamp(), to, Granularity::Minute);
assert_eq!(
prefixes,
vec![
"p/2026-06-09/1910",
"p/2026-06-09/1911",
"p/2026-06-09/1912"
]
);
}
#[test]
fn selected_service_prefixes_are_exact_full_minutes() {
let from = OffsetDateTime::from_unix_timestamp(1_781_032_200).unwrap(); let to = from.unix_timestamp() + 2 * 60; let (prefixes, truncated) =
service_time_prefixes("traces", from.unix_timestamp(), to, "api");
assert!(!truncated);
assert_eq!(
prefixes,
vec![
"traces/2026-06-09/1910/api/",
"traces/2026-06-09/1911/api/",
"traces/2026-06-09/1912/api/",
]
);
}
#[test]
fn scope_prefix_already_under_default_is_not_duplicated() {
assert_eq!(resolve_base(Some("traces"), Some("traces")), "traces");
assert_eq!(
resolve_base(Some("traces"), Some("traces/team-a")),
"traces/team-a"
);
assert_eq!(
resolve_base(Some("traces"), Some("team-a")),
"traces/team-a"
);
}
#[test]
fn service_normalization_trims_and_treats_empty_as_absent() {
assert_eq!(normalize_service(None).unwrap(), None);
assert_eq!(normalize_service(Some("")).unwrap(), None);
assert_eq!(normalize_service(Some(" \t ")).unwrap(), None);
assert_eq!(normalize_service(Some(" api ")).unwrap(), Some("api"));
}
#[test]
fn service_normalization_rejects_non_segments() {
for service in ["api/worker", "api\nworker"] {
let error = normalize_service(Some(service)).unwrap_err();
assert_eq!(error.0, StatusCode::BAD_REQUEST, "{service:?}");
}
for service in [r"api\worker", ".", ".."] {
assert_eq!(normalize_service(Some(service)).unwrap(), Some(service));
}
}
#[test]
fn known_layout_service_is_exact() {
assert_eq!(
key_service("root/2026-06-09/1910/api/host/boot/1-0.bin.gz"),
Some("api")
);
assert_eq!(
key_service("root/2026-06-09/1910/api-worker/host/boot/1-0.bin.gz"),
Some("api-worker")
);
assert_eq!(key_service("boot/trace.0.bin"), None);
}
#[test]
fn known_layout_host_is_exact() {
assert_eq!(
key_host("root/2026-06-09/1910/api/host-a/boot/1-0.bin.gz"),
Some("host-a")
);
assert_eq!(key_host("boot/trace.0.bin"), None);
}
#[test]
fn empty_base_has_no_leading_slash() {
let (prefixes, _) = time_prefixes("", 1_781_032_200, 1_781_032_200, Granularity::TenMinute);
assert_eq!(prefixes, vec!["2026-06-09/191"]);
}
#[test]
fn oversized_range_truncates() {
let (prefixes, truncated) =
time_prefixes("", 0, MAX_PREFIXES as i64 * 600 * 2, Granularity::TenMinute);
assert!(truncated);
assert_eq!(prefixes.len(), MAX_PREFIXES);
}
}