use glaredb_core::expr;
use glaredb_core::functions::table::TableFunctionInput;
use glaredb_core::functions::table::scan::ScanContext;
use glaredb_core::optimizer::expr_rewrite::ExpressionRewriteRule;
use glaredb_core::optimizer::expr_rewrite::const_fold::ConstFold;
use glaredb_core::runtime::filesystem::file_provider::{MultiFileData, MultiFileProvider};
use glaredb_core::runtime::filesystem::{FileSystemWithState, OpenFlags};
use glaredb_error::{DbError, Result, ResultExt};
use crate::table::spec;
#[derive(Debug)]
pub struct TableState {
pub root: String,
pub fs: FileSystemWithState,
pub metadata: spec::Metadata,
pub manifest_list: Option<spec::ManifestList>,
pub manifests: Option<Vec<spec::Manifest>>,
}
impl TableState {
pub async fn open_root_with_inputs(
scan_context: ScanContext<'_>,
mut input: TableFunctionInput,
) -> Result<Self> {
let root = ConstFold::rewrite(input.positional[0].clone())?
.try_into_scalar()?
.try_into_string()?;
let version = match input.named.get("version") {
Some(version) => Some(
ConstFold::rewrite(version.clone())?
.try_into_scalar()?
.try_into_string()?,
),
None => None,
};
let glob = format_glob_for_metadata(&root, version.as_deref());
input.positional[0] = expr::lit(glob).into();
let (mut mf_prov, fs) =
MultiFileProvider::try_new_from_inputs(scan_context, &input).await?;
let mut mf_data = MultiFileData::empty();
mf_prov.expand_all(&mut mf_data).await?;
let metadata_path = match mf_data.expanded().iter().max() {
Some(max) => max,
None => {
return Err(DbError::new(
"Could not find any metadata files in the table root",
));
}
};
let mut file = fs.open(OpenFlags::READ, metadata_path).await?;
let mut read_buf = vec![0; file.call_size() as usize];
file.call_read_exact(&mut read_buf).await?;
let metadata: spec::Metadata = serde_json::from_slice(&read_buf)
.context_fn(|| format!("Failed to read metadata from {metadata_path}"))?;
Ok(TableState {
root,
fs,
metadata,
manifest_list: None,
manifests: None,
})
}
pub async fn load_manifest_list(&mut self) -> Result<&spec::ManifestList> {
let curr_snap = self.current_snapshot()?;
let manifest_list_path = curr_snap
.manifest_list
.as_ref()
.ok_or_else(|| DbError::new("Missing manifest list location for current snapshot"))?;
let manifest_list_rel = relative_path(&self.metadata.location, manifest_list_path);
let path = format!("{}/{}", self.root, manifest_list_rel);
let mut file = self
.fs
.open(OpenFlags::READ, &path)
.await
.context_fn(|| format!("Failed to open manifest list at '{path}'"))?;
let mut read_buf = vec![0; file.call_size() as usize];
file.call_read_exact(&mut read_buf).await?;
let list = spec::ManifestList::from_raw_avro(&read_buf)?;
self.manifest_list = Some(list);
Ok(self.manifest_list.as_ref().unwrap())
}
pub async fn load_manifests(&mut self) -> Result<&[spec::Manifest]> {
let list = match self.manifest_list.as_ref() {
Some(list) => list,
None => {
return Err(DbError::new(
"Manifest list must be loaded before reading manifests",
));
}
};
let mut manifests = Vec::with_capacity(list.entries.len());
let mut read_buf = Vec::new();
for ent in &list.entries {
let manifest_rel = relative_path(&self.metadata.location, &ent.manifest_path);
let path = format!("{}/{}", self.root, manifest_rel);
let mut file = self.fs.open(OpenFlags::READ, &path).await?;
read_buf.resize(file.call_size() as usize, 0);
file.call_read_exact(&mut read_buf).await?;
let manifest = spec::Manifest::from_raw_avro(&read_buf)?;
manifests.push(manifest);
}
self.manifests = Some(manifests);
Ok(self.manifests.as_ref().unwrap())
}
fn current_snapshot(&self) -> Result<&spec::Snapshot> {
let curr_id = self
.metadata
.current_snapshot_id
.ok_or_else(|| DbError::new("Missing current snapshot id"))?;
let curr_snap = self
.metadata
.snapshots
.iter()
.find(|s| s.snapshot_id == curr_id)
.ok_or_else(|| DbError::new(format!("Missing snapshot for id {curr_id}")))?;
Ok(curr_snap)
}
}
fn format_glob_for_metadata(root: &str, version: Option<&str>) -> String {
let root = root.trim_end_matches("/");
match version {
Some(version) => format!("{root}/metadata/{version}*.metadata.json"),
None => {
format!("{root}/metadata/*.metadata.json")
}
}
}
fn relative_path<'a>(root: &str, path: &'a str) -> &'a str {
let metadata_location = root.trim_start_matches("./");
path.trim_start_matches(metadata_location).trim_matches('/')
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn format_glob_cases() {
struct TestCase {
root: &'static str,
version: Option<&'static str>,
expected: &'static str,
}
let cases = [
TestCase {
root: "wh/default.db",
version: None,
expected: "wh/default.db/metadata/*.metadata.json",
},
TestCase {
root: "wh/default.db",
version: Some("00001"),
expected: "wh/default.db/metadata/00001*.metadata.json",
},
TestCase {
root: "wh/default.db/", version: Some("00001"),
expected: "wh/default.db/metadata/00001*.metadata.json",
},
];
for case in cases {
let got = format_glob_for_metadata(case.root, case.version);
assert_eq!(case.expected, got);
}
}
#[test]
fn relative_path_cases() {
struct TestCase {
root: &'static str,
input: &'static str,
expected: &'static str,
}
let test_cases = vec![
TestCase {
root: "out/iceberg_table",
input: "out/iceberg_table/metadata/snap-4160073268445560424-1-095d0ad9-385f-406f-b29c-966a6e222e58.avro",
expected: "metadata/snap-4160073268445560424-1-095d0ad9-385f-406f-b29c-966a6e222e58.avro",
},
TestCase {
root: "./out/iceberg_table",
input: "out/iceberg_table/metadata/snap-4160073268445560424-1-095d0ad9-385f-406f-b29c-966a6e222e58.avro",
expected: "metadata/snap-4160073268445560424-1-095d0ad9-385f-406f-b29c-966a6e222e58.avro",
},
TestCase {
root: "/Users/sean/Code/github.com/glaredb/glaredb/testdata/iceberg/tables/lineitem_versioned",
input: "/Users/sean/Code/github.com/glaredb/glaredb/testdata/iceberg/tables/lineitem_versioned/metadata/snap-2591356646088336681-1-481f5867-e369-4c1c-a9ba-6c9e04030958.avro",
expected: "metadata/snap-2591356646088336681-1-481f5867-e369-4c1c-a9ba-6c9e04030958.avro",
},
TestCase {
root: "s3://testdata/iceberg/tables/lineitem_versioned",
input: "s3://testdata/iceberg/tables/lineitem_versioned/metadata/snap-2591356646088336681-1-481f5867-e369-4c1c-a9ba-6c9e04030958.avro",
expected: "metadata/snap-2591356646088336681-1-481f5867-e369-4c1c-a9ba-6c9e04030958.avro",
},
];
for case in test_cases {
let out = relative_path(case.root, case.input);
assert_eq!(case.expected, out, "root: {}", case.root,);
}
}
}