Skip to main content

timeseries_table_format/table/operations/
state_access.rs

1//! Reading and refreshing a table handle's committed state.
2
3use snafu::{ResultExt, Snafu};
4
5use crate::{
6    table::{TableError, TimeSeriesTable},
7    transaction_log::{CommitError, IndexSpec, TableKind, TableState},
8};
9
10/// Errors owned by table state reads and refreshes.
11#[derive(Debug, Snafu)]
12#[snafu(module, visibility(pub(crate)))]
13#[non_exhaustive]
14pub enum TableStateAccessError {
15    /// Latest metadata no longer describes a time-series table.
16    #[snafu(display("Latest table kind is {kind:?}, expected a time-series table"))]
17    NotTimeSeries {
18        /// Rejected table kind.
19        kind: TableKind,
20    },
21
22    /// Reading or replaying the transaction log failed.
23    #[snafu(context(false), display("Table state transaction-log error: {source}"))]
24    Commit {
25        /// Complete transaction-log failure.
26        #[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    /// Load the current log version without mutating the in-memory state.
40    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    /// Rebuild and return the latest time-series table state.
49    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    /// Refresh in-memory state if the transaction log has advanced.
64    #[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(&current_path)?;
215        std::fs::create_dir(&current_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}