use std::sync::Arc;
use std::time::{SystemTime, UNIX_EPOCH};
use anyhow::{bail, Context, Result};
use object_store::path::Path as ObjPath;
use object_store::{ObjectStore, ObjectStoreExt, PutMode, PutOptions, UpdateVersion};
use crate::snapshot::BackupTarget;
pub const UNFENCED_EPOCH: u64 = 0;
pub const FIRST_EPOCH: u64 = 1;
const CLAIM_LEAF: &str = "latest.owner-claim";
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ClaimRecord {
pub epoch: u64,
pub owner: String,
pub claimed_at_nanos: u128,
}
impl ClaimRecord {
fn render(&self) -> String {
format!("{} {} {}\n", self.epoch, self.claimed_at_nanos, self.owner)
}
fn parse(text: &str) -> Result<Self> {
let mut fields = text.split_whitespace();
let epoch: u64 = fields
.next()
.context("owner-claim sidecar is empty; expected '<epoch> <nanos> <owner>'")?
.parse()
.context("owner-claim epoch")?;
if epoch == UNFENCED_EPOCH {
bail!(
"owner-claim sidecar records epoch 0, which means unfenced and is never \
written by acquire(); the sidecar has been tampered with or hand-edited"
);
}
let claimed_at_nanos: u128 = fields
.next()
.context("owner-claim sidecar must be '<epoch> <nanos> <owner>'")?
.parse()
.context("owner-claim claimed_at_nanos")?;
let owner = fields
.next()
.context("owner-claim sidecar must be '<epoch> <nanos> <owner>'")?
.to_string();
Ok(ClaimRecord {
epoch,
claimed_at_nanos,
owner,
})
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ClaimOutcome {
Granted {
claim: ClaimRecord,
previous: Option<ClaimRecord>,
},
Lost { current: ClaimRecord },
}
impl ClaimOutcome {
pub fn granted_epoch(&self) -> Option<u64> {
match self {
ClaimOutcome::Granted { claim, .. } => Some(claim.epoch),
ClaimOutcome::Lost { .. } => None,
}
}
}
fn claim_key(target: &BackupTarget) -> ObjPath {
join_key(&target.prefix, CLAIM_LEAF)
}
async fn read_claim_versioned(
store: &Arc<dyn ObjectStore>,
key: &ObjPath,
) -> Result<Option<(ClaimRecord, UpdateVersion)>> {
match store.get(key).await {
Ok(res) => {
let version = UpdateVersion {
e_tag: res.meta.e_tag.clone(),
version: res.meta.version.clone(),
};
let bytes = res.bytes().await.context("reading owner-claim sidecar")?;
let record = ClaimRecord::parse(&String::from_utf8_lossy(&bytes))
.with_context(|| format!("parsing owner-claim sidecar {key}"))?;
Ok(Some((record, version)))
}
Err(object_store::Error::NotFound { .. }) => Ok(None),
Err(e) => Err(e).context("fetching owner-claim sidecar"),
}
}
pub async fn read_claim(target: &BackupTarget) -> Result<Option<ClaimRecord>> {
Ok(read_claim_versioned(&target.store, &claim_key(target))
.await?
.map(|(record, _)| record))
}
pub async fn acquire(target: &BackupTarget, owner: &str) -> Result<ClaimOutcome> {
if owner.is_empty() || owner.split_whitespace().count() != 1 {
bail!(
"claim owner label {owner:?} must be one non-empty whitespace-free token; it is the \
last field of the sidecar line and would not parse back"
);
}
let key = claim_key(target);
let existing = read_claim_versioned(&target.store, &key).await?;
let (next_epoch, previous, mode) = match &existing {
Some((record, version)) => (
record
.epoch
.checked_add(1)
.context("owner-claim epoch would overflow u64")?,
Some(record.clone()),
PutMode::Update(version.clone()),
),
None => (FIRST_EPOCH, None, PutMode::Create),
};
let claim = ClaimRecord {
epoch: next_epoch,
owner: owner.to_string(),
claimed_at_nanos: unix_nanos(),
};
let opts = PutOptions {
mode,
..Default::default()
};
match target
.store
.put_opts(&key, claim.render().into_bytes().into(), opts)
.await
{
Ok(_) => Ok(ClaimOutcome::Granted { claim, previous }),
Err(object_store::Error::Precondition { .. })
| Err(object_store::Error::AlreadyExists { .. }) => {
let (current, _) = read_claim_versioned(&target.store, &key)
.await?
.context(
"lost the owner-claim compare-and-swap, but the sidecar is now absent; \
something outside turso-backup is deleting claims",
)?;
Ok(ClaimOutcome::Lost { current })
}
Err(e) => Err(e).with_context(|| format!("writing owner-claim sidecar {key}")),
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum ClaimLost {
Superseded { current: ClaimRecord, ours: u64 },
Vanished { ours: u64 },
}
impl std::fmt::Display for ClaimLost {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
ClaimLost::Superseded { current, ours } => write!(
f,
"epoch {ours} was superseded by epoch {} (owner {})",
current.epoch, current.owner
),
ClaimLost::Vanished { ours } => write!(
f,
"epoch {ours} can no longer be verified — the owner-claim sidecar is absent"
),
}
}
}
pub async fn assert_holds(
target: &BackupTarget,
epoch: u64,
) -> Result<std::result::Result<(), ClaimLost>> {
match read_claim(target).await? {
Some(current) if current.epoch == epoch => Ok(Ok(())),
Some(current) => Ok(Err(ClaimLost::Superseded {
current,
ours: epoch,
})),
None => Ok(Err(ClaimLost::Vanished { ours: epoch })),
}
}
fn join_key(prefix: &str, leaf: &str) -> ObjPath {
let prefix = prefix.trim_matches('/');
if prefix.is_empty() {
ObjPath::from(leaf)
} else {
ObjPath::from(format!("{prefix}/{leaf}"))
}
}
fn unix_nanos() -> u128 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.map(|d| d.as_nanos())
.unwrap_or(0)
}
#[cfg(test)]
mod tests {
use super::*;
use object_store::memory::InMemory;
fn target(prefix: &str) -> BackupTarget {
BackupTarget {
store: Arc::new(InMemory::new()),
prefix: prefix.to_string(),
}
}
fn shared(prefix: &str) -> (BackupTarget, BackupTarget) {
let store: Arc<dyn ObjectStore> = Arc::new(InMemory::new());
(
BackupTarget {
store: store.clone(),
prefix: prefix.to_string(),
},
BackupTarget {
store,
prefix: prefix.to_string(),
},
)
}
#[tokio::test]
async fn a_virgin_prefix_is_unclaimed_and_the_first_acquire_mints_epoch_one() {
let t = target("acct");
assert_eq!(read_claim(&t).await.unwrap(), None);
let out = acquire(&t, "us-east-001").await.unwrap();
let ClaimOutcome::Granted { claim, previous } = out else {
panic!("first acquire on a virgin prefix must be granted, got {out:?}");
};
assert_eq!(claim.epoch, FIRST_EPOCH);
assert_eq!(claim.owner, "us-east-001");
assert_eq!(previous, None);
assert_eq!(read_claim(&t).await.unwrap(), Some(claim));
}
#[tokio::test]
async fn a_takeover_advances_the_epoch_and_names_who_it_displaced() {
let (a, b) = shared("acct");
acquire(&a, "us-east-001").await.unwrap();
let out = acquire(&b, "us-west-003").await.unwrap();
let ClaimOutcome::Granted { claim, previous } = out else {
panic!("a sequentially-later acquire is a takeover and must be granted, got {out:?}");
};
assert_eq!(claim.epoch, FIRST_EPOCH + 1);
assert_eq!(claim.owner, "us-west-003");
assert_eq!(previous.unwrap().owner, "us-east-001");
}
#[tokio::test]
async fn the_displaced_owner_fails_assert_holds_and_is_told_who_won() {
let (a, b) = shared("acct");
let ours = acquire(&a, "us-east-001")
.await
.unwrap()
.granted_epoch()
.unwrap();
assert!(assert_holds(&a, ours).await.unwrap().is_ok());
acquire(&b, "us-west-003").await.unwrap();
let lost = assert_holds(&a, ours).await.unwrap().unwrap_err();
let ClaimLost::Superseded { current, ours: had } = lost else {
panic!("expected Superseded, got {lost:?}");
};
assert_eq!(had, ours);
assert_eq!(current.epoch, ours + 1);
assert_eq!(current.owner, "us-west-003");
}
#[tokio::test]
async fn assert_holds_ignores_the_owner_label() {
let t = target("acct");
let e = acquire(&t, "node-a").await.unwrap().granted_epoch().unwrap();
assert!(assert_holds(&t, e).await.unwrap().is_ok());
}
#[tokio::test]
async fn two_bootstrappers_racing_a_fresh_prefix_produce_exactly_one_winner() {
let (a, b) = shared("acct");
let key = claim_key(&a);
acquire(&a, "node-a").await.unwrap();
let stale_create = b
.store
.put_opts(
&key,
b"1 0 node-b".as_slice().into(),
PutOptions {
mode: PutMode::Create,
..Default::default()
},
)
.await;
assert!(
stale_create.is_err(),
"PutMode::Create must be refused once the sidecar exists — if this passes, the \
store has degraded to unconditional writes and the fence is not real"
);
assert_eq!(read_claim(&a).await.unwrap().unwrap().owner, "node-a");
}
#[tokio::test]
async fn a_lost_compare_and_swap_reports_the_winner_rather_than_minting() {
let (a, b) = shared("acct");
acquire(&a, "node-a").await.unwrap();
let key = claim_key(&b);
let (_, v1) = read_claim_versioned(&b.store, &key).await.unwrap().unwrap();
acquire(&a, "node-a").await.unwrap();
let bounced = b
.store
.put_opts(
&key,
b"2 0 node-b".as_slice().into(),
PutOptions {
mode: PutMode::Update(v1),
..Default::default()
},
)
.await;
assert!(bounced.is_err(), "stale-version update must be refused");
let out = acquire(&b, "node-b").await.unwrap();
assert_eq!(out.granted_epoch(), Some(3));
}
#[tokio::test]
async fn a_garbled_sidecar_is_an_error_not_an_unclaimed_prefix() {
let t = target("acct");
t.store
.put(&claim_key(&t), b"not a claim at all".as_slice().into())
.await
.unwrap();
let err = read_claim(&t).await.unwrap_err();
assert!(
format!("{err:#}").contains("owner-claim"),
"error should name the sidecar: {err:#}"
);
}
#[tokio::test]
async fn epoch_zero_in_a_sidecar_is_rejected() {
let t = target("acct");
t.store
.put(&claim_key(&t), b"0 123 node-a".as_slice().into())
.await
.unwrap();
let err = read_claim(&t).await.unwrap_err();
assert!(format!("{err:#}").contains("epoch 0"), "{err:#}");
}
#[tokio::test]
async fn an_owner_label_with_whitespace_is_refused_before_anything_is_written() {
let t = target("acct");
let err = acquire(&t, "us east 001").await.unwrap_err();
assert!(format!("{err:#}").contains("whitespace-free"), "{err:#}");
assert_eq!(
read_claim(&t).await.unwrap(),
None,
"a refused label must not leave a sidecar behind"
);
}
#[tokio::test]
async fn an_empty_prefix_claims_at_the_bucket_root() {
let t = target("");
acquire(&t, "node-a").await.unwrap();
assert_eq!(claim_key(&t), ObjPath::from(CLAIM_LEAF));
assert!(read_claim(&t).await.unwrap().is_some());
}
#[tokio::test]
async fn assert_holds_on_a_deleted_sidecar_is_vanished_not_held() {
let t = target("acct");
let e = acquire(&t, "node-a").await.unwrap().granted_epoch().unwrap();
t.store.delete(&claim_key(&t)).await.unwrap();
assert_eq!(
assert_holds(&t, e).await.unwrap().unwrap_err(),
ClaimLost::Vanished { ours: e }
);
}
}