use snafu::{ResultExt, Snafu};
use crate::{
table::{TableError, TimeSeriesTable},
transaction_log::{CommitError, IndexSpec, TableKind, TableState},
};
#[derive(Debug, Snafu)]
#[snafu(module, visibility(pub(crate)))]
#[non_exhaustive]
pub enum TableStateAccessError {
#[snafu(display("Latest table kind is {kind:?}, expected a time-series table"))]
NotTimeSeries {
kind: TableKind,
},
#[snafu(context(false), display("Table state transaction-log error: {source}"))]
Commit {
#[snafu(source, backtrace)]
source: CommitError,
},
}
fn time_series_index_from_state(state: &TableState) -> Result<IndexSpec, TableStateAccessError> {
match &state.table_meta.kind {
TableKind::TimeSeries(index) => Ok(index.clone()),
kind => Err(TableStateAccessError::NotTimeSeries { kind: kind.clone() }),
}
}
impl TimeSeriesTable {
pub async fn current_version(&self) -> Result<u64, TableError> {
self.log
.load_current_version()
.await
.map_err(TableStateAccessError::from)
.context(crate::table::error::StateAccessSnafu)
}
pub async fn load_latest_state(&self) -> Result<TableState, TableError> {
let result: Result<TableState, TableStateAccessError> = async {
let state = self
.log
.rebuild_table_state()
.await
.map_err(TableStateAccessError::from)?;
time_series_index_from_state(&state)?;
Ok(state)
}
.await;
result.context(crate::table::error::StateAccessSnafu)
}
#[tracing::instrument(
name = "table.refresh",
target = "timeseries_table_format::table",
level = "debug",
skip_all,
fields(
previous_version = tracing::field::Empty,
observed_version = tracing::field::Empty,
refreshed = tracing::field::Empty,
new_version = tracing::field::Empty,
outcome = tracing::field::Empty
)
)]
pub async fn refresh(&mut self) -> Result<bool, TableError> {
tracing::Span::current().record("previous_version", self.state.version);
let result: Result<bool, TableStateAccessError> = async {
let current = self
.log
.load_current_version()
.await
.map_err(TableStateAccessError::from)?;
tracing::Span::current().record("observed_version", current);
if current == self.state.version {
return Ok(false);
}
let state = self
.log
.rebuild_table_state()
.await
.map_err(TableStateAccessError::from)?;
let index = time_series_index_from_state(&state)?;
self.state = state;
self.index = index;
Ok(true)
}
.await;
let span = tracing::Span::current();
match &result {
Ok(refreshed) => {
span.record("refreshed", *refreshed);
if *refreshed {
span.record("new_version", self.state.version);
span.record("outcome", "succeeded");
} else {
span.record("outcome", "no_change");
}
}
Err(_) => {
span.record("outcome", "failed");
}
}
result.context(crate::table::error::StateAccessSnafu)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::{
coverage::EntityValue,
storage::{TableLocation, layout},
table::{
OptimizeError,
test_util::{
TestResult, TraceCapture, assert_capture_excludes, assert_debug_span,
assert_no_event, captured_span, make_basic_table_meta, utc_datetime,
},
},
transaction_log::{
CommitError, IndexKind, LogAction, TableProtocolError, TimeIndexGranularity,
TransactionLogStore,
},
};
use futures::StreamExt;
use tempfile::TempDir;
#[tokio::test]
async fn refresh_reports_no_change_and_applies_a_new_index() -> TestResult {
let tmp = TempDir::new()?;
let location = TableLocation::local(tmp.path());
let meta = make_basic_table_meta();
let mut table = TimeSeriesTable::create(location.clone(), meta.clone()).await?;
let no_change_capture = TraceCapture::default();
assert!(!no_change_capture.run(table.refresh()).await?);
assert_debug_span(
&no_change_capture,
"table.refresh",
&[
("previous_version", Some("1")),
("observed_version", Some("1")),
("refreshed", Some("false")),
("new_version", None),
("outcome", Some("no_change")),
],
);
assert_eq!(
captured_span(&no_change_capture, "table.refresh").target,
"timeseries_table_format::table"
);
let mut updated_meta = meta;
let TableKind::TimeSeries(index) = &mut updated_meta.kind else {
unreachable!("test metadata is time-series");
};
index.kind = IndexKind::Timestamp {
index_granularity: TimeIndexGranularity::Minutes(5),
timezone: None,
};
TransactionLogStore::new(location)
.commit_with_expected_version(1, vec![LogAction::UpdateTableMeta(updated_meta)])
.await?;
let update_capture = TraceCapture::default();
assert!(update_capture.run(table.refresh()).await?);
assert_eq!(table.state().version, 2);
assert!(matches!(
table.index_spec().kind,
IndexKind::Timestamp {
index_granularity: TimeIndexGranularity::Minutes(5),
..
}
));
assert_debug_span(
&update_capture,
"table.refresh",
&[
("previous_version", Some("1")),
("observed_version", Some("2")),
("refreshed", Some("true")),
("new_version", Some("2")),
("outcome", Some("succeeded")),
],
);
assert_no_event(&update_capture, "table.refresh");
assert_capture_excludes(&update_capture, &[&tmp.path().display().to_string()]);
Ok(())
}
#[tokio::test]
async fn state_access_preserves_commit_failures_without_mutating_state() -> TestResult {
let current_tmp = TempDir::new()?;
let current_table = TimeSeriesTable::create(
TableLocation::local(current_tmp.path()),
make_basic_table_meta(),
)
.await?;
let current_path = current_tmp.path().join(layout::current_rel_path());
std::fs::remove_file(¤t_path)?;
std::fs::create_dir(¤t_path)?;
assert!(matches!(
current_table
.current_version()
.await
.expect_err("unreadable CURRENT must fail"),
TableError::StateAccess {
source: TableStateAccessError::Commit {
source: CommitError::Storage { .. }
}
}
));
let refresh_tmp = TempDir::new()?;
let mut table = TimeSeriesTable::create(
TableLocation::local(refresh_tmp.path()),
make_basic_table_meta(),
)
.await?;
let state_before = table.state().clone();
std::fs::write(
refresh_tmp.path().join(layout::commit_rel_path(2)),
b"not json",
)?;
std::fs::write(refresh_tmp.path().join(layout::current_rel_path()), b"2\n")?;
assert!(matches!(
table
.load_latest_state()
.await
.expect_err("corrupt commit must fail"),
TableError::StateAccess {
source: TableStateAccessError::Commit {
source: CommitError::CommitDeserialization { .. }
}
}
));
assert!(table.refresh().await.is_err());
assert_eq!(table.state(), &state_before);
Ok(())
}
#[tokio::test]
async fn state_access_rejects_a_generic_update_without_mutating_state() -> TestResult {
let tmp = TempDir::new()?;
let location = TableLocation::local(tmp.path());
let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
let state_before = table.state().clone();
let mut generic_meta = make_basic_table_meta();
generic_meta.kind = TableKind::Generic;
TransactionLogStore::new(location)
.commit_with_expected_version(1, vec![LogAction::UpdateTableMeta(generic_meta)])
.await?;
assert!(matches!(
table
.load_latest_state()
.await
.expect_err("generic update must fail"),
TableError::StateAccess {
source: TableStateAccessError::NotTimeSeries {
kind: TableKind::Generic
}
}
));
assert!(matches!(
table.refresh().await.expect_err("generic update must fail"),
TableError::StateAccess {
source: TableStateAccessError::NotTimeSeries {
kind: TableKind::Generic
}
}
));
assert_eq!(table.state(), &state_before);
Ok(())
}
#[tokio::test]
async fn refresh_applies_reader_and_writer_requirements_by_operation() -> TestResult {
let tmp = TempDir::new()?;
let location = TableLocation::local(tmp.path());
let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
let mut writer_meta = table.state().table_meta.clone();
writer_meta
.required_writer_features
.insert("future_writer".to_string());
TransactionLogStore::new(location.clone())
.commit_with_expected_version(1, vec![LogAction::UpdateTableMeta(writer_meta.clone())])
.await?;
assert!(table.refresh().await?);
assert_eq!(table.state().version, 2);
let start = utc_datetime(2025, 1, 1, 0, 0, 0);
let end = utc_datetime(2025, 1, 1, 1, 0, 0);
let mut scan = table.scan_range(start, end).await?;
assert!(scan.next().await.is_none());
assert_eq!(
table
.coverage_ratio_for_entity_range(&[("symbol", EntityValue::from("A"))], start, end,)
.await?,
0.0
);
assert!(matches!(
table
.optimize()
.await
.expect_err("unknown writer feature must reject optimize after refresh"),
TableError::Optimize {
source: OptimizeError::Protocol {
source: TableProtocolError::UnsupportedWriterFeatures { features },
..
}
} if features == ["future_writer"]
));
let state_before_reader_upgrade = table.state().clone();
let mut reader_meta = writer_meta;
reader_meta
.required_reader_features
.insert("future_reader".to_string());
TransactionLogStore::new(location)
.commit_with_expected_version(2, vec![LogAction::UpdateTableMeta(reader_meta)])
.await?;
assert!(matches!(
table
.refresh()
.await
.expect_err("unknown reader feature must reject refresh"),
TableError::StateAccess {
source: TableStateAccessError::Commit {
source: CommitError::Protocol {
source: TableProtocolError::UnsupportedReaderFeatures { features },
..
}
}
} if features == ["future_reader"]
));
assert_eq!(table.state(), &state_before_reader_upgrade);
Ok(())
}
}