Skip to main content

log_store/object_store_wal/
batch.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//! Accumulation of admitted entries into the batch that becomes the next
16//! object, and the entry id scheme.
17
18use std::collections::HashMap;
19use std::time::Instant;
20
21use snafu::ensure;
22use store_api::logstore::EntryId;
23use store_api::logstore::entry::Entry;
24use store_api::storage::RegionId;
25
26use crate::error::{Result, WalEntryPositionExhaustedSnafu};
27
28/// Bits of an entry id that hold the position of the entry among the entries
29/// of its region inside the object; the remaining high bits hold the object
30/// sequence.
31pub(crate) const POSITION_BITS: u32 = 20;
32/// Positions are `1..POSITION_LIMIT`, so an object holds at most 2^20 - 1
33/// entries of one region and a batch seals before a region reaches the limit.
34/// This is a theoretical bound: an object is sealed by size long before.
35pub(crate) const POSITION_LIMIT: u64 = 1 << POSITION_BITS;
36/// Object sequences are `0..OBJECT_SEQ_LIMIT`, the 44 bits an entry id leaves
37/// above the position.
38pub(crate) const OBJECT_SEQ_LIMIT: u64 = 1 << (u64::BITS - POSITION_BITS);
39
40/// Returns the id of the entry at `position` among the entries of its region
41/// in the object `object_seq`.
42///
43/// Ids are object-sequence-major: `entry_id >> 20` names the object that holds
44/// the entry, and the low bits are its position in that object, which starts
45/// at one so that id zero, the watermark of a region without entries, is never
46/// assigned. A region's ids increase with the object sequence and have gaps
47/// wherever other regions or other positions took the sequence.
48pub fn entry_id(object_seq: u64, position: u64) -> EntryId {
49    debug_assert!(object_seq < OBJECT_SEQ_LIMIT && (1..POSITION_LIMIT).contains(&position));
50    (object_seq << POSITION_BITS) | position
51}
52
53/// Returns the smallest object sequence whose entry ids are all greater than
54/// `entry_id`. Zero names no entry, so it needs no floor.
55///
56/// An id assigned under the earlier contiguous scheme carries no object
57/// information, but the floor keeps every new id above it just the same.
58pub(crate) fn sequence_floor(entry_id: EntryId) -> u64 {
59    if entry_id == 0 {
60        0
61    } else {
62        (entry_id >> POSITION_BITS) + 1
63    }
64}
65
66/// Entries admitted since the last seal, together with the position of the
67/// last entry admitted per region.
68///
69/// Every entry is assigned its id at admission from the sequence the batch
70/// takes when it is sealed, so a batch that is rolled back and admitted again
71/// under the same sequence hands out the same ids.
72#[derive(Debug)]
73pub(crate) struct OpenBatch {
74    max_bytes: usize,
75    entries: Vec<Entry>,
76    estimated_bytes: usize,
77    /// When the first entry of the batch was admitted.
78    first_admitted_at: Option<Instant>,
79    positions: HashMap<RegionId, u64>,
80}
81
82impl OpenBatch {
83    pub(crate) fn new(max_bytes: usize) -> Self {
84        Self {
85            max_bytes,
86            entries: Vec::new(),
87            estimated_bytes: 0,
88            first_admitted_at: None,
89            positions: HashMap::new(),
90        }
91    }
92
93    /// Returns true when admitting `entries` would take a region past the
94    /// position range, so that `entries` need a batch of their own.
95    pub(crate) fn would_exhaust_positions(&self, entries: &[Entry]) -> bool {
96        self.check_positions(entries).is_err()
97    }
98
99    /// Admits `entries` into the batch that becomes the object `object_seq`,
100    /// assigning each the next position of its region, and returns the last id
101    /// assigned to every region in `entries`. Nothing is admitted when a
102    /// region would run past the position range.
103    pub(crate) fn admit(
104        &mut self,
105        object_seq: u64,
106        mut entries: Vec<Entry>,
107    ) -> Result<HashMap<RegionId, EntryId>> {
108        self.check_positions(&entries)?;
109        let mut last_entry_ids = HashMap::new();
110        for entry in &mut entries {
111            let region_id = entry.region_id();
112            let position = self.positions.entry(region_id).or_insert(0);
113            *position += 1;
114            let entry_id = entry_id(object_seq, *position);
115            entry.set_entry_id(entry_id);
116            last_entry_ids.insert(region_id, entry_id);
117        }
118        self.estimated_bytes += entries.iter().map(Entry::estimated_size).sum::<usize>();
119        if !entries.is_empty() {
120            self.first_admitted_at.get_or_insert_with(Instant::now);
121        }
122        self.entries.extend(entries);
123        Ok(last_entry_ids)
124    }
125
126    /// Checks that every region in `entries` stays inside the position range
127    /// once they are admitted.
128    fn check_positions(&self, entries: &[Entry]) -> Result<()> {
129        let mut positions = HashMap::new();
130        for entry in entries {
131            let region_id = entry.region_id();
132            let position = positions
133                .entry(region_id)
134                .or_insert_with(|| self.positions.get(&region_id).copied().unwrap_or(0));
135            *position += 1;
136            ensure!(
137                *position < POSITION_LIMIT,
138                WalEntryPositionExhaustedSnafu { region_id }
139            );
140        }
141        Ok(())
142    }
143
144    pub(crate) fn is_empty(&self) -> bool {
145        self.entries.is_empty()
146    }
147
148    pub(crate) fn holds_region(&self, region_id: RegionId) -> bool {
149        self.positions.contains_key(&region_id)
150    }
151
152    /// Returns the estimated size of the admitted entries.
153    pub(crate) fn estimated_bytes(&self) -> usize {
154        self.estimated_bytes
155    }
156
157    /// Returns when the first admitted entry was admitted, if any.
158    pub(crate) fn first_admitted_at(&self) -> Option<Instant> {
159        self.first_admitted_at
160    }
161
162    /// Returns true once the admitted entries reach the size limit.
163    pub(crate) fn should_seal(&self) -> bool {
164        !self.is_empty() && self.estimated_bytes >= self.max_bytes
165    }
166
167    /// Takes the admitted entries out of the batch together with the time the
168    /// first of them was admitted. The next admission starts at position one
169    /// again, under the next sequence.
170    pub(crate) fn seal(&mut self) -> (Vec<Entry>, Instant) {
171        self.estimated_bytes = 0;
172        self.positions.clear();
173        let first_admitted_at = self.first_admitted_at.take().unwrap_or_else(Instant::now);
174        (std::mem::take(&mut self.entries), first_admitted_at)
175    }
176
177    /// Drops the admitted entries, so the next admission under the same
178    /// sequence hands out the same ids again.
179    pub(crate) fn reset(&mut self) {
180        let _ = self.seal();
181    }
182}
183
184#[cfg(test)]
185mod tests {
186    use store_api::logstore::entry::NaiveEntry;
187    use store_api::logstore::provider::Provider;
188
189    use super::*;
190    use crate::error::Error;
191
192    fn entry(region_id: RegionId, payload_len: usize) -> Entry {
193        Entry::Naive(NaiveEntry {
194            provider: Provider::object_store_provider(region_id, "wal".to_string()),
195            region_id,
196            entry_id: 0,
197            data: vec![0; payload_len],
198        })
199    }
200
201    fn entry_ids(entries: &[Entry]) -> Vec<(RegionId, EntryId)> {
202        entries
203            .iter()
204            .map(|entry| (entry.region_id(), entry.entry_id()))
205            .collect()
206    }
207
208    #[test]
209    fn test_entry_id_scheme() {
210        assert_eq!(1, entry_id(0, 1));
211        assert_eq!(0x10_0001, entry_id(1, 1));
212        // 7 << 52 | 3
213        assert_eq!(31_525_197_391_593_475, entry_id(0x7_0000_0000, 3));
214        assert_eq!(u64::MAX, entry_id(OBJECT_SEQ_LIMIT - 1, POSITION_LIMIT - 1));
215        assert_eq!(1 << 44, OBJECT_SEQ_LIMIT);
216
217        assert_eq!(0, sequence_floor(0));
218        assert_eq!(1, sequence_floor(1));
219        assert_eq!(1, sequence_floor(POSITION_LIMIT - 1));
220        assert_eq!(2, sequence_floor(entry_id(1, 1)));
221        assert_eq!(8, sequence_floor(entry_id(7, POSITION_LIMIT - 1)));
222        // A contiguous id well past the first object is placed by its high bits.
223        assert_eq!(5, sequence_floor(5_000_000));
224        assert_eq!(OBJECT_SEQ_LIMIT, sequence_floor(u64::MAX));
225    }
226
227    #[test]
228    fn test_batch_assigns_positions_per_region_under_the_object_sequence() {
229        let region_a = RegionId::new(1, 1);
230        let region_b = RegionId::new(1, 2);
231        let mut batch = OpenBatch::new(usize::MAX);
232
233        let first = batch
234            .admit(5, vec![entry(region_a, 1), entry(region_b, 1)])
235            .unwrap();
236        assert_eq!(
237            HashMap::from([(region_a, entry_id(5, 1)), (region_b, entry_id(5, 1))]),
238            first
239        );
240        let second = batch
241            .admit(
242                5,
243                vec![entry(region_b, 1), entry(region_a, 1), entry(region_b, 1)],
244            )
245            .unwrap();
246        assert_eq!(
247            HashMap::from([(region_a, entry_id(5, 2)), (region_b, entry_id(5, 3))]),
248            second
249        );
250
251        assert_eq!(
252            vec![
253                (region_a, entry_id(5, 1)),
254                (region_b, entry_id(5, 1)),
255                (region_b, entry_id(5, 2)),
256                (region_a, entry_id(5, 2)),
257                (region_b, entry_id(5, 3)),
258            ],
259            entry_ids(&batch.seal().0)
260        );
261        assert!(batch.is_empty());
262        assert_eq!(
263            HashMap::from([(region_a, entry_id(6, 1))]),
264            batch.admit(6, vec![entry(region_a, 1)]).unwrap()
265        );
266        assert_eq!(POSITION_LIMIT, entry_id(6, 1) - entry_id(5, 1));
267    }
268
269    #[test]
270    fn test_batch_seals_at_size_limit() {
271        let region_id = RegionId::new(1, 1);
272        let first = entry(region_id, 8);
273        let second = entry(region_id, 8);
274        let max_bytes = first.estimated_size() + second.estimated_size();
275        let mut batch = OpenBatch::new(max_bytes);
276
277        assert!(!batch.should_seal());
278        batch.admit(0, vec![first]).unwrap();
279        assert!(!batch.should_seal());
280        batch.admit(0, vec![second]).unwrap();
281        assert!(batch.should_seal());
282
283        assert_eq!(2, batch.seal().0.len());
284        assert!(!batch.should_seal());
285    }
286
287    #[test]
288    fn test_batch_admission_clock_starts_with_the_first_entry() {
289        let region_id = RegionId::new(1, 1);
290        let mut batch = OpenBatch::new(usize::MAX);
291
292        assert!(batch.admit(0, Vec::new()).unwrap().is_empty());
293        assert!(batch.is_empty());
294        assert_eq!(None, batch.first_admitted_at());
295
296        let before = Instant::now();
297        batch.admit(0, vec![entry(region_id, 1)]).unwrap();
298        let first_admitted_at = batch.first_admitted_at().unwrap();
299        assert!(first_admitted_at >= before);
300        batch.admit(0, vec![entry(region_id, 1)]).unwrap();
301        assert_eq!(Some(first_admitted_at), batch.first_admitted_at());
302        assert_eq!(first_admitted_at, batch.seal().1);
303        assert_eq!(None, batch.first_admitted_at());
304    }
305
306    #[test]
307    fn test_batch_reset_hands_out_the_same_ids_again() {
308        let region_id = RegionId::new(1, 1);
309        let mut batch = OpenBatch::new(usize::MAX);
310
311        assert_eq!(
312            HashMap::from([(region_id, entry_id(3, 1))]),
313            batch.admit(3, vec![entry(region_id, 1)]).unwrap()
314        );
315        batch.reset();
316        assert!(batch.is_empty());
317        assert_eq!(
318            HashMap::from([(region_id, entry_id(3, 1))]),
319            batch.admit(3, vec![entry(region_id, 1)]).unwrap()
320        );
321    }
322
323    #[test]
324    fn test_batch_refuses_a_region_past_the_position_range() {
325        let region_a = RegionId::new(1, 1);
326        let region_b = RegionId::new(1, 2);
327        let mut batch = OpenBatch::new(usize::MAX);
328        let entries =
329            |region_id, count: u64| (0..count).map(|_| entry(region_id, 0)).collect::<Vec<_>>();
330
331        // An append that alone runs past the range fits no object.
332        let error = batch
333            .admit(0, entries(region_a, POSITION_LIMIT))
334            .unwrap_err();
335        assert!(
336            matches!(error, Error::WalEntryPositionExhausted { region_id, .. } if region_id == region_a),
337            "unexpected error: {error:?}"
338        );
339        assert!(batch.is_empty());
340
341        // The range holds one entry fewer; the next entry of that region needs
342        // a new batch, while another region still fits.
343        let last = batch
344            .admit(0, entries(region_a, POSITION_LIMIT - 1))
345            .unwrap();
346        assert_eq!(
347            HashMap::from([(region_a, entry_id(0, POSITION_LIMIT - 1))]),
348            last
349        );
350        assert!(batch.would_exhaust_positions(&entries(region_a, 1)));
351        assert!(!batch.would_exhaust_positions(&entries(region_b, 1)));
352        let error = batch.admit(0, entries(region_a, 1)).unwrap_err();
353        assert!(
354            matches!(error, Error::WalEntryPositionExhausted { .. }),
355            "unexpected error: {error:?}"
356        );
357        assert_eq!(POSITION_LIMIT as usize - 1, batch.seal().0.len());
358        assert_eq!(
359            HashMap::from([(region_a, entry_id(1, 1))]),
360            batch.admit(1, entries(region_a, 1)).unwrap()
361        );
362    }
363}