timeseries_table_format/metadata/
protocol.rs1use std::collections::BTreeSet;
4
5use serde::{Deserialize, Deserializer};
6use snafu::prelude::*;
7
8use crate::metadata::table::TableMeta;
9
10pub 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#[derive(Debug, Snafu, PartialEq, Eq)]
65#[non_exhaustive]
66pub enum TableProtocolError {
67 #[snafu(display("unsupported table protocol version: expected {expected}, found {found}"))]
69 UnsupportedVersion {
70 expected: u32,
72 found: u64,
74 },
75
76 #[snafu(display("unsupported table reader features: {features:?}"))]
78 UnsupportedReaderFeatures {
79 features: Vec<String>,
81 },
82
83 #[snafu(display("unsupported table writer features: {features:?}"))]
85 UnsupportedWriterFeatures {
86 features: Vec<String>,
88 },
89
90 #[snafu(display("table metadata removed required reader features: {features:?}"))]
92 ReaderFeaturesRemoved {
93 features: Vec<String>,
95 },
96
97 #[snafu(display("table metadata removed required writer features: {features:?}"))]
99 WriterFeaturesRemoved {
100 features: Vec<String>,
102 },
103
104 #[snafu(display("table metadata decreased protocol version from {previous} to {next}"))]
106 ProtocolVersionDecreased {
107 previous: u32,
109 next: u32,
111 },
112}
113
114impl TableMeta {
115 pub fn protocol_version(&self) -> u32 {
117 self.protocol_version
118 }
119
120 pub fn required_reader_features(&self) -> &BTreeSet<String> {
122 &self.required_reader_features
123 }
124
125 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}