1use std::collections::{vec_deque::IterMut, VecDeque};
5
6use allocative::Allocative;
7#[cfg(with_metrics)]
8use linera_base::prometheus_util::MeasureLatency as _;
9use serde::{de::DeserializeOwned, Deserialize, Serialize};
10
11use crate::{
12 batch::Batch,
13 common::{from_bytes_option, from_bytes_option_or_default, HasherOutput},
14 context::Context,
15 hashable_wrapper::WrappedHashableContainerView,
16 historical_hash_wrapper::HistoricallyHashableView,
17 store::ReadableKeyValueStore as _,
18 views::{ClonableView, HashableView, Hasher, View, ViewError, MIN_VIEW_TAG},
19};
20
21#[cfg(with_metrics)]
22mod metrics {
23 use std::sync::LazyLock;
24
25 use linera_base::prometheus_util::{exponential_bucket_latencies, register_histogram_vec};
26 use prometheus::HistogramVec;
27
28 pub static BUCKET_QUEUE_VIEW_HASH_RUNTIME: LazyLock<HistogramVec> = LazyLock::new(|| {
30 register_histogram_vec(
31 "bucket_queue_view_hash_runtime",
32 "BucketQueueView hash runtime",
33 &[],
34 exponential_bucket_latencies(5.0),
35 )
36 });
37}
38
39#[repr(u8)]
41enum KeyTag {
42 Front = MIN_VIEW_TAG,
44 Store,
46 Index,
48}
49
50#[derive(Clone, Debug, Default, Serialize, Deserialize)]
52struct BucketStore {
53 descriptions: Vec<BucketDescription>,
56 front_position: usize,
58}
59
60#[derive(Copy, Clone, Debug, Default, Serialize, Deserialize)]
62struct BucketDescription {
63 length: usize,
65 index: usize,
67}
68
69impl BucketStore {
70 fn len(&self) -> usize {
71 self.descriptions.len()
72 }
73}
74
75#[derive(Copy, Clone, Debug, Allocative)]
77struct Cursor {
78 offset: usize,
80 position: usize,
82}
83
84#[derive(Clone, Debug, Allocative)]
86enum State<T> {
87 Loaded { data: Vec<T> },
88 NotLoaded { length: usize },
89}
90
91impl<T> Bucket<T> {
92 fn len(&self) -> usize {
93 match &self.state {
94 State::Loaded { data } => data.len(),
95 State::NotLoaded { length } => *length,
96 }
97 }
98
99 fn is_loaded(&self) -> bool {
100 match self.state {
101 State::Loaded { .. } => true,
102 State::NotLoaded { .. } => false,
103 }
104 }
105
106 fn to_description(&self) -> BucketDescription {
107 BucketDescription {
108 length: self.len(),
109 index: self.index,
110 }
111 }
112}
113
114#[derive(Clone, Debug, Allocative)]
116struct Bucket<T> {
117 index: usize,
119 state: State<T>,
121}
122
123#[derive(Debug, Allocative)]
129#[allocative(bound = "C, T: Allocative, const N: usize")]
130pub struct BucketQueueView<C, T, const N: usize> {
131 #[allocative(skip)]
133 context: C,
134 stored_buckets: VecDeque<Bucket<T>>,
136 new_back_values: VecDeque<T>,
138 stored_front_position: usize,
140 cursor: Option<Cursor>,
143 delete_storage_first: bool,
145}
146
147impl<C, T, const N: usize> View for BucketQueueView<C, T, N>
148where
149 C: Context,
150 T: Send + Sync + Clone + Serialize + DeserializeOwned,
151{
152 const NUM_INIT_KEYS: usize = 2;
153
154 type Context = C;
155
156 fn context(&self) -> &C {
157 &self.context
158 }
159
160 fn pre_load(context: &C) -> Result<Vec<Vec<u8>>, ViewError> {
161 let key1 = context.base_key().base_tag(KeyTag::Front as u8);
162 let key2 = context.base_key().base_tag(KeyTag::Store as u8);
163 Ok(vec![key1, key2])
164 }
165
166 fn post_load(context: C, values: &[Option<Vec<u8>>]) -> Result<Self, ViewError> {
167 let value1 = values.first().ok_or(ViewError::PostLoadValuesError)?;
168 let value2 = values.get(1).ok_or(ViewError::PostLoadValuesError)?;
169 let front = from_bytes_option::<Vec<T>>(value1)?;
170 let mut stored_buckets = VecDeque::from(match front {
171 Some(data) => {
172 let bucket = Bucket {
173 index: 0,
174 state: State::Loaded { data },
175 };
176 vec![bucket]
177 }
178 None => {
179 vec![]
180 }
181 });
182 let bucket_store = from_bytes_option_or_default::<BucketStore>(value2)?;
183 for i in 1..bucket_store.len() {
186 let length = bucket_store.descriptions[i].length;
187 let index = bucket_store.descriptions[i].index;
188 stored_buckets.push_back(Bucket {
189 index,
190 state: State::NotLoaded { length },
191 });
192 }
193 let cursor = if bucket_store.descriptions.is_empty() {
194 None
195 } else {
196 Some(Cursor {
197 offset: 0,
198 position: bucket_store.front_position,
199 })
200 };
201 Ok(Self {
202 context,
203 stored_buckets,
204 stored_front_position: bucket_store.front_position,
205 new_back_values: VecDeque::new(),
206 cursor,
207 delete_storage_first: false,
208 })
209 }
210
211 fn rollback(&mut self) {
212 self.delete_storage_first = false;
213 self.cursor = if self.stored_buckets.is_empty() {
214 None
215 } else {
216 Some(Cursor {
217 offset: 0,
218 position: self.stored_front_position,
219 })
220 };
221 self.new_back_values.clear();
222 }
223
224 async fn has_pending_changes(&self) -> bool {
225 if self.delete_storage_first {
226 return true;
227 }
228 if !self.stored_buckets.is_empty() {
229 let Some(cursor) = self.cursor else {
230 return true;
231 };
232 if cursor.offset != 0 || cursor.position != self.stored_front_position {
233 return true;
234 }
235 }
236 !self.new_back_values.is_empty()
237 }
238
239 fn pre_save(&self, batch: &mut Batch) -> Result<bool, ViewError> {
240 let mut delete_view = false;
241 let mut descriptions = Vec::new();
242 let mut stored_front_position = self.stored_front_position;
243 if self.stored_count() == 0 {
244 let key_prefix = self.context.base_key().bytes.clone();
245 batch.delete_key_prefix(key_prefix);
246 delete_view = true;
247 stored_front_position = 0;
248 } else if let Some(cursor) = self.cursor {
249 for i in 0..cursor.offset {
251 let bucket = &self.stored_buckets[i];
252 let index = bucket.index;
253 let key = self.get_bucket_key(index)?;
254 batch.delete_key(key);
255 }
256 stored_front_position = cursor.position;
257 let first_index = self.stored_buckets[cursor.offset].index;
259 let start_offset = if first_index != 0 {
260 let key = self.get_bucket_key(first_index)?;
262 batch.delete_key(key);
263 let key = self.get_bucket_key(0)?;
264 let bucket = &self.stored_buckets[cursor.offset];
265 let State::Loaded { data } = &bucket.state else {
266 unreachable!("The front bucket is always loaded.");
267 };
268 batch.put_key_value(key, data)?;
269 descriptions.push(BucketDescription {
270 length: bucket.len(),
271 index: 0,
272 });
273 cursor.offset + 1
274 } else {
275 cursor.offset
276 };
277 for bucket in self.stored_buckets.range(start_offset..) {
278 descriptions.push(bucket.to_description());
279 }
280 }
281 if !self.new_back_values.is_empty() {
282 delete_view = false;
283 let mut index = if self.stored_count() == 0 {
287 0
288 } else if let Some(last_description) = descriptions.last() {
289 last_description.index + 1
290 } else {
291 0
293 };
294 let mut start = 0;
295 while start < self.new_back_values.len() {
296 let end = std::cmp::min(start + N, self.new_back_values.len());
297 let value_chunk: Vec<_> = self.new_back_values.range(start..end).collect();
298 let key = self.get_bucket_key(index)?;
299 batch.put_key_value(key, &value_chunk)?;
300 descriptions.push(BucketDescription {
301 index,
302 length: end - start,
303 });
304 index += 1;
305 start = end;
306 }
307 }
308 if !delete_view {
309 let bucket_store = BucketStore {
310 descriptions,
311 front_position: stored_front_position,
312 };
313 let key = self.context.base_key().base_tag(KeyTag::Store as u8);
314 batch.put_key_value(key, &bucket_store)?;
315 }
316 Ok(delete_view)
317 }
318
319 fn post_save(&mut self) {
320 if self.stored_count() == 0 {
321 self.stored_buckets.clear();
322 self.stored_front_position = 0;
323 self.cursor = None;
324 } else if let Some(cursor) = self.cursor {
325 for _ in 0..cursor.offset {
326 self.stored_buckets.pop_front();
327 }
328 self.cursor = Some(Cursor {
329 offset: 0,
330 position: cursor.position,
331 });
332 self.stored_front_position = cursor.position;
333 self.stored_buckets[0].index = 0;
335 }
336 if !self.new_back_values.is_empty() {
337 let start = self
338 .stored_buckets
339 .back()
340 .map(|bucket| bucket.index + 1)
341 .unwrap_or_default();
342 let new_back_values = std::mem::take(&mut self.new_back_values);
343 let new_back_values = new_back_values.into_iter().collect::<Vec<_>>();
344 self.stored_buckets
345 .extend(
346 new_back_values
347 .chunks(N)
348 .zip(start..)
349 .map(|(value_chunk, index)| Bucket {
350 index,
351 state: State::Loaded {
352 data: value_chunk.to_vec(),
353 },
354 }),
355 );
356 if self.cursor.is_none() {
357 self.cursor = Some(Cursor {
358 offset: 0,
359 position: 0,
360 });
361 }
362 }
363 self.delete_storage_first = false;
364 }
365
366 fn clear(&mut self) {
367 self.delete_storage_first = true;
368 self.new_back_values.clear();
369 self.cursor = None;
370 }
371}
372
373impl<C: Clone, T: Clone, const N: usize> ClonableView for BucketQueueView<C, T, N>
374where
375 Self: View,
376{
377 fn clone_unchecked(&mut self) -> Result<Self, ViewError> {
378 Ok(BucketQueueView {
379 context: self.context.clone(),
380 stored_buckets: self.stored_buckets.clone(),
381 new_back_values: self.new_back_values.clone(),
382 stored_front_position: self.stored_front_position,
383 cursor: self.cursor,
384 delete_storage_first: self.delete_storage_first,
385 })
386 }
387}
388
389impl<C: Context, T, const N: usize> BucketQueueView<C, T, N> {
390 fn get_bucket_key(&self, index: usize) -> Result<Vec<u8>, ViewError> {
392 Ok(if index == 0 {
393 self.context.base_key().base_tag(KeyTag::Front as u8)
394 } else {
395 self.context
396 .base_key()
397 .derive_tag_key(KeyTag::Index as u8, &index)?
398 })
399 }
400
401 pub fn stored_count(&self) -> usize {
414 if self.delete_storage_first {
415 0
416 } else {
417 let Some(cursor) = self.cursor else {
418 return 0;
419 };
420 let mut stored_count = 0;
421 for offset in cursor.offset..self.stored_buckets.len() {
422 stored_count += self.stored_buckets[offset].len();
423 }
424 stored_count -= cursor.position;
425 stored_count
426 }
427 }
428
429 pub fn count(&self) -> usize {
442 self.stored_count() + self.new_back_values.len()
443 }
444}
445
446impl<C: Context, T: DeserializeOwned + Clone, const N: usize> BucketQueueView<C, T, N> {
447 pub fn front(&self) -> Option<&T> {
461 match self.cursor {
462 Some(Cursor { offset, position }) => {
463 let bucket = &self.stored_buckets[offset];
464 let State::Loaded { data } = &bucket.state else {
465 unreachable!();
466 };
467 Some(&data[position])
468 }
469 None => self.new_back_values.front(),
470 }
471 }
472
473 pub fn front_mut(&mut self) -> Option<&mut T> {
489 match self.cursor {
490 Some(Cursor { offset, position }) => {
491 let bucket = self
492 .stored_buckets
493 .get_mut(offset)
494 .expect("cursor.offset must be a valid index into stored_buckets");
495 let State::Loaded { data } = &mut bucket.state else {
496 unreachable!();
497 };
498 Some(
499 data.get_mut(position)
500 .expect("cursor.position must be a valid index within the front bucket"),
501 )
502 }
503 None => self.new_back_values.front_mut(),
504 }
505 }
506
507 pub async fn delete_front(&mut self) -> Result<(), ViewError> {
521 match self.cursor {
522 Some(cursor) => {
523 let mut offset = cursor.offset;
524 let mut position = cursor.position + 1;
525 if self.stored_buckets[offset].len() == position {
526 offset += 1;
527 position = 0;
528 }
529 if offset == self.stored_buckets.len() {
530 self.cursor = None;
531 } else {
532 if !self.stored_buckets[offset].is_loaded() {
533 let index = self.stored_buckets[offset].index;
534 let key = self.get_bucket_key(index)?;
535 let data = self.context.store().read_value(&key).await?;
536 let data = match data {
537 Some(value) => value,
538 None => {
539 return Err(ViewError::MissingEntries(
540 "BucketQueueView::delete_front".into(),
541 ));
542 }
543 };
544 self.stored_buckets[offset].state = State::Loaded { data };
545 }
546 self.cursor = Some(Cursor { offset, position });
547 }
548 }
549 None => {
550 self.new_back_values.pop_front();
551 }
552 }
553 Ok(())
554 }
555
556 pub fn push_back(&mut self, value: T) {
569 self.new_back_values.push_back(value);
570 }
571
572 pub async fn elements(&self) -> Result<Vec<T>, ViewError> {
586 let count = self.count();
587 self.read_context(self.cursor, count).await
588 }
589
590 pub async fn back(&mut self) -> Result<Option<T>, ViewError>
604 where
605 T: Clone,
606 {
607 if let Some(value) = self.new_back_values.back() {
608 return Ok(Some(value.clone()));
609 }
610 if self.cursor.is_none() {
611 return Ok(None);
612 }
613 let Some(bucket) = self.stored_buckets.back() else {
614 return Ok(None);
615 };
616 match &bucket.state {
617 State::Loaded { data } => Ok(Some(
618 data.last().expect("a stored bucket is never empty").clone(),
619 )),
620 State::NotLoaded { .. } => {
621 let key = self.get_bucket_key(bucket.index)?;
622 let data = self
623 .context
624 .store()
625 .read_value::<Vec<T>>(&key)
626 .await?
627 .ok_or_else(|| ViewError::MissingEntries("BucketQueueView::back".into()))?;
628 let result = data.last().expect("a stored bucket is never empty").clone();
629 self.stored_buckets
630 .back_mut()
631 .expect("stored_buckets is non-empty since we just accessed its back element")
632 .state = State::Loaded { data };
633 Ok(Some(result))
634 }
635 }
636 }
637
638 async fn read_context(
639 &self,
640 cursor: Option<Cursor>,
641 count: usize,
642 ) -> Result<Vec<T>, ViewError> {
643 if count == 0 {
644 return Ok(Vec::new());
645 }
646 let mut elements = Vec::<T>::new();
647 let mut count_remain = count;
648 if let Some(cursor) = cursor {
649 let mut keys = Vec::new();
650 let mut position = cursor.position;
651 for offset in cursor.offset..self.stored_buckets.len() {
652 let bucket = &self.stored_buckets[offset];
653 let size = bucket.len() - position;
654 if !bucket.is_loaded() {
655 let key = self.get_bucket_key(bucket.index)?;
656 keys.push(key);
657 };
658 if size >= count_remain {
659 break;
660 }
661 count_remain -= size;
662 position = 0;
663 }
664 let values = self.context.store().read_multi_values_bytes(&keys).await?;
665 let mut value_pos = 0;
666 count_remain = count;
667 let mut position = cursor.position;
668 for offset in cursor.offset..self.stored_buckets.len() {
669 let bucket = &self.stored_buckets[offset];
670 let size = bucket.len() - position;
671 let data = match &bucket.state {
672 State::Loaded { data } => data,
673 State::NotLoaded { .. } => {
674 let value = match &values[value_pos] {
675 Some(value) => value,
676 None => {
677 return Err(ViewError::MissingEntries(
678 "BucketQueueView::read_context".into(),
679 ));
680 }
681 };
682 value_pos += 1;
683 &bcs::from_bytes::<Vec<T>>(value)?
684 }
685 };
686 elements.extend(data[position..].iter().take(count_remain).cloned());
687 if size >= count_remain {
688 return Ok(elements);
689 }
690 count_remain -= size;
691 position = 0;
692 }
693 }
694 let count_read = std::cmp::min(count_remain, self.new_back_values.len());
695 elements.extend(self.new_back_values.range(0..count_read).cloned());
696 Ok(elements)
697 }
698
699 pub async fn read_front(&self, count: usize) -> Result<Vec<T>, ViewError> {
714 let count = std::cmp::min(count, self.count());
715 self.read_context(self.cursor, count).await
716 }
717
718 pub async fn read_back(&self, count: usize) -> Result<Vec<T>, ViewError> {
733 let count = std::cmp::min(count, self.count());
734 if count <= self.new_back_values.len() {
735 let start = self.new_back_values.len() - count;
736 Ok(self
737 .new_back_values
738 .range(start..)
739 .cloned()
740 .collect::<Vec<_>>())
741 } else {
742 let mut increment = self.count() - count;
743 let Some(cursor) = self.cursor else {
744 unreachable!();
745 };
746 let mut position = cursor.position;
747 for offset in cursor.offset..self.stored_buckets.len() {
748 let size = self.stored_buckets[offset].len() - position;
749 if increment < size {
750 return self
751 .read_context(
752 Some(Cursor {
753 offset,
754 position: position + increment,
755 }),
756 count,
757 )
758 .await;
759 }
760 increment -= size;
761 position = 0;
762 }
763 unreachable!();
764 }
765 }
766
767 async fn load_all(&mut self) -> Result<(), ViewError> {
768 if !self.delete_storage_first {
769 let elements = self.elements().await?;
770 self.new_back_values.clear();
771 for elt in elements {
772 self.new_back_values.push_back(elt);
773 }
774 self.cursor = None;
775 self.delete_storage_first = true;
776 }
777 Ok(())
778 }
779
780 pub async fn try_iter_mut(&mut self) -> Result<IterMut<'_, T>, ViewError> {
796 self.load_all().await?;
797 Ok(self.new_back_values.iter_mut())
798 }
799}
800
801impl<C: Context, T: Serialize + DeserializeOwned + Send + Sync + Clone, const N: usize> HashableView
802 for BucketQueueView<C, T, N>
803where
804 Self: View,
805{
806 type Hasher = sha3::Sha3_256;
807
808 async fn hash_mut(&mut self) -> Result<<Self::Hasher as Hasher>::Output, ViewError> {
809 self.hash().await
810 }
811
812 async fn hash(&self) -> Result<<Self::Hasher as Hasher>::Output, ViewError> {
813 #[cfg(with_metrics)]
814 let _hash_latency = metrics::BUCKET_QUEUE_VIEW_HASH_RUNTIME.measure_latency();
815 let elements = self.elements().await?;
816 let mut hasher = sha3::Sha3_256::default();
817 hasher.update_with_bcs_bytes(&elements)?;
818 Ok(hasher.finalize())
819 }
820}
821
822pub type HashedBucketQueueView<C, T, const N: usize> =
824 WrappedHashableContainerView<C, BucketQueueView<C, T, N>, HasherOutput>;
825
826pub type HistoricallyHashedBucketQueueView<C, T, const N: usize> =
828 HistoricallyHashableView<C, BucketQueueView<C, T, N>>;
829
830#[cfg(with_graphql)]
831mod graphql {
832 use std::borrow::Cow;
833
834 use super::BucketQueueView;
835 use crate::{
836 context::Context,
837 graphql::{hash_name, mangle},
838 };
839
840 impl<C: Send + Sync, T: async_graphql::OutputType, const N: usize> async_graphql::TypeName
841 for BucketQueueView<C, T, N>
842 {
843 fn type_name() -> Cow<'static, str> {
844 format!(
845 "BucketQueueView_{}_{:08x}",
846 mangle(T::type_name()),
847 hash_name::<T>()
848 )
849 .into()
850 }
851 }
852
853 #[async_graphql::Object(cache_control(no_cache), name_type)]
854 impl<C: Context, T: async_graphql::OutputType, const N: usize> BucketQueueView<C, T, N>
855 where
856 C: Send + Sync,
857 T: serde::ser::Serialize + serde::de::DeserializeOwned + Clone + Send + Sync,
858 {
859 #[graphql(derived(name = "count"))]
860 async fn count_(&self) -> Result<u32, async_graphql::Error> {
861 Ok(self.count() as u32)
862 }
863
864 async fn entries(&self, count: Option<usize>) -> async_graphql::Result<Vec<T>> {
865 Ok(self
866 .read_front(count.unwrap_or_else(|| self.count()))
867 .await?)
868 }
869 }
870}
871
872#[cfg(test)]
873mod tests {
874 use super::*;
875 use crate::{
876 batch::Batch,
877 context::{Context, MemoryContext},
878 store::WritableKeyValueStore as _,
879 };
880
881 #[tokio::test]
887 async fn delete_front_load_failure_preserves_invariant() -> Result<(), ViewError> {
888 let context = MemoryContext::new_for_testing(());
889 let mut view = BucketQueueView::<_, u8, 2>::load(context.clone()).await?;
890 for value in [1u8, 2, 3, 4] {
891 view.push_back(value);
892 }
893 save(&context, &mut view).await?;
894
895 let mut view = BucketQueueView::<_, u8, 2>::load(context.clone()).await?;
896
897 let bucket1_key = view.get_bucket_key(1)?;
898 let mut batch = Batch::new();
899 batch.delete_key(bucket1_key);
900 context.store().write_batch(batch).await?;
901
902 view.delete_front().await?;
903 let err = view.delete_front().await.expect_err("load should fail");
904 assert!(matches!(err, ViewError::MissingEntries(_)));
905
906 save(&context, &mut view).await?;
907
908 Ok(())
909 }
910
911 async fn save<V: View>(context: &V::Context, view: &mut V) -> Result<(), ViewError> {
912 let mut batch = Batch::new();
913 view.pre_save(&mut batch)?;
914 context.store().write_batch(batch).await?;
915 view.post_save();
916 Ok(())
917 }
918}