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