Skip to main content

mkit_server/indexed/
budget.rs

1//! Per-slice subrequest accounting and pack window reads (WP-4.8).
2//!
3//! A Workers alarm allows 1,000 subrequests and its clock does not advance
4//! during CPU work, so a slice is bounded in fixed units: every call that
5//! leaves the Durable Object (an R2 range, an index or membership shard call)
6//! is charged one unit, and the slice stops before it passes its share.
7
8use 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/// A shared call counter with a fixed limit.
19#[derive(Debug, Clone)]
20pub struct SliceBudget {
21    used: Arc<AtomicU32>,
22    limit: u32,
23}
24
25impl SliceBudget {
26    /// A budget of `limit` calls.
27    #[must_use]
28    pub fn new(limit: u32) -> Self {
29        Self {
30            used: Arc::new(AtomicU32::new(0)),
31            limit,
32        }
33    }
34
35    /// Calls charged so far.
36    #[must_use]
37    pub fn used(&self) -> u32 {
38        self.used.load(Ordering::SeqCst)
39    }
40
41    /// Calls left.
42    #[must_use]
43    pub fn remaining(&self) -> u32 {
44        self.limit.saturating_sub(self.used())
45    }
46
47    /// Charge one call, failing once the limit is spent. The failing call is
48    /// not counted, so an exhausted budget stays exhausted.
49    ///
50    /// # Errors
51    /// `StoreError::Unavailable`: a spent budget is not a CAS race.
52    pub fn charge(&self) -> Result<(), StoreError> {
53        self.charge_many(1)
54    }
55
56    /// Reserve a combined operation before any external effect.
57    /// # Errors
58    /// Exhausted shared allowance, without partially charging the reservation.
59    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/// Whether `error` is [`SliceBudget`] running out.
72#[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/// A store whose every call is one charged unit, batched reads included (one
78/// round trip however many keys), as the rollup's budgeted store charges.
79#[derive(Debug)]
80pub struct Budgeted<'a, S> {
81    inner: &'a S,
82    budget: &'a SliceBudget,
83}
84
85impl<'a, S> Budgeted<'a, S> {
86    /// `inner` charging `budget`.
87    #[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
146/// Blob reads use the same counter as namespace calls.
147impl<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        // R2's ranged BlobStore read checks metadata before fetching bytes.
160        // Reserve both backend requests even for stores that need only one.
161        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/// A window read that could not be served.
179#[derive(Debug, Clone, Copy, PartialEq, Eq)]
180pub enum WindowError {
181    /// The object no longer has the etag the job recorded: restart the job.
182    EtagChanged,
183    /// The pack is not in storage.
184    Missing,
185    /// The store failed or returned a short range.
186    Unavailable,
187}
188
189/// Bytes of one range and the etag of the object they came from.
190#[derive(Debug)]
191pub struct Window {
192    /// Exactly the bytes asked for.
193    pub bytes: Vec<u8>,
194    /// The object's etag.
195    pub etag: String,
196}
197
198/// Range reads of a stored pack, bound to one etag (SPEC-PACKFILE ยง11: keeping
199/// the source immutable across resumes is the caller's job).
200pub trait PackWindows: MaybeSend + MaybeSync {
201    /// Read `len` bytes at `offset`. With `etag`, an object that no longer has
202    /// it answers [`WindowError::EtagChanged`], never other bytes.
203    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/// Windows over any [`BlobStore`]. Packs are content-addressed and immutable,
213/// so the etag is the pack id; the reader binds the bytes to it as well.
214#[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}