extern crate alloc;
use alloc::collections::{BTreeMap, BTreeSet};
use alloc::vec::Vec;
use crate::kairos::{Clock, Kairos, TickCounter};
use crate::metis::{
Abandoned, Admission, ArrivalRefusal, Consigned, Cut, Departed, Departure, DepartureRefusal,
Dot, DotSet, Dotted, EpochAddress, EpochBootstrapError, EpochConsignmentError, EpochIdeal,
EpochProjection, EpochRefusal, EpochShadow, Epochs, Fenced, Metatheses, Retirement, Rhapsody,
SealRecord, SealedEpoch, Stability, VersionVector, Vouched,
};
use super::{Note, OldDelta};
mod data;
mod lifecycle;
mod persistence;
mod receive;
pub type Text = Dotted<Rhapsody>;
pub type Moves = Dotted<Metatheses>;
type NativeEntry = (Vec<Dot>, OldDelta, Option<Kairos>);
#[derive(Clone)]
pub struct Journal {
pub(super) generation: u64,
pub(super) base_text: Text,
pub(super) base_moves: Moves,
pub(super) epochs: Epochs,
pub(super) counter: u64,
pub(super) clock: Kairos,
pub(super) seals: Vec<(EpochAddress, EpochProjection)>,
pub(super) departed: BTreeSet<u32>,
pub(super) attested: Vec<u32>,
pub(super) notes: Vec<Note>,
}
pub(super) struct JoinerCheckpoint {
station: u32,
generation: u64,
base_text: Text,
base_moves: Moves,
counter: u64,
clock: Kairos,
seals: Vec<(EpochAddress, EpochProjection)>,
attested: Vec<u32>,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(super) enum JoinerBootstrapError {
Epoch(EpochBootstrapError),
StationOutsideRoster { station: u32 },
GenerationMismatch { checkpoint: u64, lineage: u64 },
}
impl From<EpochBootstrapError> for JoinerBootstrapError {
fn from(error: EpochBootstrapError) -> Self {
Self::Epoch(error)
}
}
enum Probe {
Data(bool),
Stability(Stability),
Retirement(Retirement),
Epochs(alloc::boxed::Box<Epochs>),
Departures(BTreeMap<u32, Departure>),
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum Decline {
Unwitnessed,
NotSelfSupporting {
dot: Dot,
},
Refused(EpochRefusal),
}
struct Transition {
adopted: crate::metis::Adopted,
shadow: EpochShadow,
native_text: Text,
native_moves: Moves,
native_counter: u64,
native_recorded: DotSet,
native_log: Vec<NativeEntry>,
}
pub struct Replica {
id: u32,
roster: Vec<u32>,
clock: Clock<TickCounter>,
epochs: Epochs,
generation: u64,
stability: Stability,
counter: u64,
text: Text,
moves: Moves,
base_text: Text,
base_moves: Moves,
log: Vec<(Vec<Dot>, OldDelta)>,
have: DotSet,
recorded: DotSet,
reported: Option<VersionVector>,
retirement: Retirement,
acked: DotSet,
confirmed: BTreeSet<EpochAddress>,
declaration_dots: Vec<Dot>,
gate: EpochIdeal<(Vec<Dot>, OldDelta)>,
gated: DotSet,
transition: Option<Transition>,
parked: Vec<Note>,
pub seals: Vec<(EpochAddress, EpochProjection)>,
journal: Journal,
volatile: Vec<Note>,
halt_on_refusal: bool,
halted: Option<EpochConsignmentError>,
certified: Vec<Consigned>,
departed: BTreeSet<u32>,
departures: BTreeMap<u32, Departure>,
attested: Vec<Departed>,
resurgences: Vec<Fenced>,
}
impl Replica {
pub const fn id(&self) -> u32 {
self.id
}
pub const fn tolerate_refusals(&mut self) {
self.halt_on_refusal = true;
}
pub fn abandon(&mut self, station: u32, out: &mut Vec<Note>) -> Abandoned {
let _ = self.departed.insert(station);
let abandoned = self.stability.abandon(station).expect("a family remains");
self.poll(out);
abandoned
}
pub fn open_departure(
&mut self,
departing: u32,
out: &mut Vec<Note>,
) -> Result<u64, DepartureRefusal> {
let (id, generation) = (self.id, self.generation);
let held = self.have.clone();
let mut successor = self.gated.clone();
if let Some(transition) = &self.transition {
successor = successor.merge(&transition.native_recorded);
}
let prefix = self
.enter_departure(departing)?
.propose(&held, &successor)?;
self.emit_mint(
out,
Note::Depart {
generation,
departing,
by: id,
prefix,
},
);
self.settle_departures();
self.poll(out);
Ok(prefix)
}
fn enter_departure(&mut self, departing: u32) -> Result<&mut Departure, DepartureRefusal> {
if !self.departures.contains_key(&departing) {
let survivors: Vec<u32> = self
.roster
.iter()
.copied()
.filter(|&station| {
station != departing
&& !self.departed.contains(&station)
&& self.attested(station).is_none()
})
.collect();
let _ = self
.departures
.insert(departing, Departure::opened(departing, self.id, survivors)?);
}
Ok(self
.departures
.get_mut(&departing)
.expect("just inserted or already present"))
}
pub(super) fn settle_departures(&mut self) {
let ready: Vec<Departed> = self
.departures
.values_mut()
.filter_map(|round| round.try_seal().cloned())
.collect();
for departed in ready {
if self.attested(departed.station()).is_some() {
continue;
}
if self.stability.abandon_attested(&departed).is_ok() {
self.attested.push(departed);
}
}
}
fn restore_departures(&mut self, departed: &BTreeSet<u32>, attested: &[u32]) {
for &station in departed {
let _ = self.departed.insert(station);
let _ = self.stability.abandon(station).expect("a family remains");
}
let coverage_only: BTreeSet<u32> = self
.epochs
.newest_sealed()
.and_then(SealedEpoch::admission)
.map(|admission| admission.joiners().collect())
.unwrap_or_default();
for &station in attested {
let held = self.have.floor().get(station);
let mut round = self
.enter_departure(station)
.expect("a checkpointed departure left a surviving family")
.clone();
round
.reinstate(held)
.expect("the base prefix cannot lower a fresh fence");
let family: Vec<u32> = round.proposals().map(|(station, _)| station).collect();
for surviving in family {
if surviving != self.id {
let word = if coverage_only.contains(&surviving) {
0
} else {
held
};
round
.fold(&Vouched::trust(surviving, word))
.expect("a survivor of the checkpointed family");
}
}
let departed = round
.try_seal()
.expect("every slot holds the base prefix")
.clone();
let _ = self
.stability
.abandon_attested(&departed)
.expect("a family remains");
self.attested.push(departed);
let _ = self.departures.insert(station, round);
}
}
pub fn attested_proposal(&self, departed: u32, member: u32) -> Option<u64> {
self.attested
.iter()
.find(|attestation| attestation.station() == departed)
.and_then(|attestation| attestation.proposal_of(member))
}
pub fn open_arrival(
&mut self,
joiners: &[u32],
out: &mut Vec<Note>,
) -> Result<(), ArrivalRefusal> {
let base = self
.epochs
.newest_sealed()
.ok_or(ArrivalRefusal::Unfounded)?;
let commitment = record_commitment(base);
let admission = Admission::new(base.declaration(), joiners.iter().copied())?;
let fresh = self.epochs.arrival().is_none();
let endorsement =
self.epochs
.open_arrival(admission, commitment, self.id, &self.stability)?;
let note = Note::Endorse {
by: self.id,
admission: endorsement.admission().clone(),
commitment,
};
if fresh {
self.emit_mint(out, note);
} else {
self.emit(out, note);
}
self.drive(out);
Ok(())
}
fn fenced(&mut self, dots: &[Dot]) -> Option<Fenced> {
let refusal = self.old_plane_refusal(dots);
self.refuse(refusal)
}
fn old_plane_refusal(&self, dots: &[Dot]) -> Option<Fenced> {
dots.iter().find_map(|&dot| {
self.departures
.get(&dot.station())
.and_then(|round| round.admits(dot).err())
})
}
fn fenced_native(&mut self, dots: &[Dot]) -> Option<Fenced> {
let refusal = self.native_plane_refusal(dots);
self.refuse(refusal)
}
fn native_plane_refusal(&self, dots: &[Dot]) -> Option<Fenced> {
dots.iter().find_map(|&dot| {
self.departures
.get(&dot.station())?
.admits_successor(dot)
.err()
})
}
fn note_refused(&self, note: &Note) -> bool {
match note {
Note::Old {
dots, generation, ..
} if *generation <= self.generation => self.old_plane_refusal(dots).is_some(),
Note::Old { dots, .. } | Note::Native { dots, .. } => {
self.native_plane_refusal(dots).is_some()
}
Note::Declare { declaration } => self.old_plane_refusal(&[declaration.dot()]).is_some(),
_ => false,
}
}
fn refuse(&mut self, refusal: Option<Fenced>) -> Option<Fenced> {
let refusal = refusal?;
if refusal.is_resurgence() && !self.resurgences.contains(&refusal) {
self.resurgences.push(refusal);
}
Some(refusal)
}
pub fn resurgences(&self) -> impl Iterator<Item = &Fenced> {
self.resurgences.iter()
}
pub fn departure_fence(&self, station: u32) -> Option<u64> {
self.departures.get(&station).and_then(Departure::fence)
}
pub fn attested(&self, station: u32) -> Option<u64> {
self.attested
.iter()
.find(|departed| departed.station() == station)
.map(Departed::bound)
}
pub fn resurgent(&self) -> impl Iterator<Item = u32> + '_ {
self.stability.resurgent()
}
pub fn watermark(&self) -> VersionVector {
self.stability.watermark()
}
pub const fn halted(&self) -> Option<&EpochConsignmentError> {
self.halted.as_ref()
}
fn vouch_restored_lineage(&mut self) {
self.certified = self
.epochs
.sealed()
.cloned()
.map(Consigned::vouch)
.collect();
}
pub fn certified(&self) -> impl Iterator<Item = &SealedEpoch> {
self.certified.iter().map(Consigned::record)
}
pub const fn generation(&self) -> u64 {
self.generation
}
pub const fn epochs(&self) -> &Epochs {
&self.epochs
}
pub const fn text(&self) -> &Text {
&self.text
}
pub const fn moves(&self) -> &Moves {
&self.moves
}
pub fn adopted(&self) -> bool {
self.transition.is_some()
}
pub fn effective_order(&self) -> Vec<Dot> {
self.transition.as_ref().map_or_else(
|| self.text.store().recension(self.moves.store()).order(),
|transition| {
transition
.native_text
.store()
.recension(transition.native_moves.store())
.order()
},
)
}
fn fence(&mut self) {
self.journal.notes.append(&mut self.volatile);
}
fn park(&mut self, note: &Note) {
if !self.parked.contains(note) {
self.parked.push(note.clone());
}
}
fn emit(&mut self, out: &mut Vec<Note>, note: Note) {
self.fence();
out.push(note);
}
fn data_note_novel(&self, note: &Note) -> bool {
match note {
Note::Old {
generation, dots, ..
} => {
if *generation < self.generation {
return false;
}
if *generation == self.generation + 1 {
return self.transition.as_ref().is_none_or(|transition| {
dots.iter()
.any(|&dot| !transition.native_recorded.contains(dot))
});
}
dots.iter().any(|&dot| !self.recorded.contains(dot))
}
Note::Native { epoch, dots, .. } => {
if let Some(transition) = &self.transition
&& transition.adopted.address() == *epoch
{
return dots
.iter()
.any(|&dot| !transition.native_recorded.contains(dot));
}
if self
.epochs
.sealed()
.any(|sealed| sealed.declaration() == *epoch)
{
return epoch.generation() + 1 == self.generation
&& dots.iter().any(|&dot| !self.recorded.contains(dot));
}
dots.iter().any(|&dot| !self.gated.contains(dot))
}
_ => unreachable!("only data notes are probed for novelty"),
}
}
fn emit_mint(&mut self, out: &mut Vec<Note>, note: Note) {
self.fence();
self.journal.notes.push(note.clone());
out.push(note);
}
fn fold_current(&mut self, delta: &OldDelta) {
match delta {
OldDelta::Text(delta) => self.text.merge_from(delta),
OldDelta::Moves(delta) => self.moves.merge_from(delta),
}
}
fn base_frontier(&self) -> Cut {
Cut::floor_of(&self.base_text.context().merge(self.base_moves.context()))
}
pub fn try_declare(&mut self, out: &mut Vec<Note>) -> Result<EpochAddress, Decline> {
let Some(cut) = self.stability.watermark_cut() else {
return Err(Decline::Unwitnessed);
};
if let Some(dot) = self.stratum_obstacle(&cut) {
return Err(Decline::NotSelfSupporting { dot });
}
let dot = Dot::from_parts(self.id, self.counter + 1)
.expect("the allocator's next counter is one past a count, hence nonzero");
let rank = self.clock.now(0u16);
let declaration = self
.epochs
.declare(dot, rank, &self.stability, &self.base_frontier())
.map_err(Decline::Refused)?;
self.counter += 1;
let _ = self.have.insert(dot);
self.declaration_dots.push(dot);
let address = declaration.address();
self.emit_mint(out, Note::Declare { declaration });
self.poll(out);
Ok(address)
}
fn stratum_obstacle(&self, cut: &Cut) -> Option<Dot> {
let (text, moves) = self.rebuild_at(cut);
let coverage = text.context().merge(moves.context());
if let Some(dot) = cut.to_have_set().difference(&coverage).next() {
return Some(dot);
}
let vector = cut.as_vector();
if let Some(dot) = text
.context()
.dots()
.chain(moves.context().dots())
.find(|dot| dot.counter() > vector.get(dot.station()))
{
return Some(dot);
}
match text.store().refound(moves.store()) {
Ok(_) => None,
Err(unsealed) => Some(unsealed.dot),
}
}
fn rebuild_at(&self, cut: &Cut) -> (Text, Moves) {
let vector = cut.as_vector();
let mut text = self.base_text.clone();
let mut moves = self.base_moves.clone();
for (dots, delta) in &self.log {
if dots
.iter()
.all(|dot| dot.counter() <= vector.get(dot.station()))
{
match delta {
OldDelta::Text(delta) => text.merge_from(delta),
OldDelta::Moves(delta) => moves.merge_from(delta),
}
}
}
(text, moves)
}
}
fn record_commitment(sealed: &SealedEpoch) -> [u8; 32] {
let frame = SealRecord::from_sealed(sealed).to_bytes();
let mut commitment = [0u8; 32];
for (lane, chunk) in (1u64..).zip(commitment.chunks_exact_mut(8)) {
let mut acc = 0xcbf2_9ce4_8422_2325u64 ^ lane.wrapping_mul(0x9E37_79B9_7F4A_7C15);
for &byte in &frame {
acc ^= u64::from(byte);
acc = acc.wrapping_mul(0x0000_0100_0000_01B3);
}
chunk.copy_from_slice(&acc.to_be_bytes());
}
commitment
}
fn observe_pair_ranks(clock: &Clock<TickCounter>, text: &Text, moves: &Moves) {
for dot in text.context().dots() {
if let Some(locus) = text.store().locus(dot) {
clock.observe(locus.rank);
}
}
for (_, testimony) in moves.store() {
clock.observe(testimony.to.rank);
}
}
fn fold_native(transition: &mut Transition, dots: &[Dot], delta: &OldDelta, stamp: Option<Kairos>) {
let novel = dots
.iter()
.any(|&dot| !transition.native_recorded.contains(dot));
for &dot in dots {
let _ = transition.native_recorded.insert(dot);
}
match delta {
OldDelta::Text(delta) => transition.native_text.merge_from(delta),
OldDelta::Moves(delta) => transition.native_moves.merge_from(delta),
}
if novel {
transition
.native_log
.push((dots.to_vec(), delta.clone(), stamp));
}
}