use super::ShardMap;
use crate::op::{Creation, Operation};
use crate::quota::NamespaceUsage;
use crate::relay::relay_watermark;
use crate::repo::{Addressing, RepoId, RepoName};
use crate::rt::Clock;
use crate::store::{
Batch, BatchOutcome, MultipartBlobStore, NamespaceStore, Partition, Precondition, StoreError,
Value, codec, keys,
};
use crate::timers::lease_sweep::lease_reference;
use crate::timers::registry::kinds;
use super::{HookSet, Pipeline, Snapshot, internal, meta_error, ms};
use crate::error::ServerError;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(super) struct LeaseWrite {
pub(super) value: codec::EpochLease,
pub(super) install: bool,
}
pub(super) enum LeaseObservation {
Usable(codec::EpochLease),
Renew(Box<CoordinatorLease>),
}
pub(super) struct CoordinatorLease {
namespace: Option<Value>,
repo: Option<Value>,
visibility: Option<Value>,
epoch: Option<Value>,
authority: Option<Value>,
authority_generation: Option<u64>,
leased_epoch: u64,
shard: Option<Value>,
observed_el: Option<codec::EpochLease>,
recovery: Option<codec::LeaseRecovery>,
recovery_value: Option<Value>,
relay_watermark_ms: u64,
quota_seed: Option<(u64, NamespaceUsage)>,
}
impl CoordinatorLease {
fn creation(&self) -> Creation {
Creation {
namespace: self.namespace.is_none(),
repo: self.repo.is_none(),
}
}
fn epoch(&self) -> u64 {
self.leased_epoch
}
}
impl LeaseObservation {
pub(super) fn quota_seed(&self) -> Option<(u64, NamespaceUsage)> {
match self {
Self::Renew(read) => read.quota_seed,
Self::Usable(_) => None,
}
}
pub(super) fn creation(&self, addressing: &Addressing) -> Creation {
match self {
Self::Renew(read) if matches!(addressing, Addressing::Multi(_)) => read.creation(),
_ => Creation::default(),
}
}
pub(super) fn epoch(&self) -> u64 {
match self {
Self::Usable(lease) => lease.epoch,
Self::Renew(read) => read.epoch(),
}
}
}
pub(super) fn observed_guard(key: crate::store::Key, value: Option<&Value>) -> Precondition {
match value {
Some(value) => Precondition::Equals(key, value.clone()),
None => Precondition::Absent(key),
}
}
fn shard_ref(p: &Partition) -> Result<&str, ServerError> {
match p {
Partition::Ref { shard_ref, .. } => Ok(shard_ref),
_ => Err(internal("epoch lease outside a ref shard")),
}
}
struct LeaseGrant {
creation: Creation,
value: codec::EpochLease,
batch: Batch,
}
const LEASE_GRANT_ATTEMPTS: usize = 8;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct LeaseParams {
pub authority_fence: bool,
pub epoch_lease_ms: u64,
pub lease_margin_ms: u64,
pub min_lease_budget_ms: u64,
}
impl Default for LeaseParams {
fn default() -> Self {
Self {
authority_fence: false,
epoch_lease_ms: 30_000,
lease_margin_ms: 5_000,
min_lease_budget_ms: 1_000,
}
}
}
impl From<&super::PipelineConfig> for LeaseParams {
fn from(cfg: &super::PipelineConfig) -> Self {
Self {
authority_fence: cfg.authority_fence.is_some(),
epoch_lease_ms: cfg.epoch_lease_ms,
lease_margin_ms: cfg.lease_margin_ms,
min_lease_budget_ms: cfg.min_lease_budget_ms,
}
}
}
#[allow(clippy::too_many_lines)] fn grant_batch(
read: &CoordinatorLease,
repo: &RepoName,
p: &Partition,
now: u64,
created_at_ms: u64,
cfg: &LeaseParams,
) -> Result<LeaseGrant, ServerError> {
let shard_ref = shard_ref(p)?;
let ls_key = keys::leased_shard(repo, shard_ref);
let reference = lease_reference(repo, shard_ref);
let epoch = read.epoch();
let old = read
.shard
.as_ref()
.map(codec::decode_leased_shard)
.transpose()
.map_err(meta_error)?;
let recovering = read
.recovery
.and_then(codec::LeaseRecovery::recovery_time)
.is_some_and(|resumed| {
now < resumed
.saturating_add(cfg.epoch_lease_ms)
.saturating_add(cfg.lease_margin_ms)
});
if !recovering && old.is_none_or(|lease| lease.expires_at_ms <= now) {
let observed_ls_expires = old.map_or(0, |lease| lease.expires_at_ms);
debug_assert!(
read.observed_el.is_none_or(|el| el.expires_at_ms
<= observed_ls_expires
.max(now)
.saturating_add(cfg.lease_margin_ms)),
"an observed el outlives every ls it could have been granted under, beyond the skew margin"
);
}
let shard = codec::LeasedShard {
epoch,
expires_at_ms: old
.map_or(0, |l| l.expires_at_ms)
.max(now.saturating_add(cfg.epoch_lease_ms)),
authority_generation: read.authority_generation,
acked_authority_generation: if old.is_some_and(|l| l.expires_at_ms > now) {
old.and_then(|l| l.acked_authority_generation)
} else {
read.authority_generation
},
acked_epoch: old
.filter(|l| l.expires_at_ms > now)
.map_or(epoch, |l| l.acked_epoch),
relay_watermark_ms: old
.map_or(0, |l| l.relay_watermark_ms)
.max(read.relay_watermark_ms),
sweep_due_ms: old
.map_or(0, |l| l.expires_at_ms)
.max(now.saturating_add(cfg.epoch_lease_ms)),
};
let creation = read.creation();
let nr_key = keys::namespace_record();
let rr_key = keys::repo_record(repo);
let namespace = match &read.namespace {
Some(value) => codec::decode_namespace_record(value).map_err(meta_error)?,
None => codec::NamespaceRecord {
created_at_ms,
config_version: 1,
},
};
let mut batch = Batch::new()
.require(if creation.namespace {
Precondition::Absent(nr_key.clone())
} else {
Precondition::Present(nr_key.clone())
})
.require(if creation.repo {
Precondition::Absent(rr_key.clone())
} else {
Precondition::Present(rr_key.clone())
})
.require(observed_guard(keys::grant_epoch(), read.epoch.as_ref()))
.require(observed_guard(ls_key.clone(), read.shard.as_ref()))
.require(observed_guard(
keys::lease_recovery(),
read.recovery_value.as_ref(),
));
if cfg.authority_fence
&& read.recovery.is_none_or(|mode| {
mode.authority_fence != Some(true) || mode.authority_ready != Some(true)
})
{
return Err(
ServerError::unavailable("authority activation pending; retry")
.with_header("Retry-After", "1"),
);
}
if read.authority_generation.is_some() {
batch = batch.require(observed_guard(
keys::authority_generation(),
read.authority.as_ref(),
));
}
if creation.namespace {
batch = batch.put(nr_key, codec::encode_namespace_record(&namespace));
}
if creation.repo {
batch.preconditions.push(observed_guard(
keys::repo_visibility(repo),
read.visibility.as_ref(),
));
super::list_repos::index_writes(&mut batch, repo, true, read.visibility.as_ref())?;
batch = batch.put(
rr_key,
codec::encode_repo_record(&codec::RepoRecord { created_at_ms }),
);
}
if let Some(old) = old {
batch = batch.delete(keys::timer(
old.sweep_due_ms,
kinds::LEASE_SWEEP.get(),
&reference,
));
}
batch = batch
.put(ls_key.clone(), codec::encode_leased_shard(&shard))
.put(
keys::timer(shard.sweep_due_ms, kinds::LEASE_SWEEP.get(), &reference),
Value::default(),
);
let value = codec::EpochLease {
authority_ready: cfg.authority_fence.then_some(true),
epoch,
authority_generation: read.authority_generation,
expires_at_ms: shard.expires_at_ms,
config_version: namespace.config_version,
};
Ok(LeaseGrant {
creation,
value,
batch,
})
}
#[allow(clippy::too_many_arguments, clippy::too_many_lines)] async fn read_lease_rows<L: NamespaceStore, M: NamespaceStore>(
source: &L,
coordinator_store: &M,
shards: &dyn ShardMap,
clock: &dyn Clock,
repo_id: &RepoId,
p: &Partition,
observed_el: Option<codec::EpochLease>,
seed_window: Option<u64>,
authority_fence: bool,
) -> Result<CoordinatorLease, ServerError> {
if !authority_fence && observed_el.is_some_and(|el| el.authority_generation.is_some()) {
return Err(ServerError::unavailable(
"persisted authority lease requires enabled executor",
));
}
let mut wanted = vec![
keys::namespace_record(),
keys::repo_record(&repo_id.name),
keys::grant_epoch(),
keys::leased_shard(&repo_id.name, shard_ref(p)?),
keys::lease_recovery(),
];
if authority_fence {
wanted.push(keys::authority_generation());
}
if let Some(window) = seed_window {
wanted.push(keys::quota_total(window));
}
wanted.push(keys::repo_visibility(&repo_id.name));
let reported = match relay_watermark(source, p, ms(clock.now_ms())).await {
Ok(value) => value,
Err(StoreError::Corrupt(reason)) => {
tracing::warn!(shard = ?p, %reason, "renewal cannot decode relay outbox; reporting zero");
0
}
Err(error) => return Err(meta_error(error)),
};
let rows = coordinator_store
.get_many(&shards.coordinator(&repo_id.namespace), &wanted)
.await
.map_err(meta_error)?;
if rows.len() != wanted.len() {
return Err(internal("lease get_many returned the wrong row count"));
}
let [namespace, repo, epoch, shard, recovery] = &rows[..5] else {
return Err(internal("lease get_many returned the wrong row count"));
};
let mode = recovery
.as_ref()
.map(codec::decode_lease_recovery)
.transpose()
.map_err(meta_error)?;
if !authority_fence
&& (mode.is_some_and(|m| m.authority_fence == Some(true))
|| shard
.as_ref()
.map(codec::decode_leased_shard)
.transpose()
.map_err(meta_error)?
.is_some_and(|row| row.authority_generation.is_some()))
{
return Err(ServerError::unavailable(
"persisted authority fence requires enabled executor",
));
}
let authority = if authority_fence {
rows[5].clone()
} else {
None
};
if authority_fence
&& mode.is_some_and(|m| m.authority_fence == Some(true))
&& authority.is_none()
{
return Err(ServerError::unavailable(
"authority generation missing from fenced namespace",
));
}
let quota_seed = seed_window
.map(|window| {
rows[5 + usize::from(authority_fence)]
.as_ref()
.map(codec::decode_namespace_usage)
.transpose()
.map(|total| (window, total.unwrap_or_default()))
.map_err(meta_error)
})
.transpose()?;
if let Some(value) = namespace {
codec::decode_namespace_record(value).map_err(meta_error)?;
}
if let Some(value) = repo {
codec::decode_repo_record(value).map_err(meta_error)?;
}
if namespace.is_none() && repo.is_some() {
return Err(internal("repository registered without a namespace"));
}
if let Some(value) = shard {
codec::decode_leased_shard(value).map_err(meta_error)?;
}
let read = CoordinatorLease {
namespace: namespace.clone(),
repo: repo.clone(),
visibility: rows.last().cloned().flatten(),
epoch: epoch.clone(),
authority: authority.clone(),
authority_generation: if authority_fence {
Some(
authority
.as_ref()
.map(codec::decode_u64)
.transpose()
.map_err(meta_error)?
.unwrap_or(0),
)
} else {
None
},
leased_epoch: epoch
.as_ref()
.map(codec::decode_u64)
.transpose()
.map_err(meta_error)?
.unwrap_or(0),
shard: shard.clone(),
observed_el,
recovery_value: recovery.clone(),
recovery: recovery
.as_ref()
.map(codec::decode_lease_recovery)
.transpose()
.map_err(meta_error)?,
relay_watermark_ms: reported,
quota_seed,
};
Ok(read)
}
const RELAY_LEASE_BUDGET_MS: u64 = 15_000;
fn relay_apply_error(
metrics: &dyn crate::telemetry::Metrics,
partition: &Partition,
error: StoreError,
) -> ServerError {
if matches!(error, StoreError::Full) {
metrics.incr(
super::METRIC_PARTITION_FULL,
&[("kind", partition.kind())],
1,
);
tracing::error!(kind = partition.kind(), "storage partition full");
ServerError::unavailable("storage partition full")
} else {
meta_error(error)
}
}
pub async fn renew_for_relay<L: NamespaceStore, M: NamespaceStore>(
local: &L,
meta: &M,
shards: &dyn ShardMap,
clock: &dyn Clock,
metrics: &dyn crate::telemetry::Metrics,
repo: &RepoId,
p: &Partition,
params: &LeaseParams,
) -> Result<Value, ServerError> {
for _ in 0..LEASE_GRANT_ATTEMPTS {
let raw = local
.get(p, &keys::epoch_lease())
.await
.map_err(meta_error)?;
let observed = raw
.as_ref()
.map(codec::decode_epoch_lease)
.transpose()
.map_err(meta_error)?;
if !params.authority_fence && observed.is_some_and(|el| el.authority_generation.is_some()) {
return Err(ServerError::unavailable(
"persisted authority lease requires enabled executor",
));
}
let now = ms(clock.now_ms());
if let (Some(raw), Some(lease)) = (&raw, observed)
&& (!params.authority_fence
|| (lease.authority_generation.is_some() && lease.authority_ready == Some(true)))
&& lease
.expires_at_ms
.checked_sub(params.lease_margin_ms)
.and_then(|end| end.checked_sub(now))
.is_some_and(|budget| budget >= RELAY_LEASE_BUDGET_MS)
{
return Ok(raw.clone());
}
let read = read_lease_rows(
local,
meta,
shards,
clock,
repo,
p,
observed,
None,
params.authority_fence,
)
.await?;
let grant = grant_batch(&read, &repo.name, p, now, now, params)?;
let coordinator = shards.coordinator(&repo.namespace);
match meta
.apply(&coordinator, grant.batch)
.await
.map_err(|error| relay_apply_error(metrics, &coordinator, error))?
{
BatchOutcome::Committed => {}
BatchOutcome::PreconditionFailed { .. } => continue,
BatchOutcome::DeadlinePassed { .. } => {
return Err(internal("lease grant had no deadline"));
}
}
let value = codec::encode_epoch_lease(&grant.value);
let install = Batch::new()
.require(Precondition::NotAfter(now.saturating_add(10_000)))
.require(observed_guard(keys::epoch_lease(), raw.as_ref()))
.put(keys::epoch_lease(), value.clone());
if matches!(
local
.apply(p, install)
.await
.map_err(|error| relay_apply_error(metrics, p, error))?,
BatchOutcome::Committed
) {
return Ok(value);
}
}
Err(ServerError::aborted_retryable(
"coordinator lease grant contention",
))
}
impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H> {
pub(super) async fn observe_lease(
&self,
op: &Operation,
p: &Partition,
ahead: Option<&Snapshot>,
) -> Result<LeaseObservation, ServerError> {
let snap = ahead.ok_or_else(|| internal("D34 lease requires an atomic snapshot"))?;
let seed_window = snap.namespace_window.filter(|window| {
let key = keys::quota_view(*window);
snap.contains(&key) && snap.get(&key).is_none()
});
let observed_el = snap
.get(&keys::epoch_lease())
.map(codec::decode_epoch_lease)
.transpose()
.map_err(meta_error)?;
if self.cfg.authority_fence.is_none()
&& observed_el.is_some_and(|el| el.authority_generation.is_some())
{
return Err(ServerError::unavailable(
"persisted authority lease requires enabled executor",
));
}
if let Some(lease) = observed_el.filter(|lease| {
self.cfg.authority_fence.is_none()
|| (lease.authority_generation.is_some() && lease.authority_ready == Some(true))
}) {
let now = ms(self.clock.now_ms());
let usable_until = lease.expires_at_ms.checked_sub(self.cfg.lease_margin_ms);
if usable_until
.and_then(|end| end.checked_sub(now))
.is_some_and(|budget| budget >= self.cfg.min_lease_budget_ms)
{
return Ok(LeaseObservation::Usable(lease));
}
}
Ok(LeaseObservation::Renew(Box::new(
self.read_lease(op, p, observed_el, seed_window).await?,
)))
}
async fn read_lease(
&self,
op: &Operation,
p: &Partition,
observed_el: Option<codec::EpochLease>,
seed_window: Option<u64>,
) -> Result<CoordinatorLease, ServerError> {
let read = read_lease_rows(
&self.meta,
&self.meta,
self.shards.as_ref(),
self.clock.as_ref(),
&op.repo,
p,
observed_el,
seed_window,
self.cfg.authority_fence.is_some(),
)
.await?;
if self.cfg.authority_fence.is_some()
&& read.recovery.is_none_or(|mode| {
mode.authority_fence != Some(true) || mode.authority_ready != Some(true)
})
{
Box::pin(self.ensure_authority_activation(&op.repo.namespace)).await?;
return read_lease_rows(
&self.meta,
&self.meta,
self.shards.as_ref(),
self.clock.as_ref(),
&op.repo,
p,
observed_el,
seed_window,
true,
)
.await;
}
Ok(read)
}
pub(super) async fn admit_lease(
&self,
op: &Operation,
p: &Partition,
observed: LeaseObservation,
skew_ms: i64,
) -> Result<(Creation, LeaseWrite), ServerError> {
let LeaseObservation::Renew(read) = observed else {
let LeaseObservation::Usable(value) = observed else {
unreachable!()
};
return Ok((
Creation::default(),
LeaseWrite {
value,
install: false,
},
));
};
let mut read = *read;
let coordinator = self.shards.coordinator(&op.repo.namespace);
for _ in 0..LEASE_GRANT_ATTEMPTS {
if op
.authz
.grant
.as_ref()
.is_some_and(|grant| grant.epoch != read.leased_epoch)
{
return Err(super::plan::epoch_moved());
}
if let Some(generation) = op.authz.authority_generation
&& Some(generation) != read.authority_generation
{
return Err(crate::authority::moved());
}
let now = ms(self.clock.now_ms());
let created_at_ms = ms(self.clock.now_ms().saturating_add(skew_ms));
let grant = grant_batch(
&read,
&op.repo.name,
p,
now,
created_at_ms,
&LeaseParams::from(&self.cfg),
)?;
let outcome = match self.meta.apply(&coordinator, grant.batch).await {
Ok(outcome) => outcome,
Err(StoreError::Full) => return Err(self.partition_full(&coordinator, None).await),
Err(error) => return Err(meta_error(error)),
};
match outcome {
BatchOutcome::Committed => {
let created = if matches!(self.cfg.addressing, Addressing::Multi(_)) {
grant.creation
} else {
Creation::default()
};
return Ok((
created,
LeaseWrite {
value: grant.value,
install: true,
},
));
}
BatchOutcome::PreconditionFailed { .. } => {
read = self
.read_lease(op, p, read.observed_el, read.quota_seed.map(|(w, _)| w))
.await?;
}
BatchOutcome::DeadlinePassed { .. } => {
return Err(internal("lease grant had no deadline"));
}
}
}
tracing::warn!(shard = ?p, attempts = LEASE_GRANT_ATTEMPTS, "coordinator lease grant did not settle");
Err(ServerError::aborted_retryable(
"coordinator lease grant contention",
))
}
}
#[cfg(test)]
#[path = "lease_model_tests.rs"]
mod model_tests;