use std::path::Path;
use arrow::array::{
Array, FixedSizeListArray, Float64Array, Int64Array, PrimitiveArray, UInt32Array,
};
use arrow::datatypes::ArrowPrimitiveType;
use arrow::datatypes::{DataType, Schema as ArrowSchema};
use arrow::record_batch::RecordBatch;
use prost::Message;
use crate::StorageError;
use crate::StorageResult;
use crate::lancefmt::pb::encodings::ColumnEncoding;
use crate::lancefmt::pb::encodings::column_encoding as ce;
use crate::lancefmt::pb::encodings21::compressive_encoding::Compression;
use crate::lancefmt::pb::encodings21::{
CompressiveEncoding, FixedSizeList as FslCompressive, Flat, MiniBlockLayout, PageLayout,
page_layout,
};
use crate::lancefmt::pb::filev2 as lfv2;
use crate::lancefmt::pb::filev2::{ColumnMetadata, column_metadata};
use crate::lancefmt::pb::table::Transaction;
use crate::lancefmt::pb::table::manifest::{DataStorageFormat, WriterVersion};
use crate::lancefmt::pb::table::transaction::{Operation, Overwrite};
use crate::lancefmt::pb::table::{DataFile, DataFragment, Manifest};
use super::schema::{SchemaMeta, to_lance_schema};
const FILE_MAJOR: u16 = 2;
const FILE_MINOR: u16 = 1;
const DATA_FORMAT_VERSION: &str = "2.1";
const CHUNK_ALIGNMENT: usize = 8;
const BUFFER_ALIGNMENT: usize = 64;
const FILL_BYTE: u8 = 0xFE;
const CHUNK_DATA_BUDGET: usize = 16 * 1024;
fn pad_to(buf: &mut Vec<u8>, alignment: usize) {
while !buf.len().is_multiple_of(alignment) {
buf.push(FILL_BYTE);
}
}
fn direct_encoding(type_url: &'static str, value: Vec<u8>) -> lfv2::Encoding {
let any = prost_types::Any {
type_url: type_url.to_string(),
value,
};
lfv2::Encoding {
location: Some(lfv2::encoding::Location::Direct(lfv2::DirectEncoding {
encoding: any.encode_to_vec(),
})),
}
}
struct Chunk {
metadata: u16,
bytes: Vec<u8>,
}
struct EncodedPage {
metadata_buffer: Vec<u8>,
data_buffer: Vec<u8>,
rows: u64,
num_items: u64,
value_compression: Compression,
}
fn log2_usize(v: usize) -> StorageResult<u16> {
u16::try_from(v.ilog2()).map_err(|_| StorageError::Overflow("log2 chunk".into()))
}
fn build_chunk(buffer: &[u8], ln: u16) -> StorageResult<Chunk> {
let mut bytes = Vec::with_capacity(buffer.len() + 16);
bytes.extend_from_slice(&0u16.to_le_bytes());
let size = u16::try_from(buffer.len()).map_err(|_| {
StorageError::Overflow(format!("chunk buffer size {} exceeds u16", buffer.len()))
})?;
bytes.extend_from_slice(&size.to_le_bytes());
pad_to(&mut bytes, CHUNK_ALIGNMENT);
bytes.extend_from_slice(buffer);
pad_to(&mut bytes, CHUNK_ALIGNMENT);
let chunk_bytes = bytes.len();
let divided = chunk_bytes / CHUNK_ALIGNMENT;
let metadata = ((u16::try_from(divided - 1).map_err(|_| {
StorageError::Overflow(format!("chunk bytes {chunk_bytes} exceed u16 metadata"))
})?) << 4)
| ln;
Ok(Chunk { metadata, bytes })
}
fn chunk_bytes_le(
values: &PrimitiveArray<arrow::datatypes::Float64Type>,
values_per_chunk: usize,
) -> StorageResult<Vec<Chunk>> {
chunk_bytes_le_impl(values, values_per_chunk, 8, |out, v| {
out.extend_from_slice(&v.to_le_bytes())
})
}
fn chunk_bytes_le_u32(
values: &PrimitiveArray<arrow::datatypes::UInt32Type>,
values_per_chunk: usize,
) -> StorageResult<Vec<Chunk>> {
chunk_bytes_le_impl(values, values_per_chunk, 4, |out, v| {
out.extend_from_slice(&v.to_le_bytes())
})
}
fn chunk_bytes_le_i64(
values: &PrimitiveArray<arrow::datatypes::Int64Type>,
values_per_chunk: usize,
) -> StorageResult<Vec<Chunk>> {
chunk_bytes_le_impl(values, values_per_chunk, 8, |out, v| {
out.extend_from_slice(&v.to_le_bytes())
})
}
fn chunk_bytes_le_impl<T: ArrowPrimitiveType>(
values: &PrimitiveArray<T>,
values_per_chunk: usize,
bytes_per_value: usize,
write_one: impl Fn(&mut Vec<u8>, T::Native),
) -> StorageResult<Vec<Chunk>> {
let mut chunks = Vec::new();
let mut start = 0usize;
let total = values.len();
while start < total {
let n = values_per_chunk.min(total - start);
let mut buffer = Vec::with_capacity(n * bytes_per_value);
for i in start..start + n {
write_one(&mut buffer, values.value(i));
}
let ln = if n == values_per_chunk {
log2_usize(n)?
} else {
0
};
chunks.push(build_chunk(&buffer, ln)?);
start += n;
}
Ok(chunks)
}
fn encode_column(batch: &RecordBatch, col: usize) -> StorageResult<EncodedPage> {
let column = batch.column(col);
let rows = column.len() as u64;
let (chunks, value_compression, num_items) = match column.data_type() {
DataType::Float64 => {
let arr = column
.as_any()
.downcast_ref::<Float64Array>()
.ok_or_else(|| StorageError::Invalid("float64 downcast failed".into()))?;
(
chunk_bytes_le(arr, 512)?,
Compression::Flat(Flat {
bits_per_value: 64,
data: None,
}),
rows,
)
}
DataType::UInt32 => {
let arr = column
.as_any()
.downcast_ref::<UInt32Array>()
.ok_or_else(|| StorageError::Invalid("uint32 downcast failed".into()))?;
(
chunk_bytes_le_u32(arr, 1024)?,
Compression::Flat(Flat {
bits_per_value: 32,
data: None,
}),
rows,
)
}
DataType::Int64 => {
let arr = column
.as_any()
.downcast_ref::<Int64Array>()
.ok_or_else(|| StorageError::Invalid("int64 downcast failed".into()))?;
(
chunk_bytes_le_i64(arr, 1024)?,
Compression::Flat(Flat {
bits_per_value: 64,
data: None,
}),
rows,
)
}
DataType::FixedSizeList(_child, dim) => {
let dim: i32 = *dim;
let list = column
.as_any()
.downcast_ref::<FixedSizeListArray>()
.ok_or_else(|| StorageError::Invalid("fsl downcast failed".into()))?;
if list.null_count() != 0 {
return Err(StorageError::UnsupportedFormat(
"lancefmt writer does not support nulls in FixedSizeList columns".into(),
));
}
let dim = dim as usize;
let values = list
.values()
.as_any()
.downcast_ref::<Float64Array>()
.ok_or_else(|| StorageError::Invalid("fsl child downcast failed".into()))?;
let bytes_per_row = dim * std::mem::size_of::<f64>();
let rows_per_chunk = (CHUNK_DATA_BUDGET / bytes_per_row).next_power_of_two();
let mut chunks = Vec::new();
let mut start = 0usize;
while start < values.len() {
let n_items = (rows_per_chunk * dim).min(values.len() - start);
let mut buffer = Vec::with_capacity(n_items * 8);
for i in start..start + n_items {
buffer.extend_from_slice(&values.value(i).to_le_bytes());
}
let ln = if n_items == rows_per_chunk * dim {
log2_usize(rows_per_chunk)?
} else {
0
};
chunks.push(build_chunk(&buffer, ln)?);
start += n_items;
}
let compression = Compression::FixedSizeList(Box::new(FslCompressive {
items_per_value: dim as u64,
has_validity: false,
values: Some(Box::new(CompressiveEncoding {
compression: Some(Compression::Flat(Flat {
bits_per_value: 64,
data: None,
})),
})),
}));
(chunks, compression, rows)
}
other => {
return Err(StorageError::UnsupportedFormat(format!(
"lancefmt writer: unsupported column type {other:?}"
)));
}
};
let mut metadata_buffer = Vec::with_capacity(chunks.len() * 2);
let mut data_buffer = Vec::new();
for chunk in &chunks {
metadata_buffer.extend_from_slice(&chunk.metadata.to_le_bytes());
data_buffer.extend_from_slice(&chunk.bytes);
}
Ok(EncodedPage {
metadata_buffer,
data_buffer,
rows,
num_items,
value_compression,
})
}
fn page_layout(page: &EncodedPage) -> lfv2::Encoding {
let layout = PageLayout {
layout: Some(page_layout::Layout::MiniBlockLayout(MiniBlockLayout {
rep_compression: None,
def_compression: None,
value_compression: Some(CompressiveEncoding {
compression: Some(page.value_compression.clone()),
}),
dictionary: None,
num_dictionary_items: 0,
layers: vec![1], num_buffers: 1,
repetition_index_depth: 0,
num_items: page.num_items,
has_large_chunk: false,
})),
};
direct_encoding("/lance.encodings21.PageLayout", layout.encode_to_vec())
}
pub fn write_dataset(batch: &RecordBatch, dir: &Path) -> StorageResult<()> {
if batch.num_rows() == 0 {
return Err(StorageError::Invalid(
"lancefmt writer: empty batches are not supported".into(),
));
}
let arrow_schema = batch.schema().as_ref().clone();
let (fields, schema_meta) = to_lance_schema(&arrow_schema)?;
std::fs::create_dir_all(dir)
.map_err(|e| StorageError::Io(format!("create dataset dir {:?}: {e}", dir)))?;
let data_dir = dir.join("data");
let versions_dir = dir.join("_versions");
let txn_dir = dir.join("_transactions");
for d in [&data_dir, &versions_dir, &txn_dir] {
std::fs::create_dir_all(d).map_err(|e| StorageError::Io(format!("create {:?}: {e}", d)))?;
}
let uuid = uuid::Uuid::new_v4();
let data_file_name = format!("{:024b}{}.lance", 0, uuid.simple());
let mut file_bytes: Vec<u8> = Vec::new();
let mut column_metadatas: Vec<(u64, Vec<u8>)> = Vec::new();
for col in 0..batch.num_columns() {
let page = encode_column(batch, col)?;
let page_encoding = page_layout(&page);
let mut buffers: Vec<&[u8]> = vec![&page.metadata_buffer, &page.data_buffer];
let mut page_buffer_offsets = Vec::with_capacity(buffers.len());
let mut page_buffer_sizes = Vec::with_capacity(buffers.len());
for b in buffers.drain(..) {
pad_to(&mut file_bytes, BUFFER_ALIGNMENT);
page_buffer_offsets.push(file_bytes.len() as u64);
page_buffer_sizes.push(b.len() as u64);
file_bytes.extend_from_slice(b);
}
let column_metadata = ColumnMetadata {
encoding: Some(direct_encoding(
"/lance.encodings.ColumnEncoding",
ColumnEncoding {
column_encoding: Some(ce::ColumnEncoding::Values(())),
}
.encode_to_vec(),
)),
pages: vec![column_metadata::Page {
buffer_offsets: page_buffer_offsets,
buffer_sizes: page_buffer_sizes,
length: page.rows,
encoding: Some(page_encoding),
priority: 0,
}],
buffer_offsets: vec![],
buffer_sizes: vec![],
};
column_metadatas.push((0, column_metadata.encode_to_vec()));
}
let schema_bytes =
super::schema::encode_schema_global(&fields, &schema_meta, batch.num_rows() as u64);
pad_to(&mut file_bytes, BUFFER_ALIGNMENT);
let global_pos = file_bytes.len() as u64;
file_bytes.extend_from_slice(&schema_bytes);
let col_meta0_pos = file_bytes.len() as u64;
let mut cmo_entries = Vec::with_capacity(column_metadatas.len());
for (_, meta) in &column_metadatas {
let pos = file_bytes.len() as u64;
file_bytes.extend_from_slice(meta);
cmo_entries.push((pos, meta.len() as u64));
}
pad_to(&mut file_bytes, CHUNK_ALIGNMENT);
let cmo_pos = file_bytes.len() as u64;
for (pos, size) in &cmo_entries {
file_bytes.extend_from_slice(&pos.to_le_bytes());
file_bytes.extend_from_slice(&size.to_le_bytes());
}
let gbo_pos = file_bytes.len() as u64;
file_bytes.extend_from_slice(&global_pos.to_le_bytes());
file_bytes.extend_from_slice(&(schema_bytes.len() as u64).to_le_bytes());
let global_entries_len = 1u32;
pad_to(&mut file_bytes, CHUNK_ALIGNMENT);
file_bytes.extend_from_slice(&col_meta0_pos.to_le_bytes());
file_bytes.extend_from_slice(&cmo_pos.to_le_bytes());
file_bytes.extend_from_slice(&gbo_pos.to_le_bytes());
file_bytes.extend_from_slice(&global_entries_len.to_le_bytes());
file_bytes.extend_from_slice(&(column_metadatas.len() as u32).to_le_bytes());
file_bytes.extend_from_slice(&FILE_MAJOR.to_le_bytes());
file_bytes.extend_from_slice(&FILE_MINOR.to_le_bytes());
file_bytes.extend_from_slice(b"LANC");
let data_path = data_dir.join(&data_file_name);
std::fs::write(&data_path, &file_bytes)
.map_err(|e| StorageError::Io(format!("write {data_path:?}: {e}")))?;
let prev_version = latest_manifest_version(&versions_dir)?;
let next_version = prev_version + 1;
let fragment_id = prev_version;
let data_file = DataFile {
path: data_file_name.clone(),
fields: fields.iter().map(|f| f.id).collect(),
column_indices: (0..fields.len() as i32).collect(),
file_major_version: FILE_MAJOR as u32,
file_minor_version: FILE_MINOR as u32,
file_size_bytes: file_bytes.len() as u64,
base_id: None,
};
let fragment = DataFragment {
id: fragment_id,
files: vec![data_file],
overlays: vec![],
deletion_file: None,
row_id_sequence: None,
last_updated_at_version_sequence: None,
created_at_version_sequence: None,
physical_rows: batch.num_rows() as u64,
};
let manifest = Manifest {
fields,
schema_metadata: schema_meta.clone(),
fragments: vec![fragment],
version: next_version,
version_aux_data: 0,
writer_version: Some(WriterVersion {
library: "genegraph-storage".to_string(),
version: env!("CARGO_PKG_VERSION").to_string(),
prerelease: None,
build_metadata: None,
}),
index_section: None,
timestamp: Some(prost_types::Timestamp {
seconds: chrono::Utc::now().timestamp(),
nanos: chrono::Utc::now().timestamp_subsec_nanos() as i32,
}),
tag: String::new(),
reader_feature_flags: 0,
writer_feature_flags: 0,
max_fragment_id: Some(fragment_id as u32),
transaction_file: format!("{prev_version}-{uuid}.txn"),
next_row_id: 0,
data_format: Some(DataStorageFormat {
file_format: "lance".to_string(),
version: DATA_FORMAT_VERSION.to_string(),
}),
config: Default::default(),
table_metadata: Default::default(),
base_paths: vec![],
branch: None,
transaction_section: None,
};
let transaction = Transaction {
read_version: prev_version,
uuid: uuid.to_string(),
tag: String::new(),
transaction_properties: Default::default(),
operation: Some(Operation::Overwrite(Overwrite {
fragments: manifest.fragments.clone(),
schema: manifest.fields.clone(),
schema_metadata: schema_meta.clone(),
config_upsert_values: Default::default(),
initial_bases: vec![],
})),
};
let txn_bytes = transaction.encode_to_vec();
std::fs::write(
txn_dir.join(format!("{prev_version}-{uuid}.txn")),
&txn_bytes,
)
.map_err(|e| StorageError::Io(format!("write txn: {e}")))?;
let manifest_bytes = manifest.encode_to_vec();
let mut out: Vec<u8> = Vec::new();
out.extend_from_slice(&(txn_bytes.len() as u32).to_le_bytes());
out.extend_from_slice(&txn_bytes);
let manifest_pos = out.len() as u64;
out.extend_from_slice(&(manifest_bytes.len() as u32).to_le_bytes());
out.extend_from_slice(&manifest_bytes);
out.extend_from_slice(&manifest_pos.to_le_bytes());
out.extend_from_slice(&FILE_MINOR.to_le_bytes());
out.extend_from_slice(&FILE_MAJOR.to_le_bytes());
out.extend_from_slice(b"LANC");
std::fs::write(versions_dir.join(format!("{next_version}.manifest")), &out)
.map_err(|e| StorageError::Io(format!("write manifest: {e}")))?;
std::fs::write(
versions_dir.join("latest_version_hint.json"),
format!("{{\"version\":{next_version}}}"),
)
.map_err(|e| StorageError::Io(format!("write hint: {e}")))?;
Ok(())
}
fn latest_manifest_version(versions_dir: &Path) -> StorageResult<u64> {
let entries = std::fs::read_dir(versions_dir)
.map_err(|e| StorageError::Io(format!("read {versions_dir:?}: {e}")))?;
let mut best = 0u64;
for entry in entries {
let entry = entry.map_err(|e| StorageError::Io(e.to_string()))?;
let name = entry.file_name();
let name_str = name.to_string_lossy().to_string();
let Some(stem) = name_str.strip_suffix(".manifest") else {
continue;
};
let v: u64 = stem.parse().map_err(|_| {
StorageError::Invalid(format!("unparseable manifest name {name_str:?}"))
})?;
best = best.max(v);
}
Ok(best)
}
#[allow(dead_code)]
fn unused(_: &ArrowSchema, _: &SchemaMeta) {}