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
78#[allow(dead_code)]
79pub(crate) fn validate_protocol(protocol: &DeltaProtocolInfo) -> Result<(), DeltaReaderError> {
80    if !matches!(protocol.min_reader_version, 1 | 2)
81        && protocol.min_reader_version != TABLE_FEATURES_READER_VERSION
82    {
83        return UnsupportedProtocolSnafu {
84            reason: "unsupported_reader_version",
85        }
86        .fail();
87    }
88
89    if protocol
90        .reader_features
91        .iter()
92        .any(|feature| !SUPPORTED_READER_FEATURES.contains(&feature.as_str()))
93    {
94        return UnsupportedProtocolSnafu {
95            reason: "unsupported_reader_feature",
96        }
97        .fail();
98    }
99
100    Ok(())
101}
102
103#[cfg(test)]
104mod tests {
105    use std::{
106        fs,
107        path::{Path, PathBuf},
108        time::{SystemTime, UNIX_EPOCH},
109    };
110
111    use super::{DeltaProtocolInfo, validate_protocol};
112    use crate::{
113        DeltaReaderError, DeltaReaderPhase, DeltaSnapshotSelection, DeltaStorageOptions,
114        snapshot::load_delta_table_snapshot_blocking,
115    };
116
117    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}}"#;
118
119    struct DeltaLogTable(PathBuf);
120
121    impl DeltaLogTable {
122        fn new(name: &str, protocol_json: &str) -> Result<Self, Box<dyn std::error::Error>> {
123            let nanos = SystemTime::now().duration_since(UNIX_EPOCH)?.as_nanos();
124            let path = Path::new("target")
125                .join("delta-arrow-reader-protocol-tests")
126                .join(format!("{}-{name}-{nanos}", std::process::id()));
127            let log_path = path.join("_delta_log");
128            fs::create_dir_all(&log_path)?;
129            fs::write(
130                log_path.join("00000000000000000000.json"),
131                format!("{protocol_json}\n{METADATA_JSON}\n"),
132            )?;
133            Ok(Self(path))
134        }
135
136        fn load(&self) -> Result<crate::snapshot::LoadedDeltaTableSnapshot, DeltaReaderError> {
137            load_delta_table_snapshot_blocking(
138                &self.0.to_string_lossy(),
139                &DeltaStorageOptions::new(),
140                DeltaSnapshotSelection::Latest,
141            )
142        }
143    }
144
145    impl Drop for DeltaLogTable {
146        fn drop(&mut self) {
147            let _ = fs::remove_dir_all(&self.0);
148        }
149    }
150
151    #[test]
152    fn extracts_type_widening_parity_features_in_kernel_order()
153    -> Result<(), Box<dyn std::error::Error>> {
154        let table = DeltaLogTable::new(
155            "type-widening",
156            r#"{"protocol":{"minReaderVersion":3,"minWriterVersion":7,"readerFeatures":["timestampNtz","typeWidening-preview"],"writerFeatures":["timestampNtz","typeWidening-preview"]}}"#,
157        )?;
158        let loaded = table.load()?;
159        let protocol = loaded.protocol_info();
160
161        assert_eq!(protocol.snapshot_version(), 0);
162        assert_eq!(protocol.min_reader_version(), 3);
163        assert_eq!(protocol.min_writer_version(), 7);
164        assert_eq!(
165            protocol.reader_features(),
166            ["timestampNtz", "typeWidening-preview"]
167        );
168        assert_eq!(
169            protocol.writer_features(),
170            ["timestampNtz", "typeWidening-preview"]
171        );
172        validate_protocol(protocol)?;
173        Ok(())
174    }
175
176    #[test]
177    fn loads_and_accepts_the_frozen_reader_feature_set() -> Result<(), Box<dyn std::error::Error>> {
178        let table = DeltaLogTable::new(
179            "all-supported-features",
180            r#"{"protocol":{"minReaderVersion":3,"minWriterVersion":7,"readerFeatures":["timestampNtz","deletionVectors","columnMapping","v2Checkpoint","vacuumProtocolCheck","typeWidening","typeWidening-preview"],"writerFeatures":["timestampNtz","deletionVectors","columnMapping","v2Checkpoint","vacuumProtocolCheck","typeWidening","typeWidening-preview"]}}"#,
181        )?;
182        let loaded = table.load()?;
183        let protocol = loaded.protocol_info();
184
185        assert_eq!(
186            protocol.reader_features(),
187            [
188                "timestampNtz",
189                "deletionVectors",
190                "columnMapping",
191                "v2Checkpoint",
192                "vacuumProtocolCheck",
193                "typeWidening",
194                "typeWidening-preview",
195            ]
196        );
197        validate_protocol(protocol)?;
198        Ok(())
199    }
200
201    #[test]
202    fn preserves_legacy_versions_and_treats_writer_only_features_as_diagnostic()
203    -> Result<(), Box<dyn std::error::Error>> {
204        for (name, protocol_json, reader_version, writer_features) in [
205            (
206                "legacy-column-mapping",
207                r#"{"protocol":{"minReaderVersion":2,"minWriterVersion":5}}"#,
208                2,
209                &[][..],
210            ),
211            (
212                "writer-only",
213                r#"{"protocol":{"minReaderVersion":3,"minWriterVersion":7,"readerFeatures":[],"writerFeatures":["inCommitTimestamp"]}}"#,
214                3,
215                &["inCommitTimestamp"][..],
216            ),
217        ] {
218            let table = DeltaLogTable::new(name, protocol_json)?;
219            let loaded = table.load()?;
220            let protocol = loaded.protocol_info();
221
222            assert_eq!(protocol.min_reader_version(), reader_version);
223            assert!(protocol.reader_features().is_empty());
224            assert_eq!(protocol.writer_features(), writer_features);
225            validate_protocol(protocol)?;
226        }
227        Ok(())
228    }
229
230    #[test]
231    fn unsupported_feature_remains_inspectable_until_validation()
232    -> Result<(), Box<dyn std::error::Error>> {
233        let table = DeltaLogTable::new(
234            "unknown-feature",
235            r#"{"protocol":{"minReaderVersion":3,"minWriterVersion":7,"readerFeatures":["madeUpFeature"],"writerFeatures":["madeUpFeature"]}}"#,
236        )?;
237        let loaded = table.load()?;
238        let protocol = loaded.protocol_info();
239
240        assert_eq!(protocol.reader_features(), ["madeUpFeature"]);
241        let error = validate_protocol(protocol).expect_err("unknown feature must fail");
242        assert!(matches!(
243            error,
244            DeltaReaderError::UnsupportedProtocol { .. }
245        ));
246        assert_eq!(error.phase(), DeltaReaderPhase::Protocol);
247        assert_eq!(
248            error.to_string(),
249            "delta reader error: phase=protocol error=unsupported_protocol reason=unsupported_reader_feature"
250        );
251        Ok(())
252    }
253
254    #[test]
255    fn unsupported_version_fails_at_the_kernel_boundary_and_policy_backstop()
256    -> Result<(), Box<dyn std::error::Error>> {
257        let table = DeltaLogTable::new(
258            "future-version",
259            r#"{"protocol":{"minReaderVersion":4,"minWriterVersion":7}}"#,
260        )?;
261        let load_error = match table.load() {
262            Ok(_) => panic!("Kernel must reject the future reader version"),
263            Err(error) => error,
264        };
265        assert!(matches!(load_error, DeltaReaderError::SnapshotLoad { .. }));
266        assert_eq!(load_error.phase(), DeltaReaderPhase::Snapshot);
267
268        let protocol = DeltaProtocolInfo {
269            snapshot_version: 0,
270            min_reader_version: 4,
271            min_writer_version: 7,
272            reader_features: Vec::new(),
273            writer_features: Vec::new(),
274        };
275        let error = validate_protocol(&protocol).expect_err("future version must fail");
276        assert_eq!(error.phase(), DeltaReaderPhase::Protocol);
277        assert_eq!(error.as_str(), "unsupported_protocol");
278        Ok(())
279    }
280}