use crate::fold::SyncState;
use crate::lease::{LeaseCoordinator, LeaseError};
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "decision", rename_all = "snake_case")]
pub enum FenceDecision {
Proceed,
AlreadyCommitted {
committed_op_id: String,
},
StaleEpoch {
claimed_epoch: u64,
current_epoch: u64,
current_holder: Option<String>,
},
NotHeld { claimed_epoch: u64 },
}
impl FenceDecision {
pub fn may_dispatch(&self) -> bool {
matches!(self, FenceDecision::Proceed)
}
}
pub fn check_dispatch(
coordinator: &mut dyn LeaseCoordinator,
state: &SyncState,
agent_id: &str,
run_id: &str,
device_id: &str,
epoch: u64,
) -> Result<FenceDecision, LeaseError> {
if let Some(committed) = state.committed_run(agent_id, run_id) {
return Ok(FenceDecision::AlreadyCommitted {
committed_op_id: committed.op_id.clone(),
});
}
match coordinator.current(agent_id)? {
None => Ok(FenceDecision::NotHeld { claimed_epoch: epoch }),
Some(lease) => {
if lease.holder == device_id && lease.epoch == epoch {
Ok(FenceDecision::Proceed)
} else {
Ok(FenceDecision::StaleEpoch {
claimed_epoch: epoch,
current_epoch: lease.epoch,
current_holder: Some(lease.holder),
})
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::fold::fold;
use crate::lease::{InMemoryLeaseCoordinator, Intent, IntentStatus};
use crate::oplog::{logical_clock, DeviceLog, Scope, Surface};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
fn coord() -> (Arc<AtomicU64>, InMemoryLeaseCoordinator) {
let t = Arc::new(AtomicU64::new(0));
let reader = t.clone();
(t, InMemoryLeaseCoordinator::new(Arc::new(move || reader.load(Ordering::SeqCst))))
}
fn state_with_commit(run: &str, epoch: u64) -> SyncState {
let mut dev = DeviceLog::new("dev-a");
dev.set_wall_clock(logical_clock());
let op = dev.append(
Scope::Personal,
Surface::Intent,
Intent::new("milo", run, epoch, IntentStatus::Committed).payload(),
);
fold(&[op])
}
#[test]
fn holder_at_current_epoch_may_dispatch() {
let (_t, mut c) = coord();
let lease = c.acquire("milo", "dev-a", 100).unwrap();
let empty = SyncState::default();
let decision =
check_dispatch(&mut c, &empty, "milo", "run-1", "dev-a", lease.epoch).unwrap();
assert_eq!(decision, FenceDecision::Proceed);
assert!(decision.may_dispatch());
}
#[test]
fn stale_epoch_holder_is_refused() {
let (t, mut c) = coord();
c.acquire("milo", "dev-a", 100).unwrap();
t.store(200, Ordering::SeqCst);
let stolen = c.acquire("milo", "dev-b", 100).unwrap();
assert_eq!(stolen.epoch, 2);
let empty = SyncState::default();
let decision = check_dispatch(&mut c, &empty, "milo", "run-1", "dev-a", 1).unwrap();
assert_eq!(
decision,
FenceDecision::StaleEpoch {
claimed_epoch: 1,
current_epoch: 2,
current_holder: Some("dev-b".to_string()),
}
);
assert!(!decision.may_dispatch(), "a stale-epoch holder must NOT dispatch");
}
#[test]
fn already_committed_run_is_not_re_executed_even_by_the_current_holder() {
let (_t, mut c) = coord();
let lease = c.acquire("milo", "dev-a", 100).unwrap();
let state = state_with_commit("run-nightly", 1);
let decision =
check_dispatch(&mut c, &state, "milo", "run-nightly", "dev-a", lease.epoch).unwrap();
match decision {
FenceDecision::AlreadyCommitted { ref committed_op_id } => {
assert_eq!(
*committed_op_id,
state.committed_run("milo", "run-nightly").unwrap().op_id
);
}
other => panic!("expected AlreadyCommitted, got {other:?}"),
}
assert!(!decision.may_dispatch());
}
#[test]
fn committed_oracle_beats_a_stale_epoch_check() {
let (t, mut c) = coord();
c.acquire("milo", "dev-a", 100).unwrap();
t.store(200, Ordering::SeqCst);
c.acquire("milo", "dev-b", 100).unwrap();
let state = state_with_commit("run-x", 2);
let decision = check_dispatch(&mut c, &state, "milo", "run-x", "dev-a", 1).unwrap();
assert!(matches!(decision, FenceDecision::AlreadyCommitted { .. }));
}
#[test]
fn unheld_agent_is_not_authorized() {
let (_t, mut c) = coord();
let empty = SyncState::default();
let decision = check_dispatch(&mut c, &empty, "milo", "run-1", "dev-a", 1).unwrap();
assert_eq!(decision, FenceDecision::NotHeld { claimed_epoch: 1 });
assert!(!decision.may_dispatch());
}
}