mkit-server 0.5.0

Runtime-agnostic core of the mkit server: operation model, errors, runtime and telemetry vocabulary
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
//! The key-level metadata contract (PRD ยง5.3): a partitioned, ordered
//! key-value store whose only write is one declarative [`Batch`].

use core::future::Future;

use bytes::Bytes;

use super::error::StoreError;
use super::keys;
use super::partition::Partition;
use crate::rt::{MaybeSend, MaybeSync};

// Portable limits. Every backend accepts batches within them and the core
// never plans beyond them. They leave headroom under Durable Object SQLite
// storage (2 MB per row or value, one `transactionSync` per batch; see
// https://developers.cloudflare.com/durable-objects/platform/limits/), so a
// batch that fits here fits every backend.

/// Longest accepted key, in bytes.
pub const MAX_KEY_BYTES: usize = 1024;
/// Longest accepted value, in bytes; below [`MAX_BATCH_BYTES`].
pub const MAX_VALUE_BYTES: usize = 512 * 1024;
/// Most preconditions plus writes in one [`Batch`].
pub const MAX_BATCH_OPS: usize = 100;
/// Most key and value bytes, summed over a [`Batch`]'s preconditions and
/// writes.
pub const MAX_BATCH_BYTES: usize = 1024 * 1024;
/// A timer value must fit Equals+Put while a guarded retry relocates two keys.
/// Four maximum keys (Equals, Absent, Delete, Put) share the existing batch cap.
pub(crate) const MAX_TIMER_VALUE_BYTES: usize = (MAX_BATCH_BYTES - 4 * MAX_KEY_BYTES) / 2;

macro_rules! bytes_newtype {
    ($(#[$doc:meta])* $name:ident) => {
        $(#[$doc])*
        #[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord, Default)]
        pub struct $name(Bytes);

        impl $name {
            /// Wrap `bytes`.
            #[must_use]
            pub fn new(bytes: impl Into<Bytes>) -> Self {
                Self(bytes.into())
            }

            /// The raw bytes.
            #[must_use]
            pub fn as_bytes(&self) -> &[u8] {
                &self.0
            }

            /// The underlying buffer.
            #[must_use]
            pub fn into_bytes(self) -> Bytes {
                self.0
            }
        }
    };
}

bytes_newtype!(
    /// A key. Keys order by raw bytes; `store::keys` owns every layout.
    Key
);
bytes_newtype!(
    /// An opaque value. A backend never interprets it.
    Value
);
bytes_newtype!(
    /// An opaque scan cursor. Callers only pass back one that a
    /// [`NamespaceStore::scan`] returned for the same range.
    Cursor
);

/// A condition a [`Batch`] checks before it writes: raw key/byte checks,
/// not the ref CAS. A planner decides a ref write with
/// [`crate::refs::evaluate_condition`] on the value it read, then guards
/// that read here (`Absent` or `Equals` on the ref key) so the decision
/// still holds at commit.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Precondition {
    /// The key holds no value. `observed` on failure: the value it holds.
    Absent(Key),
    /// The key holds a value. `observed` on failure: `None`.
    Present(Key),
    /// The key holds exactly this value. `observed` on failure: the value
    /// it holds, if any.
    Equals(Key, Value),
    /// Commit deadline (Unix ms): the batch commits only if the backend's
    /// own clock, read inside the check-and-write step, is `<=` it. It
    /// names no key, so every store accepts it. Failure is
    /// [`BatchOutcome::DeadlinePassed`], never `PreconditionFailed`. A
    /// clock reading before the epoch or otherwise invalid fails closed:
    /// it counts as `u64::MAX`.
    NotAfter(u64),
}

/// A write in a [`Batch`].
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Write {
    /// Set the key's value.
    Put(Key, Value),
    /// Remove the key; removing an absent key is not an error.
    Delete(Key),
}

/// One atomic, declarative write: preconditions, then puts and deletes in
/// order (a later write to the same key wins). A batch with no writes only
/// checks.
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct Batch {
    /// Checked in order; the first failure aborts the batch.
    pub preconditions: Vec<Precondition>,
    /// Applied in order if every precondition holds.
    pub writes: Vec<Write>,
}

impl Batch {
    /// An empty batch.
    #[must_use]
    pub fn new() -> Self {
        Self::default()
    }

    /// Append a precondition.
    #[must_use]
    pub fn require(mut self, precondition: Precondition) -> Self {
        self.preconditions.push(precondition);
        self
    }

    /// Append a put.
    #[must_use]
    pub fn put(mut self, key: Key, value: Value) -> Self {
        self.writes.push(Write::Put(key, value));
        self
    }

    /// Append a delete.
    #[must_use]
    pub fn delete(mut self, key: Key) -> Self {
        self.writes.push(Write::Delete(key));
        self
    }

    /// Whether the batch adds data (holds a put): what a full partition
    /// rejects with [`StoreError::Full`].
    #[must_use]
    pub fn has_put(&self) -> bool {
        self.writes.iter().any(|w| matches!(w, Write::Put(..)))
    }

    /// Check the batch against the size limits and a store's capabilities,
    /// before anything is read or written. Every backend calls this first
    /// in [`NamespaceStore::apply`].
    ///
    /// # Errors
    /// [`StoreError::Invalid`] for a key over [`MAX_KEY_BYTES`], a value
    /// over [`MAX_VALUE_BYTES`] (510 KiB for timer payloads), more than
    /// [`MAX_BATCH_OPS`] operations or
    /// more than [`MAX_BATCH_BYTES`] in total; [`StoreError::Unsupported`]
    /// for a key outside `caps.key_classes`, or, without
    /// `atomic_multi_key`, more than one write or a key precondition on
    /// another key than the write.
    pub fn validate(&self, caps: &StoreCapabilities) -> Result<(), StoreError> {
        if self.preconditions.len() + self.writes.len()
            > MAX_BATCH_OPS.saturating_sub(caps.reserved_batch_ops)
        {
            return Err(StoreError::Invalid("batch exceeds MAX_BATCH_OPS".into()));
        }
        let mut keys = Vec::new();
        for pre in &self.preconditions {
            match pre {
                Precondition::Absent(k) | Precondition::Present(k) => keys.push((k, None)),
                Precondition::Equals(k, v) => keys.push((k, Some(v))),
                Precondition::NotAfter(_) => {}
            }
        }
        let key_preconditions = keys.len();
        for write in &self.writes {
            match write {
                Write::Put(k, v) => keys.push((k, Some(v))),
                Write::Delete(k) => keys.push((k, None)),
            }
        }
        let total: usize = keys
            .iter()
            .map(|(k, v)| k.as_bytes().len() + v.map_or(0, |v| v.as_bytes().len()))
            .sum();
        if total > MAX_BATCH_BYTES {
            return Err(StoreError::Invalid("batch exceeds MAX_BATCH_BYTES".into()));
        }
        for (key, value) in &keys {
            if key.as_bytes().len() > MAX_KEY_BYTES {
                return Err(StoreError::Invalid("key exceeds MAX_KEY_BYTES".into()));
            }
            if value.is_some_and(|v| v.as_bytes().len() > MAX_VALUE_BYTES) {
                return Err(StoreError::Invalid("value exceeds MAX_VALUE_BYTES".into()));
            }
            if value.is_some_and(|v| v.as_bytes().len() > MAX_TIMER_VALUE_BYTES)
                && key.as_bytes().starts_with(b"w\0")
                && matches!(keys::parse(key), Some(keys::ParsedKey::Timer { .. }))
            {
                return Err(StoreError::Invalid(
                    "timer value exceeds retry batch allowance".into(),
                ));
            }
            if caps.key_classes == KeyClasses::RefsOnly && !keys::is_ref_key(key) {
                return Err(StoreError::Unsupported(
                    "this store holds only ref keys".into(),
                ));
            }
        }
        if !caps.atomic_multi_key {
            let (pre, writes) = keys.split_at(key_preconditions);
            let one_key = match (pre, writes) {
                ([], [] | [_]) | ([_], []) => true,
                ([(p, _)], [(w, _)]) => p == w,
                _ => false,
            };
            if !one_key {
                return Err(StoreError::Unsupported(
                    "this store commits at most one key per batch".into(),
                ));
            }
        }
        Ok(())
    }
}

/// The result of [`NamespaceStore::apply`].
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum BatchOutcome {
    /// Every precondition held and every write is durable.
    Committed,
    /// `preconditions[index]`, a key precondition, failed; nothing was
    /// written.
    PreconditionFailed {
        /// Index of the first failing precondition.
        index: usize,
        /// What the store saw (see [`Precondition`]).
        observed: Option<Value>,
    },
    /// A [`Precondition::NotAfter`] deadline had passed on the backend's
    /// clock (normative rule 8); nothing was written.
    DeadlinePassed {
        /// The backend's clock reading, Unix ms (`u64::MAX` if invalid).
        backend_now: u64,
    },
}

/// One page of a [`NamespaceStore::scan`].
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct ScanPage {
    /// Entries in ascending key order.
    pub entries: Vec<(Key, Value)>,
    /// Resume point, if the range may hold more entries.
    pub next: Option<Cursor>,
}

/// One ordered range in a batched scan. Its cursor belongs to this exact
/// range, just as for [`NamespaceStore::scan`].
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub struct RangeScan {
    /// Inclusive lower bound.
    pub start: Key,
    /// Exclusive upper bound.
    pub end: Key,
    /// Resume strictly after this cursor.
    pub after: Option<Cursor>,
    /// Requested maximum entries, at least one.
    pub limit: u32,
}

impl RangeScan {
    /// A range request, including its optional continuation cursor.
    #[must_use]
    pub fn new(start: Key, end: Key, after: Option<Cursor>, limit: u32) -> Self {
        Self {
            start,
            end,
            after,
            limit,
        }
    }
}

/// Most ranges accepted by one [`NamespaceStore::scan_many`] call.
pub const MAX_SCAN_RANGES: usize = 256;

/// Which key classes (`store::keys`) a store accepts.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub enum KeyClasses {
    /// Every class.
    All,
    /// Refs only (the `r` class): `FsLayoutStore`.
    RefsOnly,
}

/// How repo membership of a pack is decided (overview Q16).
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum MembershipMode {
    /// A pack in the blob store is a member (M0).
    StorePresence,
    /// Only packs recorded by an `AdvanceRefs` apply are members (M1+).
    Explicit,
}

/// What a store supports. The pipeline plans around it. Start from
/// [`StoreCapabilities::full`] or [`StoreCapabilities::refs_only`] and set
/// fields: the struct is `#[non_exhaustive]`, so later fields default
/// safely.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub struct StoreCapabilities {
    /// `false`: a batch holds at most one write, and at most one key
    /// precondition, on that same key (plus any `NotAfter`). The pipeline
    /// then issues sequential batches.
    pub atomic_multi_key: bool,
    /// Operations reserved for a `RefIndex` target-local atomic apply extension.
    pub reserved_batch_ops: usize,
    /// Which key classes the store accepts.
    pub key_classes: KeyClasses,
    /// How membership is decided.
    pub membership: MembershipMode,
    /// The layout version of a store that cannot hold the `v` key
    /// (`RefsOnly`): its format is versioned elsewhere, and planners skip
    /// the layout-version precondition. `None` for stores that hold `v`.
    pub implicit_layout_version: Option<u32>,
}

impl StoreCapabilities {
    /// A full store: atomic multi-key batches over every class.
    #[must_use]
    pub const fn full() -> Self {
        Self {
            atomic_multi_key: true,
            reserved_batch_ops: 0,
            key_classes: KeyClasses::All,
            membership: MembershipMode::StorePresence,
            implicit_layout_version: None,
        }
    }

    /// A refs-only, single-key store with an implicit layout version.
    #[must_use]
    pub const fn refs_only() -> Self {
        Self {
            atomic_multi_key: false,
            reserved_batch_ops: 0,
            key_classes: KeyClasses::RefsOnly,
            membership: MembershipMode::StorePresence,
            implicit_layout_version: Some(keys::LAYOUT_VERSION),
        }
    }
}

/// Storage used by one partition.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct PartitionStats {
    /// Bytes used; may be approximate.
    pub bytes: u64,
    /// Number of keys, if the backend knows it cheaply.
    pub keys: Option<u64>,
}

/// The metadata store: a partitioned, ordered key-value store. Any backend
/// that can do an atomic conditional multi-key write per partition and an
/// ordered range read can implement it: SQL is not required.
///
/// # Normative rules
///
/// 1. **Declarative batch.** [`Self::apply`] is the only write. A backend
///    never evaluates CAS, quota or replay logic: the pipeline plans every
///    write and guards every value it read with a precondition. Size limits
///    ([`Batch::validate`]) are checked first, and a violation writes
///    nothing.
/// 2. **Reads are `get`, `has`, `get_many`, `scan` and `scan_many`.** No other query
///    exists; every index is a key layout (`store::keys`).
/// 3. **Single writer is enough.** Nothing may assume two `apply` calls on
///    one partition run concurrently, and nothing may hold a lock across an
///    `.await` waiting for another `apply`. The pipeline handles contention
///    with an optimistic re-plan loop.
/// 4. **Cancellation safety.** Dropping an `apply` future at any `.await`
///    leaves the partition fully before or fully after the batch, and the
///    store usable: check-and-write runs in one non-yielding step (a
///    synchronous transaction that runs to completion even if the future is
///    dropped, a Durable Object `transactionSync`, a mutex held without
///    awaits). The dropped batch may still commit later (a blocking task
///    keeps running): every later observation is exactly the state before
///    or after it, never torn, and once the after state has been observed
///    the before state never reappears. A poisoned lock is recovered, never
///    propagated.
/// 5. **Durability.** `Committed` means durable to the level the backend
///    documents.
/// 6. **Atomicity scope.** One batch is one partition; nothing needs
///    atomicity across partitions. Cross-partition effects are ordered by
///    the pipeline or carried by outbox rows a core-owned relay delivers at
///    least once, idempotently. A backend needs nothing beyond this trait
///    and a way to run timers.
/// 7. **Bounded growth.** The core deletes what it no longer needs through
///    ordinary batches; a backend reclaims deleted keys and reports
///    [`Self::stats`]. At its cap it returns [`StoreError::Full`] for
///    batches that add data and keeps serving reads and deletes. A batch
///    whose writes are all deletes (preconditions allowed) never returns
///    `Full`, so on `Full` the caller retries pruning as a delete-only
///    batch.
/// 8. **Commit deadline.** [`Precondition::NotAfter`] is evaluated against
///    the **backend's own clock** (the `SQLite` host's, the Durable
///    Object's, an injected one in memory), read once inside the same
///    non-yielding check-and-write step as the key checks, never the
///    caller's clock; an invalid reading fails closed. A late batch
///    therefore cannot commit after its deadline, whenever it arrives.
///    A miss is [`BatchOutcome::DeadlinePassed`]: the pipeline answers a
///    retryable `unavailable`, never `aborted`, `deadline_exceeded` or
///    `resource_exhausted`, and may re-plan the write once within the
///    envelope's validity (SPEC-WRITE-GRANTS ยง5.5). Callers set the
///    deadline with a margin that covers the skew between their clock and
///    the backend's **plus** the longest synchronous span between reading
///    the clock and committing: on Workers `Date.now()` does not advance
///    during synchronous execution, so the backend's reading can lag real
///    time by that span.
pub trait NamespaceStore: MaybeSend + MaybeSync {
    /// What this store supports.
    fn capabilities(&self) -> StoreCapabilities;

    /// The value at `key`, if any.
    fn get(
        &self,
        p: &Partition,
        key: &Key,
    ) -> impl Future<Output = Result<Option<Value>, StoreError>> + MaybeSend;

    /// Whether `key` holds a value.
    fn has(
        &self,
        p: &Partition,
        key: &Key,
    ) -> impl Future<Output = Result<bool, StoreError>> + MaybeSend {
        async move { Ok(self.get(p, key).await?.is_some()) }
    }

    /// Several keys in one round trip; results in input order. The default
    /// issues sequential [`Self::get`] calls.
    fn get_many(
        &self,
        p: &Partition,
        keys: &[Key],
    ) -> impl Future<Output = Result<Vec<Option<Value>>, StoreError>> + MaybeSend {
        async move {
            let mut values = Vec::with_capacity(keys.len());
            for key in keys {
                values.push(self.get(p, key).await?);
            }
            Ok(values)
        }
    }

    /// Up to `limit` (at least 1) entries in `[start, end)`, ascending by
    /// key bytes. `after` resumes strictly after the cursor's position; a
    /// cursor outside `[start, end)` (forged, or from another range) is
    /// [`StoreError::Invalid`]. A page may hold fewer than `limit` entries
    /// and still return `next`: callers page until `next` is `None`.
    fn scan(
        &self,
        p: &Partition,
        start: &Key,
        end: &Key,
        after: Option<&Cursor>,
        limit: u32,
    ) -> impl Future<Output = Result<ScanPage, StoreError>> + MaybeSend;

    /// Scan a served prefix of `ranges` in order. A nonempty request returns
    /// at least its first page and at most one page per range; callers
    /// re-request any unserved suffix. Every returned page obeys [`Self::scan`].
    /// The default serves every range sequentially.
    fn scan_many(
        &self,
        p: &Partition,
        ranges: &[RangeScan],
    ) -> impl Future<Output = Result<Vec<ScanPage>, StoreError>> + MaybeSend {
        async move {
            if ranges.len() > MAX_SCAN_RANGES {
                return Err(StoreError::Invalid("too many scan ranges".into()));
            }
            let mut pages = Vec::with_capacity(ranges.len());
            for range in ranges {
                pages.push(
                    self.scan(
                        p,
                        &range.start,
                        &range.end,
                        range.after.as_ref(),
                        range.limit,
                    )
                    .await?,
                );
            }
            Ok(pages)
        }
    }

    /// One atomic, all-or-nothing batch: validate it, read the backend
    /// clock once, check every precondition in order against committed
    /// state (`NotAfter` against that reading); on the first failure return
    /// [`BatchOutcome::PreconditionFailed`] or
    /// [`BatchOutcome::DeadlinePassed`] and write nothing, otherwise apply
    /// every write and return [`BatchOutcome::Committed`].
    fn apply(
        &self,
        p: &Partition,
        batch: Batch,
    ) -> impl Future<Output = Result<BatchOutcome, StoreError>> + MaybeSend;

    /// Storage used by one partition; may be approximate or up to 60 s
    /// stale.
    fn stats(
        &self,
        p: &Partition,
    ) -> impl Future<Output = Result<PartitionStats, StoreError>> + MaybeSend;

    /// A cheap health check.
    fn probe(&self) -> impl Future<Output = Result<(), StoreError>> + MaybeSend;
}

// Share a store handle across response streams and retained finalizers.
impl<S: NamespaceStore + ?Sized> NamespaceStore for std::sync::Arc<S> {
    fn capabilities(&self) -> StoreCapabilities {
        (**self).capabilities()
    }
    async fn get(&self, p: &Partition, key: &Key) -> Result<Option<Value>, StoreError> {
        (**self).get(p, key).await
    }
    async fn has(&self, p: &Partition, key: &Key) -> Result<bool, StoreError> {
        (**self).has(p, key).await
    }
    async fn get_many(
        &self,
        p: &Partition,
        keys: &[Key],
    ) -> Result<Vec<Option<Value>>, StoreError> {
        (**self).get_many(p, keys).await
    }
    async fn scan(
        &self,
        p: &Partition,
        start: &Key,
        end: &Key,
        after: Option<&Cursor>,
        limit: u32,
    ) -> Result<ScanPage, StoreError> {
        (**self).scan(p, start, end, after, limit).await
    }
    async fn scan_many(
        &self,
        p: &Partition,
        ranges: &[RangeScan],
    ) -> Result<Vec<ScanPage>, StoreError> {
        (**self).scan_many(p, ranges).await
    }
    async fn apply(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, StoreError> {
        (**self).apply(p, batch).await
    }
    async fn stats(&self, p: &Partition) -> Result<PartitionStats, StoreError> {
        (**self).stats(p).await
    }
    async fn probe(&self) -> Result<(), StoreError> {
        (**self).probe().await
    }
}