log_store/object_store_wal/
batch.rs1use 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
28pub(crate) const POSITION_BITS: u32 = 20;
32pub(crate) const POSITION_LIMIT: u64 = 1 << POSITION_BITS;
36pub(crate) const OBJECT_SEQ_LIMIT: u64 = 1 << (u64::BITS - POSITION_BITS);
39
40pub 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
53pub(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#[derive(Debug)]
73pub(crate) struct OpenBatch {
74 max_bytes: usize,
75 entries: Vec<Entry>,
76 estimated_bytes: usize,
77 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 pub(crate) fn would_exhaust_positions(&self, entries: &[Entry]) -> bool {
96 self.check_positions(entries).is_err()
97 }
98
99 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 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(®ion_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(®ion_id)
150 }
151
152 pub(crate) fn estimated_bytes(&self) -> usize {
154 self.estimated_bytes
155 }
156
157 pub(crate) fn first_admitted_at(&self) -> Option<Instant> {
159 self.first_admitted_at
160 }
161
162 pub(crate) fn should_seal(&self) -> bool {
164 !self.is_empty() && self.estimated_bytes >= self.max_bytes
165 }
166
167 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 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 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 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 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 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}