1use crate::store::{
9 BlobBody, BlobKey, BlobMeta, BlobStore, ByteRange, Cursor, Key, NamespaceStore, Partition,
10 RangeScan, ScanPage, StoreCapabilities, StoreError, Value,
11};
12use crate::{Batch, BatchOutcome, BoxFuture, MaybeSend, MaybeSync, PartitionStats};
13use futures::StreamExt as _;
14use mkit_core::hash::Hash;
15use std::sync::Arc;
16use std::sync::atomic::{AtomicU32, Ordering};
17
18#[derive(Debug, Clone)]
20pub struct SliceBudget {
21 used: Arc<AtomicU32>,
22 limit: u32,
23}
24
25impl SliceBudget {
26 #[must_use]
28 pub fn new(limit: u32) -> Self {
29 Self {
30 used: Arc::new(AtomicU32::new(0)),
31 limit,
32 }
33 }
34
35 #[must_use]
37 pub fn used(&self) -> u32 {
38 self.used.load(Ordering::SeqCst)
39 }
40
41 #[must_use]
43 pub fn remaining(&self) -> u32 {
44 self.limit.saturating_sub(self.used())
45 }
46
47 pub fn charge(&self) -> Result<(), StoreError> {
53 self.charge_many(1)
54 }
55
56 pub fn charge_many(&self, calls: u32) -> Result<(), StoreError> {
60 self.used
61 .fetch_update(Ordering::SeqCst, Ordering::SeqCst, |used| {
62 used.checked_add(calls).filter(|total| *total <= self.limit)
63 })
64 .map(|_| ())
65 .map_err(|_| {
66 StoreError::Unavailable("verification slice subrequest budget exhausted".into())
67 })
68 }
69}
70
71#[must_use]
73pub fn is_exhausted(error: &StoreError) -> bool {
74 matches!(error, StoreError::Unavailable(reason) if reason.to_string().contains("subrequest budget"))
75}
76
77#[derive(Debug)]
80pub struct Budgeted<'a, S> {
81 inner: &'a S,
82 budget: &'a SliceBudget,
83}
84
85impl<'a, S> Budgeted<'a, S> {
86 #[must_use]
88 pub fn new(inner: &'a S, budget: &'a SliceBudget) -> Self {
89 Self { inner, budget }
90 }
91}
92
93impl<S: NamespaceStore> NamespaceStore for Budgeted<'_, S> {
94 fn capabilities(&self) -> StoreCapabilities {
95 self.inner.capabilities()
96 }
97 async fn get(&self, p: &Partition, key: &Key) -> Result<Option<Value>, StoreError> {
98 self.budget.charge()?;
99 self.inner.get(p, key).await
100 }
101 async fn has(&self, p: &Partition, key: &Key) -> Result<bool, StoreError> {
102 self.budget.charge()?;
103 self.inner.has(p, key).await
104 }
105 async fn get_many(
106 &self,
107 p: &Partition,
108 keys: &[Key],
109 ) -> Result<Vec<Option<Value>>, StoreError> {
110 self.budget.charge()?;
111 self.inner.get_many(p, keys).await
112 }
113 async fn scan_many(
114 &self,
115 p: &Partition,
116 ranges: &[RangeScan],
117 ) -> Result<Vec<ScanPage>, StoreError> {
118 self.budget.charge()?;
119 self.inner.scan_many(p, ranges).await
120 }
121 async fn scan(
122 &self,
123 p: &Partition,
124 start: &Key,
125 end: &Key,
126 after: Option<&Cursor>,
127 limit: u32,
128 ) -> Result<ScanPage, StoreError> {
129 self.budget.charge()?;
130 self.inner.scan(p, start, end, after, limit).await
131 }
132 async fn apply(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, StoreError> {
133 self.budget.charge()?;
134 self.inner.apply(p, batch).await
135 }
136 async fn stats(&self, p: &Partition) -> Result<PartitionStats, StoreError> {
137 self.budget.charge()?;
138 self.inner.stats(p).await
139 }
140 async fn probe(&self) -> Result<(), StoreError> {
141 self.budget.charge()?;
142 self.inner.probe().await
143 }
144}
145
146impl<B: BlobStore> BlobStore for Budgeted<'_, B> {
148 type Sink = B::Sink;
149
150 async fn begin(&self, key: BlobKey, len: u64) -> Result<Self::Sink, StoreError> {
151 self.inner.begin(key, len).await
152 }
153 async fn get(
154 &self,
155 key: &BlobKey,
156 range: Option<ByteRange>,
157 ) -> Result<Option<BlobBody>, StoreError> {
158 self.budget.charge()?;
159 if range.is_some() {
162 self.budget.charge()?;
163 }
164 self.inner.get(key, range).await
165 }
166 async fn head(&self, key: &BlobKey) -> Result<Option<BlobMeta>, StoreError> {
167 self.budget.charge()?;
168 self.inner.head(key).await
169 }
170 async fn probe(&self) -> Result<(), StoreError> {
171 self.inner.probe().await
172 }
173 async fn delete(&self, key: &BlobKey) -> Result<bool, StoreError> {
174 self.inner.delete(key).await
175 }
176}
177
178#[derive(Debug, Clone, Copy, PartialEq, Eq)]
180pub enum WindowError {
181 EtagChanged,
183 Missing,
185 Unavailable,
187}
188
189#[derive(Debug)]
191pub struct Window {
192 pub bytes: Vec<u8>,
194 pub etag: String,
196}
197
198pub trait PackWindows: MaybeSend + MaybeSync {
201 fn read<'a>(
204 &'a self,
205 pack: &'a Hash,
206 offset: u64,
207 len: u64,
208 etag: Option<&'a str>,
209 ) -> BoxFuture<'a, Result<Window, WindowError>>;
210}
211
212#[derive(Debug)]
215pub struct BlobWindows<'a, B>(pub &'a B);
216
217impl<B: BlobStore> PackWindows for BlobWindows<'_, B> {
218 fn read<'a>(
219 &'a self,
220 pack: &'a Hash,
221 offset: u64,
222 len: u64,
223 etag: Option<&'a str>,
224 ) -> BoxFuture<'a, Result<Window, WindowError>> {
225 Box::pin(async move {
226 let tag = mkit_core::hash::to_hex(pack);
227 if etag.is_some_and(|expected| expected != tag) {
228 return Err(WindowError::EtagChanged);
229 }
230 let end = offset
231 .checked_add(len.saturating_sub(1))
232 .ok_or(WindowError::Unavailable)?;
233 let body = self
234 .0
235 .get(
236 &BlobKey::pack(*pack),
237 Some(ByteRange {
238 start: offset,
239 end_inclusive: end,
240 }),
241 )
242 .await
243 .map_err(|_| WindowError::Unavailable)?
244 .ok_or(WindowError::Missing)?;
245 let mut bytes = Vec::new();
246 match body {
247 BlobBody::Bytes(chunk) => bytes.extend_from_slice(&chunk),
248 BlobBody::Stream { mut stream, .. } => {
249 while let Some(chunk) = stream.next().await {
250 bytes.extend_from_slice(&chunk.map_err(|_| WindowError::Unavailable)?);
251 if bytes.len() as u64 > len {
252 return Err(WindowError::Unavailable);
253 }
254 }
255 }
256 }
257 if bytes.len() as u64 != len {
258 return Err(WindowError::Unavailable);
259 }
260 Ok(Window { bytes, etag: tag })
261 })
262 }
263}
264
265#[cfg(test)]
266mod tests {
267 use super::*;
268
269 #[test]
270 fn budget_counts_calls_and_stays_exhausted() {
271 let budget = SliceBudget::new(2);
272 assert!(budget.charge().is_ok() && budget.charge().is_ok());
273 let error = budget.charge().unwrap_err();
274 assert!(is_exhausted(&error));
275 assert_eq!((budget.used(), budget.remaining()), (2, 0));
276 assert!(budget.charge().is_err());
277 assert!(!is_exhausted(&StoreError::Unavailable("x".into())));
278 }
279}