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,
};
pub const MAX_TICKETS_PER_ADVANCE: usize = 7;
pub const ADVANCE_SHARED_OPS: usize = 31;
const _: () = assert!(MAX_TICKETS_PER_ADVANCE * 9 + ADVANCE_SHARED_OPS <= MAX_BATCH_OPS);
pub const MAX_RELAY_PUTS: usize = 96;
const RELAY_DELETE_FIELD_BYTES: usize = 13; const _: () = assert!(MAX_RELAY_PUTS + 2 <= MAX_BATCH_OPS);
const _: () = assert!(MAX_VALUE_BYTES + 2 * (MAX_KEY_BYTES + 8) <= MAX_BATCH_BYTES);
#[must_use]
pub fn synthetic_reservation_id(replay_scope: &Hash) -> String {
format!("s:{}", to_hex(replay_scope))
}
#[derive(Debug, Clone)]
pub struct Terminal(ReservationV1);
impl Terminal {
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())
}
#[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 {
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();
}
}
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);
}
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);
}
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);
}
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);
}
pub fn relay_at(&mut self, now_ms: u64) {
self.relay_at_ms = Some(now_ms);
}
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);
}
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()));
};
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 {
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 };
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)
}
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(())
}
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(())
}
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;