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, ProjectedFileReader, RawFileMetadata, ReadProjection, ReaderProjection,
25 },
26 version::ConcreteFileVersion,
27 writer::{FileWriter, FileWriterOptions},
28};
29use lance_core::cache::LanceCache;
30
31pub mod v1;
32pub mod v2_0;
33pub mod v2_1;
34pub mod v2_2;
35pub mod v2_3;
36
37pub(crate) fn read_projection(version: ConcreteFileVersion) -> Result<Arc<dyn ReadProjection>> {
38 match version {
39 ConcreteFileVersion::V1 => Err(Error::internal(
40 "current reader composition received Lance v1".to_string(),
41 )),
42 ConcreteFileVersion::V2_0 => Ok(v2_0::read_projection()),
43 ConcreteFileVersion::V2_1 => Ok(v2_1::read_projection()),
44 ConcreteFileVersion::V2_2 => Ok(v2_2::read_projection()),
45 ConcreteFileVersion::V2_3 => Ok(v2_3::read_projection()),
46 }
47}
48
49pub enum OpenedFileReader {
51 V1 {
53 major_version: u16,
55 minor_version: u16,
57 },
58 Current(FileReader),
60}
61
62pub(crate) fn finish_metadata(
63 version: ConcreteFileVersion,
64 metadata: RawFileMetadata,
65) -> Result<CachedFileMetadata> {
66 match version {
67 ConcreteFileVersion::V1 => Err(Error::internal(
68 "current metadata dispatch received a Lance v1 file".to_string(),
69 )),
70 ConcreteFileVersion::V2_0 => v2_0::finish_metadata(metadata),
71 ConcreteFileVersion::V2_1 => v2_1::finish_metadata(metadata),
72 ConcreteFileVersion::V2_2 => v2_2::finish_metadata(metadata),
73 ConcreteFileVersion::V2_3 => v2_3::finish_metadata(metadata),
74 }
75}
76
77pub(crate) fn finish_metadata_index(index: FileMetadataIndex) -> Result<FileMetadataIndex> {
78 match index.version {
79 ConcreteFileVersion::V1 => Err(Error::version_conflict(
80 "Attempt to use the Lance current-format reader with a v1 metadata index".to_string(),
81 0,
82 2,
83 )),
84 ConcreteFileVersion::V2_0 => v2_0::finish_metadata_index(index),
85 ConcreteFileVersion::V2_1 => v2_1::finish_metadata_index(index),
86 ConcreteFileVersion::V2_2 => v2_2::finish_metadata_index(index),
87 ConcreteFileVersion::V2_3 => v2_3::finish_metadata_index(index),
88 }
89}
90
91pub(crate) fn decode_column_metadata(
92 version: ConcreteFileVersion,
93 column_metadatas: &[pbfile::ColumnMetadata],
94) -> Result<Vec<Arc<ColumnInfo>>> {
95 match version {
96 ConcreteFileVersion::V1 => Err(Error::not_supported(
97 "self-described batches are not part of the Lance v1 grammar".to_string(),
98 )),
99 ConcreteFileVersion::V2_0 => v2_0::decode_column_metadata(column_metadatas),
100 ConcreteFileVersion::V2_1 => v2_1::decode_column_metadata(column_metadatas),
101 ConcreteFileVersion::V2_2 => v2_2::decode_column_metadata(column_metadatas),
102 ConcreteFileVersion::V2_3 => v2_3::decode_column_metadata(column_metadatas),
103 }
104}
105
106pub(crate) fn validate_global_buffers(
107 version: ConcreteFileVersion,
108 buffers: &[BufferDescriptor],
109) -> Result<()> {
110 match version {
111 ConcreteFileVersion::V1 => Ok(()),
112 ConcreteFileVersion::V2_0 => v2_0::validate_global_buffers(buffers),
113 ConcreteFileVersion::V2_1 => v2_1::validate_global_buffers(buffers),
114 ConcreteFileVersion::V2_2 => v2_2::validate_global_buffers(buffers),
115 ConcreteFileVersion::V2_3 => v2_3::validate_global_buffers(buffers),
116 }
117}
118
119pub fn reader_projection_from_field_ids(
120 version: ConcreteFileVersion,
121 schema: &Schema,
122 field_id_to_column_index: &BTreeMap<u32, u32>,
123) -> Result<ReaderProjection> {
124 Ok(match version {
125 ConcreteFileVersion::V1 => v1::projection_from_field_ids(schema, field_id_to_column_index),
126 ConcreteFileVersion::V2_0 => {
127 v2_0::projection_from_field_ids(schema, field_id_to_column_index)
128 }
129 ConcreteFileVersion::V2_1 => {
130 v2_1::projection_from_field_ids(schema, field_id_to_column_index)
131 }
132 ConcreteFileVersion::V2_2 => {
133 v2_2::projection_from_field_ids(schema, field_id_to_column_index)
134 }
135 ConcreteFileVersion::V2_3 => {
136 v2_3::projection_from_field_ids(schema, field_id_to_column_index)
137 }
138 })
139}
140
141pub fn reader_projection_from_whole_schema(
142 schema: &Schema,
143 version: ConcreteFileVersion,
144) -> ReaderProjection {
145 match version {
146 ConcreteFileVersion::V1 => v1::projection_from_whole_schema(schema),
147 ConcreteFileVersion::V2_0 => v2_0::projection_from_whole_schema(schema),
148 ConcreteFileVersion::V2_1 => v2_1::projection_from_whole_schema(schema),
149 ConcreteFileVersion::V2_2 => v2_2::projection_from_whole_schema(schema),
150 ConcreteFileVersion::V2_3 => v2_3::projection_from_whole_schema(schema),
151 }
152}
153
154pub fn reader_projection_from_column_names(
155 version: ConcreteFileVersion,
156 schema: &Schema,
157 column_names: &[&str],
158) -> Result<ReaderProjection> {
159 match version {
160 ConcreteFileVersion::V1 => v1::projection_from_column_names(schema, column_names),
161 ConcreteFileVersion::V2_0 => v2_0::projection_from_column_names(schema, column_names),
162 ConcreteFileVersion::V2_1 => v2_1::projection_from_column_names(schema, column_names),
163 ConcreteFileVersion::V2_2 => v2_2::projection_from_column_names(schema, column_names),
164 ConcreteFileVersion::V2_3 => v2_3::projection_from_column_names(schema, column_names),
165 }
166}
167
168pub fn physical_column_count(
170 version: ConcreteFileVersion,
171 field: &lance_core::datatypes::Field,
172) -> usize {
173 match version {
174 ConcreteFileVersion::V1 => v1::physical_column_count(field),
175 ConcreteFileVersion::V2_0 => v2_0::physical_column_count(field),
176 ConcreteFileVersion::V2_1 => v2_1::physical_column_count(field),
177 ConcreteFileVersion::V2_2 => v2_2::physical_column_count(field),
178 ConcreteFileVersion::V2_3 => v2_3::physical_column_count(field),
179 }
180}
181
182pub fn data_file_columns(version: ConcreteFileVersion, schema: &Schema) -> (Vec<i32>, Vec<i32>) {
184 match version {
185 ConcreteFileVersion::V1 => v1::data_file_columns(schema),
186 ConcreteFileVersion::V2_0 => v2_0::data_file_columns(schema),
187 ConcreteFileVersion::V2_1 => v2_1::data_file_columns(schema),
188 ConcreteFileVersion::V2_2 => v2_2::data_file_columns(schema),
189 ConcreteFileVersion::V2_3 => v2_3::data_file_columns(schema),
190 }
191}
192
193pub async fn copy_external_metadata_column<Copy, CopyFuture>(
199 version: ConcreteFileVersion,
200 schema: &Schema,
201 column_index: usize,
202 has_existing_pages: bool,
203 copy: Copy,
204) -> Result<()>
205where
206 Copy: FnOnce() -> CopyFuture + Send,
207 CopyFuture: Future<Output = Result<()>> + Send,
208{
209 match version {
210 ConcreteFileVersion::V1 => Err(Error::not_supported(
211 "binary-copy metadata operations are not supported for Lance v1".to_string(),
212 )),
213 ConcreteFileVersion::V2_0 => {
214 if v2_0::should_copy_external_metadata_column(schema, column_index, has_existing_pages)
215 {
216 copy().await
217 } else {
218 Ok(())
219 }
220 }
221 ConcreteFileVersion::V2_1 | ConcreteFileVersion::V2_2 | ConcreteFileVersion::V2_3 => {
222 copy().await
223 }
224 }
225}
226
227pub fn finalize_external_metadata_column(
229 version: ConcreteFileVersion,
230 schema: &Schema,
231 column_index: usize,
232 pages: &mut Vec<PageInfo>,
233 num_rows: u64,
234) -> Result<()> {
235 match version {
236 ConcreteFileVersion::V1 => Err(Error::not_supported(
237 "binary-copy metadata operations are not supported for Lance v1".to_string(),
238 )),
239 ConcreteFileVersion::V2_0 => {
240 v2_0::finalize_external_metadata_column(schema, column_index, pages, num_rows);
241 Ok(())
242 }
243 ConcreteFileVersion::V2_1 | ConcreteFileVersion::V2_2 | ConcreteFileVersion::V2_3 => Ok(()),
244 }
245}
246
247pub async fn open_projected_reader<OpenIndexed, IndexedFuture, OpenFull, FullFuture>(
254 version: ConcreteFileVersion,
255 projection: &ReaderProjection,
256 prefer_indexed: bool,
257 open_indexed: OpenIndexed,
258 open_full: OpenFull,
259) -> Result<ProjectedFileReader>
260where
261 OpenIndexed: FnOnce() -> IndexedFuture + Send,
262 IndexedFuture: Future<Output = Result<Option<ProjectedFileReader>>> + Send,
263 OpenFull: FnOnce() -> FullFuture + Send,
264 FullFuture: Future<Output = Result<ProjectedFileReader>> + Send,
265{
266 match version {
267 ConcreteFileVersion::V1 => Err(Error::not_supported(
268 "projected current-format readers cannot open Lance v1 files".to_string(),
269 )),
270 ConcreteFileVersion::V2_0 => open_full().await,
271 ConcreteFileVersion::V2_1 | ConcreteFileVersion::V2_2 | ConcreteFileVersion::V2_3 => {
272 if prefer_indexed
273 && FileMetadataProvider::projection_matches_indexed_metadata(projection)
274 && let Some(reader) = open_indexed().await?
275 {
276 return Ok(reader);
277 }
278 open_full().await
279 }
280 }
281}
282
283pub async fn open_self_described_reader(
288 scheduler: FileScheduler,
289 decoder_plugins: Arc<DecoderPlugins>,
290 cache: &LanceCache,
291 options: FileReaderOptions,
292) -> Result<OpenedFileReader> {
293 FileReader::try_open_for_dispatch(scheduler, None, decoder_plugins, cache, options).await
294}
295
296pub fn create_writer(
301 version: ConcreteFileVersion,
302 object_writer: Box<dyn Writer>,
303 schema: Schema,
304 options: FileWriterOptions,
305) -> Result<FileWriter> {
306 match version {
307 ConcreteFileVersion::V1 => Err(Error::not_supported(
308 "Lance v1 files must be created with versions::v1::writer::FileWriter".to_string(),
309 )),
310 ConcreteFileVersion::V2_0 => {
311 v2_0::create_writer(object_writer, schema, options).map(Into::into)
312 }
313 ConcreteFileVersion::V2_1 => {
314 v2_1::create_writer(object_writer, schema, options).map(Into::into)
315 }
316 ConcreteFileVersion::V2_2 => {
317 v2_2::create_writer(object_writer, schema, options).map(Into::into)
318 }
319 ConcreteFileVersion::V2_3 => {
320 v2_3::create_writer(object_writer, schema, options).map(Into::into)
321 }
322 }
323}
324
325pub fn create_lazy_writer(
327 version: ConcreteFileVersion,
328 object_writer: Box<dyn Writer>,
329 options: FileWriterOptions,
330) -> Result<FileWriter> {
331 match version {
332 ConcreteFileVersion::V1 => Err(Error::not_supported(
333 "legacy v1 files require an explicit schema and manifest provider".to_string(),
334 )),
335 ConcreteFileVersion::V2_0 => Ok(v2_0::create_lazy_writer(object_writer, options).into()),
336 ConcreteFileVersion::V2_1 => Ok(v2_1::create_lazy_writer(object_writer, options).into()),
337 ConcreteFileVersion::V2_2 => Ok(v2_2::create_lazy_writer(object_writer, options).into()),
338 ConcreteFileVersion::V2_3 => Ok(v2_3::create_lazy_writer(object_writer, options).into()),
339 }
340}
341
342pub fn encode_self_described_batch(
344 version: ConcreteFileVersion,
345 batch: &EncodedBatch,
346) -> Result<Bytes> {
347 match version {
348 ConcreteFileVersion::V1 => Err(Error::not_supported(
349 "Lance v1 does not support self-described current-format batches".to_string(),
350 )),
351 ConcreteFileVersion::V2_0 => v2_0::encode_self_described_batch(batch),
352 ConcreteFileVersion::V2_1 => v2_1::encode_self_described_batch(batch),
353 ConcreteFileVersion::V2_2 => v2_2::encode_self_described_batch(batch),
354 ConcreteFileVersion::V2_3 => v2_3::encode_self_described_batch(batch),
355 }
356}
357
358pub fn encode_mini_batch(version: ConcreteFileVersion, batch: &EncodedBatch) -> Result<Bytes> {
360 match version {
361 ConcreteFileVersion::V1 => Err(Error::not_supported(
362 "Lance v1 does not support mini-lance current-format batches".to_string(),
363 )),
364 ConcreteFileVersion::V2_0 => v2_0::encode_mini_batch(batch),
365 ConcreteFileVersion::V2_1 => v2_1::encode_mini_batch(batch),
366 ConcreteFileVersion::V2_2 => v2_2::encode_mini_batch(batch),
367 ConcreteFileVersion::V2_3 => v2_3::encode_mini_batch(batch),
368 }
369}
370
371#[cfg(test)]
372mod tests {
373 use std::sync::Arc;
374
375 use super::*;
376
377 #[test]
378 fn v1_rejects_current_format_embedded_batches() {
379 let batch = EncodedBatch {
380 data: Bytes::new(),
381 page_table: Vec::new(),
382 schema: Arc::new(Schema::default()),
383 top_level_columns: Vec::new(),
384 num_rows: 0,
385 };
386
387 assert!(matches!(
388 encode_self_described_batch(ConcreteFileVersion::V1, &batch),
389 Err(Error::NotSupported { .. })
390 ));
391 assert!(matches!(
392 encode_mini_batch(ConcreteFileVersion::V1, &batch),
393 Err(Error::NotSupported { .. })
394 ));
395 }
396}