1use bytes::Bytes;
12use lance_core::{Error, Result, datatypes::Schema};
13use lance_encoding::{
14 decoder::{ColumnInfo, DecoderPlugins, PageInfo},
15 encoder::EncodedBatch,
16};
17use lance_io::{scheduler::FileScheduler, traits::Writer};
18use std::{collections::BTreeMap, future::Future, sync::Arc};
19
20use crate::{
21 format::pbfile,
22 reader::{
23 BufferDescriptor, CachedFileMetadata, FileMetadataIndex, FileMetadataProvider, FileReader,
24 FileReaderOptions, PreparedProjection, ProjectedFileReader, RawFileMetadata,
25 ReadProjection, ReaderProjection,
26 },
27 version::ConcreteFileVersion,
28 writer::{FileWriter, FileWriterOptions},
29};
30use lance_core::cache::LanceCache;
31
32pub mod v1;
33pub mod v2_0;
34pub mod v2_1;
35pub mod v2_2;
36pub mod v2_3;
37
38pub(crate) fn read_projection(version: ConcreteFileVersion) -> Result<Arc<dyn ReadProjection>> {
39 match version {
40 ConcreteFileVersion::V1 => Err(Error::internal(
41 "current reader composition received Lance v1".to_string(),
42 )),
43 ConcreteFileVersion::V2_0 => Ok(v2_0::read_projection()),
44 ConcreteFileVersion::V2_1 => Ok(v2_1::read_projection()),
45 ConcreteFileVersion::V2_2 => Ok(v2_2::read_projection()),
46 ConcreteFileVersion::V2_3 => Ok(v2_3::read_projection()),
47 }
48}
49
50pub enum OpenedFileReader {
52 V1 {
54 major_version: u16,
56 minor_version: u16,
58 },
59 Current(FileReader),
61}
62
63pub(crate) fn finish_metadata(
64 version: ConcreteFileVersion,
65 metadata: RawFileMetadata,
66) -> Result<CachedFileMetadata> {
67 match version {
68 ConcreteFileVersion::V1 => Err(Error::internal(
69 "current metadata dispatch received a Lance v1 file".to_string(),
70 )),
71 ConcreteFileVersion::V2_0 => v2_0::finish_metadata(metadata),
72 ConcreteFileVersion::V2_1 => v2_1::finish_metadata(metadata),
73 ConcreteFileVersion::V2_2 => v2_2::finish_metadata(metadata),
74 ConcreteFileVersion::V2_3 => v2_3::finish_metadata(metadata),
75 }
76}
77
78pub(crate) fn validate_external_metadata(
81 version: ConcreteFileVersion,
82 schema: &Schema,
83 metadata: &CachedFileMetadata,
84) -> Result<u64> {
85 let projection = reader_projection_from_whole_schema(schema, version);
86 if projection.column_indices.len() != metadata.column_infos.len() {
87 return Err(Error::invalid_input(format!(
88 "schema requires {} physical columns but file metadata contains {}",
89 projection.column_indices.len(),
90 metadata.column_infos.len()
91 )));
92 }
93 FileReader::validate_projection(&projection, metadata)?;
94 for (expected_index, column) in metadata.column_infos.iter().enumerate() {
95 if column.index != expected_index as u32 {
96 return Err(Error::invalid_input(format!(
97 "physical column {} reports index {}",
98 expected_index, column.index
99 )));
100 }
101 }
102 let prepared = PreparedProjection {
103 column_infos: metadata.column_infos.clone(),
104 decoder_projection: projection,
105 };
106 read_projection(version)?.read_length(&prepared)
107}
108
109pub(crate) fn finish_metadata_index(index: FileMetadataIndex) -> Result<FileMetadataIndex> {
110 match index.version {
111 ConcreteFileVersion::V1 => Err(Error::version_conflict(
112 "Attempt to use the Lance current-format reader with a v1 metadata index".to_string(),
113 0,
114 2,
115 )),
116 ConcreteFileVersion::V2_0 => v2_0::finish_metadata_index(index),
117 ConcreteFileVersion::V2_1 => v2_1::finish_metadata_index(index),
118 ConcreteFileVersion::V2_2 => v2_2::finish_metadata_index(index),
119 ConcreteFileVersion::V2_3 => v2_3::finish_metadata_index(index),
120 }
121}
122
123pub(crate) fn decode_column_metadata(
124 version: ConcreteFileVersion,
125 column_metadatas: &[pbfile::ColumnMetadata],
126) -> Result<Vec<Arc<ColumnInfo>>> {
127 match version {
128 ConcreteFileVersion::V1 => Err(Error::not_supported(
129 "self-described batches are not part of the Lance v1 grammar".to_string(),
130 )),
131 ConcreteFileVersion::V2_0 => v2_0::decode_column_metadata(column_metadatas),
132 ConcreteFileVersion::V2_1 => v2_1::decode_column_metadata(column_metadatas),
133 ConcreteFileVersion::V2_2 => v2_2::decode_column_metadata(column_metadatas),
134 ConcreteFileVersion::V2_3 => v2_3::decode_column_metadata(column_metadatas),
135 }
136}
137
138pub(crate) fn validate_global_buffers(
139 version: ConcreteFileVersion,
140 buffers: &[BufferDescriptor],
141) -> Result<()> {
142 match version {
143 ConcreteFileVersion::V1 => Ok(()),
144 ConcreteFileVersion::V2_0 => v2_0::validate_global_buffers(buffers),
145 ConcreteFileVersion::V2_1 => v2_1::validate_global_buffers(buffers),
146 ConcreteFileVersion::V2_2 => v2_2::validate_global_buffers(buffers),
147 ConcreteFileVersion::V2_3 => v2_3::validate_global_buffers(buffers),
148 }
149}
150
151pub fn reader_projection_from_field_ids(
152 version: ConcreteFileVersion,
153 schema: &Schema,
154 field_id_to_column_index: &BTreeMap<u32, u32>,
155) -> Result<ReaderProjection> {
156 Ok(match version {
157 ConcreteFileVersion::V1 => v1::projection_from_field_ids(schema, field_id_to_column_index),
158 ConcreteFileVersion::V2_0 => {
159 v2_0::projection_from_field_ids(schema, field_id_to_column_index)
160 }
161 ConcreteFileVersion::V2_1 => {
162 v2_1::projection_from_field_ids(schema, field_id_to_column_index)
163 }
164 ConcreteFileVersion::V2_2 => {
165 v2_2::projection_from_field_ids(schema, field_id_to_column_index)
166 }
167 ConcreteFileVersion::V2_3 => {
168 v2_3::projection_from_field_ids(schema, field_id_to_column_index)
169 }
170 })
171}
172
173pub fn reader_projection_from_whole_schema(
174 schema: &Schema,
175 version: ConcreteFileVersion,
176) -> ReaderProjection {
177 match version {
178 ConcreteFileVersion::V1 => v1::projection_from_whole_schema(schema),
179 ConcreteFileVersion::V2_0 => v2_0::projection_from_whole_schema(schema),
180 ConcreteFileVersion::V2_1 => v2_1::projection_from_whole_schema(schema),
181 ConcreteFileVersion::V2_2 => v2_2::projection_from_whole_schema(schema),
182 ConcreteFileVersion::V2_3 => v2_3::projection_from_whole_schema(schema),
183 }
184}
185
186pub fn reader_projection_from_column_names(
187 version: ConcreteFileVersion,
188 schema: &Schema,
189 column_names: &[&str],
190) -> Result<ReaderProjection> {
191 match version {
192 ConcreteFileVersion::V1 => v1::projection_from_column_names(schema, column_names),
193 ConcreteFileVersion::V2_0 => v2_0::projection_from_column_names(schema, column_names),
194 ConcreteFileVersion::V2_1 => v2_1::projection_from_column_names(schema, column_names),
195 ConcreteFileVersion::V2_2 => v2_2::projection_from_column_names(schema, column_names),
196 ConcreteFileVersion::V2_3 => v2_3::projection_from_column_names(schema, column_names),
197 }
198}
199
200pub fn physical_column_count(
202 version: ConcreteFileVersion,
203 field: &lance_core::datatypes::Field,
204) -> usize {
205 match version {
206 ConcreteFileVersion::V1 => v1::physical_column_count(field),
207 ConcreteFileVersion::V2_0 => v2_0::physical_column_count(field),
208 ConcreteFileVersion::V2_1 => v2_1::physical_column_count(field),
209 ConcreteFileVersion::V2_2 => v2_2::physical_column_count(field),
210 ConcreteFileVersion::V2_3 => v2_3::physical_column_count(field),
211 }
212}
213
214pub fn data_file_columns(version: ConcreteFileVersion, schema: &Schema) -> (Vec<i32>, Vec<i32>) {
216 match version {
217 ConcreteFileVersion::V1 => v1::data_file_columns(schema),
218 ConcreteFileVersion::V2_0 => v2_0::data_file_columns(schema),
219 ConcreteFileVersion::V2_1 => v2_1::data_file_columns(schema),
220 ConcreteFileVersion::V2_2 => v2_2::data_file_columns(schema),
221 ConcreteFileVersion::V2_3 => v2_3::data_file_columns(schema),
222 }
223}
224
225pub async fn copy_external_metadata_column<Copy, CopyFuture>(
231 version: ConcreteFileVersion,
232 schema: &Schema,
233 column_index: usize,
234 has_existing_pages: bool,
235 copy: Copy,
236) -> Result<()>
237where
238 Copy: FnOnce() -> CopyFuture + Send,
239 CopyFuture: Future<Output = Result<()>> + Send,
240{
241 match version {
242 ConcreteFileVersion::V1 => Err(Error::not_supported(
243 "binary-copy metadata operations are not supported for Lance v1".to_string(),
244 )),
245 ConcreteFileVersion::V2_0 => {
246 if v2_0::should_copy_external_metadata_column(schema, column_index, has_existing_pages)
247 {
248 copy().await
249 } else {
250 Ok(())
251 }
252 }
253 ConcreteFileVersion::V2_1 | ConcreteFileVersion::V2_2 | ConcreteFileVersion::V2_3 => {
254 copy().await
255 }
256 }
257}
258
259pub fn finalize_external_metadata_column(
261 version: ConcreteFileVersion,
262 schema: &Schema,
263 column_index: usize,
264 pages: &mut Vec<PageInfo>,
265 num_rows: u64,
266) -> Result<()> {
267 match version {
268 ConcreteFileVersion::V1 => Err(Error::not_supported(
269 "binary-copy metadata operations are not supported for Lance v1".to_string(),
270 )),
271 ConcreteFileVersion::V2_0 => {
272 v2_0::finalize_external_metadata_column(schema, column_index, pages, num_rows);
273 Ok(())
274 }
275 ConcreteFileVersion::V2_1 | ConcreteFileVersion::V2_2 | ConcreteFileVersion::V2_3 => Ok(()),
276 }
277}
278
279pub async fn open_projected_reader<OpenIndexed, IndexedFuture, OpenFull, FullFuture>(
286 version: ConcreteFileVersion,
287 projection: &ReaderProjection,
288 prefer_indexed: bool,
289 open_indexed: OpenIndexed,
290 open_full: OpenFull,
291) -> Result<ProjectedFileReader>
292where
293 OpenIndexed: FnOnce() -> IndexedFuture + Send,
294 IndexedFuture: Future<Output = Result<Option<ProjectedFileReader>>> + Send,
295 OpenFull: FnOnce() -> FullFuture + Send,
296 FullFuture: Future<Output = Result<ProjectedFileReader>> + Send,
297{
298 match version {
299 ConcreteFileVersion::V1 => Err(Error::not_supported(
300 "projected current-format readers cannot open Lance v1 files".to_string(),
301 )),
302 ConcreteFileVersion::V2_0 => open_full().await,
303 ConcreteFileVersion::V2_1 | ConcreteFileVersion::V2_2 | ConcreteFileVersion::V2_3 => {
304 if prefer_indexed
305 && FileMetadataProvider::projection_matches_indexed_metadata(projection)
306 && let Some(reader) = open_indexed().await?
307 {
308 return Ok(reader);
309 }
310 open_full().await
311 }
312 }
313}
314
315pub async fn open_self_described_reader(
320 scheduler: FileScheduler,
321 decoder_plugins: Arc<DecoderPlugins>,
322 cache: &LanceCache,
323 options: FileReaderOptions,
324) -> Result<OpenedFileReader> {
325 FileReader::try_open_for_dispatch(scheduler, None, decoder_plugins, cache, options).await
326}
327
328pub fn create_writer(
333 version: ConcreteFileVersion,
334 object_writer: Box<dyn Writer>,
335 schema: Schema,
336 options: FileWriterOptions,
337) -> Result<FileWriter> {
338 match version {
339 ConcreteFileVersion::V1 => Err(Error::not_supported(
340 "Lance v1 files must be created with versions::v1::writer::FileWriter".to_string(),
341 )),
342 ConcreteFileVersion::V2_0 => {
343 v2_0::create_writer(object_writer, schema, options).map(Into::into)
344 }
345 ConcreteFileVersion::V2_1 => {
346 v2_1::create_writer(object_writer, schema, options).map(Into::into)
347 }
348 ConcreteFileVersion::V2_2 => {
349 v2_2::create_writer(object_writer, schema, options).map(Into::into)
350 }
351 ConcreteFileVersion::V2_3 => {
352 v2_3::create_writer(object_writer, schema, options).map(Into::into)
353 }
354 }
355}
356
357pub fn create_lazy_writer(
359 version: ConcreteFileVersion,
360 object_writer: Box<dyn Writer>,
361 options: FileWriterOptions,
362) -> Result<FileWriter> {
363 match version {
364 ConcreteFileVersion::V1 => Err(Error::not_supported(
365 "legacy v1 files require an explicit schema and manifest provider".to_string(),
366 )),
367 ConcreteFileVersion::V2_0 => Ok(v2_0::create_lazy_writer(object_writer, options).into()),
368 ConcreteFileVersion::V2_1 => Ok(v2_1::create_lazy_writer(object_writer, options).into()),
369 ConcreteFileVersion::V2_2 => Ok(v2_2::create_lazy_writer(object_writer, options).into()),
370 ConcreteFileVersion::V2_3 => Ok(v2_3::create_lazy_writer(object_writer, options).into()),
371 }
372}
373
374pub fn encode_self_described_batch(
376 version: ConcreteFileVersion,
377 batch: &EncodedBatch,
378) -> Result<Bytes> {
379 match version {
380 ConcreteFileVersion::V1 => Err(Error::not_supported(
381 "Lance v1 does not support self-described current-format batches".to_string(),
382 )),
383 ConcreteFileVersion::V2_0 => v2_0::encode_self_described_batch(batch),
384 ConcreteFileVersion::V2_1 => v2_1::encode_self_described_batch(batch),
385 ConcreteFileVersion::V2_2 => v2_2::encode_self_described_batch(batch),
386 ConcreteFileVersion::V2_3 => v2_3::encode_self_described_batch(batch),
387 }
388}
389
390pub fn encode_mini_batch(version: ConcreteFileVersion, batch: &EncodedBatch) -> Result<Bytes> {
392 match version {
393 ConcreteFileVersion::V1 => Err(Error::not_supported(
394 "Lance v1 does not support mini-lance current-format batches".to_string(),
395 )),
396 ConcreteFileVersion::V2_0 => v2_0::encode_mini_batch(batch),
397 ConcreteFileVersion::V2_1 => v2_1::encode_mini_batch(batch),
398 ConcreteFileVersion::V2_2 => v2_2::encode_mini_batch(batch),
399 ConcreteFileVersion::V2_3 => v2_3::encode_mini_batch(batch),
400 }
401}
402
403#[cfg(test)]
404mod tests {
405 use std::sync::Arc;
406
407 use super::*;
408
409 #[test]
410 fn v1_rejects_current_format_embedded_batches() {
411 let batch = EncodedBatch {
412 data: Bytes::new(),
413 page_table: Vec::new(),
414 schema: Arc::new(Schema::default()),
415 top_level_columns: Vec::new(),
416 num_rows: 0,
417 };
418
419 assert!(matches!(
420 encode_self_described_batch(ConcreteFileVersion::V1, &batch),
421 Err(Error::NotSupported { .. })
422 ));
423 assert!(matches!(
424 encode_mini_batch(ConcreteFileVersion::V1, &batch),
425 Err(Error::NotSupported { .. })
426 ));
427 }
428}