Skip to main content

common_datasource/
packed_snapshot.rs

1// Copyright 2023 Greptime Team
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Shared on-disk index for one packed snapshot schema chunk.
16
17use 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/// Versioned index. Object paths are generated direct children of the chunk.
27#[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/// One complete independent Parquet stream, including its footer.
51#[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    /// Validates all references and complete, nonoverlapping object coverage.
63    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; // Checked against object length above.
118            }
119            if end != object.length {
120                return Err(invalid("object ranges do not cover the full object"));
121            }
122        }
123        Ok(())
124    }
125
126    /// Requires exactly the data-bearing tables selected by snapshot DDL.
127    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}