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}