Skip to main content

linera_views/views/
queue_view.rs

1// Copyright (c) Zefchain Labs, Inc.
2// SPDX-License-Identifier: Apache-2.0
3
4use std::{
5    collections::{vec_deque::IterMut, VecDeque},
6    ops::Range,
7};
8
9use allocative::Allocative;
10#[cfg(with_metrics)]
11use linera_base::prometheus_util::MeasureLatency as _;
12use linera_base::{data_types::ArithmeticError, visit_allocative_simple};
13use serde::{de::DeserializeOwned, Serialize};
14
15use crate::{
16    batch::Batch,
17    common::{from_bytes_option_or_default, HasherOutput},
18    context::Context,
19    hashable_wrapper::WrappedHashableContainerView,
20    historical_hash_wrapper::HistoricallyHashableView,
21    store::ReadableKeyValueStore as _,
22    views::{ClonableView, HashableView, Hasher, View, ViewError, MIN_VIEW_TAG},
23};
24
25#[cfg(with_metrics)]
26pub(crate) mod metrics {
27    use linera_base::prometheus_util::{exponential_bucket_latencies, register_histogram_vec};
28    use prometheus::HistogramVec;
29
30    linera_base::declare_metrics! {
31        /// The runtime of hash computation
32        pub static QUEUE_VIEW_HASH_RUNTIME: HistogramVec =
33            register_histogram_vec(
34                "queue_view_hash_runtime",
35                "QueueView hash runtime",
36                &[],
37                exponential_bucket_latencies(5.0),
38            );
39    }
40}
41
42/// Key tags to create the sub-keys of a `QueueView` on top of the base key.
43#[repr(u8)]
44enum KeyTag {
45    /// Prefix for the storing of the variable `stored_indices`.
46    Store = MIN_VIEW_TAG,
47    /// Prefix for the indices of the log.
48    Index,
49}
50
51/// A view that supports a FIFO queue for values of type `T`.
52#[derive(Debug, Allocative)]
53#[allocative(bound = "C, T: Allocative")]
54pub struct QueueView<C, T> {
55    /// The view context.
56    #[allocative(skip)]
57    context: C,
58    /// The range of indices for entries persisted in storage.
59    #[allocative(visit = visit_allocative_simple)]
60    stored_indices: Range<u32>,
61    /// The number of entries to delete from the front.
62    front_delete_count: u32,
63    /// Whether to clear storage before applying updates.
64    delete_storage_first: bool,
65    /// New values added to the back, not yet persisted to storage.
66    new_back_values: VecDeque<T>,
67}
68
69impl<C, T> View for QueueView<C, T>
70where
71    C: Context,
72    T: Serialize + Send + Sync,
73{
74    const NUM_INIT_KEYS: usize = 1;
75
76    type Context = C;
77
78    fn context(&self) -> C {
79        self.context.clone()
80    }
81
82    fn pre_load(context: &C) -> Result<Vec<Vec<u8>>, ViewError> {
83        Ok(vec![context.base_key().base_tag(KeyTag::Store as u8)])
84    }
85
86    fn post_load(context: C, values: &[Option<Vec<u8>>]) -> Result<Self, ViewError> {
87        let stored_indices =
88            from_bytes_option_or_default(values.first().ok_or(ViewError::PostLoadValuesError)?)?;
89        Ok(Self {
90            context,
91            stored_indices,
92            front_delete_count: 0,
93            delete_storage_first: false,
94            new_back_values: VecDeque::new(),
95        })
96    }
97
98    fn rollback(&mut self) {
99        self.delete_storage_first = false;
100        self.front_delete_count = 0;
101        self.new_back_values.clear();
102    }
103
104    async fn has_pending_changes(&self) -> bool {
105        if self.delete_storage_first {
106            return true;
107        }
108        if self.front_delete_count > 0 {
109            return true;
110        }
111        !self.new_back_values.is_empty()
112    }
113
114    fn pre_save(&self, batch: &mut Batch) -> Result<bool, ViewError> {
115        let mut delete_view = false;
116        if self.delete_storage_first {
117            batch.delete_key_prefix(self.context.base_key().bytes.clone());
118            delete_view = true;
119        }
120        let mut new_stored_indices = self.stored_indices.clone();
121        if self.stored_count() == 0 {
122            let key_prefix = self.context.base_key().base_tag(KeyTag::Index as u8);
123            batch.delete_key_prefix(key_prefix);
124            new_stored_indices = Range::default();
125        } else if self.front_delete_count > 0 {
126            let deletion_end = new_stored_indices.start + self.front_delete_count;
127            for index in new_stored_indices.start..deletion_end {
128                let key = self
129                    .context
130                    .base_key()
131                    .derive_tag_key(KeyTag::Index as u8, &index)?;
132                batch.delete_key(key);
133            }
134            new_stored_indices.start = deletion_end;
135        }
136        if !self.new_back_values.is_empty() {
137            delete_view = false;
138            let new_back_len =
139                u32::try_from(self.new_back_values.len()).map_err(|_| ArithmeticError::Overflow)?;
140            new_stored_indices.end = new_stored_indices
141                .end
142                .checked_add(new_back_len)
143                .ok_or(ArithmeticError::Overflow)?;
144            let start = new_stored_indices.end - new_back_len;
145            for (index, value) in (start..).zip(&self.new_back_values) {
146                let key = self
147                    .context
148                    .base_key()
149                    .derive_tag_key(KeyTag::Index as u8, &index)?;
150                batch.put_key_value(key, value)?;
151            }
152        }
153        if !self.delete_storage_first || !new_stored_indices.is_empty() {
154            let key = self.context.base_key().base_tag(KeyTag::Store as u8);
155            batch.put_key_value(key, &new_stored_indices)?;
156        }
157        Ok(delete_view)
158    }
159
160    fn post_save(&mut self) {
161        if self.stored_count() == 0 {
162            self.stored_indices = Range::default();
163        } else if self.front_delete_count > 0 {
164            self.stored_indices.start += self.front_delete_count;
165        }
166        if !self.new_back_values.is_empty() {
167            self.stored_indices.end +=
168                u32::try_from(self.new_back_values.len()).expect("verified in pre_save");
169            self.new_back_values.clear();
170        }
171        self.front_delete_count = 0;
172        self.delete_storage_first = false;
173    }
174
175    fn clear(&mut self) {
176        self.delete_storage_first = true;
177        self.new_back_values.clear();
178    }
179}
180
181impl<C, T> ClonableView for QueueView<C, T>
182where
183    C: Context,
184    T: Clone + Send + Sync + Serialize,
185{
186    fn clone_unchecked(&mut self) -> Result<Self, ViewError> {
187        Ok(QueueView {
188            context: self.context.clone(),
189            stored_indices: self.stored_indices.clone(),
190            front_delete_count: self.front_delete_count,
191            delete_storage_first: self.delete_storage_first,
192            new_back_values: self.new_back_values.clone(),
193        })
194    }
195}
196
197impl<C, T> QueueView<C, T> {
198    fn stored_count(&self) -> u32 {
199        if self.delete_storage_first {
200            0
201        } else {
202            (self.stored_indices.end - self.stored_indices.start) - self.front_delete_count
203        }
204    }
205}
206
207impl<'a, C, T> QueueView<C, T>
208where
209    C: Context,
210    T: Send + Sync + Clone + Serialize + DeserializeOwned,
211{
212    async fn get(&self, index: u32) -> Result<Option<T>, ViewError> {
213        let key = self
214            .context
215            .base_key()
216            .derive_tag_key(KeyTag::Index as u8, &index)?;
217        Ok(self.context.store().read_value(&key).await?)
218    }
219
220    /// Reads the front value, if any.
221    /// ```rust
222    /// # tokio_test::block_on(async {
223    /// # use linera_views::context::MemoryContext;
224    /// # use linera_views::queue_view::QueueView;
225    /// # use linera_views::views::View;
226    /// # let context = MemoryContext::new_for_testing(());
227    /// let mut queue = QueueView::load(context).await.unwrap();
228    /// queue.push_back(34);
229    /// queue.push_back(42);
230    /// assert_eq!(queue.front().await.unwrap(), Some(34));
231    /// # })
232    /// ```
233    pub async fn front(&self) -> Result<Option<T>, ViewError> {
234        let stored_remainder = self.stored_count();
235        let value = if stored_remainder > 0 {
236            self.get(self.stored_indices.end - stored_remainder).await?
237        } else {
238            self.new_back_values.front().cloned()
239        };
240        Ok(value)
241    }
242
243    /// Reads the back value, if any.
244    /// ```rust
245    /// # tokio_test::block_on(async {
246    /// # use linera_views::context::MemoryContext;
247    /// # use linera_views::queue_view::QueueView;
248    /// # use linera_views::views::View;
249    /// # let context = MemoryContext::new_for_testing(());
250    /// let mut queue = QueueView::load(context).await.unwrap();
251    /// queue.push_back(34);
252    /// queue.push_back(42);
253    /// assert_eq!(queue.back().await.unwrap(), Some(42));
254    /// # })
255    /// ```
256    pub async fn back(&self) -> Result<Option<T>, ViewError> {
257        Ok(match self.new_back_values.back() {
258            Some(value) => Some(value.clone()),
259            None if self.stored_count() > 0 => self.get(self.stored_indices.end - 1).await?,
260            _ => None,
261        })
262    }
263
264    /// Deletes the front value, if any.
265    /// ```rust
266    /// # tokio_test::block_on(async {
267    /// # use linera_views::context::MemoryContext;
268    /// # use linera_views::queue_view::QueueView;
269    /// # use linera_views::views::View;
270    /// # let context = MemoryContext::new_for_testing(());
271    /// let mut queue = QueueView::load(context).await.unwrap();
272    /// queue.push_back(34 as u128);
273    /// queue.delete_front();
274    /// assert_eq!(queue.elements().await.unwrap(), Vec::<u128>::new());
275    /// # })
276    /// ```
277    pub fn delete_front(&mut self) {
278        if self.stored_count() > 0 {
279            self.front_delete_count += 1;
280        } else {
281            self.new_back_values.pop_front();
282        }
283    }
284
285    /// Pushes a value to the end of the queue.
286    /// ```rust
287    /// # tokio_test::block_on(async {
288    /// # use linera_views::context::MemoryContext;
289    /// # use linera_views::queue_view::QueueView;
290    /// # use linera_views::views::View;
291    /// # let context = MemoryContext::new_for_testing(());
292    /// let mut queue = QueueView::load(context).await.unwrap();
293    /// queue.push_back(34);
294    /// queue.push_back(37);
295    /// assert_eq!(queue.elements().await.unwrap(), vec![34, 37]);
296    /// # })
297    /// ```
298    pub fn push_back(&mut self, value: T) {
299        self.new_back_values.push_back(value);
300    }
301
302    /// Reads the size of the queue.
303    /// ```rust
304    /// # tokio_test::block_on(async {
305    /// # use linera_views::context::MemoryContext;
306    /// # use linera_views::queue_view::QueueView;
307    /// # use linera_views::views::View;
308    /// # let context = MemoryContext::new_for_testing(());
309    /// let mut queue = QueueView::load(context).await.unwrap();
310    /// queue.push_back(34);
311    /// assert_eq!(queue.count(), 1);
312    /// # })
313    /// ```
314    pub fn count(&self) -> usize {
315        self.stored_count() as usize + self.new_back_values.len()
316    }
317
318    /// Obtains the extra data.
319    pub fn extra(&self) -> &C::Extra {
320        self.context.extra()
321    }
322
323    async fn read_context(&self, range: Range<u32>) -> Result<Vec<T>, ViewError> {
324        let count = range.len();
325        let mut keys = Vec::with_capacity(count);
326        for index in range {
327            let key = self
328                .context
329                .base_key()
330                .derive_tag_key(KeyTag::Index as u8, &index)?;
331            keys.push(key)
332        }
333        let mut values = Vec::with_capacity(count);
334        for entry in self.context.store().read_multi_values(&keys).await? {
335            match entry {
336                None => {
337                    return Err(ViewError::MissingEntries("QueueView".into()));
338                }
339                Some(value) => values.push(value),
340            }
341        }
342        Ok(values)
343    }
344
345    /// Reads the `count` next values in the queue (including staged ones).
346    /// ```rust
347    /// # tokio_test::block_on(async {
348    /// # use linera_views::context::MemoryContext;
349    /// # use linera_views::queue_view::QueueView;
350    /// # use linera_views::views::View;
351    /// # let context = MemoryContext::new_for_testing(());
352    /// let mut queue = QueueView::load(context).await.unwrap();
353    /// queue.push_back(34);
354    /// queue.push_back(42);
355    /// assert_eq!(queue.read_front(1).await.unwrap(), vec![34]);
356    /// # })
357    /// ```
358    pub async fn read_front(&self, mut count: usize) -> Result<Vec<T>, ViewError> {
359        if count > self.count() {
360            count = self.count();
361        }
362        if count == 0 {
363            return Ok(Vec::new());
364        }
365        let mut values = Vec::with_capacity(count);
366        if !self.delete_storage_first {
367            let stored_remainder = self.stored_count();
368            let start = self.stored_indices.end - stored_remainder;
369            if count <= stored_remainder as usize {
370                let count = u32::try_from(count).map_err(|_| ArithmeticError::Overflow)?;
371                values.extend(self.read_context(start..start + count).await?);
372            } else {
373                values.extend(self.read_context(start..self.stored_indices.end).await?);
374                values.extend(
375                    self.new_back_values
376                        .range(0..count - stored_remainder as usize)
377                        .cloned(),
378                );
379            }
380        } else {
381            values.extend(self.new_back_values.range(0..count).cloned());
382        }
383        Ok(values)
384    }
385
386    /// Reads the `count` last values in the queue (including staged ones).
387    /// ```rust
388    /// # tokio_test::block_on(async {
389    /// # use linera_views::context::MemoryContext;
390    /// # use linera_views::queue_view::QueueView;
391    /// # use linera_views::views::View;
392    /// # let context = MemoryContext::new_for_testing(());
393    /// let mut queue = QueueView::load(context).await.unwrap();
394    /// queue.push_back(34);
395    /// queue.push_back(42);
396    /// assert_eq!(queue.read_back(1).await.unwrap(), vec![42]);
397    /// # })
398    /// ```
399    pub async fn read_back(&self, mut count: usize) -> Result<Vec<T>, ViewError> {
400        if count > self.count() {
401            count = self.count();
402        }
403        if count == 0 {
404            return Ok(Vec::new());
405        }
406        let mut values = Vec::with_capacity(count);
407        let new_back_len = self.new_back_values.len();
408        if count <= new_back_len || self.delete_storage_first {
409            values.extend(
410                self.new_back_values
411                    .range((new_back_len - count)..new_back_len)
412                    .cloned(),
413            );
414        } else {
415            let stored_consumed =
416                u32::try_from(count - new_back_len).map_err(|_| ArithmeticError::Underflow)?;
417            let start = self.stored_indices.end - stored_consumed;
418            values.extend(self.read_context(start..self.stored_indices.end).await?);
419            values.extend(self.new_back_values.iter().cloned());
420        }
421        Ok(values)
422    }
423
424    /// Reads all the elements
425    /// ```rust
426    /// # tokio_test::block_on(async {
427    /// # use linera_views::context::MemoryContext;
428    /// # use linera_views::queue_view::QueueView;
429    /// # use linera_views::views::View;
430    /// # let context = MemoryContext::new_for_testing(());
431    /// let mut queue = QueueView::load(context).await.unwrap();
432    /// queue.push_back(34);
433    /// queue.push_back(37);
434    /// assert_eq!(queue.elements().await.unwrap(), vec![34, 37]);
435    /// # })
436    /// ```
437    pub async fn elements(&self) -> Result<Vec<T>, ViewError> {
438        let count = self.count();
439        self.read_front(count).await
440    }
441
442    async fn load_all(&mut self) -> Result<(), ViewError> {
443        if !self.delete_storage_first {
444            let stored_remainder = self.stored_count();
445            let start = self.stored_indices.end - stored_remainder;
446            let elements = self.read_context(start..self.stored_indices.end).await?;
447            for elt in elements {
448                self.new_back_values.push_back(elt);
449            }
450            self.new_back_values.rotate_right(stored_remainder as usize);
451            // All indices are being deleted at the next flush. This is because they are deleted either:
452            // * Because a self.front_delete_count forces them to be removed
453            // * Or because loading them means that their value can be changed which invalidates
454            //   the entries on storage
455            self.delete_storage_first = true;
456        }
457        Ok(())
458    }
459
460    /// Gets a mutable iterator on the entries of the queue
461    /// ```rust
462    /// # tokio_test::block_on(async {
463    /// # use linera_views::context::MemoryContext;
464    /// # use linera_views::queue_view::QueueView;
465    /// # use linera_views::views::View;
466    /// # let context = MemoryContext::new_for_testing(());
467    /// let mut queue = QueueView::load(context).await.unwrap();
468    /// queue.push_back(34);
469    /// let mut iter = queue.try_iter_mut().await.unwrap();
470    /// let value = iter.next().unwrap();
471    /// *value = 42;
472    /// assert_eq!(queue.elements().await.unwrap(), vec![42]);
473    /// # })
474    /// ```
475    pub async fn try_iter_mut(&'a mut self) -> Result<IterMut<'a, T>, ViewError> {
476        self.load_all().await?;
477        Ok(self.new_back_values.iter_mut())
478    }
479}
480
481impl<C, T> HashableView for QueueView<C, T>
482where
483    C: Context,
484    T: Send + Sync + Clone + Serialize + DeserializeOwned,
485{
486    type Hasher = sha3::Sha3_256;
487
488    async fn hash_mut(&mut self) -> Result<<Self::Hasher as Hasher>::Output, ViewError> {
489        self.hash().await
490    }
491
492    async fn hash(&self) -> Result<<Self::Hasher as Hasher>::Output, ViewError> {
493        #[cfg(with_metrics)]
494        let _hash_latency = metrics::QUEUE_VIEW_HASH_RUNTIME.measure_latency();
495        let elements = self.elements().await?;
496        let mut hasher = sha3::Sha3_256::default();
497        hasher.update_with_bcs_bytes(&elements)?;
498        Ok(hasher.finalize())
499    }
500}
501
502/// Type wrapping `QueueView` while memoizing the hash.
503pub type HashedQueueView<C, T> = WrappedHashableContainerView<C, QueueView<C, T>, HasherOutput>;
504
505/// Wrapper around `QueueView` to compute hashes based on the history of changes.
506pub type HistoricallyHashedQueueView<C, T> = HistoricallyHashableView<C, QueueView<C, T>>;
507
508#[cfg(with_graphql)]
509mod graphql {
510    use std::borrow::Cow;
511
512    use linera_base::data_types::ArithmeticError;
513
514    use super::QueueView;
515    use crate::{
516        context::Context,
517        graphql::{hash_name, mangle},
518    };
519
520    impl<C: Send + Sync, T: async_graphql::OutputType> async_graphql::TypeName for QueueView<C, T> {
521        fn type_name() -> Cow<'static, str> {
522            format!(
523                "QueueView_{}_{:08x}",
524                mangle(T::type_name()),
525                hash_name::<T>()
526            )
527            .into()
528        }
529    }
530
531    #[async_graphql::Object(cache_control(no_cache), name_type)]
532    impl<C: Context, T: async_graphql::OutputType> QueueView<C, T>
533    where
534        T: serde::ser::Serialize + serde::de::DeserializeOwned + Clone + Send + Sync,
535    {
536        #[graphql(derived(name = "count"))]
537        async fn count_(&self) -> Result<u32, async_graphql::Error> {
538            Ok(u32::try_from(self.count()).map_err(|_| ArithmeticError::Overflow)?)
539        }
540
541        async fn entries(&self, count: Option<usize>) -> async_graphql::Result<Vec<T>> {
542            Ok(self
543                .read_front(count.unwrap_or_else(|| self.count()))
544                .await?)
545        }
546    }
547}