Skip to main content

log_store/object_store_wal/
catalog.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//! In-memory index over the footers of the objects of one WAL prefix.
16
17use std::collections::BTreeMap;
18use std::ops::Bound::{Excluded, Unbounded};
19
20use snafu::{OptionExt, ensure};
21use store_api::storage::RegionId;
22
23use crate::error::{
24    CorruptedWalObjectSnafu, InvalidWalEntryRangeSnafu, Result, WalObjectSequenceExhaustedSnafu,
25};
26use crate::object_store_wal::batch::{OBJECT_SEQ_LIMIT, sequence_floor};
27use crate::object_store_wal::format::FooterEntry;
28
29/// Indexes objects by sequence and, per region, the objects that hold entries
30/// of that region.
31#[derive(Debug, Default)]
32pub(crate) struct ObjectCatalog {
33    objects: BTreeMap<u64, Vec<FooterEntry>>,
34    regions: BTreeMap<RegionId, BTreeMap<u64, FooterEntry>>,
35}
36
37impl ObjectCatalog {
38    /// Indexes the footer of the object `object_seq`. Objects may be inserted
39    /// in any order, which lets recovery index them as it discovers them.
40    /// Inserting a sequence that is already indexed is rejected, whether or not
41    /// the footer matches the indexed one. An object without segments holds no
42    /// entries but still takes its sequence.
43    pub(crate) fn insert_object(
44        &mut self,
45        object_seq: u64,
46        mut footer: Vec<FooterEntry>,
47    ) -> Result<()> {
48        footer.sort_unstable_by_key(|entry| entry.region_id);
49        for entries in footer.windows(2) {
50            ensure!(
51                entries[0].region_id != entries[1].region_id,
52                CorruptedWalObjectSnafu {
53                    reason: format!(
54                        "object {object_seq} has duplicate footer entries for region {}",
55                        entries[0].region_id
56                    ),
57                }
58            );
59        }
60        for entry in &footer {
61            ensure!(
62                entry.entry_count > 0 && entry.min_entry_id <= entry.max_entry_id,
63                CorruptedWalObjectSnafu {
64                    reason: format!(
65                        "object {object_seq} has invalid entry range {}..={} with {} entries for region {}",
66                        entry.min_entry_id, entry.max_entry_id, entry.entry_count, entry.region_id
67                    ),
68                }
69            );
70        }
71
72        // Object sequences are unique: recovery indexes every listed key once, and
73        // accepting a repeated insertion would hide a caller that lost track of it.
74        // Retrying an identical write stays an object store concern.
75        ensure!(
76            !self.objects.contains_key(&object_seq),
77            CorruptedWalObjectSnafu {
78                reason: format!("object {object_seq} is already indexed"),
79            }
80        );
81
82        // Validate every region before mutating either index so insertion is atomic.
83        for entry in &footer {
84            let Some(region_objects) = self.regions.get(&entry.region_id) else {
85                continue;
86            };
87            if let Some((&previous_seq, previous)) = region_objects.range(..object_seq).next_back()
88            {
89                ensure!(
90                    previous.max_entry_id < entry.min_entry_id,
91                    CorruptedWalObjectSnafu {
92                        reason: out_of_order_reason(
93                            entry.region_id,
94                            previous_seq,
95                            previous.max_entry_id,
96                            object_seq,
97                            entry.min_entry_id
98                        ),
99                    }
100                );
101            }
102            if let Some((&next_seq, next)) = region_objects
103                .range((Excluded(object_seq), Unbounded))
104                .next()
105            {
106                ensure!(
107                    entry.max_entry_id < next.min_entry_id,
108                    CorruptedWalObjectSnafu {
109                        reason: out_of_order_reason(
110                            entry.region_id,
111                            object_seq,
112                            entry.max_entry_id,
113                            next_seq,
114                            next.min_entry_id
115                        ),
116                    }
117                );
118            }
119        }
120
121        for entry in &footer {
122            self.regions
123                .entry(entry.region_id)
124                .or_default()
125                .insert(object_seq, entry.clone());
126        }
127        self.objects.insert(object_seq, footer);
128        Ok(())
129    }
130
131    /// Returns the objects that hold entries of `region_id` overlapping
132    /// `start_entry_id..=end_entry_id`, ordered by object sequence.
133    pub(crate) fn objects_for_entry_range(
134        &self,
135        region_id: RegionId,
136        start_entry_id: u64,
137        end_entry_id: u64,
138    ) -> Result<Vec<(u64, &FooterEntry)>> {
139        ensure!(
140            start_entry_id <= end_entry_id,
141            InvalidWalEntryRangeSnafu {
142                region_id,
143                start_entry_id,
144                end_entry_id,
145            }
146        );
147
148        let Some(objects) = self.regions.get(&region_id) else {
149            return Ok(Vec::new());
150        };
151        Ok(objects
152            .iter()
153            .filter(|(_, entry)| {
154                entry.max_entry_id >= start_entry_id && entry.min_entry_id <= end_entry_id
155            })
156            .map(|(&object_seq, entry)| (object_seq, entry))
157            .collect())
158    }
159
160    /// Returns the largest entry id indexed for `region_id`.
161    pub(crate) fn region_max_entry_id(&self, region_id: RegionId) -> Option<u64> {
162        self.regions
163            .get(&region_id)?
164            .last_key_value()
165            .map(|(_, entry)| entry.max_entry_id)
166    }
167
168    /// Returns the sequence to assign to the next object written after recovery.
169    ///
170    /// An empty catalog starts at zero, so the first object of a prefix always
171    /// takes sequence zero. Otherwise the sequence continues after the largest
172    /// indexed one, which recovery discovers regardless of insertion order, and
173    /// is raised further when the largest entry id of a region lies at or above
174    /// the ids that sequence would assign: ids assigned under the earlier
175    /// contiguous scheme carry no object information, and every new id of a
176    /// region must be greater than every id it already has. A sequence at or
177    /// above [`OBJECT_SEQ_LIMIT`] does not fit an entry id and is rejected.
178    pub(crate) fn next_object_seq(&self) -> Result<u64> {
179        let after_last = match self.objects.last_key_value() {
180            None => 0,
181            Some((&last_object_seq, _)) => last_object_seq
182                .checked_add(1)
183                .context(WalObjectSequenceExhaustedSnafu { last_object_seq })?,
184        };
185        let floor = self
186            .regions
187            .keys()
188            .filter_map(|region_id| self.region_max_entry_id(*region_id))
189            .map(sequence_floor)
190            .max()
191            .unwrap_or(0);
192        let next_object_seq = after_last.max(floor);
193        ensure!(
194            next_object_seq < OBJECT_SEQ_LIMIT,
195            WalObjectSequenceExhaustedSnafu {
196                last_object_seq: next_object_seq - 1,
197            }
198        );
199        Ok(next_object_seq)
200    }
201
202    /// Iterates over the indexed objects ordered by object sequence.
203    pub(crate) fn objects_in_order(&self) -> impl Iterator<Item = (u64, &[FooterEntry])> + '_ {
204        self.objects
205            .iter()
206            .map(|(&object_seq, footer)| (object_seq, footer.as_slice()))
207    }
208}
209
210fn out_of_order_reason(
211    region_id: RegionId,
212    lower_object_seq: u64,
213    lower_max_entry_id: u64,
214    upper_object_seq: u64,
215    upper_min_entry_id: u64,
216) -> String {
217    format!(
218        "entry ranges of region {region_id} are not strictly increasing, object {lower_object_seq} ends at {lower_max_entry_id}, object {upper_object_seq} starts at {upper_min_entry_id}"
219    )
220}
221
222#[cfg(test)]
223mod tests {
224    use super::*;
225    use crate::error::Error;
226    use crate::object_store_wal::batch::entry_id;
227
228    #[test]
229    fn test_catalog_indexes_objects_and_queries_ranges() {
230        let region_one = RegionId::new(1, 1);
231        let region_two = RegionId::new(2, 1);
232        let mut catalog = ObjectCatalog::default();
233
234        // Recovery may discover objects out of order.
235        catalog
236            .insert_object(
237                2,
238                vec![
239                    footer_entry(region_one, 4, 6),
240                    footer_entry(region_two, 8, 9),
241                ],
242            )
243            .unwrap();
244        catalog
245            .insert_object(1, vec![footer_entry(region_one, 1, 3)])
246            .unwrap();
247        catalog
248            .insert_object(4, vec![footer_entry(region_one, 10, 12)])
249            .unwrap();
250
251        assert_eq!(Some(12), catalog.region_max_entry_id(region_one));
252        assert_eq!(Some(9), catalog.region_max_entry_id(region_two));
253        assert_eq!(None, catalog.region_max_entry_id(RegionId::new(3, 1)));
254
255        let objects = catalog.objects_for_entry_range(region_one, 3, 10).unwrap();
256        assert_eq!(vec![1, 2, 4], object_seqs(&objects));
257        assert_eq!(
258            vec![1, 2, 4],
259            catalog
260                .objects_in_order()
261                .map(|(object_seq, _)| object_seq)
262                .collect::<Vec<_>>()
263        );
264    }
265
266    #[test]
267    fn test_catalog_rejects_duplicate_object_sequences() {
268        let region_id = RegionId::new(1, 1);
269        let mut catalog = ObjectCatalog::default();
270        let first = footer_entry(region_id, 1, 2);
271        let second = footer_entry(RegionId::new(2, 1), 4, 5);
272
273        catalog
274            .insert_object(1, vec![first.clone(), second.clone()])
275            .unwrap();
276
277        // An identical footer is rejected just like a conflicting one.
278        assert_corrupted(
279            catalog.insert_object(1, vec![second, first.clone()]),
280            "object 1 is already indexed",
281        );
282
283        let mut conflicting = first;
284        conflicting.segment_offset += 1;
285        assert_corrupted(
286            catalog.insert_object(1, vec![conflicting]),
287            "object 1 is already indexed",
288        );
289        assert_eq!(1, catalog.objects_in_order().count());
290    }
291
292    #[test]
293    fn test_catalog_resumes_object_sequence_after_recovery() {
294        let region_id = RegionId::new(1, 1);
295        let mut catalog = ObjectCatalog::default();
296        assert_eq!(0, catalog.next_object_seq().unwrap());
297
298        // Recovery may discover objects out of order.
299        catalog
300            .insert_object(4, vec![footer_entry(region_id, 10, 12)])
301            .unwrap();
302        catalog
303            .insert_object(1, vec![footer_entry(region_id, 1, 3)])
304            .unwrap();
305
306        assert_eq!(5, catalog.next_object_seq().unwrap());
307
308        // An object without segments takes its sequence all the same.
309        catalog.insert_object(7, Vec::new()).unwrap();
310        assert_eq!(8, catalog.next_object_seq().unwrap());
311        assert_eq!(Some(12), catalog.region_max_entry_id(region_id));
312    }
313
314    #[test]
315    fn test_catalog_raises_object_sequence_above_existing_entry_ids() {
316        let region_one = RegionId::new(1, 1);
317        let region_two = RegionId::new(1, 2);
318        let mut catalog = ObjectCatalog::default();
319
320        // Ids that fit below the ids of the next sequence leave it alone.
321        catalog
322            .insert_object(0, vec![footer_entry(region_one, 1, 3)])
323            .unwrap();
324        catalog
325            .insert_object(
326                1,
327                vec![footer_entry(region_two, entry_id(1, 1), entry_id(1, 2))],
328            )
329            .unwrap();
330        assert_eq!(2, catalog.next_object_seq().unwrap());
331
332        // A contiguous id past them names a later object: the sequence
333        // resumes above it, whichever region holds it.
334        catalog
335            .insert_object(2, vec![footer_entry(region_one, 5_000_000, 5_000_000)])
336            .unwrap();
337        assert_eq!(5, catalog.next_object_seq().unwrap());
338        catalog
339            .insert_object(
340                3,
341                vec![footer_entry(region_two, entry_id(7, 4), entry_id(7, 4))],
342            )
343            .unwrap();
344        assert_eq!(8, catalog.next_object_seq().unwrap());
345    }
346
347    #[test]
348    fn test_catalog_rejects_exhausted_object_sequence() {
349        let region_id = RegionId::new(1, 1);
350        let assert_exhausted = |catalog: &ObjectCatalog| {
351            let error = catalog.next_object_seq().unwrap_err();
352            assert!(
353                error.to_string().contains("object sequence is exhausted"),
354                "unexpected error: {error}"
355            );
356        };
357
358        // The last sequence that fits an entry id is indexed.
359        let mut catalog = ObjectCatalog::default();
360        catalog
361            .insert_object(OBJECT_SEQ_LIMIT - 2, vec![footer_entry(region_id, 1, 2)])
362            .unwrap();
363        assert_eq!(OBJECT_SEQ_LIMIT - 1, catalog.next_object_seq().unwrap());
364        catalog
365            .insert_object(OBJECT_SEQ_LIMIT - 1, vec![footer_entry(region_id, 3, 4)])
366            .unwrap();
367        assert_exhausted(&catalog);
368
369        // A sequence that does not fit was written by an earlier scheme.
370        let mut catalog = ObjectCatalog::default();
371        catalog
372            .insert_object(u64::MAX, vec![footer_entry(region_id, 1, 2)])
373            .unwrap();
374        assert_exhausted(&catalog);
375
376        // An entry id that leaves no sequence above it.
377        let mut catalog = ObjectCatalog::default();
378        catalog
379            .insert_object(0, vec![footer_entry(region_id, u64::MAX, u64::MAX)])
380            .unwrap();
381        assert_exhausted(&catalog);
382    }
383
384    #[test]
385    fn test_catalog_rejects_overlapping_or_reversed_region_ranges() {
386        let region_id = RegionId::new(1, 1);
387        let mut catalog = ObjectCatalog::default();
388        catalog
389            .insert_object(2, vec![footer_entry(region_id, 10, 20)])
390            .unwrap();
391
392        assert_corrupted(
393            catalog.insert_object(3, vec![footer_entry(region_id, 20, 30)]),
394            "are not strictly increasing",
395        );
396        assert_corrupted(
397            catalog.insert_object(3, vec![footer_entry(region_id, 5, 9)]),
398            "are not strictly increasing",
399        );
400        assert_corrupted(
401            catalog.insert_object(1, vec![footer_entry(region_id, 15, 19)]),
402            "are not strictly increasing",
403        );
404        assert_eq!(1, catalog.objects_in_order().count());
405    }
406
407    #[test]
408    fn test_catalog_rejects_duplicate_region_and_invalid_ranges() {
409        let region_id = RegionId::new(1, 1);
410        let mut catalog = ObjectCatalog::default();
411        assert_corrupted(
412            catalog.insert_object(
413                1,
414                vec![footer_entry(region_id, 1, 1), footer_entry(region_id, 2, 2)],
415            ),
416            "duplicate footer entries",
417        );
418        let mut invalid = footer_entry(region_id, 2, 1);
419        invalid.entry_count = 0;
420        assert_corrupted(
421            catalog.insert_object(1, vec![invalid]),
422            "has invalid entry range 2..=1",
423        );
424
425        let error = catalog
426            .objects_for_entry_range(region_id, 2, 1)
427            .unwrap_err();
428        assert!(
429            matches!(error, Error::InvalidWalEntryRange { start_entry_id, end_entry_id, .. } if start_entry_id == 2 && end_entry_id == 1),
430            "unexpected error: {error:?}"
431        );
432    }
433
434    fn footer_entry(region_id: RegionId, min_entry_id: u64, max_entry_id: u64) -> FooterEntry {
435        FooterEntry {
436            region_id,
437            min_entry_id,
438            max_entry_id,
439            entry_count: max_entry_id
440                .checked_sub(min_entry_id)
441                .and_then(|count| count.checked_add(1))
442                .unwrap_or(0) as u32,
443            segment_offset: min_entry_id.wrapping_mul(100),
444            segment_len: 100,
445            segment_crc32: min_entry_id as u32,
446        }
447    }
448
449    fn object_seqs(objects: &[(u64, &FooterEntry)]) -> Vec<u64> {
450        objects.iter().map(|(object_seq, _)| *object_seq).collect()
451    }
452
453    fn assert_corrupted(result: Result<()>, expected_reason: &str) {
454        match result {
455            Err(Error::CorruptedWalObject { reason, .. }) => assert!(
456                reason.contains(expected_reason),
457                "expected reason to contain {expected_reason:?}, actual {reason:?}"
458            ),
459            other => panic!("expected a corrupted object error, actual {other:?}"),
460        }
461    }
462}