Skip to main content

mkit_server/store/
kv.rs

1//! The key-level metadata contract (PRD §5.3): a partitioned, ordered
2//! key-value store whose only write is one declarative [`Batch`].
3
4use core::future::Future;
5
6use bytes::Bytes;
7
8use super::error::StoreError;
9use super::keys;
10use super::partition::Partition;
11use crate::rt::{MaybeSend, MaybeSync};
12
13// Portable limits. Every backend accepts batches within them and the core
14// never plans beyond them. They leave headroom under Durable Object SQLite
15// storage (2 MB per row or value, one `transactionSync` per batch; see
16// https://developers.cloudflare.com/durable-objects/platform/limits/), so a
17// batch that fits here fits every backend.
18
19/// Longest accepted key, in bytes.
20pub const MAX_KEY_BYTES: usize = 1024;
21/// Longest accepted value, in bytes; below [`MAX_BATCH_BYTES`].
22pub const MAX_VALUE_BYTES: usize = 512 * 1024;
23/// Most preconditions plus writes in one [`Batch`].
24pub const MAX_BATCH_OPS: usize = 100;
25/// Most key and value bytes, summed over a [`Batch`]'s preconditions and
26/// writes.
27pub const MAX_BATCH_BYTES: usize = 1024 * 1024;
28/// A timer value must fit Equals+Put while a guarded retry relocates two keys.
29/// Four maximum keys (Equals, Absent, Delete, Put) share the existing batch cap.
30pub(crate) const MAX_TIMER_VALUE_BYTES: usize = (MAX_BATCH_BYTES - 4 * MAX_KEY_BYTES) / 2;
31
32macro_rules! bytes_newtype {
33    ($(#[$doc:meta])* $name:ident) => {
34        $(#[$doc])*
35        #[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord, Default)]
36        pub struct $name(Bytes);
37
38        impl $name {
39            /// Wrap `bytes`.
40            #[must_use]
41            pub fn new(bytes: impl Into<Bytes>) -> Self {
42                Self(bytes.into())
43            }
44
45            /// The raw bytes.
46            #[must_use]
47            pub fn as_bytes(&self) -> &[u8] {
48                &self.0
49            }
50
51            /// The underlying buffer.
52            #[must_use]
53            pub fn into_bytes(self) -> Bytes {
54                self.0
55            }
56        }
57    };
58}
59
60bytes_newtype!(
61    /// A key. Keys order by raw bytes; `store::keys` owns every layout.
62    Key
63);
64bytes_newtype!(
65    /// An opaque value. A backend never interprets it.
66    Value
67);
68bytes_newtype!(
69    /// An opaque scan cursor. Callers only pass back one that a
70    /// [`NamespaceStore::scan`] returned for the same range.
71    Cursor
72);
73
74/// A condition a [`Batch`] checks before it writes: raw key/byte checks,
75/// not the ref CAS. A planner decides a ref write with
76/// [`crate::refs::evaluate_condition`] on the value it read, then guards
77/// that read here (`Absent` or `Equals` on the ref key) so the decision
78/// still holds at commit.
79#[derive(Debug, Clone, PartialEq, Eq)]
80pub enum Precondition {
81    /// The key holds no value. `observed` on failure: the value it holds.
82    Absent(Key),
83    /// The key holds a value. `observed` on failure: `None`.
84    Present(Key),
85    /// The key holds exactly this value. `observed` on failure: the value
86    /// it holds, if any.
87    Equals(Key, Value),
88    /// Commit deadline (Unix ms): the batch commits only if the backend's
89    /// own clock, read inside the check-and-write step, is `<=` it. It
90    /// names no key, so every store accepts it. Failure is
91    /// [`BatchOutcome::DeadlinePassed`], never `PreconditionFailed`. A
92    /// clock reading before the epoch or otherwise invalid fails closed:
93    /// it counts as `u64::MAX`.
94    NotAfter(u64),
95}
96
97/// A write in a [`Batch`].
98#[derive(Debug, Clone, PartialEq, Eq)]
99pub enum Write {
100    /// Set the key's value.
101    Put(Key, Value),
102    /// Remove the key; removing an absent key is not an error.
103    Delete(Key),
104}
105
106/// One atomic, declarative write: preconditions, then puts and deletes in
107/// order (a later write to the same key wins). A batch with no writes only
108/// checks.
109#[derive(Debug, Clone, Default, PartialEq, Eq)]
110pub struct Batch {
111    /// Checked in order; the first failure aborts the batch.
112    pub preconditions: Vec<Precondition>,
113    /// Applied in order if every precondition holds.
114    pub writes: Vec<Write>,
115}
116
117impl Batch {
118    /// An empty batch.
119    #[must_use]
120    pub fn new() -> Self {
121        Self::default()
122    }
123
124    /// Append a precondition.
125    #[must_use]
126    pub fn require(mut self, precondition: Precondition) -> Self {
127        self.preconditions.push(precondition);
128        self
129    }
130
131    /// Append a put.
132    #[must_use]
133    pub fn put(mut self, key: Key, value: Value) -> Self {
134        self.writes.push(Write::Put(key, value));
135        self
136    }
137
138    /// Append a delete.
139    #[must_use]
140    pub fn delete(mut self, key: Key) -> Self {
141        self.writes.push(Write::Delete(key));
142        self
143    }
144
145    /// Whether the batch adds data (holds a put): what a full partition
146    /// rejects with [`StoreError::Full`].
147    #[must_use]
148    pub fn has_put(&self) -> bool {
149        self.writes.iter().any(|w| matches!(w, Write::Put(..)))
150    }
151
152    /// Check the batch against the size limits and a store's capabilities,
153    /// before anything is read or written. Every backend calls this first
154    /// in [`NamespaceStore::apply`].
155    ///
156    /// # Errors
157    /// [`StoreError::Invalid`] for a key over [`MAX_KEY_BYTES`], a value
158    /// over [`MAX_VALUE_BYTES`] (510 KiB for timer payloads), more than
159    /// [`MAX_BATCH_OPS`] operations or
160    /// more than [`MAX_BATCH_BYTES`] in total; [`StoreError::Unsupported`]
161    /// for a key outside `caps.key_classes`, or, without
162    /// `atomic_multi_key`, more than one write or a key precondition on
163    /// another key than the write.
164    pub fn validate(&self, caps: &StoreCapabilities) -> Result<(), StoreError> {
165        if self.preconditions.len() + self.writes.len()
166            > MAX_BATCH_OPS.saturating_sub(caps.reserved_batch_ops)
167        {
168            return Err(StoreError::Invalid("batch exceeds MAX_BATCH_OPS".into()));
169        }
170        let mut keys = Vec::new();
171        for pre in &self.preconditions {
172            match pre {
173                Precondition::Absent(k) | Precondition::Present(k) => keys.push((k, None)),
174                Precondition::Equals(k, v) => keys.push((k, Some(v))),
175                Precondition::NotAfter(_) => {}
176            }
177        }
178        let key_preconditions = keys.len();
179        for write in &self.writes {
180            match write {
181                Write::Put(k, v) => keys.push((k, Some(v))),
182                Write::Delete(k) => keys.push((k, None)),
183            }
184        }
185        let total: usize = keys
186            .iter()
187            .map(|(k, v)| k.as_bytes().len() + v.map_or(0, |v| v.as_bytes().len()))
188            .sum();
189        if total > MAX_BATCH_BYTES {
190            return Err(StoreError::Invalid("batch exceeds MAX_BATCH_BYTES".into()));
191        }
192        for (key, value) in &keys {
193            if key.as_bytes().len() > MAX_KEY_BYTES {
194                return Err(StoreError::Invalid("key exceeds MAX_KEY_BYTES".into()));
195            }
196            if value.is_some_and(|v| v.as_bytes().len() > MAX_VALUE_BYTES) {
197                return Err(StoreError::Invalid("value exceeds MAX_VALUE_BYTES".into()));
198            }
199            if value.is_some_and(|v| v.as_bytes().len() > MAX_TIMER_VALUE_BYTES)
200                && key.as_bytes().starts_with(b"w\0")
201                && matches!(keys::parse(key), Some(keys::ParsedKey::Timer { .. }))
202            {
203                return Err(StoreError::Invalid(
204                    "timer value exceeds retry batch allowance".into(),
205                ));
206            }
207            if caps.key_classes == KeyClasses::RefsOnly && !keys::is_ref_key(key) {
208                return Err(StoreError::Unsupported(
209                    "this store holds only ref keys".into(),
210                ));
211            }
212        }
213        if !caps.atomic_multi_key {
214            let (pre, writes) = keys.split_at(key_preconditions);
215            let one_key = match (pre, writes) {
216                ([], [] | [_]) | ([_], []) => true,
217                ([(p, _)], [(w, _)]) => p == w,
218                _ => false,
219            };
220            if !one_key {
221                return Err(StoreError::Unsupported(
222                    "this store commits at most one key per batch".into(),
223                ));
224            }
225        }
226        Ok(())
227    }
228}
229
230/// The result of [`NamespaceStore::apply`].
231#[derive(Debug, Clone, PartialEq, Eq)]
232pub enum BatchOutcome {
233    /// Every precondition held and every write is durable.
234    Committed,
235    /// `preconditions[index]`, a key precondition, failed; nothing was
236    /// written.
237    PreconditionFailed {
238        /// Index of the first failing precondition.
239        index: usize,
240        /// What the store saw (see [`Precondition`]).
241        observed: Option<Value>,
242    },
243    /// A [`Precondition::NotAfter`] deadline had passed on the backend's
244    /// clock (normative rule 8); nothing was written.
245    DeadlinePassed {
246        /// The backend's clock reading, Unix ms (`u64::MAX` if invalid).
247        backend_now: u64,
248    },
249}
250
251/// One page of a [`NamespaceStore::scan`].
252#[derive(Debug, Clone, PartialEq, Eq, Default)]
253pub struct ScanPage {
254    /// Entries in ascending key order.
255    pub entries: Vec<(Key, Value)>,
256    /// Resume point, if the range may hold more entries.
257    pub next: Option<Cursor>,
258}
259
260/// One ordered range in a batched scan. Its cursor belongs to this exact
261/// range, just as for [`NamespaceStore::scan`].
262#[derive(Debug, Clone, PartialEq, Eq)]
263#[non_exhaustive]
264pub struct RangeScan {
265    /// Inclusive lower bound.
266    pub start: Key,
267    /// Exclusive upper bound.
268    pub end: Key,
269    /// Resume strictly after this cursor.
270    pub after: Option<Cursor>,
271    /// Requested maximum entries, at least one.
272    pub limit: u32,
273}
274
275impl RangeScan {
276    /// A range request, including its optional continuation cursor.
277    #[must_use]
278    pub fn new(start: Key, end: Key, after: Option<Cursor>, limit: u32) -> Self {
279        Self {
280            start,
281            end,
282            after,
283            limit,
284        }
285    }
286}
287
288/// Most ranges accepted by one [`NamespaceStore::scan_many`] call.
289pub const MAX_SCAN_RANGES: usize = 256;
290
291/// Which key classes (`store::keys`) a store accepts.
292#[derive(Debug, Clone, Copy, PartialEq, Eq)]
293#[non_exhaustive]
294pub enum KeyClasses {
295    /// Every class.
296    All,
297    /// Refs only (the `r` class): `FsLayoutStore`.
298    RefsOnly,
299}
300
301/// How repo membership of a pack is decided (overview Q16).
302#[derive(Debug, Clone, Copy, PartialEq, Eq)]
303pub enum MembershipMode {
304    /// A pack in the blob store is a member (M0).
305    StorePresence,
306    /// Only packs recorded by an `AdvanceRefs` apply are members (M1+).
307    Explicit,
308}
309
310/// What a store supports. The pipeline plans around it. Start from
311/// [`StoreCapabilities::full`] or [`StoreCapabilities::refs_only`] and set
312/// fields: the struct is `#[non_exhaustive]`, so later fields default
313/// safely.
314#[derive(Debug, Clone, Copy, PartialEq, Eq)]
315#[non_exhaustive]
316pub struct StoreCapabilities {
317    /// `false`: a batch holds at most one write, and at most one key
318    /// precondition, on that same key (plus any `NotAfter`). The pipeline
319    /// then issues sequential batches.
320    pub atomic_multi_key: bool,
321    /// Operations reserved for a `RefIndex` target-local atomic apply extension.
322    pub reserved_batch_ops: usize,
323    /// Which key classes the store accepts.
324    pub key_classes: KeyClasses,
325    /// How membership is decided.
326    pub membership: MembershipMode,
327    /// The layout version of a store that cannot hold the `v` key
328    /// (`RefsOnly`): its format is versioned elsewhere, and planners skip
329    /// the layout-version precondition. `None` for stores that hold `v`.
330    pub implicit_layout_version: Option<u32>,
331}
332
333impl StoreCapabilities {
334    /// A full store: atomic multi-key batches over every class.
335    #[must_use]
336    pub const fn full() -> Self {
337        Self {
338            atomic_multi_key: true,
339            reserved_batch_ops: 0,
340            key_classes: KeyClasses::All,
341            membership: MembershipMode::StorePresence,
342            implicit_layout_version: None,
343        }
344    }
345
346    /// A refs-only, single-key store with an implicit layout version.
347    #[must_use]
348    pub const fn refs_only() -> Self {
349        Self {
350            atomic_multi_key: false,
351            reserved_batch_ops: 0,
352            key_classes: KeyClasses::RefsOnly,
353            membership: MembershipMode::StorePresence,
354            implicit_layout_version: Some(keys::LAYOUT_VERSION),
355        }
356    }
357}
358
359/// Storage used by one partition.
360#[derive(Debug, Clone, Copy, PartialEq, Eq)]
361pub struct PartitionStats {
362    /// Bytes used; may be approximate.
363    pub bytes: u64,
364    /// Number of keys, if the backend knows it cheaply.
365    pub keys: Option<u64>,
366}
367
368/// The metadata store: a partitioned, ordered key-value store. Any backend
369/// that can do an atomic conditional multi-key write per partition and an
370/// ordered range read can implement it: SQL is not required.
371///
372/// # Normative rules
373///
374/// 1. **Declarative batch.** [`Self::apply`] is the only write. A backend
375///    never evaluates CAS, quota or replay logic: the pipeline plans every
376///    write and guards every value it read with a precondition. Size limits
377///    ([`Batch::validate`]) are checked first, and a violation writes
378///    nothing.
379/// 2. **Reads are `get`, `has`, `get_many`, `scan` and `scan_many`.** No other query
380///    exists; every index is a key layout (`store::keys`).
381/// 3. **Single writer is enough.** Nothing may assume two `apply` calls on
382///    one partition run concurrently, and nothing may hold a lock across an
383///    `.await` waiting for another `apply`. The pipeline handles contention
384///    with an optimistic re-plan loop.
385/// 4. **Cancellation safety.** Dropping an `apply` future at any `.await`
386///    leaves the partition fully before or fully after the batch, and the
387///    store usable: check-and-write runs in one non-yielding step (a
388///    synchronous transaction that runs to completion even if the future is
389///    dropped, a Durable Object `transactionSync`, a mutex held without
390///    awaits). The dropped batch may still commit later (a blocking task
391///    keeps running): every later observation is exactly the state before
392///    or after it, never torn, and once the after state has been observed
393///    the before state never reappears. A poisoned lock is recovered, never
394///    propagated.
395/// 5. **Durability.** `Committed` means durable to the level the backend
396///    documents.
397/// 6. **Atomicity scope.** One batch is one partition; nothing needs
398///    atomicity across partitions. Cross-partition effects are ordered by
399///    the pipeline or carried by outbox rows a core-owned relay delivers at
400///    least once, idempotently. A backend needs nothing beyond this trait
401///    and a way to run timers.
402/// 7. **Bounded growth.** The core deletes what it no longer needs through
403///    ordinary batches; a backend reclaims deleted keys and reports
404///    [`Self::stats`]. At its cap it returns [`StoreError::Full`] for
405///    batches that add data and keeps serving reads and deletes. A batch
406///    whose writes are all deletes (preconditions allowed) never returns
407///    `Full`, so on `Full` the caller retries pruning as a delete-only
408///    batch.
409/// 8. **Commit deadline.** [`Precondition::NotAfter`] is evaluated against
410///    the **backend's own clock** (the `SQLite` host's, the Durable
411///    Object's, an injected one in memory), read once inside the same
412///    non-yielding check-and-write step as the key checks, never the
413///    caller's clock; an invalid reading fails closed. A late batch
414///    therefore cannot commit after its deadline, whenever it arrives.
415///    A miss is [`BatchOutcome::DeadlinePassed`]: the pipeline answers a
416///    retryable `unavailable`, never `aborted`, `deadline_exceeded` or
417///    `resource_exhausted`, and may re-plan the write once within the
418///    envelope's validity (SPEC-WRITE-GRANTS §5.5). Callers set the
419///    deadline with a margin that covers the skew between their clock and
420///    the backend's **plus** the longest synchronous span between reading
421///    the clock and committing: on Workers `Date.now()` does not advance
422///    during synchronous execution, so the backend's reading can lag real
423///    time by that span.
424pub trait NamespaceStore: MaybeSend + MaybeSync {
425    /// What this store supports.
426    fn capabilities(&self) -> StoreCapabilities;
427
428    /// The value at `key`, if any.
429    fn get(
430        &self,
431        p: &Partition,
432        key: &Key,
433    ) -> impl Future<Output = Result<Option<Value>, StoreError>> + MaybeSend;
434
435    /// Whether `key` holds a value.
436    fn has(
437        &self,
438        p: &Partition,
439        key: &Key,
440    ) -> impl Future<Output = Result<bool, StoreError>> + MaybeSend {
441        async move { Ok(self.get(p, key).await?.is_some()) }
442    }
443
444    /// Several keys in one round trip; results in input order. The default
445    /// issues sequential [`Self::get`] calls.
446    fn get_many(
447        &self,
448        p: &Partition,
449        keys: &[Key],
450    ) -> impl Future<Output = Result<Vec<Option<Value>>, StoreError>> + MaybeSend {
451        async move {
452            let mut values = Vec::with_capacity(keys.len());
453            for key in keys {
454                values.push(self.get(p, key).await?);
455            }
456            Ok(values)
457        }
458    }
459
460    /// Up to `limit` (at least 1) entries in `[start, end)`, ascending by
461    /// key bytes. `after` resumes strictly after the cursor's position; a
462    /// cursor outside `[start, end)` (forged, or from another range) is
463    /// [`StoreError::Invalid`]. A page may hold fewer than `limit` entries
464    /// and still return `next`: callers page until `next` is `None`.
465    fn scan(
466        &self,
467        p: &Partition,
468        start: &Key,
469        end: &Key,
470        after: Option<&Cursor>,
471        limit: u32,
472    ) -> impl Future<Output = Result<ScanPage, StoreError>> + MaybeSend;
473
474    /// Scan a served prefix of `ranges` in order. A nonempty request returns
475    /// at least its first page and at most one page per range; callers
476    /// re-request any unserved suffix. Every returned page obeys [`Self::scan`].
477    /// The default serves every range sequentially.
478    fn scan_many(
479        &self,
480        p: &Partition,
481        ranges: &[RangeScan],
482    ) -> impl Future<Output = Result<Vec<ScanPage>, StoreError>> + MaybeSend {
483        async move {
484            if ranges.len() > MAX_SCAN_RANGES {
485                return Err(StoreError::Invalid("too many scan ranges".into()));
486            }
487            let mut pages = Vec::with_capacity(ranges.len());
488            for range in ranges {
489                pages.push(
490                    self.scan(
491                        p,
492                        &range.start,
493                        &range.end,
494                        range.after.as_ref(),
495                        range.limit,
496                    )
497                    .await?,
498                );
499            }
500            Ok(pages)
501        }
502    }
503
504    /// One atomic, all-or-nothing batch: validate it, read the backend
505    /// clock once, check every precondition in order against committed
506    /// state (`NotAfter` against that reading); on the first failure return
507    /// [`BatchOutcome::PreconditionFailed`] or
508    /// [`BatchOutcome::DeadlinePassed`] and write nothing, otherwise apply
509    /// every write and return [`BatchOutcome::Committed`].
510    fn apply(
511        &self,
512        p: &Partition,
513        batch: Batch,
514    ) -> impl Future<Output = Result<BatchOutcome, StoreError>> + MaybeSend;
515
516    /// Storage used by one partition; may be approximate or up to 60 s
517    /// stale.
518    fn stats(
519        &self,
520        p: &Partition,
521    ) -> impl Future<Output = Result<PartitionStats, StoreError>> + MaybeSend;
522
523    /// A cheap health check.
524    fn probe(&self) -> impl Future<Output = Result<(), StoreError>> + MaybeSend;
525}
526
527// Share a store handle across response streams and retained finalizers.
528impl<S: NamespaceStore + ?Sized> NamespaceStore for std::sync::Arc<S> {
529    fn capabilities(&self) -> StoreCapabilities {
530        (**self).capabilities()
531    }
532    async fn get(&self, p: &Partition, key: &Key) -> Result<Option<Value>, StoreError> {
533        (**self).get(p, key).await
534    }
535    async fn has(&self, p: &Partition, key: &Key) -> Result<bool, StoreError> {
536        (**self).has(p, key).await
537    }
538    async fn get_many(
539        &self,
540        p: &Partition,
541        keys: &[Key],
542    ) -> Result<Vec<Option<Value>>, StoreError> {
543        (**self).get_many(p, keys).await
544    }
545    async fn scan(
546        &self,
547        p: &Partition,
548        start: &Key,
549        end: &Key,
550        after: Option<&Cursor>,
551        limit: u32,
552    ) -> Result<ScanPage, StoreError> {
553        (**self).scan(p, start, end, after, limit).await
554    }
555    async fn scan_many(
556        &self,
557        p: &Partition,
558        ranges: &[RangeScan],
559    ) -> Result<Vec<ScanPage>, StoreError> {
560        (**self).scan_many(p, ranges).await
561    }
562    async fn apply(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, StoreError> {
563        (**self).apply(p, batch).await
564    }
565    async fn stats(&self, p: &Partition) -> Result<PartitionStats, StoreError> {
566        (**self).stats(p).await
567    }
568    async fn probe(&self) -> Result<(), StoreError> {
569        (**self).probe().await
570    }
571}