1use 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#[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 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 ensure!(
76 !self.objects.contains_key(&object_seq),
77 CorruptedWalObjectSnafu {
78 reason: format!("object {object_seq} is already indexed"),
79 }
80 );
81
82 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 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(®ion_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 pub(crate) fn region_max_entry_id(&self, region_id: RegionId) -> Option<u64> {
162 self.regions
163 .get(®ion_id)?
164 .last_key_value()
165 .map(|(_, entry)| entry.max_entry_id)
166 }
167
168 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 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 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 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 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 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 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 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 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 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 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}