use std::any::Any;
use std::fmt;
use std::sync::Arc;
use async_trait::async_trait;
use datafusion::catalog::{Session, TableProvider};
use datafusion::common::{DataFusionError, Result as DfResult};
use datafusion::datasource::TableType;
use datafusion::datasource::memory::MemTable;
use datafusion::logical_expr::{Expr, TableProviderFilterPushDown};
use datafusion::physical_plan::ExecutionPlan;
use datafusion::physical_plan::union::UnionExec;
use time::OffsetDateTime;
use tracing::debug;
use crate::arrow::array::RecordBatch;
use crate::arrow::datatypes::SchemaRef;
use crate::planner::predicate;
use crate::planner::split::{TierSplit, TimeRange, split};
use crate::tiering::store::{ColdStore, HotStore};
use crate::watermark::TieringWatermark;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum SnapshotSelector {
Id(i64),
Timestamp(OffsetDateTime),
}
impl fmt::Display for SnapshotSelector {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::Id(id) => write!(f, "snapshot {id}"),
Self::Timestamp(at) => write!(f, "as of {at}"),
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum ReadMode {
#[default]
Unified,
Historical,
Operational,
AsOf {
snapshot: SnapshotSelector,
max_version: Option<crate::version::Version>,
},
AsKnownAt(OffsetDateTime),
}
impl ReadMode {
fn apply(self, split: TierSplit) -> TierSplit {
match self {
Self::Unified | Self::AsKnownAt(_) => split,
Self::Historical | Self::AsOf { .. } => TierSplit { hot: None, ..split },
Self::Operational => TierSplit {
cold: None,
..split
},
}
}
pub const fn as_of(self) -> Option<(SnapshotSelector, Option<crate::version::Version>)> {
match self {
Self::AsOf {
snapshot,
max_version,
} => Some((snapshot, max_version)),
_ => None,
}
}
pub const fn recorded_at_ceiling(self) -> Option<OffsetDateTime> {
match self {
Self::AsKnownAt(at) => Some(at),
_ => None,
}
}
pub const fn forces_resolution(self) -> bool {
matches!(self, Self::AsOf { .. } | Self::AsKnownAt(_))
}
}
pub struct TieredTableProvider {
table: String,
schema: SchemaRef,
hot: Arc<dyn HotStore>,
cold: Arc<dyn ColdStore>,
cold_provider: Arc<dyn TableProvider>,
mode: ReadMode,
spec: crate::tiering::store::ScanSpec,
}
impl fmt::Debug for TieredTableProvider {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("TieredTableProvider")
.field("table", &self.table)
.field("mode", &self.mode)
.finish_non_exhaustive()
}
}
impl TieredTableProvider {
pub fn new(
table: impl Into<String>,
hot: Arc<dyn HotStore>,
cold: Arc<dyn ColdStore>,
cold_provider: Arc<dyn TableProvider>,
) -> Self {
let schema = cold_provider.schema();
let core = crate::encode::schema::storage_schema(&[]);
let extra: Vec<String> = schema
.fields()
.iter()
.filter(|f| core.field_with_name(f.name()).is_err())
.map(|f| f.name().clone())
.collect();
Self {
table: table.into(),
schema,
hot,
cold,
cold_provider,
mode: ReadMode::default(),
spec: crate::tiering::store::ScanSpec::new(
crate::encode::schema::MERGE_KEY
.iter()
.map(|s| (*s).to_string())
.collect(),
extra,
),
}
}
pub fn with_scan_spec(mut self, spec: crate::tiering::store::ScanSpec) -> Self {
self.spec = spec;
self
}
pub fn with_mode(mut self, mode: ReadMode) -> Self {
self.mode = mode;
self
}
pub const fn mode(&self) -> ReadMode {
self.mode
}
pub async fn watermark(&self) -> DfResult<TieringWatermark> {
self.cold.watermark(&self.table).await.map_err(external)
}
pub fn split_at(&self, watermark: TieringWatermark, filters: &[Expr]) -> TierSplit {
let range = predicate::time_range(filters);
self.mode.apply(split(range, watermark))
}
pub async fn plan_split(&self, filters: &[Expr]) -> DfResult<(TieringWatermark, TierSplit)> {
let watermark = self.watermark().await?;
Ok((watermark, self.split_at(watermark, filters)))
}
fn scan_hot(
&self,
range: TimeRange,
projection: Option<&Vec<usize>>,
limit: Option<usize>,
) -> DfResult<Arc<dyn ExecutionPlan>> {
Ok(Arc::new(HotScanExec::new(
self.schema.clone(),
Arc::clone(&self.hot),
self.table.clone(),
range,
self.spec.clone(),
projection.cloned(),
limit,
)?))
}
async fn scan_cold(
&self,
state: &dyn Session,
range: TimeRange,
projection: Option<&Vec<usize>>,
filters: &[Expr],
limit: Option<usize>,
) -> DfResult<Arc<dyn ExecutionPlan>> {
let mut all = filters.to_vec();
all.extend(predicate::range_filters(range));
let Some((snapshot, max_version)) = self.mode.as_of() else {
return self
.cold_provider
.scan(state, projection, &all, limit)
.await;
};
let ceiling = max_version.map(predicate::version_ceiling);
if let Some(expr) = &ceiling {
all.push(expr.clone());
}
let pinned = self
.cold
.snapshot_provider(&self.table, snapshot)
.await
.map_err(external)?;
if pinned.schema() != self.schema {
return Err(DataFusionError::Plan(format!(
"cannot read {} {snapshot}: its schema differs from the current one, so a \
reproducible read would silently change shape. Pin an explicit snapshot id \
from before the schema change, or query the raw table directly.",
self.table,
)));
}
let Some(ceiling) = ceiling else {
return pinned.scan(state, projection, &all, limit).await;
};
let scanned = pinned.scan(state, None, &all, None).await?;
let df_schema = datafusion::common::DFSchema::try_from(self.schema.as_ref().clone())?;
let physical = state.create_physical_expr(ceiling, &df_schema)?;
let filtered: Arc<dyn ExecutionPlan> = Arc::new(
datafusion::physical_plan::filter::FilterExec::try_new(physical, scanned)?,
);
let projected: Arc<dyn ExecutionPlan> = match projection {
None => filtered,
Some(indices) => {
let exprs = indices
.iter()
.map(|i| {
let field = self.schema.field(*i);
let column: Arc<dyn datafusion::physical_expr::PhysicalExpr> = Arc::new(
datafusion::physical_expr::expressions::Column::new(field.name(), *i),
);
(column, field.name().clone())
})
.collect::<Vec<_>>();
Arc::new(
datafusion::physical_plan::projection::ProjectionExec::try_new(
exprs, filtered,
)?,
)
}
};
Ok(match limit {
None => projected,
Some(fetch) => Arc::new(datafusion::physical_plan::limit::GlobalLimitExec::new(
projected,
0,
Some(fetch),
)),
})
}
pub async fn scan_at(
&self,
state: &dyn Session,
watermark: TieringWatermark,
projection: Option<&Vec<usize>>,
filters: &[Expr],
limit: Option<usize>,
) -> DfResult<Arc<dyn ExecutionPlan>> {
let started = std::time::Instant::now();
let split = self.split_at(watermark, filters);
debug!(
table = %self.table,
%watermark,
cold = split.cold.is_some(),
hot = split.hot.is_some(),
"tier split"
);
let mut plans: Vec<Arc<dyn ExecutionPlan>> = Vec::with_capacity(2);
if let Some(range) = split.cold {
plans.push(
self.scan_cold(state, range, projection, filters, limit)
.await?,
);
}
if let Some(range) = split.hot {
plans.push(self.scan_hot(range, projection, limit)?);
}
crate::observe::metrics().plan_duration.record(
started.elapsed().as_secs_f64(),
&crate::observe::table(&self.table),
);
Ok(match plans.len() {
0 => {
let empty = MemTable::try_new(self.schema.clone(), vec![vec![]])?;
empty.scan(state, projection, filters, limit).await?
}
1 => plans.pop().expect("length checked"),
_ => UnionExec::try_new(plans)?,
})
}
}
pub(crate) struct HotScanExec {
schema: SchemaRef,
source_schema: SchemaRef,
hot: Arc<dyn HotStore>,
table: String,
range: TimeRange,
spec: crate::tiering::store::ScanSpec,
projection: Option<Vec<usize>>,
limit: Option<usize>,
properties: Arc<datafusion::physical_plan::PlanProperties>,
}
impl HotScanExec {
#[allow(clippy::too_many_arguments)]
fn new(
source_schema: SchemaRef,
hot: Arc<dyn HotStore>,
table: String,
range: TimeRange,
spec: crate::tiering::store::ScanSpec,
projection: Option<Vec<usize>>,
limit: Option<usize>,
) -> DfResult<Self> {
use datafusion::physical_expr::{EquivalenceProperties, Partitioning};
use datafusion::physical_plan::PlanProperties;
use datafusion::physical_plan::execution_plan::{Boundedness, EmissionType};
let schema = match &projection {
Some(indices) => Arc::new(source_schema.project(indices)?),
None => source_schema.clone(),
};
let properties = Arc::new(PlanProperties::new(
EquivalenceProperties::new(schema.clone()),
Partitioning::UnknownPartitioning(1),
EmissionType::Incremental,
Boundedness::Bounded,
));
Ok(Self {
schema,
source_schema,
hot,
table,
range,
spec,
projection,
limit,
properties,
})
}
}
impl fmt::Debug for HotScanExec {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "HotScanExec")
}
}
impl datafusion::physical_plan::DisplayAs for HotScanExec {
fn fmt_as(
&self,
_t: datafusion::physical_plan::DisplayFormatType,
f: &mut fmt::Formatter<'_>,
) -> fmt::Result {
write!(f, "HotScanExec: streaming")?;
if let Some(limit) = self.limit {
write!(f, ", limit={limit}")?;
}
Ok(())
}
}
impl ExecutionPlan for HotScanExec {
fn name(&self) -> &str {
"HotScanExec"
}
fn as_any(&self) -> &dyn Any {
self
}
fn properties(&self) -> &Arc<datafusion::physical_plan::PlanProperties> {
&self.properties
}
fn children(&self) -> Vec<&Arc<dyn ExecutionPlan>> {
vec![]
}
fn with_new_children(
self: Arc<Self>,
_: Vec<Arc<dyn ExecutionPlan>>,
) -> DfResult<Arc<dyn ExecutionPlan>> {
Ok(self)
}
fn execute(
&self,
_partition: usize,
_context: Arc<datafusion::execution::TaskContext>,
) -> DfResult<datafusion::physical_plan::SendableRecordBatchStream> {
use datafusion::physical_plan::stream::RecordBatchStreamAdapter;
use futures::StreamExt;
let hot = Arc::clone(&self.hot);
let table = self.table.clone();
let range = self.range;
let spec = self.spec.clone();
let table_name = self.table.clone();
let source_schema = self.source_schema.clone();
let projection = self.projection.clone();
let mut remaining = self.limit.unwrap_or(usize::MAX);
let batches = async_stream::stream! {
let started = std::time::Instant::now();
let mut inner = match hot.scan_range(&table, range, &spec).await {
Ok(stream) => stream,
Err(e) => {
yield Err(external(e));
return;
}
};
while let Some(batch) = inner.next().await {
if remaining == 0 {
break;
}
let mapped = batch
.map_err(external)
.and_then(|b| cast_to(&b, &source_schema))
.and_then(|b| match &projection {
Some(indices) => Ok(b.project(indices)?),
None => Ok(b),
});
if let Ok(ref b) = mapped {
remaining = remaining.saturating_sub(b.num_rows());
crate::observe::metrics()
.rows_scanned
.add(b.num_rows() as u64, &crate::observe::table_tier(&table_name, "hot"));
}
yield mapped;
}
crate::observe::metrics().scan_duration.record(
started.elapsed().as_secs_f64(),
&crate::observe::table_tier(&table_name, "hot"),
);
};
Ok(Box::pin(RecordBatchStreamAdapter::new(
self.schema.clone(),
batches,
)))
}
}
fn cast_to(batch: &RecordBatch, schema: &SchemaRef) -> DfResult<RecordBatch> {
if batch.schema() == *schema {
return Ok(batch.clone());
}
let columns = batch
.columns()
.iter()
.zip(schema.fields())
.map(|(array, field)| crate::arrow::compute::cast(array, field.data_type()))
.collect::<std::result::Result<Vec<_>, _>>()?;
Ok(RecordBatch::try_new(schema.clone(), columns)?)
}
fn external(e: crate::Error) -> DataFusionError {
DataFusionError::External(Box::new(e))
}
#[async_trait]
impl TableProvider for TieredTableProvider {
fn as_any(&self) -> &dyn Any {
self
}
fn schema(&self) -> SchemaRef {
self.schema.clone()
}
fn table_type(&self) -> TableType {
TableType::Base
}
async fn scan(
&self,
state: &dyn Session,
projection: Option<&Vec<usize>>,
filters: &[Expr],
limit: Option<usize>,
) -> DfResult<Arc<dyn ExecutionPlan>> {
let watermark = self.watermark().await?;
self.scan_at(state, watermark, projection, filters, limit)
.await
}
fn supports_filters_pushdown(
&self,
filters: &[&Expr],
) -> DfResult<Vec<TableProviderFilterPushDown>> {
Ok(vec![TableProviderFilterPushDown::Inexact; filters.len()])
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::encode::schema;
use crate::tiering::store::{CommitInfo, PartitionId};
use crate::watermark::ArchivalWindow;
use datafusion::logical_expr::{col as df_col, lit};
use datafusion::prelude::SessionContext;
use std::sync::Mutex;
use time::macros::datetime;
use time::{Duration, OffsetDateTime};
const BOUNDARY: OffsetDateTime = datetime!(2026-07-20 00:00 UTC);
#[derive(Default)]
struct Calls {
hot: Vec<TimeRange>,
}
struct FakeHot {
rows: usize,
calls: Mutex<Calls>,
}
#[async_trait]
impl HotStore for FakeHot {
async fn append_reporting(
&self,
_table: &str,
_merge_key: &[String],
_batches: &[RecordBatch],
) -> crate::error::Result<Vec<crate::session::Displacement>> {
Ok(Vec::new())
}
async fn drop_table(&self, _table: &str) -> crate::error::Result<()> {
Ok(())
}
async fn scan_range(
&self,
_table: &str,
range: TimeRange,
_spec: &crate::tiering::store::ScanSpec,
) -> crate::Result<crate::tiering::store::BatchStream> {
self.calls.lock().unwrap().hot.push(range);
let b = batch(self.rows);
Ok(Box::pin(futures::stream::once(async move { Ok(b) })))
}
async fn create_tables(
&self,
_table: &str,
_key: &[String],
_extra: &[crate::arrow::datatypes::Field],
) -> crate::Result<()> {
Ok(())
}
async fn append(
&self,
_t: &str,
_k: &[String],
batches: &[RecordBatch],
) -> crate::Result<u64> {
Ok(batches.iter().map(|b| b.num_rows() as u64).sum())
}
async fn ensure_partitions(
&self,
_t: &str,
_f: OffsetDateTime,
_u: OffsetDateTime,
_s: Duration,
) -> crate::Result<Vec<PartitionId>> {
Ok(vec![])
}
async fn partition_exists(&self, _p: &PartitionId) -> crate::Result<bool> {
Ok(true)
}
async fn detach_partition(&self, _p: &PartitionId) -> crate::Result<()> {
Ok(())
}
async fn scan_detached(
&self,
_p: &PartitionId,
_spec: &crate::tiering::store::ScanSpec,
) -> crate::Result<crate::tiering::store::BatchStream> {
Ok(crate::tiering::store::stream_of(Vec::new()))
}
async fn drop_partition(&self, _p: &PartitionId) -> crate::Result<()> {
Ok(())
}
async fn orphaned_partitions(&self, _t: &str) -> crate::Result<Vec<PartitionId>> {
Ok(vec![])
}
async fn invariant_violations(&self, _t: &str, _w: TieringWatermark) -> crate::Result<u64> {
Ok(0)
}
}
struct FakeCold {
watermark: TieringWatermark,
}
#[async_trait]
impl ColdStore for FakeCold {
async fn purge_table(&self, _table: &str) -> crate::error::Result<()> {
Ok(())
}
async fn create_tables(
&self,
_table: &str,
_identity: &[String],
_extra: &[crate::arrow::datatypes::Field],
) -> crate::Result<()> {
Ok(())
}
async fn watermark(&self, _t: &str) -> crate::Result<TieringWatermark> {
Ok(self.watermark)
}
async fn append_and_commit(
&self,
_t: &str,
_b: crate::tiering::store::BatchStream,
_h: crate::tiering::store::WriteHints,
w: ArchivalWindow,
) -> crate::Result<CommitInfo> {
Ok(CommitInfo {
snapshot_id: 1,
rows: 0,
watermark: w.resulting_watermark(),
})
}
async fn expire_snapshots(
&self,
_t: &str,
_retain_for: time::Duration,
_retain_last: usize,
_now: OffsetDateTime,
) -> crate::Result<usize> {
Ok(0)
}
async fn append_only(
&self,
_t: &str,
_b: crate::tiering::store::BatchStream,
_h: crate::tiering::store::WriteHints,
) -> crate::Result<CommitInfo> {
Ok(CommitInfo {
snapshot_id: 1,
rows: 0,
watermark: self.watermark,
})
}
}
fn batch(rows: usize) -> RecordBatch {
use crate::arrow::array::{
Decimal128Array, StringArray, TimestampMicrosecondArray, UInt8Array,
};
let s = schema::storage_schema(&[]);
let columns = s
.fields()
.iter()
.map(|f| match f.data_type() {
crate::arrow::datatypes::DataType::Utf8 => {
Arc::new(StringArray::from(vec!["x"; rows])) as _
}
crate::arrow::datatypes::DataType::UInt8 => {
Arc::new(UInt8Array::from(vec![0u8; rows])) as _
}
crate::arrow::datatypes::DataType::Timestamp(_, _) => {
Arc::new(TimestampMicrosecondArray::from(vec![0i64; rows]).with_timezone("UTC"))
as _
}
crate::arrow::datatypes::DataType::Decimal128(p, sc) => Arc::new(
Decimal128Array::from(vec![0i128; rows])
.with_precision_and_scale(*p, *sc)
.unwrap(),
) as _,
other => panic!("unhandled {other:?}"),
})
.collect();
RecordBatch::try_new(s, columns).unwrap()
}
fn cold_provider(rows: usize) -> Arc<dyn TableProvider> {
let s = schema::storage_schema(&[]);
Arc::new(MemTable::try_new(s, vec![vec![batch(rows)]]).unwrap())
}
fn provider(hot_rows: usize, cold_rows: usize) -> TieredTableProvider {
TieredTableProvider::new(
"readings",
Arc::new(FakeHot {
rows: hot_rows,
calls: Mutex::new(Calls::default()),
}),
Arc::new(FakeCold {
watermark: TieringWatermark::new(BOUNDARY),
}),
cold_provider(cold_rows),
)
}
fn ts(t: OffsetDateTime) -> Expr {
lit(datafusion::common::ScalarValue::TimestampMicrosecond(
Some((t.unix_timestamp_nanos() / 1_000) as i64),
Some("UTC".into()),
))
}
async fn row_count(p: TieredTableProvider, filters: Vec<Expr>) -> usize {
let ctx = SessionContext::new();
let plan = p
.scan(&ctx.state(), None, &filters, None)
.await
.expect("scan");
datafusion::physical_plan::collect(plan, ctx.task_ctx())
.await
.expect("collect")
.iter()
.map(|b| b.num_rows())
.sum()
}
#[tokio::test]
async fn a_historical_query_reads_only_the_cold_tier() {
let p = provider(7, 3);
let filters = vec![df_col(schema::col::FROM).lt(ts(datetime!(2026-07-10 00:00 UTC)))];
let (_, split) = p.plan_split(&filters).await.unwrap();
assert!(split.is_cold_only());
assert_eq!(row_count(p, filters).await, 3);
}
#[tokio::test]
async fn a_recent_query_reads_only_the_hot_tier() {
let p = provider(7, 3);
let filters = vec![df_col(schema::col::FROM).gt_eq(ts(datetime!(2026-07-25 00:00 UTC)))];
let (_, split) = p.plan_split(&filters).await.unwrap();
assert!(split.hot.is_some() && split.cold.is_none());
assert_eq!(row_count(p, filters).await, 7);
}
#[tokio::test]
async fn a_spanning_query_unions_both_tiers_without_deduplicating() {
let p = provider(7, 3);
let filters = vec![
df_col(schema::col::FROM).gt_eq(ts(datetime!(2026-07-10 00:00 UTC))),
df_col(schema::col::FROM).lt(ts(datetime!(2026-07-30 00:00 UTC))),
];
let (_, split) = p.plan_split(&filters).await.unwrap();
assert!(split.spans_tiers());
assert_eq!(row_count(p, filters).await, 10);
}
#[tokio::test]
async fn an_unfiltered_query_reads_both_tiers() {
let p = provider(7, 3);
assert_eq!(row_count(p, vec![]).await, 10);
}
#[tokio::test]
async fn the_hot_tier_is_asked_only_for_the_range_above_the_watermark() {
let hot = Arc::new(FakeHot {
rows: 1,
calls: Mutex::new(Calls::default()),
});
let p = TieredTableProvider::new(
"readings",
hot.clone(),
Arc::new(FakeCold {
watermark: TieringWatermark::new(BOUNDARY),
}),
cold_provider(1),
);
let filters = vec![
df_col(schema::col::FROM).gt_eq(ts(datetime!(2026-07-10 00:00 UTC))),
df_col(schema::col::FROM).lt(ts(datetime!(2026-07-30 00:00 UTC))),
];
let _ = row_count(p, filters).await;
let calls = hot.calls.lock().unwrap();
assert_eq!(calls.hot.len(), 1);
assert_eq!(calls.hot[0].start(), Some(BOUNDARY));
assert_eq!(calls.hot[0].end(), Some(datetime!(2026-07-30 00:00 UTC)));
}
#[tokio::test]
async fn an_empty_range_produces_a_plan_that_returns_nothing() {
let p = provider(7, 3);
let filters = vec![
df_col(schema::col::FROM).gt_eq(ts(datetime!(2026-08-01 00:00 UTC))),
df_col(schema::col::FROM).lt(ts(datetime!(2026-07-01 00:00 UTC))),
];
let (_, split) = p.plan_split(&filters).await.unwrap();
assert!(split.is_empty());
assert_eq!(row_count(p, filters).await, 0);
}
#[tokio::test]
async fn historical_mode_never_touches_the_hot_tier() {
let hot = Arc::new(FakeHot {
rows: 7,
calls: Mutex::new(Calls::default()),
});
let p = TieredTableProvider::new(
"readings",
hot.clone(),
Arc::new(FakeCold {
watermark: TieringWatermark::new(BOUNDARY),
}),
cold_provider(3),
)
.with_mode(ReadMode::Historical);
assert_eq!(row_count(p, vec![]).await, 3);
assert!(
hot.calls.lock().unwrap().hot.is_empty(),
"reporting queries must place no load on PostgreSQL"
);
}
#[tokio::test]
async fn operational_mode_reads_only_the_recent_window() {
let p = provider(7, 3).with_mode(ReadMode::Operational);
assert_eq!(row_count(p, vec![]).await, 7);
}
#[tokio::test]
async fn an_empty_watermark_sends_everything_to_the_hot_tier() {
let p = TieredTableProvider::new(
"readings",
Arc::new(FakeHot {
rows: 5,
calls: Mutex::new(Calls::default()),
}),
Arc::new(FakeCold {
watermark: TieringWatermark::empty(),
}),
cold_provider(3),
);
let (_, split) = p.plan_split(&[]).await.unwrap();
assert!(split.cold.is_none(), "nothing has been archived");
assert_eq!(row_count(p, vec![]).await, 5);
}
#[test]
fn no_filter_is_ever_reported_exact() {
let p = provider(1, 1);
let time = df_col(schema::col::FROM).gt_eq(ts(BOUNDARY));
let other = df_col(schema::col::MALO_ID).eq(lit("12345678901"));
let got = p.supports_filters_pushdown(&[&time, &other]).unwrap();
assert!(
got.iter()
.all(|p| *p == TableProviderFilterPushDown::Inexact)
);
}
#[test]
fn read_mode_narrows_a_split() {
let both = TierSplit {
cold: Some(TimeRange::unbounded()),
hot: Some(TimeRange::unbounded()),
};
assert!(ReadMode::Unified.apply(both).spans_tiers());
assert!(ReadMode::Historical.apply(both).is_cold_only());
assert!(ReadMode::Operational.apply(both).cold.is_none());
}
}