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
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
//! Pure reservation and outbox fragments. One guarded `o` row arbitrates
//! Ticketed or Pending -> terminal; delivery removes that row only after
//! acknowledgement.

use std::collections::{BTreeMap, BTreeSet};

use mkit_core::hash::{Hash, to_hex};

use super::codec::{self, Backlog, RelayV1, ReservationV1};
use super::{
    Batch, Key, MAX_BATCH_BYTES, MAX_BATCH_OPS, MAX_KEY_BYTES, MAX_VALUE_BYTES, Partition,
    Precondition, StoreCapabilities, StoreError, Value, Write, keys,
};

/// Each ticket costs at most nine ops: ticket guard/delete, index
/// guard/delete, reservation guard/put, pending-outcome put, membership
/// put and one relay-row share. An advance uses one signer and runs no
/// admission (the quota planner asserts this in `pipeline::plan_namespace`),
/// so `tu` and `tc` are each guarded/written once. Shared
/// overhead is at most 31: publication guard/state 2, retained value 1,
/// two published refs 2, deadline 1, lease guard/install 2, absent layout
/// version guard/install 2, absent repo-known guard/install 2, two ref CAS
/// pairs 4, replay 3, counters 4, outbox sequence/backlog 4, and relay kick
/// 1, ref-index relay rows 2, and one outcome-delivery kick 1. These figures are D34's. On Single, a grant guard replaces the lease
/// pair and there is no relay share or relay kick, so seven tickets cost
/// `8 * 7 + 29 = 85`, including durable authority-generation/mode absence guards.
/// The real maximal planner batches are tested separately. On D34, seven tickets
/// cost `9 * 7 + 31 = 94` ops before
/// opportunistic pruning.
///
/// The same constant caps an implicit transport-identity session's pending
/// packs (WP-1.15 B9): a D34 packmap write consuming all seven — one
/// membership put and one relay row each, plus WP-1.28b's ref-index
/// relay row for the packmap name — plans a 27-op batch
/// (`maximal_implicit_consume_plans_a_valid_batch`).
pub const MAX_TICKETS_PER_ADVANCE: usize = 7;
/// The advance batch's ops outside the per-ticket and per-signer ones.
pub const ADVANCE_SHARED_OPS: usize = 31;
const _: () = assert!(MAX_TICKETS_PER_ADVANCE * 9 + ADVANCE_SHARED_OPS <= MAX_BATCH_OPS);

/// Maximum operations (puts plus deletes) per relay row; two ops guard/advance rh,
/// and two remain for hooks.
pub const MAX_RELAY_PUTS: usize = 96;
const RELAY_DELETE_FIELD_BYTES: usize = 13; // ,"deletes":[]
const _: () = assert!(MAX_RELAY_PUTS + 2 <= MAX_BATCH_OPS);
// The encoded row is at most 512 KiB, leaving room for the worst rh guard/put.
const _: () = assert!(MAX_VALUE_BYTES + 2 * (MAX_KEY_BYTES + 8) <= MAX_BATCH_BYTES);

/// Reservation id for an allowance that has no deployment reservation.
#[must_use]
pub fn synthetic_reservation_id(replay_scope: &Hash) -> String {
    format!("s:{}", to_hex(replay_scope))
}

/// A validated terminal value. A Ticketed value cannot be an outcome.
#[derive(Debug, Clone)]
pub struct Terminal(ReservationV1);

impl Terminal {
    /// Validate a terminal record before planning its replacement.
    pub fn new(record: ReservationV1) -> Result<Self, StoreError> {
        if matches!(
            record,
            ReservationV1::Ticketed { .. } | ReservationV1::Pending { .. }
        ) {
            return Err(StoreError::Invalid("outcome must be terminal".into()));
        }
        codec::decode_reservation(&codec::encode_reservation(&record))?;
        Ok(Self(record))
    }

    fn occurred_at_ms(&self) -> u64 {
        match &self.0 {
            ReservationV1::Committed { occurred_at_ms, .. }
            | ReservationV1::Aborted { occurred_at_ms, .. }
            | ReservationV1::Expired { occurred_at_ms, .. }
            | ReservationV1::ReadServed { occurred_at_ms, .. } => *occurred_at_ms,
            _ => 0,
        }
    }
}

pub(crate) fn guard(key: Key, prior: Option<&Value>) -> Precondition {
    match prior {
        Some(value) => Precondition::Equals(key, value.clone()),
        None => Precondition::Absent(key),
    }
}

fn corrupt(message: &'static str) -> StoreError {
    StoreError::Corrupt(message.into())
}

/// One batch's outbox edits, using a single snapshot of os and oc.
///
/// The fixed infallible fragment methods defer malformed inputs/overflow
/// until finish. `try_finish` reports these errors without changing its
/// output vectors. `finish` instead appends mutually exclusive guards,
/// making the entire caller batch fail closed with no writes applied.
/// Use one builder per batch, finishing it before any acknowledgements.
#[derive(Debug)]
pub struct OutboxBuilder {
    os: Option<Value>,
    oc: Option<Value>,
    seq: u64,
    backlog: Backlog,
    sequence_touched: bool,
    backlog_touched: bool,
    pre: Vec<Precondition>,
    writes: Vec<Write>,
    relays: BTreeMap<Partition, BTreeMap<Key, Option<Value>>>,
    reservations: BTreeSet<String>,
    error: Option<StoreError>,
    relay_at_ms: Option<u64>,
    kick_at_ms: Option<u64>,
}

impl OutboxBuilder {
    /// Read sequence/backlog once. Missing rows mean zero.
    pub fn new(os: Option<&Value>, oc: Option<&Value>) -> Result<Self, StoreError> {
        let seq = os.map(codec::decode_u64).transpose()?.unwrap_or(0);
        if os.is_some() && seq == 0 {
            return Err(corrupt("outbox sequence is zero"));
        }
        Ok(Self {
            os: os.cloned(),
            oc: oc.cloned(),
            seq,
            backlog: oc
                .map(codec::decode_backlog)
                .transpose()?
                .unwrap_or(Backlog { rows: 0, bytes: 0 }),
            sequence_touched: false,
            backlog_touched: false,
            pre: Vec::new(),
            writes: Vec::new(),
            relays: BTreeMap::new(),
            reservations: BTreeSet::new(),
            error: None,
            relay_at_ms: None,
            kick_at_ms: None,
        })
    }

    fn allocate(&mut self) -> Result<u64, StoreError> {
        self.seq = self
            .seq
            .checked_add(1)
            .ok_or_else(|| corrupt("outbox sequence overflow"))?;
        self.sequence_touched = true;
        Ok(self.seq)
    }

    fn remember(&mut self, result: Result<(), StoreError>) {
        if self.error.is_none() {
            self.error = result.err();
        }
    }

    /// Reserve a unique id. An existing row fails Absent at commit even
    /// when it contains the same ticket; callers resolve replay beforehand.
    pub fn reserve(&mut self, rid: &str, ticket_id: [u8; 32], prior: Option<&Value>) {
        let result = (|| {
            let key = keys::reservation(rid)?;
            if !self.reservations.insert(rid.to_owned()) {
                return Err(StoreError::Invalid("duplicate reservation in batch".into()));
            }
            if let Some(value) = prior
                && !matches!(
                    codec::decode_reservation(value)?,
                    ReservationV1::Pending {
                        op: codec::PendingOp::Write,
                        ..
                    }
                )
            {
                return Err(StoreError::Invalid("reservation id already in use".into()));
            }
            self.pre.push(guard(key.clone(), prior));
            self.writes.push(Write::Put(
                key,
                codec::encode_reservation(&ReservationV1::Ticketed { ticket_id }),
            ));
            Ok(())
        })();
        self.remember(result);
    }

    /// Durably record an admitted reservation before any guarded apply.
    /// A present prior is an invalid admission decision, not a replay.
    pub fn pending(&mut self, rid: &str, prior: Option<&Value>, record: &ReservationV1) {
        let result = (|| {
            if prior.is_some() {
                return Err(StoreError::Invalid("reservation id already in use".into()));
            }
            let ReservationV1::Pending {
                reconcile_at_ms, ..
            } = record
            else {
                return Err(StoreError::Invalid(
                    "pending requires Pending record".into(),
                ));
            };
            let key = keys::reservation(rid)?;
            if !self.reservations.insert(rid.to_owned()) {
                return Err(StoreError::Invalid("duplicate reservation in batch".into()));
            }
            let value = codec::encode_reservation(record);
            codec::decode_reservation(&value)?;
            self.pre.push(Precondition::Absent(key.clone()));
            self.writes.push(Write::Put(key, value));
            self.writes.push(Write::Put(
                keys::timer(
                    *reconcile_at_ms,
                    crate::timers::registry::kinds::RESERVATION_RECONCILE.get(),
                    rid.as_bytes(),
                ),
                Value::default(),
            ));
            Ok(())
        })();
        self.remember(result);
    }

    /// Replace still-Ticketed or Pending with exactly one terminal outcome,
    /// queued for delivery. A terminal prior is rejected rather than replaced.
    pub fn outcome(&mut self, rid: &str, prior: &Value, terminal: Terminal) {
        let occurred_at_ms = terminal.occurred_at_ms();
        let record = terminal.0;
        let result = (|| {
            let key = keys::reservation(rid)?;
            let permitted = match (codec::decode_reservation(prior)?, &record) {
                (
                    ReservationV1::Ticketed { .. },
                    ReservationV1::Committed { .. }
                    | ReservationV1::Aborted { .. }
                    | ReservationV1::Expired { .. },
                ) => true,
                (
                    ReservationV1::Pending {
                        repository: prior_repo,
                        op: codec::PendingOp::Write,
                        ..
                    },
                    ReservationV1::Committed { repository, .. }
                    | ReservationV1::Aborted { repository, .. },
                )
                | (
                    ReservationV1::Pending {
                        repository: prior_repo,
                        op: codec::PendingOp::Read,
                        ..
                    },
                    ReservationV1::ReadServed { repository, .. }
                    | ReservationV1::Aborted { repository, .. },
                ) => prior_repo == repository.as_str(),
                _ => false,
            };
            if !permitted {
                return Err(corrupt("illegal reservation transition"));
            }
            if !self.reservations.insert(rid.to_owned()) {
                return Err(StoreError::Invalid("duplicate reservation in batch".into()));
            }
            let value = codec::encode_reservation(&record);
            let size = (key.as_bytes().len() + value.as_bytes().len()) as u64;
            self.backlog.rows = self
                .backlog
                .rows
                .checked_add(1)
                .ok_or_else(|| corrupt("backlog rows overflow"))?;
            self.backlog.bytes = self
                .backlog
                .bytes
                .checked_add(size)
                .ok_or_else(|| corrupt("backlog bytes overflow"))?;
            let seq = self.allocate()?;
            if self.oc.is_none() && self.kick_at_ms.is_none() {
                self.kick_at_ms = Some(occurred_at_ms);
            }
            self.backlog_touched = true;
            self.pre
                .push(Precondition::Equals(key.clone(), prior.clone()));
            self.writes.push(Write::Put(key, value));
            self.writes.push(Write::Put(
                keys::outcome_pending(seq, rid)?,
                Value::default(),
            ));
            Ok(())
        })();
        self.remember(result);
    }

    /// Fail a ticketless streaming reservation in one guarded Absent unit.
    pub fn abort_direct(&mut self, rid: &str, terminal: Terminal) {
        let occurred_at_ms = terminal.occurred_at_ms();
        let record = terminal.0;
        let result = (|| {
            if !matches!(record, ReservationV1::Aborted { .. }) {
                return Err(StoreError::Invalid("direct abort requires Aborted".into()));
            }
            let key = keys::reservation(rid)?;
            if !self.reservations.insert(rid.to_owned()) {
                return Err(StoreError::Invalid("duplicate reservation in batch".into()));
            }
            let value = codec::encode_reservation(&record);
            let size = (key.as_bytes().len() + value.as_bytes().len()) as u64;
            self.backlog.rows = self
                .backlog
                .rows
                .checked_add(1)
                .ok_or_else(|| corrupt("backlog rows overflow"))?;
            self.backlog.bytes = self
                .backlog
                .bytes
                .checked_add(size)
                .ok_or_else(|| corrupt("backlog bytes overflow"))?;
            let seq = self.allocate()?;
            if self.oc.is_none() && self.kick_at_ms.is_none() {
                self.kick_at_ms = Some(occurred_at_ms);
            }
            self.backlog_touched = true;
            self.pre.push(Precondition::Absent(key.clone()));
            self.writes.push(Write::Put(key, value));
            self.writes.push(Write::Put(
                keys::outcome_pending(seq, rid)?,
                Value::default(),
            ));
            Ok(())
        })();
        self.remember(result);
    }

    /// Stamp relay rows and schedule their immediate source-side kick.
    pub fn relay_at(&mut self, now_ms: u64) {
        self.relay_at_ms = Some(now_ms);
    }

    /// Group idempotent upserts by target, sorting keys deterministically.
    /// Conflicting values for one target/key invalidate the whole fragment.
    pub fn relay(&mut self, target: &Partition, puts: Vec<(Key, Value)>) {
        let result = (|| {
            target.encode()?;
            let group = self.relays.entry(target.clone()).or_default();
            for (key, value) in puts {
                if group
                    .get(&key)
                    .is_some_and(|old| old.as_ref() != Some(&value))
                {
                    return Err(StoreError::Invalid("conflicting relay upserts".into()));
                }
                group.insert(key, Some(value));
            }
            Ok(())
        })();
        self.remember(result);
    }

    /// Group idempotent deletes by target. A key cannot be both put and
    /// deleted in the same source batch.
    pub fn relay_delete(&mut self, target: &Partition, keys: Vec<Key>) {
        let result = (|| {
            target.encode()?;
            let group = self.relays.entry(target.clone()).or_default();
            for key in keys {
                if group.get(&key).is_some_and(Option::is_some) {
                    return Err(StoreError::Invalid("relay put/delete overlap".into()));
                }
                group.insert(key, None);
            }
            Ok(())
        })();
        self.remember(result);
    }

    fn plan_relay_rows(
        &mut self,
        target: Partition,
        operations: BTreeMap<Key, Option<Value>>,
    ) -> Result<u64, StoreError> {
        let at_ms = self
            .relay_at_ms
            .ok_or_else(|| StoreError::Invalid("relay rows need relay_at".into()))?;
        let mut row = RelayV1 {
            at_ms,
            target,
            puts: Vec::new(),
            deletes: Vec::new(),
        };
        let base_bytes = codec::encode_relay(&row)?.as_bytes().len();
        let mut encoded_bytes = base_bytes;
        let (puts, deletes): (Vec<_>, Vec<_>) = operations
            .into_iter()
            .partition(|(_, value)| value.is_some());
        for (key, value) in puts {
            let Some(value) = value else {
                return Err(StoreError::Invalid("missing relay upsert value".into()));
            };
            // JSON uses hex strings: ["key","value"], plus a comma after the first.
            let bytes = 7 + 2 * (key.as_bytes().len() + value.as_bytes().len());
            if key.as_bytes().len() > MAX_KEY_BYTES
                || value.as_bytes().len() > MAX_VALUE_BYTES
                || base_bytes + bytes > MAX_VALUE_BYTES
            {
                return Err(StoreError::Invalid(
                    "relay upsert cannot fit one row".into(),
                ));
            }
            let addition = bytes + usize::from(!row.puts.is_empty());
            if row.puts.len() == MAX_RELAY_PUTS
                || encoded_bytes + addition > MAX_VALUE_BYTES
                || encoded_bytes + addition + 2 * (MAX_KEY_BYTES + 8) > MAX_BATCH_BYTES
            {
                self.push_relay(&row)?;
                row.puts.clear();
                encoded_bytes = base_bytes;
            }
            encoded_bytes += bytes + usize::from(!row.puts.is_empty());
            row.puts.push((key, value));
        }
        for (key, _) in deletes {
            // First delete adds the optional JSON field. Every entry
            // costs at most a comma, quotes, and two hex chars per byte.
            let bytes = 3 + 2 * key.as_bytes().len();
            if key.as_bytes().len() > MAX_KEY_BYTES
                || base_bytes + RELAY_DELETE_FIELD_BYTES + bytes > MAX_VALUE_BYTES
            {
                return Err(StoreError::Invalid(
                    "relay delete cannot fit one row".into(),
                ));
            }
            let addition = bytes
                + if row.deletes.is_empty() {
                    RELAY_DELETE_FIELD_BYTES
                } else {
                    1 // comma between delete keys
                };
            if row.puts.len() + row.deletes.len() == MAX_RELAY_PUTS
                || encoded_bytes + addition > MAX_VALUE_BYTES
                || encoded_bytes + addition + 2 * (MAX_KEY_BYTES + 8) > MAX_BATCH_BYTES
            {
                self.push_relay(&row)?;
                row.puts.clear();
                row.deletes.clear();
                encoded_bytes = base_bytes;
            }
            encoded_bytes += bytes
                + if row.deletes.is_empty() {
                    RELAY_DELETE_FIELD_BYTES
                } else {
                    1
                };
            row.deletes.push(key);
        }
        self.push_relay(&row)?;
        Ok(at_ms)
    }

    /// Finish with error reporting. Errors leave output vectors unchanged.
    pub fn try_finish(
        mut self,
        pre: &mut Vec<Precondition>,
        writes: &mut Vec<Write>,
    ) -> Result<(), StoreError> {
        if let Some(error) = self.error.take() {
            return Err(error);
        }
        for rid in &self.reservations {
            require_unplanned(&keys::reservation(rid)?, pre, writes)?;
        }
        let mut relay_due = None;
        for (target, operations) in std::mem::take(&mut self.relays) {
            if operations.is_empty() {
                continue;
            }
            relay_due = Some(self.plan_relay_rows(target, operations)?);
        }
        if let Some(due) = relay_due {
            self.writes.push(Write::Put(
                keys::timer(due, crate::timers::registry::kinds::RELAY.get(), b""),
                Value::default(),
            ));
        }
        if self.sequence_touched {
            require_unplanned(&keys::outbox_sequence(), pre, writes)?;
            self.pre
                .push(guard(keys::outbox_sequence(), self.os.as_ref()));
            self.writes.push(Write::Put(
                keys::outbox_sequence(),
                codec::encode_u64(self.seq),
            ));
        }
        if self.backlog_touched {
            require_unplanned(&keys::outcome_backlog(), pre, writes)?;
            self.pre
                .push(guard(keys::outcome_backlog(), self.oc.as_ref()));
            self.writes.push(Write::Put(
                keys::outcome_backlog(),
                codec::encode_backlog(&self.backlog),
            ));
        }
        if let Some(at) = self.kick_at_ms {
            self.writes.push(Write::Put(
                keys::timer(
                    at,
                    crate::timers::registry::kinds::OUTCOME_DELIVERY.get(),
                    b"",
                ),
                Value::default(),
            ));
        }
        let batch = Batch {
            preconditions: self.pre,
            writes: self.writes,
        };
        batch.validate(&StoreCapabilities::full())?;
        pre.extend(batch.preconditions);
        writes.extend(batch.writes);
        Ok(())
    }

    fn push_relay(&mut self, row: &RelayV1) -> Result<(), StoreError> {
        let value = codec::encode_relay(row)?;
        codec::decode_relay(&value)?;
        let seq = self.allocate()?;
        self.writes.push(Write::Put(keys::relay(seq), value));
        Ok(())
    }

    /// Finish the fixed fragment API; malformed input makes the caller's
    /// complete batch uncommittable, which looks like a retryable conflict.
    /// Wiring code (WP-1.9, 1.10, 1.14, 3.3) MUST call `try_finish` so the
    /// error is reported instead.
    pub fn finish(self, pre: &mut Vec<Precondition>, writes: &mut Vec<Write>) {
        if self.try_finish(pre, writes).is_err() {
            let key = keys::outbox_sequence();
            pre.extend([
                Precondition::Absent(key.clone()),
                Precondition::Present(key),
            ]);
        }
    }
}

fn require_unplanned(key: &Key, pre: &[Precondition], writes: &[Write]) -> Result<(), StoreError> {
    let guarded = pre.iter().any(|p| match p {
        Precondition::Equals(k, _) | Precondition::Absent(k) | Precondition::Present(k) => k == key,
        Precondition::NotAfter(_) => false,
    });
    let written = writes
        .iter()
        .any(|w| matches!(w, Write::Put(k, _) | Write::Delete(k) if k == key));
    if guarded || written {
        return Err(StoreError::Invalid(
            "outbox key already planned in batch".into(),
        ));
    }
    Ok(())
}

/// Acknowledge exactly this terminal row and pending index. Counting the
/// raw stored bytes makes the decrement independent of codec re-encoding.
pub fn plan_ack(
    rid: &str,
    value: &Value,
    oq_seq: u64,
    oc: Option<&Value>,
    pre: &mut Vec<Precondition>,
    writes: &mut Vec<Write>,
) -> Result<(), StoreError> {
    if oq_seq == 0
        || matches!(
            codec::decode_reservation(value)?,
            ReservationV1::Ticketed { .. } | ReservationV1::Pending { .. }
        )
    {
        return Err(corrupt("ack requires a terminal indexed outcome"));
    }
    let key = keys::reservation(rid)?;
    let pending = keys::outcome_pending(oq_seq, rid)?;
    if writes
        .iter()
        .any(|w| matches!(w, Write::Put(k, _) | Write::Delete(k) if k == &key))
    {
        return Err(StoreError::Invalid(
            "outcome already planned in batch".into(),
        ));
    }
    let observed = oc
        .map(codec::decode_backlog)
        .transpose()?
        .ok_or_else(|| corrupt("missing outcome backlog"))?;
    let backlog_key = keys::outcome_backlog();
    let expected = guard(backlog_key.clone(), oc);
    let existing_guard = pre.iter().find(|p| match p {
        Precondition::Equals(k, _) | Precondition::Absent(k) | Precondition::Present(k) => {
            k == &backlog_key
        }
        Precondition::NotAfter(_) => false,
    });
    if existing_guard.is_some_and(|p| p != &expected) {
        return Err(StoreError::Invalid("inconsistent backlog snapshots".into()));
    }
    let position = writes
        .iter()
        .rposition(|w| matches!(w, Write::Put(k, _) | Write::Delete(k) if k == &backlog_key));
    let mut backlog = match position.map(|i| &writes[i]) {
        Some(Write::Put(_, value)) => codec::decode_backlog(value)?,
        Some(Write::Delete(_)) => Backlog::default(),
        None => observed,
    };
    backlog.rows = backlog
        .rows
        .checked_sub(1)
        .ok_or_else(|| corrupt("backlog rows underflow"))?;
    backlog.bytes = backlog
        .bytes
        .checked_sub((key.as_bytes().len() + value.as_bytes().len()) as u64)
        .ok_or_else(|| corrupt("backlog bytes underflow"))?;
    if (backlog.rows == 0) != (backlog.bytes == 0) {
        return Err(corrupt("inconsistent outcome backlog"));
    }
    let needs_guard = existing_guard.is_none();
    pre.extend([
        Precondition::Equals(key.clone(), value.clone()),
        Precondition::Equals(pending.clone(), Value::default()),
    ]);
    writes.extend([Write::Delete(key), Write::Delete(pending)]);
    if needs_guard {
        pre.push(expected);
    }
    let write = if backlog.rows == 0 {
        Write::Delete(backlog_key)
    } else {
        Write::Put(backlog_key, codec::encode_backlog(&backlog))
    };
    if let Some(i) = position {
        writes[i] = write;
    } else {
        writes.push(write);
    }
    Ok(())
}

#[cfg(test)]
#[path = "outbox_tests.rs"]
mod tests;