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
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
//! `SqlKvStore`: the [`NamespaceStore`] contract over one `kv` table.

use std::collections::HashMap;
use std::fmt;
use std::sync::atomic::{AtomicU32, Ordering};
use std::sync::{Mutex, PoisonError};

use super::{Capacity, Row, SqlConn, SqlError, SqlValue, TxFn, blob, count, schema};
use crate::store::{
    Batch, BatchOutcome, Cursor, Key, NamespaceStore, Partition, PartitionStats, Precondition,
    ScanPage, StoreCapabilities, StoreError, StoreMaintenance, Value, Write, codec, keys,
};

/// Keys per `get_many` statement: one partition parameter plus this many
/// keys stays far below [`super::MAX_BOUND_PARAMS`].
pub const GET_MANY_CHUNK: usize = 64;

/// How long a partition's [`NamespaceStore::stats`] is cached, ms.
const STATS_TTL_MS: u64 = 60_000;
/// Most partitions whose stats are cached at once.
const STATS_CACHE_MAX: usize = 4096;

pub(super) const GET: &str = "SELECT value FROM kv WHERE part = ?1 AND key = ?2";
pub(super) const PUT: &str = "INSERT INTO kv (part, key, value) VALUES (?1, ?2, ?3) \
     ON CONFLICT (part, key) DO UPDATE SET value = excluded.value";
pub(super) const DELETE: &str = "DELETE FROM kv WHERE part = ?1 AND key = ?2";
pub(super) const SCAN_FROM: &str = "SELECT key, value FROM kv \
     WHERE part = ?1 AND key >= ?2 AND key < ?3 ORDER BY key LIMIT ?4";
pub(super) const SCAN_AFTER: &str = "SELECT key, value FROM kv \
     WHERE part = ?1 AND key > ?2 AND key < ?3 ORDER BY key LIMIT ?4";
pub(super) const STATS: &str =
    "SELECT COUNT(*), SUM(length(key) + length(value)) FROM kv WHERE part = ?1";
/// Bounded first window over the existing timer index, including raw malformed keys.
pub const TIMER_WINDOW_START: &str = "SELECT part, key FROM kv INDEXED BY kv_timers \
     WHERE key >= x'7700' AND key < x'7701' ORDER BY key, part LIMIT ?1";
/// Resume strictly after an index position; no aggregation or implicit wrap.
// Keep the composite cursor first: SQLite otherwise chooses the fixed prefix
// lower bound and filters every preceding timer instead of seeking the cursor.
pub const TIMER_WINDOW_AFTER: &str = "SELECT part, key FROM kv INDEXED BY kv_timers \
     WHERE (key, part) > (?1, ?2) AND key >= x'7700' AND key < x'7701' \
     ORDER BY key, part LIMIT ?3";

/// One raw timer-index row and the exclusive position for its successor.
///
/// Advancing by the raw key, including unregistered kinds and malformed timer
/// keys, prevents bounded enumeration from restarting at the same unknown row.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TimerCursor {
    /// Raw timer key; parsing belongs to the driver, after accounting the row.
    pub key: Key,
    /// Logical partition sharing the physical SQL database.
    pub partition: Partition,
}

pub(super) const PROBE: &str = "SELECT 1";

/// The `get_many` statement for `n` keys (`?2` … `?{n+1}`).
pub(super) fn get_many_sql(n: usize) -> String {
    let marks: Vec<String> = (2..n + 2).map(|i| format!("?{i}")).collect();
    format!(
        "SELECT key, value FROM kv WHERE part = ?1 AND key IN ({})",
        marks.join(", ")
    )
}

/// The [`NamespaceStore`] contract (and [`StoreMaintenance`]) over any
/// [`SqlConn`]: full capabilities, every batch one [`SqlConn::transaction`].
///
/// With a [`Capacity`], an ordinary batch holding a put returns
/// [`StoreError::Full`] once the database uses [`Capacity::soft_limit`] bytes
/// or more, checked inside its transaction. Delete-only batches and bounded,
/// guarded relay-scan checkpoints, relay timer reschedules, and guarded timer
/// retry moves may use the reserve above that limit. These timer exceptions
/// preserve immediate relay rescheduling and persisted infrastructure backoff
/// on a full shard.
/// An engine `Full` on a delete-only batch is reported as
/// [`StoreError::Unavailable`], never `Full`.
///
/// Every method is synchronous inside: its future completes on first poll.
/// A host that drives a blocking engine should run each call on a
/// blocking thread.
pub struct SqlKvStore<C> {
    conn: C,
    schema_version: AtomicU32,
    stats: Mutex<HashMap<Vec<u8>, (u64, PartitionStats)>>,
    capacity: Option<Capacity>,
}

impl<C> fmt::Debug for SqlKvStore<C> {
    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
        f.debug_struct("SqlKvStore")
            .field(
                "schema_version",
                &self.schema_version.load(Ordering::Relaxed),
            )
            .field("capacity", &self.capacity)
            .finish_non_exhaustive()
    }
}

impl<C: SqlConn> SqlKvStore<C> {
    /// A store over `conn`, after migrating its schema
    /// ([`schema::migrate`]).
    ///
    /// # Errors
    /// [`StoreError::Unsupported`] for a database with a newer schema than
    /// this binary's; the engine's error otherwise.
    pub fn open(conn: C) -> Result<Self, StoreError> {
        let version = schema::migrate(&conn)?;
        Ok(Self::from_checked_conn(conn, version))
    }

    /// Open an existing, matching schema without migrating or writing it.
    /// Used by native export so a backup operation cannot mutate its source.
    ///
    /// # Errors
    /// The existing schema must exactly match this binary's version.
    pub fn open_existing(conn: C) -> Result<Self, StoreError> {
        let version = schema::require_current(&conn)?;
        Ok(Self::from_checked_conn(conn, version))
    }

    fn from_checked_conn(conn: C, version: u32) -> Self {
        Self {
            conn,
            schema_version: AtomicU32::new(version),
            stats: Mutex::default(),
            capacity: None,
        }
    }

    /// [`Self::open`], capped at `capacity`: the connection enforces the
    /// hard cap ([`SqlConn::set_size_limit`]) and the store the soft limit.
    ///
    /// # Errors
    /// As [`Self::open`].
    pub fn open_with_capacity(conn: C, capacity: Capacity) -> Result<Self, StoreError> {
        let mut store = Self::open(conn)?;
        store.conn.set_size_limit(capacity.cap_bytes())?;
        store.capacity = Some(capacity);
        Ok(store)
    }

    /// Apply with a synchronous target-local extension inside the guarded transaction.
    /// The extension reads local keys and appends writes only after all input guards hold.
    ///
    /// # Errors
    /// Invalid batches, extension errors, or backend failures roll back every effect.
    pub fn apply_extended<F>(
        &self,
        p: &Partition,
        batch: Batch,
        extend: F,
    ) -> Result<BatchOutcome, StoreError>
    where
        F: FnOnce(
                &dyn Fn(&Key) -> Result<Option<Value>, SqlError>,
                &mut Batch,
                u64,
            ) -> Result<(), SqlError>
            + 'static,
    {
        batch.validate(&self.capabilities())?;
        let part = part(p)?;
        let adds = batch.has_put();
        let soft_limit = self.capacity.map(|c| c.soft_limit());
        match self.conn.transaction(Box::new(move |conn| {
            check_and_write_extended(&conn, &part, batch, soft_limit, extend)
        })) {
            Ok(outcome) => Ok(outcome),
            Err(SqlError::Full) if !adds => Err(StoreError::unavailable(
                "database full during a delete-only batch",
            )),
            Err(error) => Err(error.into()),
        }
    }

    /// Read at most `limit` raw timer rows in `(key, partition)` index order.
    ///
    /// No row is fetched beyond the allowance, and several rows may name the
    /// same partition. An empty window marks the end; the caller decides when
    /// to wrap by passing `None`. The cursor need not still exist in storage.
    ///
    /// # Errors
    /// Zero limit, failed queries, corrupt partition encodings or result columns.
    pub fn timer_window(
        &self,
        after: Option<&TimerCursor>,
        limit: u32,
    ) -> Result<Vec<TimerCursor>, StoreError> {
        if limit == 0 {
            return Err(StoreError::Invalid(
                "timer window limit must be positive".into(),
            ));
        }
        let (sql, params) = match after {
            None => (
                TIMER_WINDOW_START,
                vec![SqlValue::Integer(i64::from(limit))],
            ),
            Some(cursor) => (
                TIMER_WINDOW_AFTER,
                vec![
                    SqlValue::Blob(cursor.key.as_bytes().to_vec()),
                    SqlValue::Blob(cursor.partition.encode()?.to_vec()),
                    SqlValue::Integer(i64::from(limit)),
                ],
            ),
        };
        self.conn
            .query(sql, &params)?
            .into_iter()
            .map(|mut row| {
                Ok(TimerCursor {
                    partition: Partition::decode(&blob(&mut row, 0)?)?,
                    key: Key::new(blob(&mut row, 1)?),
                })
            })
            .collect()
    }

    /// The cap, if any.
    #[must_use]
    pub fn capacity(&self) -> Option<Capacity> {
        self.capacity
    }

    /// The connection.
    #[must_use]
    pub fn conn(&self) -> &C {
        &self.conn
    }

    /// Forget cached [`NamespaceStore::stats`], so the next call reads the
    /// table.
    pub fn clear_stats_cache(&self) {
        self.stats
            .lock()
            .unwrap_or_else(PoisonError::into_inner)
            .clear();
    }
}

fn part(p: &Partition) -> Result<SqlValue, StoreError> {
    Ok(SqlValue::Blob(p.encode()?.to_vec()))
}

fn key_param(key: &Key) -> SqlValue {
    SqlValue::Blob(key.as_bytes().to_vec())
}

fn read<C: SqlConn>(conn: &C, part: &SqlValue, key: &Key) -> Result<Option<Value>, SqlError> {
    let mut rows = conn.query(GET, &[part.clone(), key_param(key)])?;
    rows.first_mut()
        .map(|row| blob(row, 0).map(Value::new))
        .transpose()
}

/// The transaction body of `apply`: read the backend clock once, check the
/// preconditions in order, refuse a batch with a put at the soft limit,
/// then write.
fn check_and_write<C: SqlConn>(
    conn: &C,
    part: &SqlValue,
    batch: Batch,
    soft_limit: Option<u64>,
) -> Result<BatchOutcome, SqlError> {
    check_and_write_extended(conn, part, batch, soft_limit, |_, _, _| Ok(()))
}

fn check_and_write_extended<C: SqlConn, F>(
    conn: &C,
    part: &SqlValue,
    mut batch: Batch,
    soft_limit: Option<u64>,
    extend: F,
) -> Result<BatchOutcome, SqlError>
where
    F: FnOnce(
        &dyn Fn(&Key) -> Result<Option<Value>, SqlError>,
        &mut Batch,
        u64,
    ) -> Result<(), SqlError>,
{
    // Rule 8: the backend's own clock, read inside the transaction.
    let now = conn.now_ms();
    for (index, pre) in batch.preconditions.iter().enumerate() {
        let (holds, observed) = match pre {
            Precondition::NotAfter(deadline) if now > *deadline => {
                return Ok(BatchOutcome::DeadlinePassed { backend_now: now });
            }
            Precondition::NotAfter(_) => continue,
            Precondition::Absent(key) => {
                let current = read(conn, part, key)?;
                (current.is_none(), current)
            }
            Precondition::Present(key) => (read(conn, part, key)?.is_some(), None),
            Precondition::Equals(key, value) => {
                let current = read(conn, part, key)?;
                (current.as_ref() == Some(value), current)
            }
        };
        if !holds {
            return Ok(BatchOutcome::PreconditionFailed { index, observed });
        }
    }
    extend(&|key| read(conn, part, key), &mut batch, now)?;
    batch
        .validate(&StoreCapabilities::full())
        .map_err(|_| SqlError::Corrupt("extended batch exceeds limits"))?;
    if let Some(limit) = soft_limit
        && batch.has_put()
        && !is_relay_scan_checkpoint(&batch)
        && !is_relay_timer_reschedule(&batch)
        && !is_timer_retry_move(&batch)
        && conn.size_bytes()? >= limit
    {
        return Err(SqlError::Full);
    }
    for write in batch.writes {
        match write {
            Write::Put(key, value) => {
                let value = SqlValue::Blob(value.into_bytes().into());
                conn.exec(PUT, &[part.clone(), key_param(&key), value])?
            }
            Write::Delete(key) => conn.exec(DELETE, &[part.clone(), key_param(&key)])?,
        };
    }
    Ok(BatchOutcome::Committed)
}

// The SQL reserve also belongs to source cleanup. Relay scan progress is a
// bounded control row, guarded by its prior value and committed with only
// relay-row deletions. Exclude exactly that shape from the soft put cutoff;
// malformed or unrelated writes still receive Full. The hard engine limit
// remains enforced by SqlConn for every batch.
fn is_relay_scan_checkpoint(batch: &Batch) -> bool {
    let scan_key = keys::relay_scan();
    let guarded = batch.preconditions.iter().any(|pre| {
        matches!(pre, Precondition::Absent(key) | Precondition::Equals(key, _) if key == &scan_key)
    });
    let mut scan_puts = 0;
    let only_cleanup = batch.writes.iter().all(|write| match write {
        Write::Put(key, value) if key == &scan_key => {
            scan_puts += 1;
            codec::decode_relay_scan(value).is_ok()
        }
        Write::Delete(key) => matches!(keys::parse(key), Some(keys::ParsedKey::Relay(_))),
        Write::Put(..) => false,
    });
    guarded && only_cleanup && scan_puts == 1
}

// Moving a relay timer after progress on a full source keeps its next fire
// immediate instead of taking the runner's 5-second retry backoff. The
// guarded old row is deleted before the empty replacement; this exception
// cannot create another timer or write data.
pub(super) fn is_relay_timer_reschedule(batch: &Batch) -> bool {
    let (
        [
            Precondition::Equals(old_key, old_value),
            Precondition::Absent(absent),
        ],
        [Write::Delete(deleted), Write::Put(new_key, new_value)],
    ) = (batch.preconditions.as_slice(), batch.writes.as_slice())
    else {
        return false;
    };
    if old_key != deleted
        || absent != new_key
        || !old_value.as_bytes().is_empty()
        || !new_value.as_bytes().is_empty()
    {
        return false;
    }
    let (
        Some(keys::ParsedKey::Timer {
            due_at_ms: old_due,
            kind: old_kind,
            reference: old_ref,
        }),
        Some(keys::ParsedKey::Timer {
            due_at_ms: new_due,
            kind: new_kind,
            reference: new_ref,
        }),
    ) = (keys::parse(old_key), keys::parse(new_key))
    else {
        return false;
    };
    old_kind == crate::timers::registry::kinds::RELAY.get()
        && new_kind == old_kind
        && old_ref == new_ref
        && new_due > old_due
}

// A guarded infrastructure retry only relocates one opaque timer payload.
// The destination guard prevents collision, including unknown timer kinds.
pub(super) fn is_timer_retry_move(batch: &Batch) -> bool {
    let (
        [
            Precondition::Equals(old_key, old_value),
            Precondition::Absent(absent),
        ],
        [Write::Delete(deleted), Write::Put(new_key, new_value)],
    ) = (batch.preconditions.as_slice(), batch.writes.as_slice())
    else {
        return false;
    };
    if old_key != deleted || absent != new_key || old_value != new_value {
        return false;
    }
    let (
        Some(keys::ParsedKey::Timer {
            due_at_ms: old_due,
            kind: old_kind,
            reference: old_ref,
        }),
        Some(keys::ParsedKey::Timer {
            due_at_ms: new_due,
            kind: new_kind,
            reference: new_ref,
        }),
        Some((old_original, old_attempt)),
        Some((new_original, new_attempt)),
    ) = (
        keys::parse(old_key),
        keys::parse(new_key),
        keys::timer_retry_state(old_key),
        keys::timer_retry_state(new_key),
    )
    else {
        return false;
    };
    new_due > old_due
        && new_kind == old_kind
        && new_ref == old_ref
        && new_original == old_original
        && new_attempt
            == old_attempt
                .saturating_add(1)
                .min(keys::MAX_TIMER_RETRY_ATTEMPT)
}

fn entry(mut row: Row) -> Result<(Key, Value), SqlError> {
    Ok((Key::new(blob(&mut row, 0)?), Value::new(blob(&mut row, 1)?)))
}

impl<C: SqlConn> NamespaceStore for SqlKvStore<C> {
    fn capabilities(&self) -> StoreCapabilities {
        StoreCapabilities::full()
    }

    async fn get(&self, p: &Partition, key: &Key) -> Result<Option<Value>, StoreError> {
        Ok(read(&self.conn, &part(p)?, key)?)
    }

    async fn get_many(
        &self,
        p: &Partition,
        keys: &[Key],
    ) -> Result<Vec<Option<Value>>, StoreError> {
        let part = part(p)?;
        let mut found: HashMap<Vec<u8>, Value> = HashMap::new();
        for chunk in keys.chunks(GET_MANY_CHUNK) {
            let params: Vec<SqlValue> = core::iter::once(part.clone())
                .chain(chunk.iter().map(key_param))
                .collect();
            for row in self.conn.query(&get_many_sql(chunk.len()), &params)? {
                let (key, value) = entry(row)?;
                found.insert(key.into_bytes().into(), value);
            }
        }
        Ok(keys
            .iter()
            .map(|key| found.get(key.as_bytes()).cloned())
            .collect())
    }

    async fn scan(
        &self,
        p: &Partition,
        start: &Key,
        end: &Key,
        after: Option<&Cursor>,
        limit: u32,
    ) -> Result<ScanPage, StoreError> {
        if limit == 0 {
            return Err(StoreError::Invalid("scan limit must be at least 1".into()));
        }
        let (sql, lower) = match after {
            None => (SCAN_FROM, start.as_bytes()),
            // Every cursor this range returns is one of its keys.
            Some(c) if start.as_bytes() <= c.as_bytes() && c.as_bytes() < end.as_bytes() => {
                (SCAN_AFTER, c.as_bytes())
            }
            Some(_) => {
                return Err(StoreError::Invalid(
                    "scan cursor outside the scanned range".into(),
                ));
            }
        };
        if lower >= end.as_bytes() {
            return Ok(ScanPage::default());
        }
        let params = [
            part(p)?,
            SqlValue::Blob(lower.to_vec()),
            key_param(end),
            SqlValue::Integer(i64::from(limit) + 1),
        ];
        let mut entries = self
            .conn
            .query(sql, &params)?
            .into_iter()
            .map(entry)
            .collect::<Result<Vec<_>, _>>()?;
        let want = usize::try_from(limit).unwrap_or(usize::MAX);
        let next = (entries.len() > want).then(|| {
            entries.truncate(want);
            Cursor::new(entries[want - 1].0.clone().into_bytes())
        });
        Ok(ScanPage { entries, next })
    }

    async fn apply(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, StoreError> {
        batch.validate(&self.capabilities())?;
        let part = part(p)?;
        let adds = batch.has_put();
        let soft_limit = self.capacity.map(|c| c.soft_limit());
        // The batch and the partition move into the owned transaction body.
        let body: TxFn<C, BatchOutcome> =
            Box::new(move |conn: C| check_and_write(&conn, &part, batch, soft_limit));
        match self.conn.transaction(body) {
            Ok(outcome) => Ok(outcome),
            // Rule 7: a delete-only batch never reports `Full`. Reaching
            // the engine limit here means the reserve was too small; the
            // batch rolled back.
            Err(SqlError::Full) if !adds => Err(StoreError::unavailable(
                "database full during a delete-only batch",
            )),
            Err(e) => Err(e.into()),
        }
    }

    async fn stats(&self, p: &Partition) -> Result<PartitionStats, StoreError> {
        let name = p.encode()?.to_vec();
        let now = self.conn.now_ms();
        let fresh = |at: u64| at <= now && now - at < STATS_TTL_MS;
        {
            let cache = self.stats.lock().unwrap_or_else(PoisonError::into_inner);
            if let Some((at, stats)) = cache.get(&name)
                && fresh(*at)
            {
                return Ok(*stats);
            }
        }
        let rows = self.conn.query(STATS, &[SqlValue::Blob(name.clone())])?;
        let row = rows
            .first()
            .ok_or(SqlError::Corrupt("stats returned no row"))?;
        let stats = PartitionStats {
            bytes: count(row, 1)?,
            keys: Some(count(row, 0)?),
        };
        let mut cache = self.stats.lock().unwrap_or_else(PoisonError::into_inner);
        if cache.len() >= STATS_CACHE_MAX {
            cache.retain(|_, (at, _)| fresh(*at));
            if cache.len() >= STATS_CACHE_MAX {
                cache.clear();
            }
        }
        cache.insert(name, (now, stats));
        Ok(stats)
    }

    async fn probe(&self) -> Result<(), StoreError> {
        self.conn.query(PROBE, &[])?;
        Ok(())
    }
}

impl<C: SqlConn> StoreMaintenance for SqlKvStore<C> {
    fn layout_version(&self) -> u32 {
        self.schema_version.load(Ordering::Relaxed)
    }

    async fn migrate(&self) -> Result<u32, StoreError> {
        let version = schema::migrate(&self.conn)?;
        self.schema_version.store(version, Ordering::Relaxed);
        Ok(version)
    }

    /// The connection's backup ([`SqlConn::backup_to`]):
    /// [`StoreError::Unsupported`] unless the backend has one.
    async fn backup_to(&self, dest: &str) -> Result<(), StoreError> {
        self.conn.backup_to(dest)
    }
}