Skip to main content

delta_arrow_reader/
protocol.rs

1//! Delta protocol metadata and reader compatibility policy.
2
3use crate::{
4    DeltaReaderError,
5    error::UnsupportedProtocolSnafu,
6    kernel::{
7        DeltaKernelProtocol, KernelSnapshot, TABLE_FEATURES_READER_VERSION,
8        snapshot_protocol_report,
9    },
10};
11
12#[allow(dead_code)]
13const SUPPORTED_READER_FEATURES: &[&str] = &[
14    "timestampNtz",
15    "deletionVectors",
16    "columnMapping",
17    "v2Checkpoint",
18    "vacuumProtocolCheck",
19    "typeWidening",
20    "typeWidening-preview",
21];
22
23/// Protocol metadata captured from one immutable Delta snapshot.
24#[derive(Debug, Clone, PartialEq, Eq)]
25pub struct DeltaProtocolInfo {
26    snapshot_version: u64,
27    min_reader_version: i32,
28    min_writer_version: i32,
29    reader_features: Vec<String>,
30    writer_features: Vec<String>,
31}
32
33impl DeltaProtocolInfo {
34    pub(crate) fn from_snapshot(snapshot: &KernelSnapshot) -> Self {
35        let version = snapshot.version();
36        let DeltaKernelProtocol {
37            min_reader_version,
38            min_writer_version,
39            reader_features,
40            writer_features,
41        } = snapshot_protocol_report(snapshot);
42
43        Self {
44            snapshot_version: version,
45            min_reader_version,
46            min_writer_version,
47            reader_features,
48            writer_features,
49        }
50    }
51
52    /// Returns the loaded snapshot version.
53    pub const fn snapshot_version(&self) -> u64 {
54        self.snapshot_version
55    }
56
57    /// Returns the minimum Delta reader protocol version.
58    pub const fn min_reader_version(&self) -> i32 {
59        self.min_reader_version
60    }
61
62    /// Returns the minimum Delta writer protocol version.
63    pub const fn min_writer_version(&self) -> i32 {
64        self.min_writer_version
65    }
66
67    /// Returns the required reader feature names in deterministic Kernel order.
68    pub fn reader_features(&self) -> &[String] {
69        &self.reader_features
70    }
71
72    /// Returns the writer feature names in deterministic Kernel order.
73    pub fn writer_features(&self) -> &[String] {
74        &self.writer_features
75    }
76
77    /// Returns the first required reader feature unsupported by this crate.
78    pub fn first_unsupported_reader_feature(&self) -> Option<&str> {
79        self.reader_features
80            .iter()
81            .map(String::as_str)
82            .find(|feature| !SUPPORTED_READER_FEATURES.contains(feature))
83    }
84}
85
86#[allow(dead_code)]
87pub(crate) fn validate_protocol(protocol: &DeltaProtocolInfo) -> Result<(), DeltaReaderError> {
88    if !matches!(protocol.min_reader_version, 1 | 2)
89        && protocol.min_reader_version != TABLE_FEATURES_READER_VERSION
90    {
91        return UnsupportedProtocolSnafu {
92            reason: "unsupported_reader_version",
93        }
94        .fail();
95    }
96
97    if protocol.first_unsupported_reader_feature().is_some() {
98        return UnsupportedProtocolSnafu {
99            reason: "unsupported_reader_feature",
100        }
101        .fail();
102    }
103
104    Ok(())
105}
106
107#[cfg(test)]
108mod tests {
109    use std::{
110        fs,
111        path::{Path, PathBuf},
112        time::{SystemTime, UNIX_EPOCH},
113    };
114
115    use super::{DeltaProtocolInfo, validate_protocol};
116    use crate::{
117        DeltaReaderError, DeltaReaderPhase, DeltaSnapshotSelection, DeltaStorageOptions,
118        snapshot::load_delta_table_snapshot_blocking,
119    };
120
121    const METADATA_JSON: &str = r#"{"metaData":{"id":"delta-arrow-reader-test","format":{"provider":"parquet","options":{}},"schemaString":"{\"type\":\"struct\",\"fields\":[{\"name\":\"id\",\"type\":\"integer\",\"nullable\":true,\"metadata\":{}}]}","partitionColumns":[],"configuration":{},"createdTime":1587968585495}}"#;
122
123    struct DeltaLogTable(PathBuf);
124
125    impl DeltaLogTable {
126        fn new(name: &str, protocol_json: &str) -> Result<Self, Box<dyn std::error::Error>> {
127            let nanos = SystemTime::now().duration_since(UNIX_EPOCH)?.as_nanos();
128            let path = Path::new("target")
129                .join("delta-arrow-reader-protocol-tests")
130                .join(format!("{}-{name}-{nanos}", std::process::id()));
131            let log_path = path.join("_delta_log");
132            fs::create_dir_all(&log_path)?;
133            fs::write(
134                log_path.join("00000000000000000000.json"),
135                format!("{protocol_json}\n{METADATA_JSON}\n"),
136            )?;
137            Ok(Self(path))
138        }
139
140        fn load(&self) -> Result<crate::snapshot::LoadedDeltaTableSnapshot, DeltaReaderError> {
141            load_delta_table_snapshot_blocking(
142                &self.0.to_string_lossy(),
143                &DeltaStorageOptions::new(),
144                DeltaSnapshotSelection::Latest,
145            )
146        }
147    }
148
149    impl Drop for DeltaLogTable {
150        fn drop(&mut self) {
151            let _ = fs::remove_dir_all(&self.0);
152        }
153    }
154
155    #[test]
156    fn extracts_type_widening_parity_features_in_kernel_order()
157    -> Result<(), Box<dyn std::error::Error>> {
158        let table = DeltaLogTable::new(
159            "type-widening",
160            r#"{"protocol":{"minReaderVersion":3,"minWriterVersion":7,"readerFeatures":["timestampNtz","typeWidening-preview"],"writerFeatures":["timestampNtz","typeWidening-preview"]}}"#,
161        )?;
162        let loaded = table.load()?;
163        let protocol = loaded.protocol_info();
164
165        assert_eq!(protocol.snapshot_version(), 0);
166        assert_eq!(protocol.min_reader_version(), 3);
167        assert_eq!(protocol.min_writer_version(), 7);
168        assert_eq!(
169            protocol.reader_features(),
170            ["timestampNtz", "typeWidening-preview"]
171        );
172        assert_eq!(
173            protocol.writer_features(),
174            ["timestampNtz", "typeWidening-preview"]
175        );
176        validate_protocol(protocol)?;
177        Ok(())
178    }
179
180    #[test]
181    fn loads_and_accepts_the_frozen_reader_feature_set() -> Result<(), Box<dyn std::error::Error>> {
182        let table = DeltaLogTable::new(
183            "all-supported-features",
184            r#"{"protocol":{"minReaderVersion":3,"minWriterVersion":7,"readerFeatures":["timestampNtz","deletionVectors","columnMapping","v2Checkpoint","vacuumProtocolCheck","typeWidening","typeWidening-preview"],"writerFeatures":["timestampNtz","deletionVectors","columnMapping","v2Checkpoint","vacuumProtocolCheck","typeWidening","typeWidening-preview"]}}"#,
185        )?;
186        let loaded = table.load()?;
187        let protocol = loaded.protocol_info();
188
189        assert_eq!(
190            protocol.reader_features(),
191            [
192                "timestampNtz",
193                "deletionVectors",
194                "columnMapping",
195                "v2Checkpoint",
196                "vacuumProtocolCheck",
197                "typeWidening",
198                "typeWidening-preview",
199            ]
200        );
201        validate_protocol(protocol)?;
202        Ok(())
203    }
204
205    #[test]
206    fn preserves_legacy_versions_and_treats_writer_only_features_as_diagnostic()
207    -> Result<(), Box<dyn std::error::Error>> {
208        for (name, protocol_json, reader_version, writer_features) in [
209            (
210                "legacy-column-mapping",
211                r#"{"protocol":{"minReaderVersion":2,"minWriterVersion":5}}"#,
212                2,
213                &[][..],
214            ),
215            (
216                "writer-only",
217                r#"{"protocol":{"minReaderVersion":3,"minWriterVersion":7,"readerFeatures":[],"writerFeatures":["inCommitTimestamp"]}}"#,
218                3,
219                &["inCommitTimestamp"][..],
220            ),
221        ] {
222            let table = DeltaLogTable::new(name, protocol_json)?;
223            let loaded = table.load()?;
224            let protocol = loaded.protocol_info();
225
226            assert_eq!(protocol.min_reader_version(), reader_version);
227            assert!(protocol.reader_features().is_empty());
228            assert_eq!(protocol.writer_features(), writer_features);
229            validate_protocol(protocol)?;
230        }
231        Ok(())
232    }
233
234    #[test]
235    fn unsupported_feature_remains_inspectable_until_validation()
236    -> Result<(), Box<dyn std::error::Error>> {
237        let table = DeltaLogTable::new(
238            "unknown-feature",
239            r#"{"protocol":{"minReaderVersion":3,"minWriterVersion":7,"readerFeatures":["madeUpFeature"],"writerFeatures":["madeUpFeature"]}}"#,
240        )?;
241        let loaded = table.load()?;
242        let protocol = loaded.protocol_info();
243
244        assert_eq!(protocol.reader_features(), ["madeUpFeature"]);
245        assert_eq!(
246            protocol.first_unsupported_reader_feature(),
247            Some("madeUpFeature")
248        );
249        let error = validate_protocol(protocol).expect_err("unknown feature must fail");
250        assert!(matches!(
251            error,
252            DeltaReaderError::UnsupportedProtocol { .. }
253        ));
254        assert_eq!(error.phase(), DeltaReaderPhase::Protocol);
255        assert_eq!(
256            error.to_string(),
257            "delta reader error: phase=protocol error=unsupported_protocol reason=unsupported_reader_feature"
258        );
259        Ok(())
260    }
261
262    #[test]
263    fn unsupported_version_fails_at_the_kernel_boundary_and_policy_backstop()
264    -> Result<(), Box<dyn std::error::Error>> {
265        let table = DeltaLogTable::new(
266            "future-version",
267            r#"{"protocol":{"minReaderVersion":4,"minWriterVersion":7}}"#,
268        )?;
269        let load_error = match table.load() {
270            Ok(_) => panic!("Kernel must reject the future reader version"),
271            Err(error) => error,
272        };
273        assert!(matches!(load_error, DeltaReaderError::SnapshotLoad { .. }));
274        assert_eq!(load_error.phase(), DeltaReaderPhase::Snapshot);
275
276        let protocol = DeltaProtocolInfo {
277            snapshot_version: 0,
278            min_reader_version: 4,
279            min_writer_version: 7,
280            reader_features: Vec::new(),
281            writer_features: Vec::new(),
282        };
283        let error = validate_protocol(&protocol).expect_err("future version must fail");
284        assert_eq!(error.phase(), DeltaReaderPhase::Protocol);
285        assert_eq!(error.as_str(), "unsupported_protocol");
286        Ok(())
287    }
288}