Skip to main content

timeseries_table_format/metadata/
protocol.rs

1//! Table protocol compatibility and feature negotiation.
2
3use std::collections::BTreeSet;
4
5use serde::{Deserialize, Deserializer};
6use snafu::prelude::*;
7
8use crate::metadata::table::TableMeta;
9
10/// Current table protocol version written by new tables.
11///
12/// This changes only when the core metadata or commit-log envelope can no
13/// longer be decoded under the current protocol. Optional capabilities are
14/// declared separately in the reader and writer feature sets.
15pub const TABLE_PROTOCOL_VERSION: u32 = 7;
16
17const SUPPORTED_READER_FEATURES: &[&str] = &[];
18const SUPPORTED_WRITER_FEATURES: &[&str] = &[];
19
20#[derive(Deserialize)]
21pub(crate) struct RawTableProtocolRequirements {
22    protocol_version: u64,
23    #[serde(deserialize_with = "deserialize_required_features")]
24    required_reader_features: BTreeSet<String>,
25    #[serde(
26        rename = "required_writer_features",
27        deserialize_with = "deserialize_required_features"
28    )]
29    _required_writer_features: BTreeSet<String>,
30}
31
32pub(crate) fn deserialize_required_features<'de, D>(
33    deserializer: D,
34) -> Result<BTreeSet<String>, D::Error>
35where
36    D: Deserializer<'de>,
37{
38    use serde::de::Error as _;
39
40    let features = Vec::<String>::deserialize(deserializer)?;
41    let mut unique = BTreeSet::new();
42    for feature in features {
43        if !is_valid_feature_name(&feature) {
44            return Err(D::Error::custom(format!(
45                "invalid table feature {feature:?}; expected [a-z][a-z0-9_]*"
46            )));
47        }
48        if !unique.insert(feature.clone()) {
49            return Err(D::Error::custom(format!(
50                "duplicate table feature {feature:?}"
51            )));
52        }
53    }
54    Ok(unique)
55}
56
57fn is_valid_feature_name(feature: &str) -> bool {
58    let mut bytes = feature.bytes();
59    matches!(bytes.next(), Some(b'a'..=b'z'))
60        && bytes.all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'_')
61}
62
63/// A table protocol requirement or transition is not supported.
64#[derive(Debug, Snafu, PartialEq, Eq)]
65#[non_exhaustive]
66pub enum TableProtocolError {
67    /// The table uses a core protocol version this client cannot decode.
68    #[snafu(display("unsupported table protocol version: expected {expected}, found {found}"))]
69    UnsupportedVersion {
70        /// Core protocol version supported by this client.
71        expected: u32,
72        /// Core protocol version required by the table.
73        found: u64,
74    },
75
76    /// The table requires reader features this client does not support.
77    #[snafu(display("unsupported table reader features: {features:?}"))]
78    UnsupportedReaderFeatures {
79        /// Unsupported feature identifiers in canonical order.
80        features: Vec<String>,
81    },
82
83    /// The table requires writer features this client does not support.
84    #[snafu(display("unsupported table writer features: {features:?}"))]
85    UnsupportedWriterFeatures {
86        /// Unsupported feature identifiers in canonical order.
87        features: Vec<String>,
88    },
89
90    /// A metadata update removed previously required reader features.
91    #[snafu(display("table metadata removed required reader features: {features:?}"))]
92    ReaderFeaturesRemoved {
93        /// Removed feature identifiers in canonical order.
94        features: Vec<String>,
95    },
96
97    /// A metadata update removed previously required writer features.
98    #[snafu(display("table metadata removed required writer features: {features:?}"))]
99    WriterFeaturesRemoved {
100        /// Removed feature identifiers in canonical order.
101        features: Vec<String>,
102    },
103
104    /// A metadata update decreased the core protocol version.
105    #[snafu(display("table metadata decreased protocol version from {previous} to {next}"))]
106    ProtocolVersionDecreased {
107        /// Protocol version before the metadata update.
108        previous: u32,
109        /// Protocol version after the metadata update.
110        next: u32,
111    },
112}
113
114impl TableMeta {
115    /// Returns the on-disk table protocol version.
116    pub fn protocol_version(&self) -> u32 {
117        self.protocol_version
118    }
119
120    /// Returns the features a client must support to read the table.
121    pub fn required_reader_features(&self) -> &BTreeSet<String> {
122        &self.required_reader_features
123    }
124
125    /// Returns the features a client must support to write the table.
126    pub fn required_writer_features(&self) -> &BTreeSet<String> {
127        &self.required_writer_features
128    }
129
130    pub(crate) fn ensure_read_compatible(&self) -> Result<(), TableProtocolError> {
131        self.ensure_read_compatible_with(SUPPORTED_READER_FEATURES)
132    }
133
134    pub(crate) fn ensure_write_compatible(&self) -> Result<(), TableProtocolError> {
135        self.ensure_write_compatible_with(SUPPORTED_READER_FEATURES, SUPPORTED_WRITER_FEATURES)
136    }
137
138    fn ensure_read_compatible_with(
139        &self,
140        supported_reader_features: &[&str],
141    ) -> Result<(), TableProtocolError> {
142        ensure_read_compatible(
143            u64::from(self.protocol_version),
144            &self.required_reader_features,
145            supported_reader_features,
146        )
147    }
148
149    fn ensure_write_compatible_with(
150        &self,
151        supported_reader_features: &[&str],
152        supported_writer_features: &[&str],
153    ) -> Result<(), TableProtocolError> {
154        self.ensure_read_compatible_with(supported_reader_features)?;
155
156        let unsupported =
157            unsupported_features(&self.required_writer_features, supported_writer_features);
158        if !unsupported.is_empty() {
159            return Err(TableProtocolError::UnsupportedWriterFeatures {
160                features: unsupported,
161            });
162        }
163        Ok(())
164    }
165
166    pub(crate) fn ensure_valid_transition_to(&self, next: &Self) -> Result<(), TableProtocolError> {
167        if next.protocol_version < self.protocol_version {
168            return Err(TableProtocolError::ProtocolVersionDecreased {
169                previous: self.protocol_version,
170                next: next.protocol_version,
171            });
172        }
173
174        let removed_reader_features = self
175            .required_reader_features
176            .difference(&next.required_reader_features)
177            .cloned()
178            .collect::<Vec<_>>();
179        if !removed_reader_features.is_empty() {
180            return Err(TableProtocolError::ReaderFeaturesRemoved {
181                features: removed_reader_features,
182            });
183        }
184
185        let removed_writer_features = self
186            .required_writer_features
187            .difference(&next.required_writer_features)
188            .cloned()
189            .collect::<Vec<_>>();
190        if !removed_writer_features.is_empty() {
191            return Err(TableProtocolError::WriterFeaturesRemoved {
192                features: removed_writer_features,
193            });
194        }
195        Ok(())
196    }
197}
198
199impl RawTableProtocolRequirements {
200    pub(crate) fn ensure_read_compatible(&self) -> Result<(), TableProtocolError> {
201        ensure_read_compatible(
202            self.protocol_version,
203            &self.required_reader_features,
204            SUPPORTED_READER_FEATURES,
205        )
206    }
207}
208
209fn ensure_read_compatible(
210    protocol_version: u64,
211    required_reader_features: &BTreeSet<String>,
212    supported_reader_features: &[&str],
213) -> Result<(), TableProtocolError> {
214    if protocol_version != u64::from(TABLE_PROTOCOL_VERSION) {
215        return Err(TableProtocolError::UnsupportedVersion {
216            expected: TABLE_PROTOCOL_VERSION,
217            found: protocol_version,
218        });
219    }
220
221    let unsupported = unsupported_features(required_reader_features, supported_reader_features);
222    if !unsupported.is_empty() {
223        return Err(TableProtocolError::UnsupportedReaderFeatures {
224            features: unsupported,
225        });
226    }
227    Ok(())
228}
229
230fn unsupported_features(required: &BTreeSet<String>, supported: &[&str]) -> Vec<String> {
231    required
232        .iter()
233        .filter(|feature| !supported.contains(&feature.as_str()))
234        .cloned()
235        .collect()
236}
237
238#[cfg(test)]
239mod tests {
240    use crate::metadata::{
241        index::{IndexKind, IndexSpec, TimeIndexGranularity},
242        table::TableMeta,
243    };
244
245    use super::*;
246
247    fn sample_table_meta() -> TableMeta {
248        TableMeta::new_time_series(IndexSpec {
249            column: "ts".to_string(),
250            entity_columns: vec!["symbol".to_string()],
251            kind: IndexKind::Timestamp {
252                index_granularity: TimeIndexGranularity::Minutes(1),
253                timezone: None,
254            },
255        })
256    }
257
258    #[test]
259    fn table_meta_requires_valid_unique_feature_lists() {
260        let valid = serde_json::to_value(sample_table_meta()).unwrap();
261        let invalid_lists = [
262            serde_json::Value::Null,
263            serde_json::json!(["duplicate", "duplicate"]),
264            serde_json::json!(["Uppercase"]),
265            serde_json::json!([""]),
266        ];
267
268        for invalid in invalid_lists {
269            let mut json = valid.clone();
270            json["required_reader_features"] = invalid;
271            assert!(serde_json::from_value::<TableMeta>(json).is_err());
272        }
273
274        for field in ["required_reader_features", "required_writer_features"] {
275            let mut json = valid.clone();
276            json.as_object_mut().unwrap().remove(field);
277            assert!(serde_json::from_value::<TableMeta>(json).is_err());
278        }
279    }
280
281    #[test]
282    fn table_protocol_compatibility_is_operation_specific() {
283        let mut meta = sample_table_meta();
284        meta.required_reader_features
285            .extend(["reader_a".to_string(), "reader_b".to_string()]);
286        meta.required_writer_features.insert("writer_a".to_string());
287
288        assert!(
289            meta.ensure_read_compatible_with(&["reader_a", "reader_b"])
290                .is_ok()
291        );
292        assert_eq!(
293            meta.ensure_read_compatible_with(&["reader_b"]),
294            Err(TableProtocolError::UnsupportedReaderFeatures {
295                features: vec!["reader_a".to_string()],
296            })
297        );
298        assert_eq!(
299            meta.ensure_write_compatible_with(&["reader_a", "reader_b"], &[]),
300            Err(TableProtocolError::UnsupportedWriterFeatures {
301                features: vec!["writer_a".to_string()],
302            })
303        );
304        assert!(
305            meta.ensure_write_compatible_with(&["reader_a", "reader_b"], &["writer_a"])
306                .is_ok()
307        );
308        assert_eq!(
309            meta.ensure_write_compatible(),
310            Err(TableProtocolError::UnsupportedReaderFeatures {
311                features: vec!["reader_a".to_string(), "reader_b".to_string()],
312            })
313        );
314    }
315
316    #[test]
317    fn table_protocol_transition_requirements_are_monotonic() {
318        let mut previous = sample_table_meta();
319        previous
320            .required_reader_features
321            .extend(["reader_a".to_string(), "reader_b".to_string()]);
322        previous
323            .required_writer_features
324            .insert("writer_a".to_string());
325
326        let mut additive = previous.clone();
327        additive
328            .required_writer_features
329            .insert("writer_b".to_string());
330        assert!(previous.ensure_valid_transition_to(&additive).is_ok());
331
332        let mut removed = previous.clone();
333        removed.required_reader_features.clear();
334        assert_eq!(
335            previous.ensure_valid_transition_to(&removed),
336            Err(TableProtocolError::ReaderFeaturesRemoved {
337                features: vec!["reader_a".to_string(), "reader_b".to_string()],
338            })
339        );
340
341        let mut removed = previous.clone();
342        removed.required_writer_features.clear();
343        assert_eq!(
344            previous.ensure_valid_transition_to(&removed),
345            Err(TableProtocolError::WriterFeaturesRemoved {
346                features: vec!["writer_a".to_string()],
347            })
348        );
349
350        let mut decreased = previous.clone();
351        decreased.protocol_version -= 1;
352        assert_eq!(
353            previous.ensure_valid_transition_to(&decreased),
354            Err(TableProtocolError::ProtocolVersionDecreased {
355                previous: TABLE_PROTOCOL_VERSION,
356                next: TABLE_PROTOCOL_VERSION - 1,
357            })
358        );
359    }
360}