use std::collections::HashSet;
use std::time::Duration;
use crate::identity::ProducerId;
use crate::time::{LocalInstant, RobotInstant, TimelineMismatch};
#[derive(Clone, Copy, Debug, PartialEq, Eq, thiserror::Error)]
pub enum LeaseRejection {
#[error("sequence {observed} does not follow {accepted} for the accepted producer")]
StaleSequence {
accepted: u64,
observed: u64,
},
#[error("producer {superseded} was superseded by {accepted}")]
SupersededProducer {
superseded: ProducerId,
accepted: ProducerId,
},
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum LeaseDecision {
Renewed,
ProducerReplaced {
superseded: ProducerId,
},
Rejected(LeaseRejection),
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct ProducerFence {
accepted: Option<(ProducerId, u64)>,
superseded: HashSet<ProducerId>,
}
impl ProducerFence {
pub fn new() -> Self {
ProducerFence::default()
}
pub fn accepted_producer(&self) -> Option<ProducerId> {
self.accepted.map(|(producer, _)| producer)
}
pub fn admit(&mut self, producer: ProducerId, sequence: u64) -> LeaseDecision {
if self.superseded.contains(&producer) {
return LeaseDecision::Rejected(LeaseRejection::SupersededProducer {
superseded: producer,
accepted: self.accepted.map_or(producer, |(accepted, _)| accepted),
});
}
match self.accepted {
None => {
self.accepted = Some((producer, sequence));
LeaseDecision::Renewed
}
Some((accepted, last)) if accepted == producer => {
if sequence > last {
self.accepted = Some((producer, sequence));
LeaseDecision::Renewed
} else {
LeaseDecision::Rejected(LeaseRejection::StaleSequence {
accepted: last,
observed: sequence,
})
}
}
Some((superseded, _)) => {
self.fence(superseded);
self.accepted = Some((producer, sequence));
LeaseDecision::ProducerReplaced { superseded }
}
}
}
pub fn permits(&self, producer: ProducerId) -> bool {
!self.superseded.contains(&producer)
&& self
.accepted
.is_none_or(|(accepted, _)| accepted == producer)
}
pub fn clear(&mut self) {
self.accepted = None;
}
fn fence(&mut self, producer: ProducerId) {
self.superseded.insert(producer);
}
}
pub const LEASE_TRACE_TARGET: &str = "phoxal.lease";
#[derive(Clone, Debug, PartialEq, Eq)]
pub struct Lease<B> {
input: &'static str,
silence: Duration,
hold: Duration,
fence: ProducerFence,
held: Option<Held<B>>,
observations: u64,
}
#[derive(Clone, Debug, PartialEq, Eq)]
struct Held<B> {
body: B,
observed_at: LocalInstant,
producer: ProducerId,
sequence: u64,
observation: u64,
accepted_at: Option<RobotInstant>,
}
impl<B> Lease<B> {
pub fn new(input: &'static str, silence: Duration, hold: Duration) -> Self {
Lease {
input,
silence,
hold,
fence: ProducerFence::new(),
held: None,
observations: 0,
}
}
pub const fn observations(&self) -> u64 {
self.observations
}
pub const fn silence(&self) -> Duration {
self.silence
}
pub const fn hold(&self) -> Duration {
self.hold
}
pub fn producer(&self) -> Option<ProducerId> {
self.fence.accepted_producer()
}
pub fn offer(
&mut self,
producer: ProducerId,
sequence: u64,
observed_at: LocalInstant,
body: B,
) -> LeaseDecision {
let decision = self.fence.admit(producer, sequence);
self.observations = self.observations.saturating_add(1);
let observation = self.observations;
match decision {
LeaseDecision::Renewed => {
tracing::debug!(
target: LEASE_TRACE_TARGET,
input = self.input,
producer = %producer,
sequence,
observation,
decision = "renewed",
"command renewed the lease"
);
}
LeaseDecision::ProducerReplaced { superseded } => {
tracing::info!(
target: LEASE_TRACE_TARGET,
input = self.input,
producer = %producer,
sequence,
observation,
superseded = %superseded,
decision = "producer_replaced",
"a replacement producer took the lease; the previous one is fenced"
);
}
LeaseDecision::Rejected(rejection) => {
tracing::info!(
target: LEASE_TRACE_TARGET,
input = self.input,
producer = %producer,
sequence,
observation,
decision = "rejected",
reason = %rejection,
"command rejected"
);
}
}
match decision {
LeaseDecision::Renewed | LeaseDecision::ProducerReplaced { .. } => {
self.held = Some(Held {
body,
observed_at,
producer,
sequence,
observation,
accepted_at: None,
});
}
LeaseDecision::Rejected(_) => {}
}
decision
}
pub fn live(&mut self, now: LocalInstant, step: RobotInstant) -> Option<&B> {
let input = self.input;
let silence = self.silence;
let hold = self.hold;
let held = self.held.as_mut()?;
let (producer, sequence, observation) = (held.producer, held.sequence, held.observation);
if now.saturating_duration_since(held.observed_at) >= silence {
tracing::info!(
target: LEASE_TRACE_TARGET,
input,
producer = %producer,
sequence,
observation,
decision = "expired_silence",
"the lease expired on host-monotonic silence"
);
self.held = None;
return None;
}
let anchor = match held.accepted_at {
Some(anchor) => anchor,
None => {
tracing::debug!(
target: LEASE_TRACE_TARGET,
input,
producer = %producer,
sequence,
observation,
step = %step,
decision = "applied",
"the held command was applied at this step and anchors the hold horizon"
);
*held.accepted_at.insert(step)
}
};
match step.duration_since(anchor) {
Ok(elapsed) if elapsed < hold => Some(&self.held.as_ref()?.body),
Ok(_) => {
tracing::info!(
target: LEASE_TRACE_TARGET,
input,
producer = %producer,
sequence,
observation,
step = %step,
decision = "expired_hold",
"the lease expired on its logical hold horizon"
);
self.held = None;
None
}
Err(TimelineMismatch { .. }) => {
tracing::info!(
target: LEASE_TRACE_TARGET,
input,
producer = %producer,
sequence,
observation,
step = %step,
decision = "timeline_replaced",
"the world was replaced under the held command, so it was dropped"
);
self.held = None;
None
}
}
}
pub fn clear(&mut self) {
self.held = None;
self.fence.clear();
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::identity::TimelineId;
const SILENCE: Duration = Duration::from_millis(150);
const HOLD: Duration = Duration::from_millis(500);
fn step(line: TimelineId, ticks: u64) -> RobotInstant {
RobotInstant::new(line, ticks)
}
fn ms(value: u64) -> u64 {
value * 1_000_000
}
#[test]
fn a_sequence_that_does_not_increase_is_rejected_without_touching_the_lease() {
let producer = ProducerId::mint();
let start = LocalInstant::from_boot_ns(1_000);
let mut lease = Lease::new("test/input", SILENCE, HOLD);
assert_eq!(
lease.offer(producer, 1, start, "go"),
LeaseDecision::Renewed
);
assert_eq!(
lease.offer(producer, 1, start, "replay"),
LeaseDecision::Rejected(LeaseRejection::StaleSequence {
accepted: 1,
observed: 1
})
);
assert_eq!(
lease.offer(producer, 0, start, "reorder"),
LeaseDecision::Rejected(LeaseRejection::StaleSequence {
accepted: 1,
observed: 0
})
);
let line = TimelineId::mint();
assert_eq!(lease.live(start, step(line, 0)), Some(&"go"));
}
#[test]
fn a_fresh_producer_starting_at_zero_supersedes_and_fences_the_previous_one() {
let first = ProducerId::mint();
let second = ProducerId::mint();
let start = LocalInstant::from_boot_ns(1_000);
let mut lease = Lease::new("test/input", SILENCE, HOLD);
lease.offer(first, 9, start, "old");
assert_eq!(
lease.offer(second, 0, start, "new"),
LeaseDecision::ProducerReplaced { superseded: first }
);
assert!(!lease.fence.permits(first));
assert!(lease.fence.permits(second));
assert_eq!(
lease.offer(first, 10, start, "zombie"),
LeaseDecision::Rejected(LeaseRejection::SupersededProducer {
superseded: first,
accepted: second,
})
);
assert_eq!(lease.producer(), Some(second));
assert_eq!(
lease.live(start, step(TimelineId::mint(), 0)),
Some(&"new"),
"the zombie must not have replaced the held command either"
);
}
#[test]
fn clearing_a_lease_does_not_unfence_a_superseded_producer() {
let first = ProducerId::mint();
let second = ProducerId::mint();
let start = LocalInstant::from_boot_ns(0);
let mut lease = Lease::new("test/input", SILENCE, HOLD);
lease.offer(first, 1, start, "old");
lease.offer(second, 0, start, "new");
lease.clear();
assert!(matches!(
lease.offer(first, 2, start, "zombie"),
LeaseDecision::Rejected(LeaseRejection::SupersededProducer { .. })
));
assert_eq!(
lease.offer(second, 1, start, "live"),
LeaseDecision::Renewed
);
}
#[test]
fn a_superseded_producer_cannot_renew_an_actuator_permit() {
let first = ProducerId::mint();
let second = ProducerId::mint();
let mut fence = ProducerFence::new();
fence.admit(first, 1);
fence.admit(second, 0);
assert_eq!(fence.accepted_producer(), Some(second));
assert!(
!fence.permits(first),
"the fenced authority must not be able to hold a permit open"
);
assert!(matches!(
fence.admit(first, 2),
LeaseDecision::Rejected(LeaseRejection::SupersededProducer { .. })
));
assert_eq!(
fence.accepted_producer(),
Some(second),
"a rejected zombie must not become the accepted producer"
);
}
#[test]
fn a_superseded_producer_stays_fenced_after_many_replacements() {
let first = ProducerId::mint();
let mut fence = ProducerFence::new();
fence.admit(first, 1);
for _ in 0..64 {
fence.admit(ProducerId::mint(), 0);
}
assert!(matches!(
fence.admit(first, 2),
LeaseDecision::Rejected(LeaseRejection::SupersededProducer { .. })
));
assert!(!fence.permits(first));
}
#[test]
fn host_silence_expires_the_lease_independently_of_logical_time() {
let producer = ProducerId::mint();
let line = TimelineId::mint();
let start = LocalInstant::from_boot_ns(0);
let mut lease = Lease::new("test/input", SILENCE, HOLD);
lease.offer(producer, 1, start, "go");
let quiet = start.saturating_add(SILENCE);
assert_eq!(lease.live(quiet, step(line, 0)), None);
assert_eq!(
lease.live(start, step(line, 0)),
None,
"an expired lease stays expired until a new command renews it"
);
}
#[test]
fn the_logical_horizon_expires_the_lease_independently_of_host_time() {
let producer = ProducerId::mint();
let line = TimelineId::mint();
let start = LocalInstant::from_boot_ns(0);
let mut lease = Lease::new("test/input", SILENCE, HOLD);
lease.offer(producer, 1, start, "go");
assert_eq!(lease.live(start, step(line, 0)), Some(&"go"));
assert_eq!(lease.live(start, step(line, ms(499))), Some(&"go"));
assert_eq!(lease.live(start, step(line, ms(500))), None);
}
#[test]
fn a_lease_that_expired_while_paused_does_not_apply_on_the_first_resumed_step() {
let producer = ProducerId::mint();
let line = TimelineId::mint();
let start = LocalInstant::from_boot_ns(0);
let mut lease = Lease::new("test/input", SILENCE, HOLD);
lease.offer(producer, 1, start, "go");
assert_eq!(lease.live(start, step(line, 0)), Some(&"go"));
let resumed = start.saturating_add(SILENCE);
assert_eq!(lease.live(resumed, step(line, 1)), None);
}
#[test]
fn a_replacement_timeline_drops_a_held_command_rather_than_comparing_across_worlds() {
let producer = ProducerId::mint();
let first_line = TimelineId::mint();
let second_line = TimelineId::mint();
let start = LocalInstant::from_boot_ns(0);
let mut lease = Lease::new("test/input", SILENCE, HOLD);
lease.offer(producer, 1, start, "go");
assert_eq!(lease.live(start, step(first_line, 0)), Some(&"go"));
assert_eq!(lease.live(start, step(second_line, 0)), None);
}
}