timeseries_table_format/table/operations/
state_access.rs1use snafu::{ResultExt, Snafu};
4
5use crate::{
6 table::{TableError, TimeSeriesTable},
7 transaction_log::{CommitError, IndexSpec, TableKind, TableState},
8};
9
10#[derive(Debug, Snafu)]
12#[snafu(module, visibility(pub(crate)))]
13#[non_exhaustive]
14pub enum TableStateAccessError {
15 #[snafu(display("Latest table kind is {kind:?}, expected a time-series table"))]
17 NotTimeSeries {
18 kind: TableKind,
20 },
21
22 #[snafu(context(false), display("Table state transaction-log error: {source}"))]
24 Commit {
25 #[snafu(source, backtrace)]
27 source: CommitError,
28 },
29}
30
31fn time_series_index_from_state(state: &TableState) -> Result<IndexSpec, TableStateAccessError> {
32 match &state.table_meta.kind {
33 TableKind::TimeSeries(index) => Ok(index.clone()),
34 kind => Err(TableStateAccessError::NotTimeSeries { kind: kind.clone() }),
35 }
36}
37
38impl TimeSeriesTable {
39 pub async fn current_version(&self) -> Result<u64, TableError> {
41 self.log
42 .load_current_version()
43 .await
44 .map_err(TableStateAccessError::from)
45 .context(crate::table::error::StateAccessSnafu)
46 }
47
48 pub async fn load_latest_state(&self) -> Result<TableState, TableError> {
50 let result: Result<TableState, TableStateAccessError> = async {
51 let state = self
52 .log
53 .rebuild_table_state()
54 .await
55 .map_err(TableStateAccessError::from)?;
56 time_series_index_from_state(&state)?;
57 Ok(state)
58 }
59 .await;
60 result.context(crate::table::error::StateAccessSnafu)
61 }
62
63 #[tracing::instrument(
65 name = "table.refresh",
66 target = "timeseries_table_format::table",
67 level = "debug",
68 skip_all,
69 fields(
70 previous_version = tracing::field::Empty,
71 observed_version = tracing::field::Empty,
72 refreshed = tracing::field::Empty,
73 new_version = tracing::field::Empty,
74 outcome = tracing::field::Empty
75 )
76 )]
77 pub async fn refresh(&mut self) -> Result<bool, TableError> {
78 tracing::Span::current().record("previous_version", self.state.version);
79 let result: Result<bool, TableStateAccessError> = async {
80 let current = self
81 .log
82 .load_current_version()
83 .await
84 .map_err(TableStateAccessError::from)?;
85 tracing::Span::current().record("observed_version", current);
86 if current == self.state.version {
87 return Ok(false);
88 }
89
90 let state = self
91 .log
92 .rebuild_table_state()
93 .await
94 .map_err(TableStateAccessError::from)?;
95 let index = time_series_index_from_state(&state)?;
96 self.state = state;
97 self.index = index;
98 Ok(true)
99 }
100 .await;
101
102 let span = tracing::Span::current();
103 match &result {
104 Ok(refreshed) => {
105 span.record("refreshed", *refreshed);
106 if *refreshed {
107 span.record("new_version", self.state.version);
108 span.record("outcome", "succeeded");
109 } else {
110 span.record("outcome", "no_change");
111 }
112 }
113 Err(_) => {
114 span.record("outcome", "failed");
115 }
116 }
117 result.context(crate::table::error::StateAccessSnafu)
118 }
119}
120
121#[cfg(test)]
122mod tests {
123 use super::*;
124 use crate::{
125 coverage::EntityValue,
126 storage::{TableLocation, layout},
127 table::{
128 OptimizeError,
129 test_util::{
130 TestResult, TraceCapture, assert_capture_excludes, assert_debug_span,
131 assert_no_event, captured_span, make_basic_table_meta, utc_datetime,
132 },
133 },
134 transaction_log::{
135 CommitError, IndexKind, LogAction, TableProtocolError, TimeIndexGranularity,
136 TransactionLogStore,
137 },
138 };
139 use futures::StreamExt;
140 use tempfile::TempDir;
141
142 #[tokio::test]
143 async fn refresh_reports_no_change_and_applies_a_new_index() -> TestResult {
144 let tmp = TempDir::new()?;
145 let location = TableLocation::local(tmp.path());
146 let meta = make_basic_table_meta();
147 let mut table = TimeSeriesTable::create(location.clone(), meta.clone()).await?;
148 let no_change_capture = TraceCapture::default();
149
150 assert!(!no_change_capture.run(table.refresh()).await?);
151 assert_debug_span(
152 &no_change_capture,
153 "table.refresh",
154 &[
155 ("previous_version", Some("1")),
156 ("observed_version", Some("1")),
157 ("refreshed", Some("false")),
158 ("new_version", None),
159 ("outcome", Some("no_change")),
160 ],
161 );
162 assert_eq!(
163 captured_span(&no_change_capture, "table.refresh").target,
164 "timeseries_table_format::table"
165 );
166
167 let mut updated_meta = meta;
168 let TableKind::TimeSeries(index) = &mut updated_meta.kind else {
169 unreachable!("test metadata is time-series");
170 };
171 index.kind = IndexKind::Timestamp {
172 index_granularity: TimeIndexGranularity::Minutes(5),
173 timezone: None,
174 };
175 TransactionLogStore::new(location)
176 .commit_with_expected_version(1, vec![LogAction::UpdateTableMeta(updated_meta)])
177 .await?;
178 let update_capture = TraceCapture::default();
179
180 assert!(update_capture.run(table.refresh()).await?);
181 assert_eq!(table.state().version, 2);
182 assert!(matches!(
183 table.index_spec().kind,
184 IndexKind::Timestamp {
185 index_granularity: TimeIndexGranularity::Minutes(5),
186 ..
187 }
188 ));
189 assert_debug_span(
190 &update_capture,
191 "table.refresh",
192 &[
193 ("previous_version", Some("1")),
194 ("observed_version", Some("2")),
195 ("refreshed", Some("true")),
196 ("new_version", Some("2")),
197 ("outcome", Some("succeeded")),
198 ],
199 );
200 assert_no_event(&update_capture, "table.refresh");
201 assert_capture_excludes(&update_capture, &[&tmp.path().display().to_string()]);
202 Ok(())
203 }
204
205 #[tokio::test]
206 async fn state_access_preserves_commit_failures_without_mutating_state() -> TestResult {
207 let current_tmp = TempDir::new()?;
208 let current_table = TimeSeriesTable::create(
209 TableLocation::local(current_tmp.path()),
210 make_basic_table_meta(),
211 )
212 .await?;
213 let current_path = current_tmp.path().join(layout::current_rel_path());
214 std::fs::remove_file(¤t_path)?;
215 std::fs::create_dir(¤t_path)?;
216 assert!(matches!(
217 current_table
218 .current_version()
219 .await
220 .expect_err("unreadable CURRENT must fail"),
221 TableError::StateAccess {
222 source: TableStateAccessError::Commit {
223 source: CommitError::Storage { .. }
224 }
225 }
226 ));
227
228 let refresh_tmp = TempDir::new()?;
229 let mut table = TimeSeriesTable::create(
230 TableLocation::local(refresh_tmp.path()),
231 make_basic_table_meta(),
232 )
233 .await?;
234 let state_before = table.state().clone();
235 std::fs::write(
236 refresh_tmp.path().join(layout::commit_rel_path(2)),
237 b"not json",
238 )?;
239 std::fs::write(refresh_tmp.path().join(layout::current_rel_path()), b"2\n")?;
240 assert!(matches!(
241 table
242 .load_latest_state()
243 .await
244 .expect_err("corrupt commit must fail"),
245 TableError::StateAccess {
246 source: TableStateAccessError::Commit {
247 source: CommitError::CommitDeserialization { .. }
248 }
249 }
250 ));
251 assert!(table.refresh().await.is_err());
252 assert_eq!(table.state(), &state_before);
253 Ok(())
254 }
255
256 #[tokio::test]
257 async fn state_access_rejects_a_generic_update_without_mutating_state() -> TestResult {
258 let tmp = TempDir::new()?;
259 let location = TableLocation::local(tmp.path());
260 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
261 let state_before = table.state().clone();
262 let mut generic_meta = make_basic_table_meta();
263 generic_meta.kind = TableKind::Generic;
264 TransactionLogStore::new(location)
265 .commit_with_expected_version(1, vec![LogAction::UpdateTableMeta(generic_meta)])
266 .await?;
267
268 assert!(matches!(
269 table
270 .load_latest_state()
271 .await
272 .expect_err("generic update must fail"),
273 TableError::StateAccess {
274 source: TableStateAccessError::NotTimeSeries {
275 kind: TableKind::Generic
276 }
277 }
278 ));
279 assert!(matches!(
280 table.refresh().await.expect_err("generic update must fail"),
281 TableError::StateAccess {
282 source: TableStateAccessError::NotTimeSeries {
283 kind: TableKind::Generic
284 }
285 }
286 ));
287 assert_eq!(table.state(), &state_before);
288 Ok(())
289 }
290
291 #[tokio::test]
292 async fn refresh_applies_reader_and_writer_requirements_by_operation() -> TestResult {
293 let tmp = TempDir::new()?;
294 let location = TableLocation::local(tmp.path());
295 let mut table = TimeSeriesTable::create(location.clone(), make_basic_table_meta()).await?;
296
297 let mut writer_meta = table.state().table_meta.clone();
298 writer_meta
299 .required_writer_features
300 .insert("future_writer".to_string());
301 TransactionLogStore::new(location.clone())
302 .commit_with_expected_version(1, vec![LogAction::UpdateTableMeta(writer_meta.clone())])
303 .await?;
304
305 assert!(table.refresh().await?);
306 assert_eq!(table.state().version, 2);
307
308 let start = utc_datetime(2025, 1, 1, 0, 0, 0);
309 let end = utc_datetime(2025, 1, 1, 1, 0, 0);
310 let mut scan = table.scan_range(start, end).await?;
311 assert!(scan.next().await.is_none());
312 assert_eq!(
313 table
314 .coverage_ratio_for_entity_range(&[("symbol", EntityValue::from("A"))], start, end,)
315 .await?,
316 0.0
317 );
318 assert!(matches!(
319 table
320 .optimize()
321 .await
322 .expect_err("unknown writer feature must reject optimize after refresh"),
323 TableError::Optimize {
324 source: OptimizeError::Protocol {
325 source: TableProtocolError::UnsupportedWriterFeatures { features },
326 ..
327 }
328 } if features == ["future_writer"]
329 ));
330
331 let state_before_reader_upgrade = table.state().clone();
332 let mut reader_meta = writer_meta;
333 reader_meta
334 .required_reader_features
335 .insert("future_reader".to_string());
336 TransactionLogStore::new(location)
337 .commit_with_expected_version(2, vec![LogAction::UpdateTableMeta(reader_meta)])
338 .await?;
339
340 assert!(matches!(
341 table
342 .refresh()
343 .await
344 .expect_err("unknown reader feature must reject refresh"),
345 TableError::StateAccess {
346 source: TableStateAccessError::Commit {
347 source: CommitError::Protocol {
348 source: TableProtocolError::UnsupportedReaderFeatures { features },
349 ..
350 }
351 }
352 } if features == ["future_reader"]
353 ));
354 assert_eq!(table.state(), &state_before_reader_upgrade);
355 Ok(())
356 }
357}