use crate::metadata_store::{RabsMetadataStore, StoreError};
use rabs_protocol::result_identity::TypedDigest;
use crate::publication::authority_digest;
use rabs_protocol::authority::CoordinatorAuthority;
pub const PUBLICATION_PIN_CLASS: &str = "action-publication";
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Releaser {
Coordinator(CoordinatorAuthority),
Worker(String),
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ReleaseOutcome {
Released,
AlreadyReleased,
UnknownPin,
RefusedWorkerOnPublicationRoot,
RefusedNotOwner,
RefusedNotActiveAuthority,
}
pub fn release_pin_scoped(
store: &mut dyn RabsMetadataStore,
pin_id: u128,
releaser: &Releaser,
) -> Result<ReleaseOutcome, StoreError> {
let Some(row) = store.pin_row(pin_id)? else {
return Ok(ReleaseOutcome::UnknownPin);
};
if row.released {
return Ok(ReleaseOutcome::AlreadyReleased);
}
match releaser {
Releaser::Worker(worker_owner) => {
if row.class == PUBLICATION_PIN_CLASS {
return Ok(ReleaseOutcome::RefusedWorkerOnPublicationRoot);
}
if row.owner != *worker_owner {
return Ok(ReleaseOutcome::RefusedNotOwner);
}
store.release_pin(pin_id, worker_owner)?;
Ok(ReleaseOutcome::Released)
}
Releaser::Coordinator(authority) => {
let presented: TypedDigest = authority_digest(authority);
match store.active_authority()? {
Some(active) if active.digest == presented => {
store.release_pin(pin_id, &row.owner)?;
Ok(ReleaseOutcome::Released)
}
_ => Ok(ReleaseOutcome::RefusedNotActiveAuthority),
}
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum PinProtection {
Protecting,
GraceProtecting,
Expired,
}
pub fn pin_protection(
store: &mut dyn RabsMetadataStore,
pin_id: u128,
now_seq: u64,
reconciliation_confirmed: bool,
grace_windows: u64,
) -> Result<Option<PinProtection>, StoreError> {
let Some(row) = store.pin_row(pin_id)? else {
return Ok(None);
};
if row.released {
return Ok(Some(PinProtection::Expired));
}
let Some(expires_at_seq) = row.expires_at_seq else {
return Ok(Some(PinProtection::Protecting));
};
if now_seq <= expires_at_seq {
return Ok(Some(PinProtection::Protecting));
}
if !reconciliation_confirmed {
return Ok(Some(PinProtection::GraceProtecting));
}
if now_seq <= expires_at_seq.saturating_add(grace_windows) {
return Ok(Some(PinProtection::GraceProtecting));
}
Ok(Some(PinProtection::Expired))
}
#[cfg(test)]
mod tests {
use super::*;
use crate::metadata_store::{
AuthorityRow, FsqliteEngine, RusqliteEngine, SqlEngine, SqlMetadataStore, StoreError,
};
use rabs_protocol::authority::{ClusterId, CoordinatorIncarnationId};
use rabs_protocol::result_identity::DigestAlgorithm;
use std::sync::atomic::{AtomicU64, Ordering};
static DB_COUNTER: AtomicU64 = AtomicU64::new(0);
fn fresh_path(tag: &str) -> std::path::PathBuf {
let n = DB_COUNTER.fetch_add(1, Ordering::SeqCst);
std::env::temp_dir().join(format!("rabs-h041-{}-{}-{}.db", std::process::id(), tag, n))
}
fn digest(tag: u8) -> TypedDigest {
TypedDigest {
algorithm: DigestAlgorithm::Sha256V1,
domain: "rabs.object.sha256.v1",
bytes: [tag; 32],
}
}
fn authority(term: u64) -> CoordinatorAuthority {
CoordinatorAuthority {
cluster_id: ClusterId("c".to_owned()),
credential_generation: 1,
term,
incarnation_id: CoordinatorIncarnationId(9),
}
}
fn install<E: SqlEngine>(store: &mut SqlMetadataStore<E>) {
store
.acquire_authority(&AuthorityRow {
digest: authority_digest(&authority(3)),
cluster_id: "c".to_owned(),
incarnation: 9,
term: 3,
acquired_seq: 1,
})
.unwrap();
store
.create_pin(
1,
&digest(1),
"coordinator",
PUBLICATION_PIN_CLASS,
None,
None,
true,
"publication root",
)
.unwrap();
store
.create_pin(
2,
&digest(2),
"worker-a",
"materialization",
None,
None,
false,
"worker hold",
)
.unwrap();
store
.create_pin(
3,
&digest(3),
"worker-a",
"transfer",
Some(100),
None,
false,
"expiring hold",
)
.unwrap();
}
#[test]
fn h041_worker_can_never_release_a_publication_root() {
let mut store = SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap();
install(&mut store);
for claimed in ["worker-a", "coordinator"] {
assert_eq!(
release_pin_scoped(&mut store, 1, &Releaser::Worker(claimed.to_owned())).unwrap(),
ReleaseOutcome::RefusedWorkerOnPublicationRoot
);
}
assert!(!store.pin_row(1).unwrap().unwrap().released);
}
#[test]
fn h041_release_is_authority_scoped_and_idempotent() {
let mut store = SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap();
install(&mut store);
assert_eq!(
release_pin_scoped(&mut store, 1, &Releaser::Coordinator(authority(2))).unwrap(),
ReleaseOutcome::RefusedNotActiveAuthority
);
assert_eq!(
release_pin_scoped(&mut store, 1, &Releaser::Coordinator(authority(3))).unwrap(),
ReleaseOutcome::Released
);
assert_eq!(
release_pin_scoped(&mut store, 1, &Releaser::Coordinator(authority(3))).unwrap(),
ReleaseOutcome::AlreadyReleased
);
assert_eq!(
release_pin_scoped(&mut store, 2, &Releaser::Worker("worker-b".to_owned())).unwrap(),
ReleaseOutcome::RefusedNotOwner
);
assert_eq!(
release_pin_scoped(&mut store, 2, &Releaser::Worker("worker-a".to_owned())).unwrap(),
ReleaseOutcome::Released
);
assert_eq!(
release_pin_scoped(&mut store, 2, &Releaser::Worker("worker-a".to_owned())).unwrap(),
ReleaseOutcome::AlreadyReleased
);
assert_eq!(
release_pin_scoped(&mut store, 99, &Releaser::Worker("worker-a".to_owned())).unwrap(),
ReleaseOutcome::UnknownPin
);
}
#[test]
fn h041_expiry_fails_toward_retention_through_grace() {
let mut store = SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap();
install(&mut store);
let judge = |store: &mut SqlMetadataStore<RusqliteEngine>,
now: u64,
confirmed: bool|
-> PinProtection {
pin_protection(store, 3, now, confirmed, 10)
.unwrap()
.unwrap()
};
assert_eq!(judge(&mut store, 100, true), PinProtection::Protecting);
assert_eq!(
judge(&mut store, 101, false),
PinProtection::GraceProtecting
);
assert_eq!(
judge(&mut store, 10_000, false),
PinProtection::GraceProtecting
);
assert_eq!(judge(&mut store, 110, true), PinProtection::GraceProtecting);
assert_eq!(judge(&mut store, 111, true), PinProtection::Expired);
assert_eq!(
pin_protection(&mut store, 2, u64::MAX, true, 0)
.unwrap()
.unwrap(),
PinProtection::Protecting
);
release_pin_scoped(&mut store, 2, &Releaser::Worker("worker-a".to_owned())).unwrap();
assert_eq!(
pin_protection(&mut store, 2, 0, false, 0).unwrap().unwrap(),
PinProtection::Expired
);
assert_eq!(pin_protection(&mut store, 99, 0, false, 0).unwrap(), None);
}
#[test]
fn h041_contradictory_lease_renewals_are_inert() {
let mut store = SqlMetadataStore::open(RusqliteEngine::open_in_memory().unwrap()).unwrap();
install(&mut store);
store.renew_pin(3, 5).unwrap();
assert_eq!(
store.renew_pin(3, 5),
Err(StoreError::NonMonotonicPinRenewal)
);
assert_eq!(store.renew_pin(98, 6), Err(StoreError::UnknownPin));
assert_eq!(store.pin_row(3).unwrap().unwrap().renewal_seq, 5);
}
#[test]
fn h041_unresolved_pin_survives_restart() {
let path = fresh_path("restart");
{
let mut store = SqlMetadataStore::open(RusqliteEngine::open(&path).unwrap()).unwrap();
install(&mut store);
}
let mut store = SqlMetadataStore::open(RusqliteEngine::open(&path).unwrap()).unwrap();
assert_eq!(
pin_protection(&mut store, 3, 10_000, false, 10)
.unwrap()
.unwrap(),
PinProtection::GraceProtecting
);
}
#[test]
fn h041_differential_reference_vs_frankensqlite() {
fn scenario<E: SqlEngine>(store: &mut SqlMetadataStore<E>) -> Vec<String> {
install(store);
assert_eq!(
release_pin_scoped(store, 1, &Releaser::Worker("worker-a".to_owned())).unwrap(),
ReleaseOutcome::RefusedWorkerOnPublicationRoot
);
assert_eq!(
release_pin_scoped(store, 2, &Releaser::Worker("worker-a".to_owned())).unwrap(),
ReleaseOutcome::Released
);
assert_eq!(
release_pin_scoped(store, 1, &Releaser::Coordinator(authority(3))).unwrap(),
ReleaseOutcome::Released
);
store.differential_snapshot().unwrap()
}
let mut reference =
SqlMetadataStore::open(RusqliteEngine::open(&fresh_path("ref41")).unwrap()).unwrap();
let mut candidate =
SqlMetadataStore::open(FsqliteEngine::open(&fresh_path("fsq41")).unwrap()).unwrap();
assert_eq!(scenario(&mut reference), scenario(&mut candidate));
}
}