diff --git a/crates/paimon/src/spec/core_options.rs b/crates/paimon/src/spec/core_options.rs index eec1ca57..fce04169 100644 --- a/crates/paimon/src/spec/core_options.rs +++ b/crates/paimon/src/spec/core_options.rs @@ -123,6 +123,7 @@ pub(crate) const BLOB_DESCRIPTOR_FIELD_OPTION: &str = "blob-descriptor-field"; pub(crate) const BLOB_VIEW_FIELD_OPTION: &str = "blob-view-field"; pub const BLOB_VIEW_RESOLVE_ENABLED_OPTION: &str = "blob-view.resolve.enabled"; const PK_VECTOR_INDEX_COLUMNS_OPTION: &str = "pk-vector.index.columns"; +const PK_FULL_TEXT_INDEX_COLUMNS_OPTION: &str = "pk-full-text.index.columns"; /// Merge engine for primary-key tables. /// @@ -1198,6 +1199,23 @@ impl<'a> CoreOptions<'a> { crate::vindex::pkvector::metric::VectorSearchMetric::parse(&raw)?; Ok(raw) } + + /// True when the PK full-text index column option key is present (any value). + pub fn primary_key_full_text_index_enabled(&self) -> bool { + self.options.contains_key(PK_FULL_TEXT_INDEX_COLUMNS_OPTION) + } + + /// Configured PK full-text index columns: split on ',' and trim each token, + /// mirroring Java `split(",",-1).map(trim)`. Blank tokens are PRESERVED (do NOT + /// filter them) so parsing matches Java exactly; `[]` only when the key is + /// absent. Never returns a `Result`, never errors (do not copy the fail-loud + /// shape of `primary_key_vector_index_columns`). + pub fn primary_key_full_text_index_columns(&self) -> Vec { + match self.options.get(PK_FULL_TEXT_INDEX_COLUMNS_OPTION) { + None => Vec::new(), + Some(raw) => raw.split(',').map(|c| c.trim().to_string()).collect(), + } + } } /// Parse a memory size string to bytes using binary (1024-based) semantics. @@ -2147,4 +2165,56 @@ mod tests { .primary_key_vector_index_type("e") .is_err()); } + + #[test] + fn test_pk_full_text_index_absent_is_disabled_and_empty() { + let opts = HashMap::new(); + let co = CoreOptions::new(&opts); + assert!(!co.primary_key_full_text_index_enabled()); + assert_eq!( + co.primary_key_full_text_index_columns(), + Vec::::new() + ); + } + + #[test] + fn test_pk_full_text_index_columns_split_and_trim() { + let opts = HashMap::from([( + "pk-full-text.index.columns".to_string(), + "a, b ,c".to_string(), + )]); + let co = CoreOptions::new(&opts); + assert!(co.primary_key_full_text_index_enabled()); + assert_eq!( + co.primary_key_full_text_index_columns(), + vec!["a".to_string(), "b".to_string(), "c".to_string()] + ); + } + + #[test] + fn test_pk_full_text_index_columns_preserve_blank_tokens() { + // Java `split(",",-1).map(trim)` keeps empty tokens, so "a,,b" yields + // three columns with a blank in the middle. + let opts = HashMap::from([("pk-full-text.index.columns".to_string(), "a,,b".to_string())]); + let co = CoreOptions::new(&opts); + assert!(co.primary_key_full_text_index_enabled()); + assert_eq!( + co.primary_key_full_text_index_columns(), + vec!["a".to_string(), "".to_string(), "b".to_string()] + ); + } + + #[test] + fn test_pk_full_text_index_blank_value_is_enabled_but_empty_no_error() { + // key present but only blanks: enabled (key exists), NO error (lenient, + // unlike vector). Blank tokens are PRESERVED to match Java's + // `split(",",-1).map(trim)`, so " , " yields two empty tokens (NOT []). + let opts = HashMap::from([("pk-full-text.index.columns".to_string(), " , ".to_string())]); + let co = CoreOptions::new(&opts); + assert!(co.primary_key_full_text_index_enabled()); + assert_eq!( + co.primary_key_full_text_index_columns(), + vec!["".to_string(), "".to_string()] + ); + } } diff --git a/crates/paimon/src/spec/mod.rs b/crates/paimon/src/spec/mod.rs index 13f9c3d6..2902ee72 100644 --- a/crates/paimon/src/spec/mod.rs +++ b/crates/paimon/src/spec/mod.rs @@ -64,8 +64,8 @@ pub use manifest_file_meta::*; mod index_file_meta; pub use index_file_meta::*; -mod pk_vector_source; -pub use pk_vector_source::*; +mod pk_index_source; +pub use pk_index_source::*; mod index_manifest; pub use index_manifest::{IndexManifest, IndexManifestEntry}; diff --git a/crates/paimon/src/spec/pk_vector_source.rs b/crates/paimon/src/spec/pk_index_source.rs similarity index 77% rename from crates/paimon/src/spec/pk_vector_source.rs rename to crates/paimon/src/spec/pk_index_source.rs index 802cbf7c..a821bc6d 100644 --- a/crates/paimon/src/spec/pk_vector_source.rs +++ b/crates/paimon/src/spec/pk_index_source.rs @@ -15,20 +15,30 @@ // specific language governing permissions and limitations // under the License. -//! Primary-key vector (bucket-local ANN) source metadata. +//! Shared primary-key index source metadata. //! //! Parses the `_SOURCE_META` blob embedded in [`crate::spec::GlobalIndexMeta`], -//! written by Java Paimon's `PkVectorSourceMeta` (apache/paimon#8549). The blob -//! lists the ordered source data files backing an ANN segment; a segment's -//! ordinals concatenate those files in order, so an ordinal maps to a -//! `(data file, physical row position)` pair. +//! written by Java Paimon's `PrimaryKeyIndexSourceMeta` (apache/paimon#8549). The +//! blob lists the ordered source data files backing an index payload; the +//! payload's ordinals concatenate those files in order, so an ordinal maps to a +//! `(data file, physical row position)` pair. Shared by the primary-key vector +//! and full-text read paths. //! //! Two distinct encodings are involved and must not be conflated: Avro extracts //! the opaque `_SOURCE_META` bytes (handled in the Avro decoder), and the bytes //! *inside* use Java `DataOutput` (big-endian ints/longs, modified-UTF-8 via //! `writeUTF`) — parsed here, independent of Avro. -use crate::spec::GlobalIndexMeta; +use crate::spec::{DataFileMeta, GlobalIndexMeta}; + +/// Compacted file source discriminant, matching Java `FileSource.COMPACT`. +const FILE_SOURCE_COMPACT: i32 = 1; + +/// Mirror of `PrimaryKeyIndexSourcePolicy.shouldRead`: only compacted, non-level-0 +/// files back a primary-key index; an absent file source reads as false. +pub fn should_read_pk_index_source(file: &DataFileMeta) -> bool { + matches!(file.file_source, Some(src) if src == FILE_SOURCE_COMPACT) && file.level > 0 +} fn data_invalid(message: impl Into) -> crate::Error { crate::Error::DataInvalid { @@ -37,19 +47,19 @@ fn data_invalid(message: impl Into) -> crate::Error { } } -/// The `_SOURCE_META` frame version written by Java `PkVectorSourceMeta`. +/// The `_SOURCE_META` frame version written by Java `PrimaryKeyIndexSourceMeta`. const SOURCE_META_VERSION: i32 = 1; -/// One source data file captured when a PK-vector ANN segment was built. +/// One source data file captured when a primary-key index payload was built. /// -/// Mirrors Java `org.apache.paimon.index.pkvector.PkVectorSourceFile`. +/// Mirrors Java `org.apache.paimon.index.pk.PrimaryKeyIndexSourceFile`. #[derive(Debug, Clone, PartialEq, Eq)] -pub struct PkVectorSourceFile { +pub struct PrimaryKeyIndexSourceFile { file_name: String, row_count: i64, } -impl PkVectorSourceFile { +impl PrimaryKeyIndexSourceFile { pub fn new(file_name: String, row_count: i64) -> crate::Result { if row_count < 0 { return Err(data_invalid(format!( @@ -71,24 +81,27 @@ impl PkVectorSourceFile { } } -/// Ordered source data files for a primary-key vector index payload. +/// Ordered source data files for a primary-key index payload. /// -/// Mirrors Java `org.apache.paimon.index.pkvector.PkVectorSourceMeta`. +/// Mirrors Java `org.apache.paimon.index.pk.PrimaryKeyIndexSourceMeta`. #[derive(Debug, Clone, PartialEq, Eq)] -pub struct PkVectorSourceMeta { +pub struct PrimaryKeyIndexSourceMeta { data_level: i32, - source_files: Vec, + source_files: Vec, } -impl PkVectorSourceMeta { - pub fn new(data_level: i32, source_files: Vec) -> crate::Result { +impl PrimaryKeyIndexSourceMeta { + pub fn new( + data_level: i32, + source_files: Vec, + ) -> crate::Result { if data_level <= 0 { return Err(data_invalid(format!( "source meta data level must be positive: {data_level}" ))); } if source_files.is_empty() { - return Err(data_invalid("a vector index must reference source files")); + return Err(data_invalid("an index must reference source files")); } Ok(Self { data_level, @@ -100,7 +113,7 @@ impl PkVectorSourceMeta { self.data_level } - pub fn source_files(&self) -> &[PkVectorSourceFile] { + pub fn source_files(&self) -> &[PrimaryKeyIndexSourceFile] { &self.source_files } @@ -110,7 +123,7 @@ impl PkVectorSourceMeta { let bytes = meta .source_meta .as_deref() - .ok_or_else(|| data_invalid("global index meta has no vector source metadata"))?; + .ok_or_else(|| data_invalid("global index meta has no source metadata"))?; Self::deserialize(bytes) } @@ -122,31 +135,31 @@ impl PkVectorSourceMeta { pub fn resolve(&self, ordinal: i64) -> crate::Result<(String, i64)> { if ordinal < 0 { return Err(data_invalid(format!( - "vector ordinal must not be negative: {ordinal}" + "source ordinal must not be negative: {ordinal}" ))); } let mut cumulative: i64 = 0; for file in &self.source_files { let next = cumulative .checked_add(file.row_count) - .ok_or_else(|| data_invalid("vector source row counts overflow i64"))?; + .ok_or_else(|| data_invalid("index source row counts overflow i64"))?; if ordinal < next { return Ok((file.file_name.clone(), ordinal - cumulative)); } cumulative = next; } Err(data_invalid(format!( - "vector ordinal {ordinal} is out of range (total rows {cumulative})" + "source ordinal {ordinal} is out of range (total rows {cumulative})" ))) } - /// Parse a Java `PkVectorSourceMeta`-serialized `_SOURCE_META` blob. + /// Parse a Java `PrimaryKeyIndexSourceMeta`-serialized `_SOURCE_META` blob. pub fn deserialize(bytes: &[u8]) -> crate::Result { let mut cursor = DataInputCursor::new(bytes); let version = cursor.read_i32_be()?; if version != SOURCE_META_VERSION { return Err(data_invalid(format!( - "unsupported vector source version: {version}" + "unsupported index source version: {version}" ))); } let data_level = cursor.read_i32_be()?; @@ -157,18 +170,18 @@ impl PkVectorSourceMeta { } let count = cursor.read_i32_be()?; if count <= 0 { - return Err(data_invalid("a vector index must reference source files")); + return Err(data_invalid("an index must reference source files")); } // NOT Vec::with_capacity(count): count is untrusted and may be huge. let mut source_files = Vec::new(); for _ in 0..count { let file_name = read_java_utf(&mut cursor)?; let row_count = cursor.read_i64_be()?; - source_files.push(PkVectorSourceFile::new(file_name, row_count)?); + source_files.push(PrimaryKeyIndexSourceFile::new(file_name, row_count)?); } if cursor.remaining() != 0 { return Err(data_invalid( - "unexpected trailing bytes in vector source metadata", + "unexpected trailing bytes in index source metadata", )); } Self::new(data_level, source_files) @@ -195,7 +208,7 @@ impl<'a> DataInputCursor<'a> { fn read_exact(&mut self, len: usize) -> crate::Result<&'a [u8]> { if self.remaining() < len { return Err(data_invalid(format!( - "unexpected end of vector source metadata: need {len} bytes, {} remain", + "unexpected end of index source metadata: need {len} bytes, {} remain", self.remaining() ))); } @@ -340,7 +353,7 @@ mod tests { #[test] fn deserialize_single_source_file() { let bytes = frame(1, &[("data-abc.parquet", 100)]); - let meta = PkVectorSourceMeta::deserialize(&bytes).unwrap(); + let meta = PrimaryKeyIndexSourceMeta::deserialize(&bytes).unwrap(); assert_eq!(meta.data_level(), 1); assert_eq!(meta.source_files().len(), 1); assert_eq!(meta.source_files()[0].file_name(), "data-abc.parquet"); @@ -350,7 +363,7 @@ mod tests { #[test] fn deserialize_multi_source_files() { let bytes = frame(2, &[("f0", 3), ("f1", 5)]); - let meta = PkVectorSourceMeta::deserialize(&bytes).unwrap(); + let meta = PrimaryKeyIndexSourceMeta::deserialize(&bytes).unwrap(); assert_eq!(meta.data_level(), 2); assert_eq!(meta.source_files().len(), 2); assert_eq!(meta.source_files()[1].row_count(), 5); @@ -360,13 +373,13 @@ mod tests { fn deserialize_rejects_bad_version() { let mut bytes = frame(1, &[("f0", 1)]); bytes[0..4].copy_from_slice(&2i32.to_be_bytes()); - assert!(PkVectorSourceMeta::deserialize(&bytes).is_err()); + assert!(PrimaryKeyIndexSourceMeta::deserialize(&bytes).is_err()); } #[test] fn deserialize_rejects_zero_or_negative_data_level() { - assert!(PkVectorSourceMeta::deserialize(&frame(0, &[("f0", 1)])).is_err()); - assert!(PkVectorSourceMeta::deserialize(&frame(-1, &[("f0", 1)])).is_err()); + assert!(PrimaryKeyIndexSourceMeta::deserialize(&frame(0, &[("f0", 1)])).is_err()); + assert!(PrimaryKeyIndexSourceMeta::deserialize(&frame(-1, &[("f0", 1)])).is_err()); } #[test] @@ -375,41 +388,41 @@ mod tests { out.extend_from_slice(&1i32.to_be_bytes()); out.extend_from_slice(&1i32.to_be_bytes()); out.extend_from_slice(&0i32.to_be_bytes()); - assert!(PkVectorSourceMeta::deserialize(&out).is_err()); + assert!(PrimaryKeyIndexSourceMeta::deserialize(&out).is_err()); } #[test] fn deserialize_rejects_trailing_bytes() { let mut bytes = frame(1, &[("f0", 1)]); bytes.push(0xFF); - assert!(PkVectorSourceMeta::deserialize(&bytes).is_err()); + assert!(PrimaryKeyIndexSourceMeta::deserialize(&bytes).is_err()); } #[test] fn deserialize_rejects_truncated_input() { let bytes = frame(1, &[("f0", 1)]); - assert!(PkVectorSourceMeta::deserialize(&bytes[..bytes.len() - 2]).is_err()); + assert!(PrimaryKeyIndexSourceMeta::deserialize(&bytes[..bytes.len() - 2]).is_err()); } #[test] fn deserialize_rejects_negative_row_count() { let bytes = frame(1, &[("f0", -1)]); - assert!(PkVectorSourceMeta::deserialize(&bytes).is_err()); + assert!(PrimaryKeyIndexSourceMeta::deserialize(&bytes).is_err()); } #[test] fn new_rejects_empty() { - assert!(PkVectorSourceMeta::new(1, Vec::new()).is_err()); - assert!(PkVectorSourceMeta::new( + assert!(PrimaryKeyIndexSourceMeta::new(1, Vec::new()).is_err()); + assert!(PrimaryKeyIndexSourceMeta::new( 0, - vec![PkVectorSourceFile::new("f0".to_string(), 1).unwrap()] + vec![PrimaryKeyIndexSourceFile::new("f0".to_string(), 1).unwrap()] ) .is_err()); } #[test] fn resolve_single_file() { - let meta = PkVectorSourceMeta::deserialize(&frame(1, &[("f0", 3)])).unwrap(); + let meta = PrimaryKeyIndexSourceMeta::deserialize(&frame(1, &[("f0", 3)])).unwrap(); assert_eq!(meta.resolve(0).unwrap(), ("f0".to_string(), 0)); assert_eq!(meta.resolve(2).unwrap(), ("f0".to_string(), 2)); } @@ -417,7 +430,8 @@ mod tests { #[test] fn resolve_multi_file_prefix_sum_boundaries() { // f0 owns ordinals 0..=2, f1 owns 3..=7. - let meta = PkVectorSourceMeta::deserialize(&frame(1, &[("f0", 3), ("f1", 5)])).unwrap(); + let meta = + PrimaryKeyIndexSourceMeta::deserialize(&frame(1, &[("f0", 3), ("f1", 5)])).unwrap(); assert_eq!(meta.resolve(2).unwrap(), ("f0".to_string(), 2)); // last of f0 assert_eq!(meta.resolve(3).unwrap(), ("f1".to_string(), 0)); // first of f1 assert_eq!(meta.resolve(7).unwrap(), ("f1".to_string(), 4)); // last of f1 @@ -425,13 +439,13 @@ mod tests { #[test] fn resolve_rejects_negative_ordinal() { - let meta = PkVectorSourceMeta::deserialize(&frame(1, &[("f0", 3)])).unwrap(); + let meta = PrimaryKeyIndexSourceMeta::deserialize(&frame(1, &[("f0", 3)])).unwrap(); assert!(meta.resolve(-1).is_err()); } #[test] fn resolve_rejects_ordinal_at_or_past_total() { - let meta = PkVectorSourceMeta::deserialize(&frame(1, &[("f0", 3)])).unwrap(); + let meta = PrimaryKeyIndexSourceMeta::deserialize(&frame(1, &[("f0", 3)])).unwrap(); assert!(meta.resolve(3).is_err()); // total == 3, valid range 0..=2 } @@ -445,7 +459,7 @@ mod tests { index_meta: None, source_meta: None, }; - assert!(PkVectorSourceMeta::from_global_index_meta(&meta).is_err()); + assert!(PrimaryKeyIndexSourceMeta::from_global_index_meta(&meta).is_err()); } #[test] @@ -458,7 +472,7 @@ mod tests { index_meta: None, source_meta: Some(frame(1, &[("f0", 3)])), }; - let parsed = PkVectorSourceMeta::from_global_index_meta(&meta).unwrap(); + let parsed = PrimaryKeyIndexSourceMeta::from_global_index_meta(&meta).unwrap(); assert_eq!(parsed.data_level(), 1); assert_eq!(parsed.source_files()[0].file_name(), "f0"); } @@ -467,11 +481,11 @@ mod tests { fn resolve_rejects_row_count_overflow() { // Two individually-valid row counts whose prefix sum overflows i64. // Resolving past the first file forces the checked_add on the second. - let meta = PkVectorSourceMeta::new( + let meta = PrimaryKeyIndexSourceMeta::new( 1, vec![ - PkVectorSourceFile::new("f0".to_string(), i64::MAX).unwrap(), - PkVectorSourceFile::new("f1".to_string(), 1).unwrap(), + PrimaryKeyIndexSourceFile::new("f0".to_string(), i64::MAX).unwrap(), + PrimaryKeyIndexSourceFile::new("f1".to_string(), 1).unwrap(), ], ) .unwrap(); diff --git a/crates/paimon/src/table/mod.rs b/crates/paimon/src/table/mod.rs index 5c34918d..5096f202 100644 --- a/crates/paimon/src/table/mod.rs +++ b/crates/paimon/src/table/mod.rs @@ -56,6 +56,7 @@ mod lumina_index_build_builder; pub(crate) mod merge_tree_split_generator; mod partition_filter; mod partition_stat; +mod pk_full_text_bucket_state; mod pk_vector_data_file_reader; mod pk_vector_indexed_split_read; mod pk_vector_orchestrator; diff --git a/crates/paimon/src/table/pk_full_text_bucket_state.rs b/crates/paimon/src/table/pk_full_text_bucket_state.rs new file mode 100644 index 00000000..f9122ea1 --- /dev/null +++ b/crates/paimon/src/table/pk_full_text_bucket_state.rs @@ -0,0 +1,614 @@ +// Licensed to the Apache Software Foundation (ASF) under one +// or more contributor license agreements. See the NOTICE file +// distributed with this work for additional information +// regarding copyright ownership. The ASF licenses this file +// to you under the Apache License, Version 2.0 (the +// "License"); you may not use this file except in compliance +// with the License. You may obtain a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, +// software distributed under the License is distributed on an +// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY +// KIND, either express or implied. See the License for the +// specific language governing permissions and limitations +// under the License. + +//! Per-bucket current/stale reconciliation for primary-key full-text index +//! payloads, aligned to active data files by level. +//! +//! Mirrors Java `PkFullTextBucketIndexState.fromActiveDataFiles`: a payload is +//! *current* only when it exactly covers the live source files at its level and +//! passes its own per-payload validation. Payloads that no longer match the live +//! data files, fail to parse, or fail their singleton validation (bad row count, +//! bad row range, or an overflowing source row total) are collected as *stale* +//! (not an error); payloads that are not full-text or carry no global index meta +//! are ignored. Only cross-payload conflicts across the surviving current set +//! (a duplicate payload file name, or a source data file covered by two current +//! payloads) fail loud. + +use std::collections::{BTreeMap, HashSet}; + +use indexmap::IndexMap; + +use crate::spec::{ + should_read_pk_index_source, DataFileMeta, GlobalIndexMeta, IndexFileMeta, + PrimaryKeyIndexSourceFile, PrimaryKeyIndexSourceMeta, +}; + +pub(crate) const PK_FULL_TEXT_INDEX_TYPE: &str = "full-text"; + +fn data_invalid(message: impl Into) -> crate::Error { + crate::Error::DataInvalid { + message: message.into(), + source: None, + } +} + +/// Per-payload singleton validation, mirroring Java's public +/// `PkFullTextBucketIndexState` payload constructor guarded by a `try/catch` +/// that routes any `RuntimeException` to stale: the payload's declared row count +/// must equal the sum of its source rows, and its row range must span exactly +/// `[0, total - 1]`. The sum is computed with `checked_add` so an overflowing +/// total fails the check (Java `Math.addExact` throws → caught → stale) instead +/// of panicking. Returns `false` on any defect so the caller routes the payload +/// to stale. +fn payload_matches_source( + payload: &IndexFileMeta, + global_meta: &GlobalIndexMeta, + source_meta: &PrimaryKeyIndexSourceMeta, +) -> bool { + let mut total: i64 = 0; + for source in source_meta.source_files() { + match total.checked_add(source.row_count()) { + Some(next) => total = next, + None => return false, // overflow → stale + } + } + i64::from(payload.row_count) == total + && global_meta.row_range_start == 0 + && global_meta.row_range_end == total - 1 +} + +/// Reconciled full-text payloads for a single bucket. +// The read path consumes this once it lands; until then the state has no +// production caller. +#[allow(dead_code)] +pub(crate) struct PkFullTextBucketState { + // Consumed by the read path once it lands, kept to identify the indexed field. + text_field_id: i32, + current_payloads: Vec, + stale_payloads: Vec, + // Maps each covered source data file to its current payload. Insertion order + // follows the current-payload order and, within a payload, its source-file + // order, so the read path can build deterministically ordered splits from it. + payload_by_source_file: IndexMap, +} + +#[allow(dead_code)] +impl PkFullTextBucketState { + pub(crate) fn from_active_data_files( + text_field_id: i32, + active_data_files: &[DataFileMeta], + active_payloads: Vec, + ) -> crate::Result { + // Desired live source set per level (sorted by file name). + let mut sources_by_level: BTreeMap> = BTreeMap::new(); + for file in active_data_files { + if should_read_pk_index_source(file) { + sources_by_level.entry(file.level).or_default().push( + PrimaryKeyIndexSourceFile::new(file.file_name.clone(), file.row_count)?, + ); + } + } + for sources in sources_by_level.values_mut() { + sources.sort_by(|a, b| a.file_name().cmp(b.file_name())); + } + + // Classify payloads. A payload lands in `by_level` only when it parses, + // matches the desired live source set exactly, and passes its own + // per-payload validation; every other single-payload defect routes it to + // stale (mirroring Java's `try/catch (RuntimeException → stale)`). + let mut by_level: BTreeMap> = + BTreeMap::new(); + let mut stale: Vec = Vec::new(); + for payload in active_payloads { + if payload.index_type != PK_FULL_TEXT_INDEX_TYPE { + continue; // IGNORED + } + let Some(global_meta) = payload.global_index_meta.as_ref() else { + continue; // IGNORED + }; + // Field id is checked before the source-meta is even parsed: a payload + // that indexes a different text field is stale only when it carries + // source-meta, and is otherwise skipped without ever decoding it. + if global_meta.index_field_id != text_field_id { + if global_meta.source_meta.is_some() { + stale.push(payload); // STALE (wrong field, has source-meta) + } + continue; // wrong field, no source-meta → IGNORED + } + // Matching field: parse source-meta. A parse failure or absent + // source-meta makes the payload stale (mirrors Java decoding it inside + // the per-candidate try/catch). + let source_meta = match global_meta + .source_meta + .as_ref() + .map(|b| PrimaryKeyIndexSourceMeta::deserialize(b)) + { + Some(Ok(m)) => m, + // No source-meta or a parse failure: cannot be a current payload. + Some(Err(_)) | None => { + stale.push(payload); + continue; + } + }; + // Desired-set match (exact, order-sensitive). + let matches_desired = sources_by_level + .get(&source_meta.data_level()) + .is_some_and(|desired| desired.as_slice() == source_meta.source_files()); + if !matches_desired { + stale.push(payload); // STALE (missing level or set mismatch) + continue; + } + // Per-payload singleton validation (row total via checked_add, row + // count, row range). Java runs this inside try/catch, so any failure + // — including an overflowing row total — routes the payload to stale, + // never an error. + if !payload_matches_source(&payload, global_meta, &source_meta) { + stale.push(payload); // STALE (per-payload defect) + continue; + } + by_level + .entry(source_meta.data_level()) + .or_default() + .push((payload, source_meta)); + } + + // Per level: exactly one survivor is current; >1 → all stale. + let mut current: Vec<(IndexFileMeta, PrimaryKeyIndexSourceMeta)> = Vec::new(); + for (_level, mut payloads) in by_level { + if payloads.len() == 1 { + current.push(payloads.remove(0)); + } else { + stale.extend(payloads.into_iter().map(|(payload, _)| payload)); + } + } + + // Cross-payload conflicts across the current set are the only fail-loud + // cases (the surviving payloads already passed per-payload validation): + // a duplicate payload file name, or a source data file covered twice. + let mut payload_by_source_file: IndexMap = IndexMap::new(); + let mut seen_payloads: HashSet = HashSet::new(); + let mut current_payloads: Vec = Vec::with_capacity(current.len()); + for (payload, source_meta) in current { + if !seen_payloads.insert(payload.file_name.clone()) { + return Err(data_invalid(format!( + "full-text payload {} appears more than once", + payload.file_name + ))); + } + for source in source_meta.source_files() { + if payload_by_source_file + .insert(source.file_name().to_string(), payload.clone()) + .is_some() + { + return Err(data_invalid(format!( + "source data file {} is covered by more than one full-text payload", + source.file_name() + ))); + } + } + current_payloads.push(payload); + } + + Ok(Self { + text_field_id, + current_payloads, + stale_payloads: stale, + payload_by_source_file, + }) + } + + pub(crate) fn current_payloads(&self) -> &[IndexFileMeta] { + &self.current_payloads + } + pub(crate) fn stale_payloads(&self) -> &[IndexFileMeta] { + &self.stale_payloads + } + pub(crate) fn payload_by_source_file(&self) -> &IndexMap { + &self.payload_by_source_file + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::spec::stats::BinaryTableStats; + use crate::spec::GlobalIndexMeta; + + /// COMPACT file source discriminant (matches Java `FileSource.COMPACT`). + const FILE_SOURCE_COMPACT: i32 = 1; + + /// A COMPACT data file at `level` (>0), so `should_read_pk_index_source` + /// accepts it as a live source. + fn dfm(name: &str, rows: i64, level: i32) -> DataFileMeta { + DataFileMeta { + file_name: name.into(), + file_size: 1, + row_count: rows, + min_key: Vec::new(), + max_key: Vec::new(), + key_stats: BinaryTableStats::empty(), + value_stats: BinaryTableStats::empty(), + min_sequence_number: 0, + max_sequence_number: 0, + schema_id: 1, + level, + extra_files: Vec::new(), + creation_time: None, + delete_row_count: None, + embedded_index: None, + file_source: Some(FILE_SOURCE_COMPACT), + value_stats_cols: None, + external_path: None, + first_row_id: Some(0), + write_cols: None, + } + } + + /// One Java `DataOutput#writeUTF` value (u16-BE length + modified UTF-8) for + /// the ASCII file names used in these tests. + fn java_write_utf(s: &str) -> Vec { + let mut body = Vec::new(); + for c in s.encode_utf16() { + if (0x0001..=0x007F).contains(&c) { + body.push(c as u8); + } else if c > 0x07FF { + body.push(0xE0 | (c >> 12) as u8); + body.push(0x80 | ((c >> 6) & 0x3F) as u8); + body.push(0x80 | (c & 0x3F) as u8); + } else { + body.push(0xC0 | (c >> 6) as u8); + body.push(0x80 | (c & 0x3F) as u8); + } + } + let mut out = (body.len() as u16).to_be_bytes().to_vec(); + out.extend_from_slice(&body); + out + } + + /// A valid `_SOURCE_META` frame (version 1) for the given level and files. + fn frame(data_level: i32, files: &[(&str, i64)]) -> Vec { + let mut out = Vec::new(); + out.extend_from_slice(&1i32.to_be_bytes()); // version + out.extend_from_slice(&data_level.to_be_bytes()); + out.extend_from_slice(&(files.len() as i32).to_be_bytes()); + for (name, rows) in files { + out.extend_from_slice(&java_write_utf(name)); + out.extend_from_slice(&rows.to_be_bytes()); + } + out + } + + /// Build a full-text `IndexFileMeta` payload with a `GlobalIndexMeta`. + fn payload( + file_name: &str, + index_type: &str, + row_count: i32, + global_index_meta: Option, + ) -> IndexFileMeta { + IndexFileMeta { + index_type: index_type.into(), + file_name: file_name.into(), + file_size: 1, + row_count, + deletion_vectors_ranges: None, + global_index_meta, + } + } + + /// A `GlobalIndexMeta` with a source-meta blob and matching row range. + fn gim(field_id: i32, start: i64, end: i64, source_meta: Option>) -> GlobalIndexMeta { + GlobalIndexMeta { + row_range_start: start, + row_range_end: end, + index_field_id: field_id, + extra_field_ids: None, + index_meta: None, + source_meta, + } + } + + // (a) exact match → current. + #[test] + fn exact_match_is_current() { + let files = [dfm("d0.parquet", 100, 1)]; + let payloads = vec![payload( + "ft-0", + PK_FULL_TEXT_INDEX_TYPE, + 100, + Some(gim(7, 0, 99, Some(frame(1, &[("d0.parquet", 100)])))), + )]; + let state = PkFullTextBucketState::from_active_data_files(7, &files, payloads).unwrap(); + assert_eq!(state.current_payloads().len(), 1); + assert_eq!(state.current_payloads()[0].file_name, "ft-0"); + assert!(state.stale_payloads().is_empty()); + assert_eq!(state.payload_by_source_file().len(), 1); + assert!(state.payload_by_source_file().contains_key("d0.parquet")); + } + + // (b) source drift (source set no longer matches live files) → stale. + #[test] + fn source_drift_is_stale() { + let files = [dfm("d0.parquet", 100, 1)]; + // Payload references a file that is no longer live at level 1. + let payloads = vec![payload( + "ft-0", + PK_FULL_TEXT_INDEX_TYPE, + 50, + Some(gim(7, 0, 49, Some(frame(1, &[("stale.parquet", 50)])))), + )]; + let state = PkFullTextBucketState::from_active_data_files(7, &files, payloads).unwrap(); + assert!(state.current_payloads().is_empty()); + assert_eq!(state.stale_payloads().len(), 1); + assert_eq!(state.stale_payloads()[0].file_name, "ft-0"); + } + + // (b') reorder-only drift → stale. The desired set is sorted by file name, so + // a payload whose source-meta lists the same files in a different order fails + // the order-sensitive equality gate (pins that the match is order-sensitive). + #[test] + fn reordered_source_files_are_stale() { + let files = [dfm("a.parquet", 100, 1), dfm("b.parquet", 200, 1)]; + // Same files and row counts, but listed b-before-a instead of the sorted + // a-before-b, so the exact List.equals gate rejects it. + let payloads = vec![payload( + "ft-0", + PK_FULL_TEXT_INDEX_TYPE, + 300, + Some(gim( + 7, + 0, + 299, + Some(frame(1, &[("b.parquet", 200), ("a.parquet", 100)])), + )), + )]; + let state = PkFullTextBucketState::from_active_data_files(7, &files, payloads).unwrap(); + assert!(state.current_payloads().is_empty()); + assert_eq!(state.stale_payloads().len(), 1); + assert_eq!(state.stale_payloads()[0].file_name, "ft-0"); + } + #[test] + fn two_payloads_same_level_are_both_stale() { + let files = [dfm("d0.parquet", 100, 1)]; + let source = frame(1, &[("d0.parquet", 100)]); + let payloads = vec![ + payload( + "ft-a", + PK_FULL_TEXT_INDEX_TYPE, + 100, + Some(gim(7, 0, 99, Some(source.clone()))), + ), + payload( + "ft-b", + PK_FULL_TEXT_INDEX_TYPE, + 100, + Some(gim(7, 0, 99, Some(source))), + ), + ]; + let state = PkFullTextBucketState::from_active_data_files(7, &files, payloads).unwrap(); + assert!(state.current_payloads().is_empty()); + let stale: Vec<&str> = state + .stale_payloads() + .iter() + .map(|p| p.file_name.as_str()) + .collect(); + assert_eq!(stale.len(), 2); + assert!(stale.contains(&"ft-a") && stale.contains(&"ft-b")); + } + + // (d) wrong index_type → ignored (in neither list). + #[test] + fn wrong_index_type_is_ignored() { + let files = [dfm("d0.parquet", 100, 1)]; + let payloads = vec![payload( + "gi-0", + "GLOBAL_INDEX", + 100, + Some(gim(7, 0, 99, Some(frame(1, &[("d0.parquet", 100)])))), + )]; + let state = PkFullTextBucketState::from_active_data_files(7, &files, payloads).unwrap(); + assert!(state.current_payloads().is_empty()); + assert!(state.stale_payloads().is_empty()); + } + + // (e) no global_index_meta → ignored (in neither list). + #[test] + fn no_global_index_meta_is_ignored() { + let files = [dfm("d0.parquet", 100, 1)]; + let payloads = vec![payload("ft-0", PK_FULL_TEXT_INDEX_TYPE, 100, None)]; + let state = PkFullTextBucketState::from_active_data_files(7, &files, payloads).unwrap(); + assert!(state.current_payloads().is_empty()); + assert!(state.stale_payloads().is_empty()); + } + + // (f) field-id mismatch (with source-meta present) → stale. + #[test] + fn field_id_mismatch_is_stale() { + let files = [dfm("d0.parquet", 100, 1)]; + let payloads = vec![payload( + "ft-0", + PK_FULL_TEXT_INDEX_TYPE, + 100, + // index_field_id 9 != text_field_id 7. + Some(gim(9, 0, 99, Some(frame(1, &[("d0.parquet", 100)])))), + )]; + let state = PkFullTextBucketState::from_active_data_files(7, &files, payloads).unwrap(); + assert!(state.current_payloads().is_empty()); + assert_eq!(state.stale_payloads().len(), 1); + assert_eq!(state.stale_payloads()[0].file_name, "ft-0"); + } + + // (f') field-id mismatch with NO source-meta → ignored (skipped, not stale): + // a payload indexing another field that carries no source metadata is another + // field's concern, so it is dropped without ever being decoded. + #[test] + fn field_id_mismatch_without_source_meta_is_ignored() { + let files = [dfm("d0.parquet", 100, 1)]; + let payloads = vec![payload( + "ft-other", + PK_FULL_TEXT_INDEX_TYPE, + 100, + // index_field_id 9 != text_field_id 7, and no source-meta blob. + Some(gim(9, 0, 99, None)), + )]; + let state = PkFullTextBucketState::from_active_data_files(7, &files, payloads).unwrap(); + assert!(state.current_payloads().is_empty()); + assert!(state.stale_payloads().is_empty()); + assert!(state.payload_by_source_file().is_empty()); + } + + // (g) corrupt source-meta bytes → stale. + #[test] + fn corrupt_source_meta_is_stale() { + let files = [dfm("d0.parquet", 100, 1)]; + let payloads = vec![payload( + "ft-0", + PK_FULL_TEXT_INDEX_TYPE, + 100, + Some(gim(7, 0, 99, Some(vec![0xFF, 0x00, 0x01]))), + )]; + let state = PkFullTextBucketState::from_active_data_files(7, &files, payloads).unwrap(); + assert!(state.current_payloads().is_empty()); + assert_eq!(state.stale_payloads().len(), 1); + assert_eq!(state.stale_payloads()[0].file_name, "ft-0"); + } + + // (h) desired level missing from live files → stale. + #[test] + fn missing_desired_level_is_stale() { + // Live files are all at level 1; payload's source-meta claims level 2. + let files = [dfm("d0.parquet", 100, 1)]; + let payloads = vec![payload( + "ft-0", + PK_FULL_TEXT_INDEX_TYPE, + 100, + Some(gim(7, 0, 99, Some(frame(2, &[("d0.parquet", 100)])))), + )]; + let state = PkFullTextBucketState::from_active_data_files(7, &files, payloads).unwrap(); + assert!(state.current_payloads().is_empty()); + assert_eq!(state.stale_payloads().len(), 1); + assert_eq!(state.stale_payloads()[0].file_name, "ft-0"); + } + + // (l) duplicate payload file_name across levels in the current set → Err. + #[test] + fn duplicate_payload_file_name_is_error() { + // Same file name lives at two levels, each with a matching current payload + // that happens to share the same index file name. + let files = [dfm("d0.parquet", 100, 1), dfm("d1.parquet", 200, 2)]; + let payloads = vec![ + payload( + "ft-dup", + PK_FULL_TEXT_INDEX_TYPE, + 100, + Some(gim(7, 0, 99, Some(frame(1, &[("d0.parquet", 100)])))), + ), + payload( + "ft-dup", + PK_FULL_TEXT_INDEX_TYPE, + 200, + Some(gim(7, 0, 199, Some(frame(2, &[("d1.parquet", 200)])))), + ), + ]; + assert!(PkFullTextBucketState::from_active_data_files(7, &files, payloads).is_err()); + } + + // (i) row_count mismatch (per-payload defect) → stale, NOT error. + #[test] + fn row_count_mismatch_is_stale() { + let files = [dfm("d0.parquet", 100, 1)]; + // Source rows total 100 but payload declares 99. + let payloads = vec![payload( + "ft-0", + PK_FULL_TEXT_INDEX_TYPE, + 99, + Some(gim(7, 0, 98, Some(frame(1, &[("d0.parquet", 100)])))), + )]; + let state = PkFullTextBucketState::from_active_data_files(7, &files, payloads).unwrap(); + assert!(state.current_payloads().is_empty()); + assert_eq!(state.stale_payloads().len(), 1); + assert_eq!(state.stale_payloads()[0].file_name, "ft-0"); + } + + // (j) row_range mismatch alone (row_count correct, per-payload defect) → + // stale, NOT error. + #[test] + fn row_range_mismatch_is_stale() { + let files = [dfm("d0.parquet", 100, 1)]; + // row_count matches total (100) but row range end is wrong (should be 99). + let payloads = vec![payload( + "ft-0", + PK_FULL_TEXT_INDEX_TYPE, + 100, + Some(gim(7, 0, 50, Some(frame(1, &[("d0.parquet", 100)])))), + )]; + let state = PkFullTextBucketState::from_active_data_files(7, &files, payloads).unwrap(); + assert!(state.current_payloads().is_empty()); + assert_eq!(state.stale_payloads().len(), 1); + assert_eq!(state.stale_payloads()[0].file_name, "ft-0"); + } + + // (k) source row counts overflow i64 (per-payload defect) → stale, NOT a + // panic and NOT error. + #[test] + fn source_row_count_overflow_is_stale() { + // Two live level-1 files whose row counts sum past i64::MAX; the payload + // matches the desired set exactly, so only the overflowing row total can + // reject it — and it must do so without panicking. + let files = [dfm("a.parquet", i64::MAX, 1), dfm("b.parquet", i64::MAX, 1)]; + let payloads = vec![payload( + "ft-0", + PK_FULL_TEXT_INDEX_TYPE, + 0, + Some(gim( + 7, + 0, + 0, + Some(frame( + 1, + &[("a.parquet", i64::MAX), ("b.parquet", i64::MAX)], + )), + )), + )]; + let state = PkFullTextBucketState::from_active_data_files(7, &files, payloads).unwrap(); + assert!(state.current_payloads().is_empty()); + assert_eq!(state.stale_payloads().len(), 1); + assert_eq!(state.stale_payloads()[0].file_name, "ft-0"); + } + + // (m) a source data file covered by two current payloads → Err. + #[test] + fn source_double_covered_is_error() { + // The same source file name is live at two levels; two current payloads + // (one per level) each cover it, colliding in payload_by_source_file. + let files = [dfm("shared.parquet", 100, 1), dfm("shared.parquet", 200, 2)]; + let payloads = vec![ + payload( + "ft-a", + PK_FULL_TEXT_INDEX_TYPE, + 100, + Some(gim(7, 0, 99, Some(frame(1, &[("shared.parquet", 100)])))), + ), + payload( + "ft-b", + PK_FULL_TEXT_INDEX_TYPE, + 200, + Some(gim(7, 0, 199, Some(frame(2, &[("shared.parquet", 200)])))), + ), + ]; + assert!(PkFullTextBucketState::from_active_data_files(7, &files, payloads).is_err()); + } +} diff --git a/crates/paimon/src/table/pk_vector_orchestrator.rs b/crates/paimon/src/table/pk_vector_orchestrator.rs index 99ab5108..2278b0d2 100644 --- a/crates/paimon/src/table/pk_vector_orchestrator.rs +++ b/crates/paimon/src/table/pk_vector_orchestrator.rs @@ -914,7 +914,8 @@ mod e2e_tests { use crate::io::{FileIO, FileIOBuilder}; use crate::spec::stats::BinaryTableStats; use crate::spec::{ - DataField, DataFileMeta, DataType, IntType, PkVectorSourceFile, PkVectorSourceMeta, + DataField, DataFileMeta, DataType, IntType, PrimaryKeyIndexSourceFile, + PrimaryKeyIndexSourceMeta, }; use crate::table::pk_vector_indexed_split_read::PkVectorIndexedSplitRead; use crate::table::pk_vector_position_read::{PKEY_VECTOR_POSITION_COLUMN, SEARCH_SCORE_COLUMN}; @@ -1070,11 +1071,11 @@ mod e2e_tests { fn ann_segment(sources: &[(&str, i64)]) -> BucketAnnSegment { BucketAnnSegment::for_test( - PkVectorSourceMeta::new( + PrimaryKeyIndexSourceMeta::new( 1, sources .iter() - .map(|(n, r)| PkVectorSourceFile::new((*n).to_string(), *r).unwrap()) + .map(|(n, r)| PrimaryKeyIndexSourceFile::new((*n).to_string(), *r).unwrap()) .collect(), ) .unwrap(), diff --git a/crates/paimon/src/table/pk_vector_scan.rs b/crates/paimon/src/table/pk_vector_scan.rs index c4947240..e6ba0c55 100644 --- a/crates/paimon/src/table/pk_vector_scan.rs +++ b/crates/paimon/src/table/pk_vector_scan.rs @@ -23,8 +23,8 @@ use std::collections::{BTreeMap, HashSet}; use crate::spec::{ - BinaryRow, DataFileMeta, FileKind, GlobalIndexMeta, IndexManifest, PkVectorSourceFile, - PkVectorSourceMeta, Predicate, + should_read_pk_index_source, BinaryRow, DataFileMeta, FileKind, GlobalIndexMeta, IndexManifest, + Predicate, PrimaryKeyIndexSourceFile, PrimaryKeyIndexSourceMeta, }; use crate::table::pk_vector_orchestrator::PkVectorSearchSplit; use crate::table::source::{DataSplit, DataSplitBuilder, DeletionFile}; @@ -32,7 +32,6 @@ use crate::table::Table; use crate::vindex::pkvector::bucket::{BucketActiveFile, BucketAnnSegment}; const INDEX_DIR: &str = "index"; -const FILE_SOURCE_COMPACT: i32 = 1; fn data_invalid(message: impl Into) -> crate::Error { crate::Error::DataInvalid { @@ -41,13 +40,7 @@ fn data_invalid(message: impl Into) -> crate::Error { } } -/// Mirror of `PrimaryKeyIndexSourcePolicy.shouldRead`: only compacted, non-level-0 -/// files back the PK-vector index; an absent file source reads as false. -fn should_read_pk_index_source(file: &DataFileMeta) -> bool { - matches!(file.file_source, Some(src) if src == FILE_SOURCE_COMPACT) && file.level > 0 -} - -fn source_files_unique(files: &[PkVectorSourceFile]) -> bool { +fn source_files_unique(files: &[PrimaryKeyIndexSourceFile]) -> bool { let mut seen = HashSet::new(); files.iter().all(|file| seen.insert(file.file_name())) } @@ -56,13 +49,13 @@ fn current_ann_segments( active_data_files: &[DataFileMeta], ann_segments: Vec, ) -> crate::Result> { - let mut sources_by_level: BTreeMap> = BTreeMap::new(); + let mut sources_by_level: BTreeMap> = BTreeMap::new(); for file in active_data_files { if should_read_pk_index_source(file) { sources_by_level .entry(file.level) .or_default() - .push(PkVectorSourceFile::new( + .push(PrimaryKeyIndexSourceFile::new( file.file_name.clone(), file.row_count, )?); @@ -320,7 +313,7 @@ fn plan_from_inputs( // Phase A: group ANN payloads by (partition, bucket). let mut segments_by_bucket: BTreeMap> = BTreeMap::new(); for (partition, bucket, gim, path, file_size, file_name) in index_entries { - let source_meta = PkVectorSourceMeta::from_global_index_meta(&gim) + let source_meta = PrimaryKeyIndexSourceMeta::from_global_index_meta(&gim) .map_err(|_| data_invalid(format!("index file {file_name} is not active")))?; let key = (partition.to_serialized_bytes(), bucket); segments_by_bucket @@ -435,9 +428,9 @@ mod tests { out } - /// Build a `_SOURCE_META` blob the way `PkVectorSourceMeta::deserialize` + /// Build a `_SOURCE_META` blob the way `PrimaryKeyIndexSourceMeta::deserialize` /// expects it. There is no public serializer, so we mirror the frame used by - /// `pk_vector_source.rs`'s own round-trip tests. + /// `pk_index_source.rs`'s own round-trip tests. fn source_meta_bytes(data_level: i32, files: &[(&str, i64)]) -> Vec { let mut out = Vec::new(); out.extend_from_slice(&1i32.to_be_bytes()); // version @@ -463,12 +456,12 @@ mod tests { fn ann_segment(data_level: i32, path: &str, source_files: &[(&str, i64)]) -> BucketAnnSegment { BucketAnnSegment { - source_meta: PkVectorSourceMeta::new( + source_meta: PrimaryKeyIndexSourceMeta::new( data_level, source_files .iter() .map(|(name, rows)| { - PkVectorSourceFile::new((*name).to_string(), *rows).unwrap() + PrimaryKeyIndexSourceFile::new((*name).to_string(), *rows).unwrap() }) .collect(), ) diff --git a/crates/paimon/src/table/vector_search_builder.rs b/crates/paimon/src/table/vector_search_builder.rs index 9f88558e..273dced9 100644 --- a/crates/paimon/src/table/vector_search_builder.rs +++ b/crates/paimon/src/table/vector_search_builder.rs @@ -4161,9 +4161,9 @@ mod tests { /// A `PkVectorSearchSplit` carrying a single ANN segment addressed by `path`. fn pk_split_with_segment(path: &str) -> PkVectorSearchSplit { let mut split = pk_search_split(0, vec![pk_data_file("file-a", 3, Some(0))]); - let source_meta = crate::spec::PkVectorSourceMeta::new( + let source_meta = crate::spec::PrimaryKeyIndexSourceMeta::new( 1, - vec![crate::spec::PkVectorSourceFile::new("file-a".to_string(), 3).unwrap()], + vec![crate::spec::PrimaryKeyIndexSourceFile::new("file-a".to_string(), 3).unwrap()], ) .unwrap(); let mut segment = BucketAnnSegment::for_test(source_meta); @@ -4174,9 +4174,9 @@ mod tests { fn pk_split_with_lumina_segment(path: &str, metric: &str) -> PkVectorSearchSplit { let mut split = pk_search_split(0, vec![pk_data_file("file-a", 3, Some(0))]); - let source_meta = crate::spec::PkVectorSourceMeta::new( + let source_meta = crate::spec::PrimaryKeyIndexSourceMeta::new( 1, - vec![crate::spec::PkVectorSourceFile::new("file-a".to_string(), 3).unwrap()], + vec![crate::spec::PrimaryKeyIndexSourceFile::new("file-a".to_string(), 3).unwrap()], ) .unwrap(); let mut segment = BucketAnnSegment::for_test(source_meta); diff --git a/crates/paimon/src/vindex/pkvector/ann.rs b/crates/paimon/src/vindex/pkvector/ann.rs index bbd7088d..ab50196a 100644 --- a/crates/paimon/src/vindex/pkvector/ann.rs +++ b/crates/paimon/src/vindex/pkvector/ann.rs @@ -23,7 +23,7 @@ use super::data_invalid; use super::metric::{java_float_compare, VectorSearchMetric}; use super::result::PkVectorSearchResult; use crate::deletion_vector::DeletionVector; -use crate::spec::{PkVectorSourceFile, PkVectorSourceMeta}; +use crate::spec::{PrimaryKeyIndexSourceFile, PrimaryKeyIndexSourceMeta}; use crate::vector_search::VectorSearch; /// Build the live-row-id mask for the ANN reader's `include_row_ids` filter, in @@ -46,7 +46,7 @@ use crate::vector_search::VectorSearch; /// no deletion vector is relevant — nothing to mask. Otherwise returns the masked /// live ids. pub(crate) fn build_live_row_ids( - source_files: &[PkVectorSourceFile], + source_files: &[PrimaryKeyIndexSourceFile], active_source_files: &HashSet, deletion_vectors: &HashMap>, residual_ranges: Option<&HashMap>, @@ -122,7 +122,7 @@ pub(crate) fn build_live_row_ids( /// sorted BEST_FIRST. pub(crate) fn map_ann_results( scored: &[(u64, f32)], - source_meta: &PkVectorSourceMeta, + source_meta: &PrimaryKeyIndexSourceMeta, active_source_files: &HashSet, deletion_vectors: &HashMap>, residual_ranges: Option<&HashMap>, @@ -327,12 +327,12 @@ mod tests { use super::*; use roaring::RoaringBitmap; - fn source_meta(files: &[(&str, i64)]) -> PkVectorSourceMeta { + fn source_meta(files: &[(&str, i64)]) -> PrimaryKeyIndexSourceMeta { let files = files .iter() - .map(|(name, rows)| PkVectorSourceFile::new((*name).to_string(), *rows).unwrap()) + .map(|(name, rows)| PrimaryKeyIndexSourceFile::new((*name).to_string(), *rows).unwrap()) .collect(); - PkVectorSourceMeta::new(1, files).unwrap() + PrimaryKeyIndexSourceMeta::new(1, files).unwrap() } fn dv(deleted: &[u32]) -> Arc { @@ -349,7 +349,7 @@ mod tests { #[test] fn test_build_live_row_ids_none_when_all_active_and_no_relevant_dv() { - let files = [PkVectorSourceFile::new("f0".into(), 3).unwrap()]; + let files = [PrimaryKeyIndexSourceFile::new("f0".into(), 3).unwrap()]; let active = active_set(&["f0"]); // All active + empty map -> None. assert!(build_live_row_ids(&files, &active, &HashMap::new(), None) @@ -368,8 +368,8 @@ mod tests { // f0 rows 0..3 (global 0,1,2), f1 rows 0..2 (global 3,4). f1 is inactive, // so its whole ordinal range is masked out; f0 stays fully live. No DV. let files = vec![ - PkVectorSourceFile::new("f0".into(), 3).unwrap(), - PkVectorSourceFile::new("f1".into(), 2).unwrap(), + PrimaryKeyIndexSourceFile::new("f0".into(), 3).unwrap(), + PrimaryKeyIndexSourceFile::new("f1".into(), 2).unwrap(), ]; let live = build_live_row_ids(&files, &active_set(&["f0"]), &HashMap::new(), None) .unwrap() @@ -381,8 +381,8 @@ mod tests { fn test_build_live_row_ids_masks_deleted_positions_with_file_offsets() { // f0 rows 0..3 (global 0,1,2), f1 rows 0..2 (global 3,4). let files = vec![ - PkVectorSourceFile::new("f0".into(), 3).unwrap(), - PkVectorSourceFile::new("f1".into(), 2).unwrap(), + PrimaryKeyIndexSourceFile::new("f0".into(), 3).unwrap(), + PrimaryKeyIndexSourceFile::new("f1".into(), 2).unwrap(), ]; let mut dvs = HashMap::new(); dvs.insert("f0".to_string(), dv(&[1])); // deletes global 1 @@ -498,12 +498,12 @@ mod tests { ), ); let segment = BucketAnnSegment::for_test({ - use crate::spec::{PkVectorSourceFile, PkVectorSourceMeta}; - PkVectorSourceMeta::new( + use crate::spec::{PrimaryKeyIndexSourceFile, PrimaryKeyIndexSourceMeta}; + PrimaryKeyIndexSourceMeta::new( 1, vec![ - PkVectorSourceFile::new("f0".into(), 3).unwrap(), - PkVectorSourceFile::new("f1".into(), 5).unwrap(), + PrimaryKeyIndexSourceFile::new("f0".into(), 3).unwrap(), + PrimaryKeyIndexSourceFile::new("f1".into(), 5).unwrap(), ], ) .unwrap() @@ -541,9 +541,12 @@ mod tests { }), ); let segment = BucketAnnSegment::for_test({ - use crate::spec::{PkVectorSourceFile, PkVectorSourceMeta}; - PkVectorSourceMeta::new(1, vec![PkVectorSourceFile::new("f0".into(), 1).unwrap()]) - .unwrap() + use crate::spec::{PrimaryKeyIndexSourceFile, PrimaryKeyIndexSourceMeta}; + PrimaryKeyIndexSourceMeta::new( + 1, + vec![PrimaryKeyIndexSourceFile::new("f0".into(), 1).unwrap()], + ) + .unwrap() }); let err = searcher .search( @@ -569,9 +572,12 @@ mod tests { }), ); let segment = BucketAnnSegment::for_test({ - use crate::spec::{PkVectorSourceFile, PkVectorSourceMeta}; - PkVectorSourceMeta::new(1, vec![PkVectorSourceFile::new("f0".into(), 1).unwrap()]) - .unwrap() + use crate::spec::{PrimaryKeyIndexSourceFile, PrimaryKeyIndexSourceMeta}; + PrimaryKeyIndexSourceMeta::new( + 1, + vec![PrimaryKeyIndexSourceFile::new("f0".into(), 1).unwrap()], + ) + .unwrap() }); let results = searcher .search( @@ -603,8 +609,8 @@ mod tests { // entry (empty allow). Result: f0 keeps {0} (1 is residual-allowed but // deleted, 2 not residual-allowed); f1 contributes nothing. let files = vec![ - PkVectorSourceFile::new("f0".into(), 3).unwrap(), - PkVectorSourceFile::new("f1".into(), 2).unwrap(), + PrimaryKeyIndexSourceFile::new("f0".into(), 3).unwrap(), + PrimaryKeyIndexSourceFile::new("f1".into(), 2).unwrap(), ]; let mut dvs = HashMap::new(); dvs.insert("f0".to_string(), dv(&[1])); @@ -621,8 +627,8 @@ mod tests { // f0 rows global 0,1,2; f1 rows global 3,4. residual allows f0={2}, f1={1}. // f1 physical pos 1 -> global 3 + 1 = 4. Result {2, 4}. No DV. let files = vec![ - PkVectorSourceFile::new("f0".into(), 3).unwrap(), - PkVectorSourceFile::new("f1".into(), 2).unwrap(), + PrimaryKeyIndexSourceFile::new("f0".into(), 3).unwrap(), + PrimaryKeyIndexSourceFile::new("f1".into(), 2).unwrap(), ]; let mut residual = HashMap::new(); residual.insert("f0".to_string(), treemap(&[2])); @@ -642,7 +648,7 @@ mod tests { fn test_build_live_row_ids_residual_some_returns_mask_even_when_all_active_no_dv() { // All active, no DV: without residual this returns None. With a residual // present, a mask is always required. - let files = [PkVectorSourceFile::new("f0".into(), 3).unwrap()]; + let files = [PrimaryKeyIndexSourceFile::new("f0".into(), 3).unwrap()]; let mut residual = HashMap::new(); residual.insert("f0".to_string(), treemap(&[0, 2])); let live = build_live_row_ids( diff --git a/crates/paimon/src/vindex/pkvector/bucket.rs b/crates/paimon/src/vindex/pkvector/bucket.rs index b946918f..4666d403 100644 --- a/crates/paimon/src/vindex/pkvector/bucket.rs +++ b/crates/paimon/src/vindex/pkvector/bucket.rs @@ -28,7 +28,7 @@ use super::data_invalid; use super::metric::{java_float_compare, VectorSearchMetric}; use super::result::PkVectorSearchResult; use crate::deletion_vector::DeletionVector; -use crate::spec::PkVectorSourceMeta; +use crate::spec::PrimaryKeyIndexSourceMeta; /// Search one uncovered data file for its per-query exact Top-K. Returns one /// bounded, BEST_FIRST list per query (outer index aligns to the `queries` slice @@ -44,7 +44,7 @@ pub(crate) type ExactFileSearchFuture<'a> = /// masking; the remaining fields address the segment's index file for the ANN /// scorer that reads it. pub(crate) struct BucketAnnSegment { - pub source_meta: PkVectorSourceMeta, + pub source_meta: PrimaryKeyIndexSourceMeta, /// Resolved index-file path (globally unique; the scorer's preload key). pub path: String, pub file_size: u64, @@ -55,7 +55,7 @@ pub(crate) struct BucketAnnSegment { impl BucketAnnSegment { /// Build a segment with dummy index-file fields for tests that exercise only /// `source_meta`-driven logic. - pub(crate) fn for_test(source_meta: PkVectorSourceMeta) -> Self { + pub(crate) fn for_test(source_meta: PrimaryKeyIndexSourceMeta) -> Self { Self { source_meta, path: "seg".to_string(), @@ -714,18 +714,18 @@ pub(crate) async fn bucket_search_batch( #[cfg(test)] mod tests { use super::*; - use crate::spec::PkVectorSourceFile; + use crate::spec::PrimaryKeyIndexSourceFile; use crate::vindex::pkvector::ann::PkVectorAnnSearcher; use crate::vindex::pkvector::exact::exact_search; use crate::vindex::pkvector::reader::test_support::ArrayReader; use roaring::RoaringBitmap; - fn meta(files: &[(&str, i64)]) -> PkVectorSourceMeta { - PkVectorSourceMeta::new( + fn meta(files: &[(&str, i64)]) -> PrimaryKeyIndexSourceMeta { + PrimaryKeyIndexSourceMeta::new( 1, files .iter() - .map(|(n, r)| PkVectorSourceFile::new((*n).into(), *r).unwrap()) + .map(|(n, r)| PrimaryKeyIndexSourceFile::new((*n).into(), *r).unwrap()) .collect(), ) .unwrap() diff --git a/crates/paimon/tests/pk_vector_baseline_test.rs b/crates/paimon/tests/pk_vector_baseline_test.rs index 1328cddd..058e499f 100644 --- a/crates/paimon/tests/pk_vector_baseline_test.rs +++ b/crates/paimon/tests/pk_vector_baseline_test.rs @@ -32,11 +32,11 @@ //! binaries and nothing skipped. //! //! Two constraints the primary-key read path enforces are satisfied by hand -//! (mirroring Java `PrimaryKeyIndexSourcePolicy` and `PkVectorSourceMeta`): +//! (mirroring Java `PrimaryKeyIndexSourcePolicy` and `PrimaryKeyIndexSourceMeta`): //! 1. Only a compacted (`file_source == COMPACT`), non-level-0 data file backs //! the index, so the written file's meta is cloned with `level = 1` and //! `file_source = Some(1)`. -//! 2. `GlobalIndexMeta.source_meta` must be the Java `PkVectorSourceMeta` frame +//! 2. `GlobalIndexMeta.source_meta` must be the Java `PrimaryKeyIndexSourceMeta` frame //! (big-endian ints/longs, `writeUTF` file names), assembled below. //! //! Determinism: every fixture uses `nlist = 1`, so the single IVF inverted list @@ -188,7 +188,7 @@ fn data_batch(vectors: &[[f32; DIM]]) -> RecordBatch { /// Encode one Java `DataOutput#writeUTF` value (u16-BE byte length + modified /// UTF-8). ASCII file names are the common case; multibyte handling mirrors the -/// round-trip helper in `PkVectorSourceMeta`'s own tests. +/// round-trip helper in `PrimaryKeyIndexSourceMeta`'s own tests. fn java_write_utf(s: &str) -> Vec { let mut body = Vec::new(); for c in s.encode_utf16() { @@ -208,8 +208,8 @@ fn java_write_utf(s: &str) -> Vec { out } -/// Assemble the `_SOURCE_META` frame the way Java `PkVectorSourceMeta` writes it -/// and `PkVectorSourceMeta::deserialize` expects: `i32-BE version=1`, `i32-BE +/// Assemble the `_SOURCE_META` frame the way Java `PrimaryKeyIndexSourceMeta` writes it +/// and `PrimaryKeyIndexSourceMeta::deserialize` expects: `i32-BE version=1`, `i32-BE /// data_level`, `i32-BE count`, then per source file a `writeUTF` name and an /// `i64-BE` row count. No trailing bytes. Source files are listed in global /// ordinal order. @@ -411,7 +411,7 @@ async fn build_table_with_first_row_id( assert_segment_reads_back(&bytes, query, &analytic_topk(query, vectors, k)); } - // Constraint 2: GlobalIndexMeta.source_meta must be the Java PkVectorSourceMeta + // Constraint 2: GlobalIndexMeta.source_meta must be the Java PrimaryKeyIndexSourceMeta // frame naming the backing data file(s) in ordinal order. Here one source file // owns all rows, so ordinal == physical position. let vector_field_id = schema diff --git a/crates/paimon/tests/pk_vector_java_fixture_test.rs b/crates/paimon/tests/pk_vector_java_fixture_test.rs index 579a58eb..f64584bb 100644 --- a/crates/paimon/tests/pk_vector_java_fixture_test.rs +++ b/crates/paimon/tests/pk_vector_java_fixture_test.rs @@ -22,7 +22,7 @@ //! Java writer with the production `ivf-flat` primary-key vector indexer, and //! asserts Rust reads + searches it correctly. It validates the cross-language //! contract at the table-metadata layer (snapshot / manifest / data-file -//! `level`+`file_source` / `GlobalIndexMeta` / the `PkVectorSourceMeta` frame); +//! `level`+`file_source` / `GlobalIndexMeta` / the `PrimaryKeyIndexSourceMeta` frame); //! the ANN segment interior is the same native core on both sides //! (`paimon-vindex-core` <-> its JNI wrapper), so a divergence here would be a //! real metadata-parse bug, not a fixture artifact.