Skip to main content

ailake_catalog/
puffin.rs

1// SPDX-License-Identifier: MIT OR Apache-2.0
2//! AI-Lake Puffin statistics files (Phase F).
3//!
4//! Extends the existing Puffin DV support (`delete.rs`) with two new blob types:
5//!
6//! - `ailake-vector-stats-v1`: all data-file centroids + radii for a snapshot,
7//!   stored as bincode-encoded `Vec<VectorStatEntry>`. Enables readers to fetch
8//!   geometric pruning data with a single GET instead of scanning manifest KV.
9//!
10//! - `ailake-bm25-bloom-v1`: per-file Bloom filters over BM25 indexed terms,
11//!   stored as bincode-encoded `Vec<BM25BloomEntry>`. Readers skip files where
12//!   no query term passes the filter (zero false negatives).
13//!
14//! Puffin file layout (reuses the Iceberg Puffin spec §4 format):
15//! ```text
16//! PFAc  (4 bytes magic)
17//! [blob bytes — vector stats]
18//! [blob bytes — BM25 bloom, optional]
19//! footer_JSON
20//! footer_len (4 bytes LE)
21//! PFAc  (4 bytes magic)
22//! ```
23
24use ailake_core::{AilakeError, AilakeResult};
25use bytes::Bytes;
26use serde::{Deserialize, Serialize};
27
28// ── Blob type tags ────────────────────────────────────────────────────────────
29
30/// Puffin blob type for per-snapshot vector statistics (centroid + radius per file).
31pub const BLOB_TYPE_VECTOR_STATS: &str = "ailake-vector-stats-v1";
32/// Puffin blob type for per-file BM25 Bloom filters.
33pub const BLOB_TYPE_BM25_BLOOM: &str = "ailake-bm25-bloom-v1";
34
35// ── Blob entry types ──────────────────────────────────────────────────────────
36
37/// Per-file vector statistics stored in the `ailake-vector-stats-v1` blob.
38#[derive(Debug, Clone, Serialize, Deserialize)]
39pub struct VectorStatEntry {
40    /// Data file path (relative within the table root).
41    pub path: String,
42    /// Centroid vector (F32, length = column dim).
43    pub centroid: Vec<f32>,
44    /// Maximum distance from any vector in the file to the centroid.
45    pub radius: f32,
46}
47
48/// Per-file BM25 Bloom filter stored in the `ailake-bm25-bloom-v1` blob.
49#[derive(Debug, Clone, Serialize, Deserialize)]
50pub struct BM25BloomEntry {
51    /// Data file path (relative within the table root).
52    pub path: String,
53    /// Serialized `BloomFilter::to_bytes()` output.
54    pub bloom_bytes: Vec<u8>,
55}
56
57// ── Puffin writer ─────────────────────────────────────────────────────────────
58
59const PUFFIN_MAGIC: &[u8] = b"PFAc";
60
61/// Result of `AilakePuffinWriter::write_stats`.
62pub struct PuffinStatsResult {
63    pub bytes: Bytes,
64    /// Byte length of the Puffin footer JSON (for `file-footer-size-in-bytes`).
65    pub footer_size: usize,
66    /// (offset, length) of the `ailake-vector-stats-v1` blob.
67    pub vector_stats_blob: (u64, u64),
68    /// (offset, length) of the `ailake-bm25-bloom-v1` blob, if present.
69    pub bm25_bloom_blob: Option<(u64, u64)>,
70}
71
72pub struct AilakePuffinWriter;
73
74impl AilakePuffinWriter {
75    /// Build a Puffin stats file containing vector-stats and optional BM25 bloom blobs.
76    ///
77    /// Returns `PuffinStatsResult` with the raw bytes and blob-location metadata,
78    /// suitable for constructing `IcebergStatisticsRef`.
79    pub fn write_stats(
80        vector_stats: &[VectorStatEntry],
81        bm25_blooms: &[BM25BloomEntry],
82        snap_id: i64,
83    ) -> AilakeResult<PuffinStatsResult> {
84        // Blob 1: vector stats (bincode, no compression — small random-access data)
85        let vec_blob =
86            bincode::serialize(vector_stats).map_err(|e| AilakeError::Bincode(e.to_string()))?;
87
88        // Blob 2: BM25 bloom (optional)
89        let bm25_blob: Option<Vec<u8>> = if !bm25_blooms.is_empty() {
90            Some(bincode::serialize(bm25_blooms).map_err(|e| AilakeError::Bincode(e.to_string()))?)
91        } else {
92            None
93        };
94
95        // Compute blob offsets (magic is 4 bytes at position 0)
96        let vec_offset = PUFFIN_MAGIC.len() as u64;
97        let vec_len = vec_blob.len() as u64;
98
99        let bm25_offset = vec_offset + vec_len;
100        let bm25_len = bm25_blob.as_ref().map_or(0, |b| b.len() as u64);
101
102        // Puffin footer JSON
103        let mut blobs_json = vec![serde_json::json!({
104            "type": BLOB_TYPE_VECTOR_STATS,
105            "snapshot-id": snap_id,
106            "sequence-number": 0,
107            "offset": vec_offset,
108            "length": vec_len,
109            "fields": []
110        })];
111        if bm25_blob.is_some() {
112            blobs_json.push(serde_json::json!({
113                "type": BLOB_TYPE_BM25_BLOOM,
114                "snapshot-id": snap_id,
115                "sequence-number": 0,
116                "offset": bm25_offset,
117                "length": bm25_len,
118                "fields": []
119            }));
120        }
121        let footer_json = serde_json::json!({
122            "blobs": blobs_json,
123            "properties": { "ailake.stats-version": "1" }
124        })
125        .to_string();
126        let footer_bytes = footer_json.as_bytes();
127        let footer_size = footer_bytes.len();
128        let footer_len_le = (footer_size as u32).to_le_bytes();
129
130        let mut out = Vec::with_capacity(
131            PUFFIN_MAGIC.len() * 2 + vec_blob.len() + bm25_len as usize + footer_size + 4,
132        );
133        out.extend_from_slice(PUFFIN_MAGIC);
134        out.extend_from_slice(&vec_blob);
135        if let Some(ref b) = bm25_blob {
136            out.extend_from_slice(b);
137        }
138        out.extend_from_slice(footer_bytes);
139        out.extend_from_slice(&footer_len_le);
140        out.extend_from_slice(PUFFIN_MAGIC);
141
142        Ok(PuffinStatsResult {
143            bytes: Bytes::from(out),
144            footer_size,
145            vector_stats_blob: (vec_offset, vec_len),
146            bm25_bloom_blob: bm25_blob.map(|_| (bm25_offset, bm25_len)),
147        })
148    }
149}
150
151// ── Puffin reader ─────────────────────────────────────────────────────────────
152
153pub struct AilakePuffinReader<'a> {
154    data: &'a [u8],
155}
156
157impl<'a> AilakePuffinReader<'a> {
158    pub fn new(data: &'a [u8]) -> Self {
159        Self { data }
160    }
161
162    /// Parse the Puffin footer JSON.
163    fn footer(&self) -> AilakeResult<serde_json::Value> {
164        let n = self.data.len();
165        // Minimum: magic(4) + footer_len(4) + magic(4) = 12
166        if n < 12 {
167            return Err(AilakeError::Catalog("Puffin file too short".into()));
168        }
169        let footer_len = u32::from_le_bytes(self.data[n - 8..n - 4].try_into().unwrap()) as usize;
170        let footer_start = n - 8 - footer_len;
171        serde_json::from_slice(&self.data[footer_start..footer_start + footer_len])
172            .map_err(|e| AilakeError::Catalog(format!("Puffin footer parse: {e}")))
173    }
174
175    fn blob_slice(&self, blob: &serde_json::Value) -> AilakeResult<&[u8]> {
176        let offset = blob["offset"].as_u64().unwrap_or(0) as usize;
177        let length = blob["length"].as_u64().unwrap_or(0) as usize;
178        if offset + length > self.data.len() {
179            return Err(AilakeError::Catalog(
180                "Puffin blob offset out of range".into(),
181            ));
182        }
183        Ok(&self.data[offset..offset + length])
184    }
185
186    /// Read the `ailake-vector-stats-v1` blob. Returns empty `Vec` if not present.
187    pub fn read_vector_stats(&self) -> AilakeResult<Vec<VectorStatEntry>> {
188        let footer = self.footer()?;
189        let blobs = match footer["blobs"].as_array() {
190            Some(b) => b,
191            None => return Ok(vec![]),
192        };
193        for blob in blobs {
194            if blob["type"].as_str() == Some(BLOB_TYPE_VECTOR_STATS) {
195                let slice = self.blob_slice(blob)?;
196                return bincode::deserialize(slice)
197                    .map_err(|e| AilakeError::Bincode(e.to_string()));
198            }
199        }
200        Ok(vec![])
201    }
202
203    /// Read the `ailake-bm25-bloom-v1` blob. Returns empty `Vec` if not present.
204    pub fn read_bm25_blooms(&self) -> AilakeResult<Vec<BM25BloomEntry>> {
205        let footer = self.footer()?;
206        let blobs = match footer["blobs"].as_array() {
207            Some(b) => b,
208            None => return Ok(vec![]),
209        };
210        for blob in blobs {
211            if blob["type"].as_str() == Some(BLOB_TYPE_BM25_BLOOM) {
212                let slice = self.blob_slice(blob)?;
213                return bincode::deserialize(slice)
214                    .map_err(|e| AilakeError::Bincode(e.to_string()));
215            }
216        }
217        Ok(vec![])
218    }
219}
220
221// ── Tests ─────────────────────────────────────────────────────────────────────
222
223#[cfg(test)]
224mod tests {
225    use super::*;
226
227    // Inline FNV-64a for test-only bloom construction (avoids circular dep on ailake-query).
228    fn fnv64a(data: &[u8], seed: u64) -> u64 {
229        let mut h = seed ^ 14695981039346656037u64;
230        for &b in data {
231            h ^= b as u64;
232            h = h.wrapping_mul(1099511628211u64);
233        }
234        h
235    }
236
237    fn sample_vector_stats() -> Vec<VectorStatEntry> {
238        vec![
239            VectorStatEntry {
240                path: "data/part-00001.parquet".into(),
241                centroid: vec![0.1, 0.2, 0.3],
242                radius: 0.5,
243            },
244            VectorStatEntry {
245                path: "data/part-00002.parquet".into(),
246                centroid: vec![0.9, 0.8, 0.7],
247                radius: 0.3,
248            },
249        ]
250    }
251
252    fn sample_bloom() -> Vec<BM25BloomEntry> {
253        // Build minimal bloom bytes inline (format: u64_le(num_bits) || u64 words).
254        // Use same FNV double-hashing as BloomFilter but here inline for test isolation.
255        let num_bits: usize = 1024;
256        let mut words = vec![0u64; num_bits / 64];
257        for term in &["rust", "iceberg"] {
258            let h1 = fnv64a(term.as_bytes(), 0);
259            let h2 = fnv64a(term.as_bytes(), h1);
260            for k in 0..4u64 {
261                let bit = ((h1.wrapping_add(k.wrapping_mul(h2))) as usize) % num_bits;
262                words[bit / 64] |= 1u64 << (bit % 64);
263            }
264        }
265        let mut bytes = (num_bits as u64).to_le_bytes().to_vec();
266        for w in &words {
267            bytes.extend_from_slice(&w.to_le_bytes());
268        }
269        vec![BM25BloomEntry {
270            path: "data/part-00001.parquet".into(),
271            bloom_bytes: bytes,
272        }]
273    }
274
275    #[test]
276    fn vector_stats_roundtrip() {
277        let stats = sample_vector_stats();
278        let result = AilakePuffinWriter::write_stats(&stats, &[], 42).unwrap();
279        assert!(result.bytes.starts_with(PUFFIN_MAGIC));
280        assert!(result.bytes.ends_with(PUFFIN_MAGIC));
281        assert!(result.bm25_bloom_blob.is_none());
282
283        let reader = AilakePuffinReader::new(&result.bytes);
284        let recovered = reader.read_vector_stats().unwrap();
285        assert_eq!(recovered.len(), 2);
286        assert_eq!(recovered[0].path, "data/part-00001.parquet");
287        assert!((recovered[0].radius - 0.5).abs() < 1e-6);
288        assert_eq!(recovered[1].centroid, vec![0.9f32, 0.8, 0.7]);
289    }
290
291    #[test]
292    fn bm25_bloom_roundtrip() {
293        let stats = sample_vector_stats();
294        let blooms = sample_bloom();
295        let result = AilakePuffinWriter::write_stats(&stats, &blooms, 99).unwrap();
296        assert!(result.bm25_bloom_blob.is_some());
297
298        let reader = AilakePuffinReader::new(&result.bytes);
299        let recovered_blooms = reader.read_bm25_blooms().unwrap();
300        assert_eq!(recovered_blooms.len(), 1);
301        assert_eq!(recovered_blooms[0].path, "data/part-00001.parquet");
302
303        // Verify bloom bytes roundtrip intact (content-check done in ailake-query bloom tests).
304        assert!(!recovered_blooms[0].bloom_bytes.is_empty());
305        // num_bits header must parse
306        assert!(recovered_blooms[0].bloom_bytes.len() >= 8);
307        let nb = u64::from_le_bytes(recovered_blooms[0].bloom_bytes[..8].try_into().unwrap());
308        assert_eq!(nb, 1024);
309    }
310
311    #[test]
312    fn empty_bloom_produces_no_bloom_blob() {
313        let stats = sample_vector_stats();
314        let result = AilakePuffinWriter::write_stats(&stats, &[], 1).unwrap();
315        let reader = AilakePuffinReader::new(&result.bytes);
316        let blooms = reader.read_bm25_blooms().unwrap();
317        assert!(blooms.is_empty());
318    }
319
320    #[test]
321    fn footer_size_matches_actual() {
322        let stats = sample_vector_stats();
323        let result = AilakePuffinWriter::write_stats(&stats, &[], 7).unwrap();
324        // The declared footer_size must match what the reader parses.
325        let n = result.bytes.len();
326        let footer_len_from_file =
327            u32::from_le_bytes(result.bytes[n - 8..n - 4].try_into().unwrap()) as usize;
328        assert_eq!(footer_len_from_file, result.footer_size);
329    }
330}