use std::cell::RefCell;
use std::collections::BTreeSet;
use crate::Context;
use crate::cell::Source;
#[derive(Debug, Clone)]
pub struct LeaseCore<P> {
holder: Option<P>,
expiry: u64,
fence: u64,
}
impl<P: Clone + PartialEq> Default for LeaseCore<P> {
fn default() -> Self {
Self::new()
}
}
impl<P: Clone + PartialEq> LeaseCore<P> {
pub fn new() -> Self {
Self {
holder: None,
expiry: 0,
fence: 0,
}
}
fn is_expired(&self, now: u64) -> bool {
self.holder.is_some() && now >= self.expiry
}
pub fn is_held(&self, now: u64) -> bool {
self.holder.is_some() && !self.is_expired(now)
}
pub fn holder(&self, now: u64) -> Option<P> {
if self.is_held(now) {
self.holder.clone()
} else {
None
}
}
pub fn fence(&self) -> u64 {
self.fence
}
pub fn acquire(&mut self, peer: P, now: u64, ttl: u64) -> Option<u64> {
let free = match &self.holder {
None => true,
Some(_) => self.is_expired(now),
};
if free {
self.fence += 1;
self.holder = Some(peer);
self.expiry = now + ttl;
return Some(self.fence);
}
if self.holder.as_ref() == Some(&peer) {
self.expiry = now + ttl; return Some(self.fence);
}
None
}
pub fn renew(&mut self, peer: P, now: u64, ttl: u64) -> bool {
if self.is_held(now) && self.holder.as_ref() == Some(&peer) {
self.expiry = now + ttl;
true
} else {
false
}
}
pub fn release(&mut self, peer: &P) {
if self.holder.as_ref() == Some(peer) {
self.holder = None;
}
}
pub fn tick(&mut self, now: u64) -> bool {
if self.is_expired(now) {
self.holder = None;
true
} else {
false
}
}
}
pub struct LeaseCell<P> {
core: RefCell<LeaseCore<P>>,
holder: Source<Option<P>>,
}
impl<P: Clone + PartialEq + 'static> LeaseCell<P> {
pub fn new(ctx: &Context) -> Self {
Self {
core: RefCell::new(LeaseCore::new()),
holder: ctx.source(None),
}
}
fn refresh(&self, ctx: &Context, now: u64) {
let h = self.core.borrow().holder(now);
self.holder.set(ctx, h);
}
pub fn acquire(&self, ctx: &Context, peer: P, now: u64, ttl: u64) -> Option<u64> {
let r = self.core.borrow_mut().acquire(peer, now, ttl);
self.refresh(ctx, now);
r
}
pub fn renew(&self, ctx: &Context, peer: P, now: u64, ttl: u64) -> bool {
let r = self.core.borrow_mut().renew(peer, now, ttl);
self.refresh(ctx, now);
r
}
pub fn release(&self, ctx: &Context, peer: &P, now: u64) {
self.core.borrow_mut().release(peer);
self.refresh(ctx, now);
}
pub fn tick(&self, ctx: &Context, now: u64) -> bool {
let r = self.core.borrow_mut().tick(now);
self.refresh(ctx, now);
r
}
pub fn holder(&self, now: u64) -> Option<P> {
self.core.borrow().holder(now)
}
pub fn is_held(&self, now: u64) -> bool {
self.core.borrow().is_held(now)
}
pub fn fence(&self) -> u64 {
self.core.borrow().fence()
}
pub fn holder_cell(&self) -> Source<Option<P>> {
self.holder
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum LeaderRole {
Leader,
Follower,
Candidate,
}
pub struct LeaderCell<P> {
core: RefCell<LeaseCore<P>>,
me: P,
current_leader: Source<Option<P>>,
}
impl<P: Clone + PartialEq + 'static> LeaderCell<P> {
pub fn new(ctx: &Context, me: P) -> Self {
Self {
core: RefCell::new(LeaseCore::new()),
me,
current_leader: ctx.source(None),
}
}
fn refresh(&self, ctx: &Context, now: u64) {
let l = self.core.borrow().holder(now);
self.current_leader.set(ctx, l);
}
pub fn campaign(&self, ctx: &Context, now: u64, ttl: u64) -> LeaderRole {
self.core.borrow_mut().acquire(self.me.clone(), now, ttl);
self.refresh(ctx, now);
self.role(now)
}
pub fn contend(&self, ctx: &Context, peer: P, now: u64, ttl: u64) -> LeaderRole {
self.core.borrow_mut().acquire(peer, now, ttl);
self.refresh(ctx, now);
self.role(now)
}
pub fn tick(&self, ctx: &Context, now: u64) -> LeaderRole {
self.core.borrow_mut().tick(now);
self.refresh(ctx, now);
self.role(now)
}
pub fn current_leader(&self, now: u64) -> Option<P> {
self.core.borrow().holder(now)
}
pub fn role(&self, now: u64) -> LeaderRole {
match self.core.borrow().holder(now) {
Some(h) if h == self.me => LeaderRole::Leader,
Some(_) => LeaderRole::Follower,
None => LeaderRole::Candidate,
}
}
pub fn current_leader_cell(&self) -> Source<Option<P>> {
self.current_leader
}
}
pub struct LockCell<P> {
core: RefCell<LeaseCore<P>>,
is_locked: Source<bool>,
}
impl<P: Clone + PartialEq + 'static> LockCell<P> {
pub fn new(ctx: &Context) -> Self {
Self {
core: RefCell::new(LeaseCore::new()),
is_locked: ctx.source(false),
}
}
fn refresh(&self, ctx: &Context, now: u64) {
let held = self.core.borrow().is_held(now);
self.is_locked.set(ctx, held);
}
pub fn acquire(&self, ctx: &Context, peer: P, now: u64, ttl: u64) -> Option<u64> {
let r = self.core.borrow_mut().acquire(peer, now, ttl);
self.refresh(ctx, now);
r
}
pub fn release(&self, ctx: &Context, peer: &P, now: u64) {
self.core.borrow_mut().release(peer);
self.refresh(ctx, now);
}
pub fn tick(&self, ctx: &Context, now: u64) -> bool {
let r = self.core.borrow_mut().tick(now);
self.refresh(ctx, now);
r
}
pub fn validate(&self, fence: u64) -> bool {
self.core.borrow().fence() == fence
}
pub fn is_locked(&self, now: u64) -> bool {
self.core.borrow().is_held(now)
}
pub fn fence(&self) -> u64 {
self.core.borrow().fence()
}
pub fn is_locked_cell(&self) -> Source<bool> {
self.is_locked
}
}
#[derive(Debug, Clone, Copy)]
pub struct SemaphoreCore {
capacity: u64,
acquired: u64,
}
impl SemaphoreCore {
pub fn new(capacity: u64) -> Self {
Self {
capacity,
acquired: 0,
}
}
pub fn available(&self) -> u64 {
self.capacity - self.acquired
}
pub fn acquire(&mut self) -> bool {
if self.acquired < self.capacity {
self.acquired += 1;
true
} else {
false
}
}
pub fn release(&mut self) {
if self.acquired > 0 {
self.acquired -= 1;
}
}
}
pub struct SemaphoreCell {
core: RefCell<SemaphoreCore>,
available: Source<u64>,
}
impl SemaphoreCell {
pub fn new(ctx: &Context, capacity: u64) -> Self {
Self {
core: RefCell::new(SemaphoreCore::new(capacity)),
available: ctx.source(capacity),
}
}
fn refresh(&self, ctx: &Context) {
let a = self.core.borrow().available();
self.available.set(ctx, a);
}
pub fn acquire(&self, ctx: &Context) -> bool {
let r = self.core.borrow_mut().acquire();
self.refresh(ctx);
r
}
pub fn release(&self, ctx: &Context) {
self.core.borrow_mut().release();
self.refresh(ctx);
}
pub fn permits_available(&self, ctx: &Context) -> u64 {
self.available.get(ctx)
}
pub fn permits_available_cell(&self) -> Source<u64> {
self.available
}
}
#[derive(Debug, Clone)]
pub struct BarrierCore<P> {
required: u64,
arrived: BTreeSet<P>,
}
impl<P: Ord + Clone> BarrierCore<P> {
pub fn new(required: u64) -> Self {
Self {
required,
arrived: BTreeSet::new(),
}
}
pub fn arrive(&mut self, peer: P) -> bool {
self.arrived.insert(peer);
self.is_open()
}
pub fn count(&self) -> u64 {
self.arrived.len() as u64
}
pub fn is_open(&self) -> bool {
self.count() >= self.required
}
}
pub struct BarrierCell<P> {
core: RefCell<BarrierCore<P>>,
is_open: Source<bool>,
}
impl<P: Ord + Clone + 'static> BarrierCell<P> {
pub fn new(ctx: &Context, required: u64) -> Self {
let core = BarrierCore::new(required);
let open = core.is_open();
Self {
core: RefCell::new(core),
is_open: ctx.source(open),
}
}
pub fn quorum(ctx: &Context, total: u64) -> Self {
Self::new(ctx, total / 2 + 1)
}
fn refresh(&self, ctx: &Context) {
let o = self.core.borrow().is_open();
self.is_open.set(ctx, o);
}
pub fn arrive(&self, ctx: &Context, peer: P) -> bool {
let r = self.core.borrow_mut().arrive(peer);
self.refresh(ctx);
r
}
pub fn count(&self) -> u64 {
self.core.borrow().count()
}
pub fn is_open(&self, ctx: &Context) -> bool {
self.is_open.get(ctx)
}
pub fn is_open_cell(&self) -> Source<bool> {
self.is_open
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn lease_fence_monotone_renew_keeps() {
let mut l = LeaseCore::<u64>::new();
assert_eq!(l.acquire(1, 0, 10), Some(1));
assert_eq!(l.acquire(2, 1, 10), None); assert!(l.renew(1, 5, 10));
assert_eq!(l.fence(), 1); assert!(l.tick(15)); assert_eq!(l.acquire(2, 16, 10), Some(2)); }
#[test]
fn semaphore_bounded_ops() {
let mut s = SemaphoreCore::new(2);
assert!(s.acquire());
assert!(s.acquire());
assert!(!s.acquire()); assert_eq!(s.available(), 0);
s.release();
assert_eq!(s.available(), 1);
}
#[test]
fn quorum_opens_at_majority() {
let mut b = BarrierCore::<u64>::new(5 / 2 + 1); assert!(!b.arrive(1));
assert!(!b.arrive(2));
assert!(b.arrive(3)); assert!(b.arrive(1)); assert_eq!(b.count(), 3);
}
}