1use 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 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#[repr(u8)]
44enum KeyTag {
45 Store = MIN_VIEW_TAG,
47 Index,
49}
50
51#[derive(Debug, Allocative)]
53#[allocative(bound = "C, T: Allocative")]
54pub struct QueueView<C, T> {
55 #[allocative(skip)]
57 context: C,
58 #[allocative(visit = visit_allocative_simple)]
60 stored_indices: Range<u32>,
61 front_delete_count: u32,
63 delete_storage_first: bool,
65 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 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 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 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 pub fn push_back(&mut self, value: T) {
299 self.new_back_values.push_back(value);
300 }
301
302 pub fn count(&self) -> usize {
315 self.stored_count() as usize + self.new_back_values.len()
316 }
317
318 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 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 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 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 self.delete_storage_first = true;
456 }
457 Ok(())
458 }
459
460 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
502pub type HashedQueueView<C, T> = WrappedHashableContainerView<C, QueueView<C, T>, HasherOutput>;
504
505pub 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}