1use std::collections::{BTreeMap, BTreeSet};
6
7use mkit_core::hash::{Hash, to_hex};
8
9use super::codec::{self, Backlog, RelayV1, ReservationV1};
10use super::{
11 Batch, Key, MAX_BATCH_BYTES, MAX_BATCH_OPS, MAX_KEY_BYTES, MAX_VALUE_BYTES, Partition,
12 Precondition, StoreCapabilities, StoreError, Value, Write, keys,
13};
14
15pub const MAX_TICKETS_PER_ADVANCE: usize = 7;
37pub const ADVANCE_SHARED_OPS: usize = 31;
39const _: () = assert!(MAX_TICKETS_PER_ADVANCE * 9 + ADVANCE_SHARED_OPS <= MAX_BATCH_OPS);
40
41pub const MAX_RELAY_PUTS: usize = 96;
44const RELAY_DELETE_FIELD_BYTES: usize = 13; const _: () = assert!(MAX_RELAY_PUTS + 2 <= MAX_BATCH_OPS);
46const _: () = assert!(MAX_VALUE_BYTES + 2 * (MAX_KEY_BYTES + 8) <= MAX_BATCH_BYTES);
48
49#[must_use]
51pub fn synthetic_reservation_id(replay_scope: &Hash) -> String {
52 format!("s:{}", to_hex(replay_scope))
53}
54
55#[derive(Debug, Clone)]
57pub struct Terminal(ReservationV1);
58
59impl Terminal {
60 pub fn new(record: ReservationV1) -> Result<Self, StoreError> {
62 if matches!(
63 record,
64 ReservationV1::Ticketed { .. } | ReservationV1::Pending { .. }
65 ) {
66 return Err(StoreError::Invalid("outcome must be terminal".into()));
67 }
68 codec::decode_reservation(&codec::encode_reservation(&record))?;
69 Ok(Self(record))
70 }
71
72 fn occurred_at_ms(&self) -> u64 {
73 match &self.0 {
74 ReservationV1::Committed { occurred_at_ms, .. }
75 | ReservationV1::Aborted { occurred_at_ms, .. }
76 | ReservationV1::Expired { occurred_at_ms, .. }
77 | ReservationV1::ReadServed { occurred_at_ms, .. } => *occurred_at_ms,
78 _ => 0,
79 }
80 }
81}
82
83pub(crate) fn guard(key: Key, prior: Option<&Value>) -> Precondition {
84 match prior {
85 Some(value) => Precondition::Equals(key, value.clone()),
86 None => Precondition::Absent(key),
87 }
88}
89
90fn corrupt(message: &'static str) -> StoreError {
91 StoreError::Corrupt(message.into())
92}
93
94#[derive(Debug)]
102pub struct OutboxBuilder {
103 os: Option<Value>,
104 oc: Option<Value>,
105 seq: u64,
106 backlog: Backlog,
107 sequence_touched: bool,
108 backlog_touched: bool,
109 pre: Vec<Precondition>,
110 writes: Vec<Write>,
111 relays: BTreeMap<Partition, BTreeMap<Key, Option<Value>>>,
112 reservations: BTreeSet<String>,
113 error: Option<StoreError>,
114 relay_at_ms: Option<u64>,
115 kick_at_ms: Option<u64>,
116}
117
118impl OutboxBuilder {
119 pub fn new(os: Option<&Value>, oc: Option<&Value>) -> Result<Self, StoreError> {
121 let seq = os.map(codec::decode_u64).transpose()?.unwrap_or(0);
122 if os.is_some() && seq == 0 {
123 return Err(corrupt("outbox sequence is zero"));
124 }
125 Ok(Self {
126 os: os.cloned(),
127 oc: oc.cloned(),
128 seq,
129 backlog: oc
130 .map(codec::decode_backlog)
131 .transpose()?
132 .unwrap_or(Backlog { rows: 0, bytes: 0 }),
133 sequence_touched: false,
134 backlog_touched: false,
135 pre: Vec::new(),
136 writes: Vec::new(),
137 relays: BTreeMap::new(),
138 reservations: BTreeSet::new(),
139 error: None,
140 relay_at_ms: None,
141 kick_at_ms: None,
142 })
143 }
144
145 fn allocate(&mut self) -> Result<u64, StoreError> {
146 self.seq = self
147 .seq
148 .checked_add(1)
149 .ok_or_else(|| corrupt("outbox sequence overflow"))?;
150 self.sequence_touched = true;
151 Ok(self.seq)
152 }
153
154 fn remember(&mut self, result: Result<(), StoreError>) {
155 if self.error.is_none() {
156 self.error = result.err();
157 }
158 }
159
160 pub fn reserve(&mut self, rid: &str, ticket_id: [u8; 32], prior: Option<&Value>) {
163 let result = (|| {
164 let key = keys::reservation(rid)?;
165 if !self.reservations.insert(rid.to_owned()) {
166 return Err(StoreError::Invalid("duplicate reservation in batch".into()));
167 }
168 if let Some(value) = prior
169 && !matches!(
170 codec::decode_reservation(value)?,
171 ReservationV1::Pending {
172 op: codec::PendingOp::Write,
173 ..
174 }
175 )
176 {
177 return Err(StoreError::Invalid("reservation id already in use".into()));
178 }
179 self.pre.push(guard(key.clone(), prior));
180 self.writes.push(Write::Put(
181 key,
182 codec::encode_reservation(&ReservationV1::Ticketed { ticket_id }),
183 ));
184 Ok(())
185 })();
186 self.remember(result);
187 }
188
189 pub fn pending(&mut self, rid: &str, prior: Option<&Value>, record: &ReservationV1) {
192 let result = (|| {
193 if prior.is_some() {
194 return Err(StoreError::Invalid("reservation id already in use".into()));
195 }
196 let ReservationV1::Pending {
197 reconcile_at_ms, ..
198 } = record
199 else {
200 return Err(StoreError::Invalid(
201 "pending requires Pending record".into(),
202 ));
203 };
204 let key = keys::reservation(rid)?;
205 if !self.reservations.insert(rid.to_owned()) {
206 return Err(StoreError::Invalid("duplicate reservation in batch".into()));
207 }
208 let value = codec::encode_reservation(record);
209 codec::decode_reservation(&value)?;
210 self.pre.push(Precondition::Absent(key.clone()));
211 self.writes.push(Write::Put(key, value));
212 self.writes.push(Write::Put(
213 keys::timer(
214 *reconcile_at_ms,
215 crate::timers::registry::kinds::RESERVATION_RECONCILE.get(),
216 rid.as_bytes(),
217 ),
218 Value::default(),
219 ));
220 Ok(())
221 })();
222 self.remember(result);
223 }
224
225 pub fn outcome(&mut self, rid: &str, prior: &Value, terminal: Terminal) {
228 let occurred_at_ms = terminal.occurred_at_ms();
229 let record = terminal.0;
230 let result = (|| {
231 let key = keys::reservation(rid)?;
232 let permitted = match (codec::decode_reservation(prior)?, &record) {
233 (
234 ReservationV1::Ticketed { .. },
235 ReservationV1::Committed { .. }
236 | ReservationV1::Aborted { .. }
237 | ReservationV1::Expired { .. },
238 ) => true,
239 (
240 ReservationV1::Pending {
241 repository: prior_repo,
242 op: codec::PendingOp::Write,
243 ..
244 },
245 ReservationV1::Committed { repository, .. }
246 | ReservationV1::Aborted { repository, .. },
247 )
248 | (
249 ReservationV1::Pending {
250 repository: prior_repo,
251 op: codec::PendingOp::Read,
252 ..
253 },
254 ReservationV1::ReadServed { repository, .. }
255 | ReservationV1::Aborted { repository, .. },
256 ) => prior_repo == repository.as_str(),
257 _ => false,
258 };
259 if !permitted {
260 return Err(corrupt("illegal reservation transition"));
261 }
262 if !self.reservations.insert(rid.to_owned()) {
263 return Err(StoreError::Invalid("duplicate reservation in batch".into()));
264 }
265 let value = codec::encode_reservation(&record);
266 let size = (key.as_bytes().len() + value.as_bytes().len()) as u64;
267 self.backlog.rows = self
268 .backlog
269 .rows
270 .checked_add(1)
271 .ok_or_else(|| corrupt("backlog rows overflow"))?;
272 self.backlog.bytes = self
273 .backlog
274 .bytes
275 .checked_add(size)
276 .ok_or_else(|| corrupt("backlog bytes overflow"))?;
277 let seq = self.allocate()?;
278 if self.oc.is_none() && self.kick_at_ms.is_none() {
279 self.kick_at_ms = Some(occurred_at_ms);
280 }
281 self.backlog_touched = true;
282 self.pre
283 .push(Precondition::Equals(key.clone(), prior.clone()));
284 self.writes.push(Write::Put(key, value));
285 self.writes.push(Write::Put(
286 keys::outcome_pending(seq, rid)?,
287 Value::default(),
288 ));
289 Ok(())
290 })();
291 self.remember(result);
292 }
293
294 pub fn abort_direct(&mut self, rid: &str, terminal: Terminal) {
296 let occurred_at_ms = terminal.occurred_at_ms();
297 let record = terminal.0;
298 let result = (|| {
299 if !matches!(record, ReservationV1::Aborted { .. }) {
300 return Err(StoreError::Invalid("direct abort requires Aborted".into()));
301 }
302 let key = keys::reservation(rid)?;
303 if !self.reservations.insert(rid.to_owned()) {
304 return Err(StoreError::Invalid("duplicate reservation in batch".into()));
305 }
306 let value = codec::encode_reservation(&record);
307 let size = (key.as_bytes().len() + value.as_bytes().len()) as u64;
308 self.backlog.rows = self
309 .backlog
310 .rows
311 .checked_add(1)
312 .ok_or_else(|| corrupt("backlog rows overflow"))?;
313 self.backlog.bytes = self
314 .backlog
315 .bytes
316 .checked_add(size)
317 .ok_or_else(|| corrupt("backlog bytes overflow"))?;
318 let seq = self.allocate()?;
319 if self.oc.is_none() && self.kick_at_ms.is_none() {
320 self.kick_at_ms = Some(occurred_at_ms);
321 }
322 self.backlog_touched = true;
323 self.pre.push(Precondition::Absent(key.clone()));
324 self.writes.push(Write::Put(key, value));
325 self.writes.push(Write::Put(
326 keys::outcome_pending(seq, rid)?,
327 Value::default(),
328 ));
329 Ok(())
330 })();
331 self.remember(result);
332 }
333
334 pub fn relay_at(&mut self, now_ms: u64) {
336 self.relay_at_ms = Some(now_ms);
337 }
338
339 pub fn relay(&mut self, target: &Partition, puts: Vec<(Key, Value)>) {
342 let result = (|| {
343 target.encode()?;
344 let group = self.relays.entry(target.clone()).or_default();
345 for (key, value) in puts {
346 if group
347 .get(&key)
348 .is_some_and(|old| old.as_ref() != Some(&value))
349 {
350 return Err(StoreError::Invalid("conflicting relay upserts".into()));
351 }
352 group.insert(key, Some(value));
353 }
354 Ok(())
355 })();
356 self.remember(result);
357 }
358
359 pub fn relay_delete(&mut self, target: &Partition, keys: Vec<Key>) {
362 let result = (|| {
363 target.encode()?;
364 let group = self.relays.entry(target.clone()).or_default();
365 for key in keys {
366 if group.get(&key).is_some_and(Option::is_some) {
367 return Err(StoreError::Invalid("relay put/delete overlap".into()));
368 }
369 group.insert(key, None);
370 }
371 Ok(())
372 })();
373 self.remember(result);
374 }
375
376 fn plan_relay_rows(
377 &mut self,
378 target: Partition,
379 operations: BTreeMap<Key, Option<Value>>,
380 ) -> Result<u64, StoreError> {
381 let at_ms = self
382 .relay_at_ms
383 .ok_or_else(|| StoreError::Invalid("relay rows need relay_at".into()))?;
384 let mut row = RelayV1 {
385 at_ms,
386 target,
387 puts: Vec::new(),
388 deletes: Vec::new(),
389 };
390 let base_bytes = codec::encode_relay(&row)?.as_bytes().len();
391 let mut encoded_bytes = base_bytes;
392 let (puts, deletes): (Vec<_>, Vec<_>) = operations
393 .into_iter()
394 .partition(|(_, value)| value.is_some());
395 for (key, value) in puts {
396 let Some(value) = value else {
397 return Err(StoreError::Invalid("missing relay upsert value".into()));
398 };
399 let bytes = 7 + 2 * (key.as_bytes().len() + value.as_bytes().len());
401 if key.as_bytes().len() > MAX_KEY_BYTES
402 || value.as_bytes().len() > MAX_VALUE_BYTES
403 || base_bytes + bytes > MAX_VALUE_BYTES
404 {
405 return Err(StoreError::Invalid(
406 "relay upsert cannot fit one row".into(),
407 ));
408 }
409 let addition = bytes + usize::from(!row.puts.is_empty());
410 if row.puts.len() == MAX_RELAY_PUTS
411 || encoded_bytes + addition > MAX_VALUE_BYTES
412 || encoded_bytes + addition + 2 * (MAX_KEY_BYTES + 8) > MAX_BATCH_BYTES
413 {
414 self.push_relay(&row)?;
415 row.puts.clear();
416 encoded_bytes = base_bytes;
417 }
418 encoded_bytes += bytes + usize::from(!row.puts.is_empty());
419 row.puts.push((key, value));
420 }
421 for (key, _) in deletes {
422 let bytes = 3 + 2 * key.as_bytes().len();
425 if key.as_bytes().len() > MAX_KEY_BYTES
426 || base_bytes + RELAY_DELETE_FIELD_BYTES + bytes > MAX_VALUE_BYTES
427 {
428 return Err(StoreError::Invalid(
429 "relay delete cannot fit one row".into(),
430 ));
431 }
432 let addition = bytes
433 + if row.deletes.is_empty() {
434 RELAY_DELETE_FIELD_BYTES
435 } else {
436 1 };
438 if row.puts.len() + row.deletes.len() == MAX_RELAY_PUTS
439 || encoded_bytes + addition > MAX_VALUE_BYTES
440 || encoded_bytes + addition + 2 * (MAX_KEY_BYTES + 8) > MAX_BATCH_BYTES
441 {
442 self.push_relay(&row)?;
443 row.puts.clear();
444 row.deletes.clear();
445 encoded_bytes = base_bytes;
446 }
447 encoded_bytes += bytes
448 + if row.deletes.is_empty() {
449 RELAY_DELETE_FIELD_BYTES
450 } else {
451 1
452 };
453 row.deletes.push(key);
454 }
455 self.push_relay(&row)?;
456 Ok(at_ms)
457 }
458
459 pub fn try_finish(
461 mut self,
462 pre: &mut Vec<Precondition>,
463 writes: &mut Vec<Write>,
464 ) -> Result<(), StoreError> {
465 if let Some(error) = self.error.take() {
466 return Err(error);
467 }
468 for rid in &self.reservations {
469 require_unplanned(&keys::reservation(rid)?, pre, writes)?;
470 }
471 let mut relay_due = None;
472 for (target, operations) in std::mem::take(&mut self.relays) {
473 if operations.is_empty() {
474 continue;
475 }
476 relay_due = Some(self.plan_relay_rows(target, operations)?);
477 }
478 if let Some(due) = relay_due {
479 self.writes.push(Write::Put(
480 keys::timer(due, crate::timers::registry::kinds::RELAY.get(), b""),
481 Value::default(),
482 ));
483 }
484 if self.sequence_touched {
485 require_unplanned(&keys::outbox_sequence(), pre, writes)?;
486 self.pre
487 .push(guard(keys::outbox_sequence(), self.os.as_ref()));
488 self.writes.push(Write::Put(
489 keys::outbox_sequence(),
490 codec::encode_u64(self.seq),
491 ));
492 }
493 if self.backlog_touched {
494 require_unplanned(&keys::outcome_backlog(), pre, writes)?;
495 self.pre
496 .push(guard(keys::outcome_backlog(), self.oc.as_ref()));
497 self.writes.push(Write::Put(
498 keys::outcome_backlog(),
499 codec::encode_backlog(&self.backlog),
500 ));
501 }
502 if let Some(at) = self.kick_at_ms {
503 self.writes.push(Write::Put(
504 keys::timer(
505 at,
506 crate::timers::registry::kinds::OUTCOME_DELIVERY.get(),
507 b"",
508 ),
509 Value::default(),
510 ));
511 }
512 let batch = Batch {
513 preconditions: self.pre,
514 writes: self.writes,
515 };
516 batch.validate(&StoreCapabilities::full())?;
517 pre.extend(batch.preconditions);
518 writes.extend(batch.writes);
519 Ok(())
520 }
521
522 fn push_relay(&mut self, row: &RelayV1) -> Result<(), StoreError> {
523 let value = codec::encode_relay(row)?;
524 codec::decode_relay(&value)?;
525 let seq = self.allocate()?;
526 self.writes.push(Write::Put(keys::relay(seq), value));
527 Ok(())
528 }
529
530 pub fn finish(self, pre: &mut Vec<Precondition>, writes: &mut Vec<Write>) {
535 if self.try_finish(pre, writes).is_err() {
536 let key = keys::outbox_sequence();
537 pre.extend([
538 Precondition::Absent(key.clone()),
539 Precondition::Present(key),
540 ]);
541 }
542 }
543}
544
545fn require_unplanned(key: &Key, pre: &[Precondition], writes: &[Write]) -> Result<(), StoreError> {
546 let guarded = pre.iter().any(|p| match p {
547 Precondition::Equals(k, _) | Precondition::Absent(k) | Precondition::Present(k) => k == key,
548 Precondition::NotAfter(_) => false,
549 });
550 let written = writes
551 .iter()
552 .any(|w| matches!(w, Write::Put(k, _) | Write::Delete(k) if k == key));
553 if guarded || written {
554 return Err(StoreError::Invalid(
555 "outbox key already planned in batch".into(),
556 ));
557 }
558 Ok(())
559}
560
561pub fn plan_ack(
564 rid: &str,
565 value: &Value,
566 oq_seq: u64,
567 oc: Option<&Value>,
568 pre: &mut Vec<Precondition>,
569 writes: &mut Vec<Write>,
570) -> Result<(), StoreError> {
571 if oq_seq == 0
572 || matches!(
573 codec::decode_reservation(value)?,
574 ReservationV1::Ticketed { .. } | ReservationV1::Pending { .. }
575 )
576 {
577 return Err(corrupt("ack requires a terminal indexed outcome"));
578 }
579 let key = keys::reservation(rid)?;
580 let pending = keys::outcome_pending(oq_seq, rid)?;
581 if writes
582 .iter()
583 .any(|w| matches!(w, Write::Put(k, _) | Write::Delete(k) if k == &key))
584 {
585 return Err(StoreError::Invalid(
586 "outcome already planned in batch".into(),
587 ));
588 }
589 let observed = oc
590 .map(codec::decode_backlog)
591 .transpose()?
592 .ok_or_else(|| corrupt("missing outcome backlog"))?;
593 let backlog_key = keys::outcome_backlog();
594 let expected = guard(backlog_key.clone(), oc);
595 let existing_guard = pre.iter().find(|p| match p {
596 Precondition::Equals(k, _) | Precondition::Absent(k) | Precondition::Present(k) => {
597 k == &backlog_key
598 }
599 Precondition::NotAfter(_) => false,
600 });
601 if existing_guard.is_some_and(|p| p != &expected) {
602 return Err(StoreError::Invalid("inconsistent backlog snapshots".into()));
603 }
604 let position = writes
605 .iter()
606 .rposition(|w| matches!(w, Write::Put(k, _) | Write::Delete(k) if k == &backlog_key));
607 let mut backlog = match position.map(|i| &writes[i]) {
608 Some(Write::Put(_, value)) => codec::decode_backlog(value)?,
609 Some(Write::Delete(_)) => Backlog::default(),
610 None => observed,
611 };
612 backlog.rows = backlog
613 .rows
614 .checked_sub(1)
615 .ok_or_else(|| corrupt("backlog rows underflow"))?;
616 backlog.bytes = backlog
617 .bytes
618 .checked_sub((key.as_bytes().len() + value.as_bytes().len()) as u64)
619 .ok_or_else(|| corrupt("backlog bytes underflow"))?;
620 if (backlog.rows == 0) != (backlog.bytes == 0) {
621 return Err(corrupt("inconsistent outcome backlog"));
622 }
623 let needs_guard = existing_guard.is_none();
624 pre.extend([
625 Precondition::Equals(key.clone(), value.clone()),
626 Precondition::Equals(pending.clone(), Value::default()),
627 ]);
628 writes.extend([Write::Delete(key), Write::Delete(pending)]);
629 if needs_guard {
630 pre.push(expected);
631 }
632 let write = if backlog.rows == 0 {
633 Write::Delete(backlog_key)
634 } else {
635 Write::Put(backlog_key, codec::encode_backlog(&backlog))
636 };
637 if let Some(i) = position {
638 writes[i] = write;
639 } else {
640 writes.push(write);
641 }
642 Ok(())
643}
644
645#[cfg(test)]
646#[path = "outbox_tests.rs"]
647mod tests;