1use ailake_core::{AilakeError, AilakeResult};
25use bytes::Bytes;
26use serde::{Deserialize, Serialize};
27
28pub const BLOB_TYPE_VECTOR_STATS: &str = "ailake-vector-stats-v1";
32pub const BLOB_TYPE_BM25_BLOOM: &str = "ailake-bm25-bloom-v1";
34
35#[derive(Debug, Clone, Serialize, Deserialize)]
39pub struct VectorStatEntry {
40 pub path: String,
42 pub centroid: Vec<f32>,
44 pub radius: f32,
46}
47
48#[derive(Debug, Clone, Serialize, Deserialize)]
50pub struct BM25BloomEntry {
51 pub path: String,
53 pub bloom_bytes: Vec<u8>,
55}
56
57const PUFFIN_MAGIC: &[u8] = b"PFAc";
60
61pub struct PuffinStatsResult {
63 pub bytes: Bytes,
64 pub footer_size: usize,
66 pub vector_stats_blob: (u64, u64),
68 pub bm25_bloom_blob: Option<(u64, u64)>,
70}
71
72pub struct AilakePuffinWriter;
73
74impl AilakePuffinWriter {
75 pub fn write_stats(
80 vector_stats: &[VectorStatEntry],
81 bm25_blooms: &[BM25BloomEntry],
82 snap_id: i64,
83 ) -> AilakeResult<PuffinStatsResult> {
84 let vec_blob =
86 bincode::serialize(vector_stats).map_err(|e| AilakeError::Bincode(e.to_string()))?;
87
88 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 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 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
151pub 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 fn footer(&self) -> AilakeResult<serde_json::Value> {
164 let n = self.data.len();
165 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 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 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#[cfg(test)]
224mod tests {
225 use super::*;
226
227 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 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 assert!(!recovered_blooms[0].bloom_bytes.is_empty());
305 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 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}