common_datasource/
packed_snapshot.rs1use std::collections::{HashMap, HashSet};
18
19use serde::{Deserialize, Serialize};
20
21use crate::error::{InvalidPackedSnapshotSnafu, Result};
22
23pub const PACK_INDEX_FILE: &str = "pack-index.json";
24pub const PACKED_LAYOUT: &str = "metric-parquet-packs";
25
26#[derive(Debug, Clone, Serialize, Deserialize)]
28#[serde(deny_unknown_fields)]
29pub struct PackIndex {
30 pub version: u32,
31 pub objects: Vec<PackObject>,
32 pub tables: Vec<PackTable>,
33}
34
35#[derive(Debug, Clone, Serialize, Deserialize)]
36#[serde(deny_unknown_fields)]
37pub struct PackObject {
38 pub path: String,
39 pub kind: ObjectKind,
40 pub length: u64,
41}
42
43#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
44#[serde(rename_all = "lowercase")]
45pub enum ObjectKind {
46 Pack,
47 Parquet,
48}
49
50#[derive(Debug, Clone, Serialize, Deserialize)]
52#[serde(deny_unknown_fields)]
53pub struct PackTable {
54 pub table_name: String,
55 pub object: String,
56 pub offset: u64,
57 pub length: u64,
58 pub row_count: u64,
59}
60
61impl PackIndex {
62 pub fn validate(&self) -> Result<()> {
64 let invalid = |reason: &str| InvalidPackedSnapshotSnafu { reason }.build();
65 if self.version != 1 {
66 return Err(invalid("unsupported index version"));
67 }
68 let mut objects = HashMap::new();
69 for object in &self.objects {
70 let (prefix, suffix) = match object.kind {
71 ObjectKind::Pack => ("pack-", ".bin"),
72 ObjectKind::Parquet => ("table-", ".parquet"),
73 };
74 let generated = object
75 .path
76 .strip_prefix(prefix)
77 .and_then(|s| s.strip_suffix(suffix))
78 .is_some_and(|s| !s.is_empty() && s.bytes().all(|b| b.is_ascii_digit()));
79 if !generated
80 || object.length == 0
81 || objects.insert(object.path.as_str(), object).is_some()
82 {
83 return Err(invalid("invalid or duplicate object path/length"));
84 }
85 }
86 let mut names = HashSet::new();
87 let mut ranges: HashMap<&str, Vec<&PackTable>> = HashMap::new();
88 for table in &self.tables {
89 if table.table_name.is_empty() || !names.insert(table.table_name.as_str()) {
90 return Err(invalid("empty or duplicate table name"));
91 }
92 let object = objects
93 .get(table.object.as_str())
94 .ok_or_else(|| invalid("unknown object reference"))?;
95 let end = table
96 .offset
97 .checked_add(table.length)
98 .ok_or_else(|| invalid("table range overflow"))?;
99 if table.length < 12 || end > object.length {
100 return Err(invalid("invalid Parquet stream range"));
101 }
102 ranges.entry(&table.object).or_default().push(table);
103 }
104 for object in &self.objects {
105 let entries = ranges
106 .get_mut(object.path.as_str())
107 .ok_or_else(|| invalid("unreferenced object"))?;
108 entries.sort_unstable_by_key(|t| t.offset);
109 if object.kind == ObjectKind::Parquet && entries.len() != 1 {
110 return Err(invalid("standalone object must contain one table"));
111 }
112 let mut end = 0;
113 for entry in entries {
114 if entry.offset != end {
115 return Err(invalid("object ranges have a gap or overlap"));
116 }
117 end += entry.length; }
119 if end != object.length {
120 return Err(invalid("object ranges do not cover the full object"));
121 }
122 }
123 Ok(())
124 }
125
126 pub fn validate_membership<'a>(&self, tables: impl IntoIterator<Item = &'a str>) -> Result<()> {
128 self.validate()?;
129 let expected: HashSet<_> = tables.into_iter().collect();
130 let actual: HashSet<_> = self.tables.iter().map(|t| t.table_name.as_str()).collect();
131 if expected != actual {
132 return InvalidPackedSnapshotSnafu {
133 reason: "index table membership differs from snapshot DDL",
134 }
135 .fail();
136 }
137 Ok(())
138 }
139}
140
141#[cfg(test)]
142mod tests {
143 use super::*;
144
145 fn index() -> PackIndex {
146 PackIndex {
147 version: 1,
148 objects: vec![PackObject {
149 path: "pack-000000.bin".into(),
150 kind: ObjectKind::Pack,
151 length: 24,
152 }],
153 tables: ["literal.name", "other"]
154 .into_iter()
155 .enumerate()
156 .map(|(i, name)| PackTable {
157 table_name: name.into(),
158 object: "pack-000000.bin".into(),
159 offset: i as u64 * 12,
160 length: 12,
161 row_count: 0,
162 })
163 .collect(),
164 }
165 }
166
167 #[test]
168 fn validates_complete_membership_and_ranges() {
169 index()
170 .validate_membership(["other", "literal.name"])
171 .unwrap();
172 assert!(index().validate_membership(["other"]).is_err());
173 for mutate in [
174 |i: &mut PackIndex| i.version = 2,
175 |i: &mut PackIndex| i.objects[0].path = "../pack-0.bin".into(),
176 |i: &mut PackIndex| i.objects.push(i.objects[0].clone()),
177 |i: &mut PackIndex| i.tables[1].table_name = "literal.name".into(),
178 |i: &mut PackIndex| i.tables[1].object = "missing".into(),
179 |i: &mut PackIndex| i.tables[1].offset = 11,
180 |i: &mut PackIndex| i.tables[1].length = u64::MAX,
181 |i: &mut PackIndex| i.tables.pop().map(|_| ()).unwrap(),
182 ] {
183 let mut bad = index();
184 mutate(&mut bad);
185 assert!(bad.validate().is_err(), "{bad:?}");
186 }
187 }
188}