1use super::ShardMap;
11use crate::op::{Creation, Operation};
12use crate::quota::NamespaceUsage;
13use crate::relay::relay_watermark;
14use crate::repo::{Addressing, RepoId, RepoName};
15use crate::rt::Clock;
16use crate::store::{
17 Batch, BatchOutcome, MultipartBlobStore, NamespaceStore, Partition, Precondition, StoreError,
18 Value, codec, keys,
19};
20use crate::timers::lease_sweep::lease_reference;
21use crate::timers::registry::kinds;
22
23use super::{HookSet, Pipeline, Snapshot, internal, meta_error, ms};
24use crate::error::ServerError;
25
26#[derive(Debug, Clone, Copy, PartialEq, Eq)]
28pub(super) struct LeaseWrite {
29 pub(super) value: codec::EpochLease,
30 pub(super) install: bool,
31}
32
33pub(super) enum LeaseObservation {
34 Usable(codec::EpochLease),
35 Renew(Box<CoordinatorLease>),
36}
37
38pub(super) struct CoordinatorLease {
39 namespace: Option<Value>,
40 repo: Option<Value>,
41 visibility: Option<Value>,
42 epoch: Option<Value>,
43 authority: Option<Value>,
44 authority_generation: Option<u64>,
45 leased_epoch: u64,
46 shard: Option<Value>,
47 observed_el: Option<codec::EpochLease>,
48 recovery: Option<codec::LeaseRecovery>,
49 recovery_value: Option<Value>,
50 relay_watermark_ms: u64,
51 quota_seed: Option<(u64, NamespaceUsage)>,
52}
53
54impl CoordinatorLease {
55 fn creation(&self) -> Creation {
56 Creation {
57 namespace: self.namespace.is_none(),
58 repo: self.repo.is_none(),
59 }
60 }
61
62 fn epoch(&self) -> u64 {
63 self.leased_epoch
64 }
65}
66
67impl LeaseObservation {
68 pub(super) fn quota_seed(&self) -> Option<(u64, NamespaceUsage)> {
69 match self {
70 Self::Renew(read) => read.quota_seed,
71 Self::Usable(_) => None,
72 }
73 }
74
75 pub(super) fn creation(&self, addressing: &Addressing) -> Creation {
76 match self {
77 Self::Renew(read) if matches!(addressing, Addressing::Multi(_)) => read.creation(),
78 _ => Creation::default(),
79 }
80 }
81
82 pub(super) fn epoch(&self) -> u64 {
83 match self {
84 Self::Usable(lease) => lease.epoch,
85 Self::Renew(read) => read.epoch(),
86 }
87 }
88}
89
90pub(super) fn observed_guard(key: crate::store::Key, value: Option<&Value>) -> Precondition {
91 match value {
92 Some(value) => Precondition::Equals(key, value.clone()),
93 None => Precondition::Absent(key),
94 }
95}
96
97fn shard_ref(p: &Partition) -> Result<&str, ServerError> {
98 match p {
99 Partition::Ref { shard_ref, .. } => Ok(shard_ref),
100 _ => Err(internal("epoch lease outside a ref shard")),
101 }
102}
103
104struct LeaseGrant {
105 creation: Creation,
106 value: codec::EpochLease,
107 batch: Batch,
108}
109
110const LEASE_GRANT_ATTEMPTS: usize = 8;
114
115#[derive(Debug, Clone, Copy, PartialEq, Eq)]
119pub struct LeaseParams {
120 pub authority_fence: bool,
122 pub epoch_lease_ms: u64,
124 pub lease_margin_ms: u64,
126 pub min_lease_budget_ms: u64,
128}
129
130impl Default for LeaseParams {
131 fn default() -> Self {
132 Self {
133 authority_fence: false,
134 epoch_lease_ms: 30_000,
135 lease_margin_ms: 5_000,
136 min_lease_budget_ms: 1_000,
137 }
138 }
139}
140
141impl From<&super::PipelineConfig> for LeaseParams {
142 fn from(cfg: &super::PipelineConfig) -> Self {
143 Self {
144 authority_fence: cfg.authority_fence.is_some(),
145 epoch_lease_ms: cfg.epoch_lease_ms,
146 lease_margin_ms: cfg.lease_margin_ms,
147 min_lease_budget_ms: cfg.min_lease_budget_ms,
148 }
149 }
150}
151
152#[allow(clippy::too_many_lines)] fn grant_batch(
154 read: &CoordinatorLease,
155 repo: &RepoName,
156 p: &Partition,
157 now: u64,
158 created_at_ms: u64,
159 cfg: &LeaseParams,
160) -> Result<LeaseGrant, ServerError> {
161 let shard_ref = shard_ref(p)?;
162 let ls_key = keys::leased_shard(repo, shard_ref);
163 let reference = lease_reference(repo, shard_ref);
164 let epoch = read.epoch();
165 let old = read
166 .shard
167 .as_ref()
168 .map(codec::decode_leased_shard)
169 .transpose()
170 .map_err(meta_error)?;
171 let recovering = read
174 .recovery
175 .and_then(codec::LeaseRecovery::recovery_time)
176 .is_some_and(|resumed| {
177 now < resumed
178 .saturating_add(cfg.epoch_lease_ms)
179 .saturating_add(cfg.lease_margin_ms)
180 });
181 if !recovering && old.is_none_or(|lease| lease.expires_at_ms <= now) {
182 let observed_ls_expires = old.map_or(0, |lease| lease.expires_at_ms);
183 debug_assert!(
186 read.observed_el.is_none_or(|el| el.expires_at_ms
187 <= observed_ls_expires
188 .max(now)
189 .saturating_add(cfg.lease_margin_ms)),
190 "an observed el outlives every ls it could have been granted under, beyond the skew margin"
191 );
192 }
193 let shard = codec::LeasedShard {
194 epoch,
195 expires_at_ms: old
196 .map_or(0, |l| l.expires_at_ms)
197 .max(now.saturating_add(cfg.epoch_lease_ms)),
198 authority_generation: read.authority_generation,
199 acked_authority_generation: if old.is_some_and(|l| l.expires_at_ms > now) {
200 old.and_then(|l| l.acked_authority_generation)
201 } else {
202 read.authority_generation
203 },
204 acked_epoch: old
205 .filter(|l| l.expires_at_ms > now)
206 .map_or(epoch, |l| l.acked_epoch),
207 relay_watermark_ms: old
208 .map_or(0, |l| l.relay_watermark_ms)
209 .max(read.relay_watermark_ms),
210 sweep_due_ms: old
211 .map_or(0, |l| l.expires_at_ms)
212 .max(now.saturating_add(cfg.epoch_lease_ms)),
213 };
214 let creation = read.creation();
215 let nr_key = keys::namespace_record();
216 let rr_key = keys::repo_record(repo);
217 let namespace = match &read.namespace {
218 Some(value) => codec::decode_namespace_record(value).map_err(meta_error)?,
219 None => codec::NamespaceRecord {
220 created_at_ms,
221 config_version: 1,
222 },
223 };
224 let mut batch = Batch::new()
225 .require(if creation.namespace {
226 Precondition::Absent(nr_key.clone())
227 } else {
228 Precondition::Present(nr_key.clone())
229 })
230 .require(if creation.repo {
231 Precondition::Absent(rr_key.clone())
232 } else {
233 Precondition::Present(rr_key.clone())
234 })
235 .require(observed_guard(keys::grant_epoch(), read.epoch.as_ref()))
236 .require(observed_guard(ls_key.clone(), read.shard.as_ref()))
237 .require(observed_guard(
238 keys::lease_recovery(),
239 read.recovery_value.as_ref(),
240 ));
241 if cfg.authority_fence
242 && read.recovery.is_none_or(|mode| {
243 mode.authority_fence != Some(true) || mode.authority_ready != Some(true)
244 })
245 {
246 return Err(
247 ServerError::unavailable("authority activation pending; retry")
248 .with_header("Retry-After", "1"),
249 );
250 }
251 if read.authority_generation.is_some() {
252 batch = batch.require(observed_guard(
253 keys::authority_generation(),
254 read.authority.as_ref(),
255 ));
256 }
257 if creation.namespace {
258 batch = batch.put(nr_key, codec::encode_namespace_record(&namespace));
259 }
260 if creation.repo {
261 batch.preconditions.push(observed_guard(
262 keys::repo_visibility(repo),
263 read.visibility.as_ref(),
264 ));
265 super::list_repos::index_writes(&mut batch, repo, true, read.visibility.as_ref())?;
266 batch = batch.put(
267 rr_key,
268 codec::encode_repo_record(&codec::RepoRecord { created_at_ms }),
269 );
270 }
271 if let Some(old) = old {
272 batch = batch.delete(keys::timer(
273 old.sweep_due_ms,
274 kinds::LEASE_SWEEP.get(),
275 &reference,
276 ));
277 }
278 batch = batch
279 .put(ls_key.clone(), codec::encode_leased_shard(&shard))
280 .put(
281 keys::timer(shard.sweep_due_ms, kinds::LEASE_SWEEP.get(), &reference),
282 Value::default(),
283 );
284 let value = codec::EpochLease {
285 authority_ready: cfg.authority_fence.then_some(true),
286 epoch,
287 authority_generation: read.authority_generation,
288 expires_at_ms: shard.expires_at_ms,
289 config_version: namespace.config_version,
290 };
291 Ok(LeaseGrant {
292 creation,
293 value,
294 batch,
295 })
296}
297
298#[allow(clippy::too_many_arguments, clippy::too_many_lines)] async fn read_lease_rows<L: NamespaceStore, M: NamespaceStore>(
303 source: &L,
304 coordinator_store: &M,
305 shards: &dyn ShardMap,
306 clock: &dyn Clock,
307 repo_id: &RepoId,
308 p: &Partition,
309 observed_el: Option<codec::EpochLease>,
310 seed_window: Option<u64>,
311 authority_fence: bool,
312) -> Result<CoordinatorLease, ServerError> {
313 if !authority_fence && observed_el.is_some_and(|el| el.authority_generation.is_some()) {
314 return Err(ServerError::unavailable(
315 "persisted authority lease requires enabled executor",
316 ));
317 }
318 let mut wanted = vec![
319 keys::namespace_record(),
320 keys::repo_record(&repo_id.name),
321 keys::grant_epoch(),
322 keys::leased_shard(&repo_id.name, shard_ref(p)?),
323 keys::lease_recovery(),
324 ];
325 if authority_fence {
326 wanted.push(keys::authority_generation());
327 }
328 if let Some(window) = seed_window {
329 wanted.push(keys::quota_total(window));
330 }
331 wanted.push(keys::repo_visibility(&repo_id.name));
332 let reported = match relay_watermark(source, p, ms(clock.now_ms())).await {
333 Ok(value) => value,
334 Err(StoreError::Corrupt(reason)) => {
335 tracing::warn!(shard = ?p, %reason, "renewal cannot decode relay outbox; reporting zero");
336 0
337 }
338 Err(error) => return Err(meta_error(error)),
339 };
340 let rows = coordinator_store
341 .get_many(&shards.coordinator(&repo_id.namespace), &wanted)
342 .await
343 .map_err(meta_error)?;
344 if rows.len() != wanted.len() {
345 return Err(internal("lease get_many returned the wrong row count"));
346 }
347 let [namespace, repo, epoch, shard, recovery] = &rows[..5] else {
348 return Err(internal("lease get_many returned the wrong row count"));
349 };
350 let mode = recovery
351 .as_ref()
352 .map(codec::decode_lease_recovery)
353 .transpose()
354 .map_err(meta_error)?;
355 if !authority_fence
356 && (mode.is_some_and(|m| m.authority_fence == Some(true))
357 || shard
358 .as_ref()
359 .map(codec::decode_leased_shard)
360 .transpose()
361 .map_err(meta_error)?
362 .is_some_and(|row| row.authority_generation.is_some()))
363 {
364 return Err(ServerError::unavailable(
365 "persisted authority fence requires enabled executor",
366 ));
367 }
368 let authority = if authority_fence {
369 rows[5].clone()
370 } else {
371 None
372 };
373 if authority_fence
374 && mode.is_some_and(|m| m.authority_fence == Some(true))
375 && authority.is_none()
376 {
377 return Err(ServerError::unavailable(
378 "authority generation missing from fenced namespace",
379 ));
380 }
381 let quota_seed = seed_window
382 .map(|window| {
383 rows[5 + usize::from(authority_fence)]
384 .as_ref()
385 .map(codec::decode_namespace_usage)
386 .transpose()
387 .map(|total| (window, total.unwrap_or_default()))
388 .map_err(meta_error)
389 })
390 .transpose()?;
391 if let Some(value) = namespace {
392 codec::decode_namespace_record(value).map_err(meta_error)?;
393 }
394 if let Some(value) = repo {
395 codec::decode_repo_record(value).map_err(meta_error)?;
396 }
397 if namespace.is_none() && repo.is_some() {
398 return Err(internal("repository registered without a namespace"));
399 }
400 if let Some(value) = shard {
401 codec::decode_leased_shard(value).map_err(meta_error)?;
402 }
403 let read = CoordinatorLease {
404 namespace: namespace.clone(),
405 repo: repo.clone(),
406 visibility: rows.last().cloned().flatten(),
407 epoch: epoch.clone(),
408 authority: authority.clone(),
409 authority_generation: if authority_fence {
410 Some(
411 authority
412 .as_ref()
413 .map(codec::decode_u64)
414 .transpose()
415 .map_err(meta_error)?
416 .unwrap_or(0),
417 )
418 } else {
419 None
420 },
421 leased_epoch: epoch
422 .as_ref()
423 .map(codec::decode_u64)
424 .transpose()
425 .map_err(meta_error)?
426 .unwrap_or(0),
427 shard: shard.clone(),
428 observed_el,
429 recovery_value: recovery.clone(),
430 recovery: recovery
431 .as_ref()
432 .map(codec::decode_lease_recovery)
433 .transpose()
434 .map_err(meta_error)?,
435 relay_watermark_ms: reported,
436 quota_seed,
437 };
438 Ok(read)
439}
440
441const RELAY_LEASE_BUDGET_MS: u64 = 15_000;
444
445fn relay_apply_error(
446 metrics: &dyn crate::telemetry::Metrics,
447 partition: &Partition,
448 error: StoreError,
449) -> ServerError {
450 if matches!(error, StoreError::Full) {
451 metrics.incr(
452 super::METRIC_PARTITION_FULL,
453 &[("kind", partition.kind())],
454 1,
455 );
456 tracing::error!(kind = partition.kind(), "storage partition full");
457 ServerError::unavailable("storage partition full")
458 } else {
459 meta_error(error)
460 }
461}
462
463pub async fn renew_for_relay<L: NamespaceStore, M: NamespaceStore>(
478 local: &L,
479 meta: &M,
480 shards: &dyn ShardMap,
481 clock: &dyn Clock,
482 metrics: &dyn crate::telemetry::Metrics,
483 repo: &RepoId,
484 p: &Partition,
485 params: &LeaseParams,
486) -> Result<Value, ServerError> {
487 for _ in 0..LEASE_GRANT_ATTEMPTS {
488 let raw = local
489 .get(p, &keys::epoch_lease())
490 .await
491 .map_err(meta_error)?;
492 let observed = raw
493 .as_ref()
494 .map(codec::decode_epoch_lease)
495 .transpose()
496 .map_err(meta_error)?;
497 if !params.authority_fence && observed.is_some_and(|el| el.authority_generation.is_some()) {
498 return Err(ServerError::unavailable(
499 "persisted authority lease requires enabled executor",
500 ));
501 }
502 let now = ms(clock.now_ms());
503 if let (Some(raw), Some(lease)) = (&raw, observed)
504 && (!params.authority_fence
505 || (lease.authority_generation.is_some() && lease.authority_ready == Some(true)))
506 && lease
507 .expires_at_ms
508 .checked_sub(params.lease_margin_ms)
509 .and_then(|end| end.checked_sub(now))
510 .is_some_and(|budget| budget >= RELAY_LEASE_BUDGET_MS)
511 {
512 return Ok(raw.clone());
513 }
514 let read = read_lease_rows(
515 local,
516 meta,
517 shards,
518 clock,
519 repo,
520 p,
521 observed,
522 None,
523 params.authority_fence,
524 )
525 .await?;
526 let grant = grant_batch(&read, &repo.name, p, now, now, params)?;
527 let coordinator = shards.coordinator(&repo.namespace);
528 match meta
529 .apply(&coordinator, grant.batch)
530 .await
531 .map_err(|error| relay_apply_error(metrics, &coordinator, error))?
532 {
533 BatchOutcome::Committed => {}
534 BatchOutcome::PreconditionFailed { .. } => continue,
535 BatchOutcome::DeadlinePassed { .. } => {
536 return Err(internal("lease grant had no deadline"));
537 }
538 }
539 let value = codec::encode_epoch_lease(&grant.value);
540 let install = Batch::new()
541 .require(Precondition::NotAfter(now.saturating_add(10_000)))
542 .require(observed_guard(keys::epoch_lease(), raw.as_ref()))
543 .put(keys::epoch_lease(), value.clone());
544 if matches!(
545 local
546 .apply(p, install)
547 .await
548 .map_err(|error| relay_apply_error(metrics, p, error))?,
549 BatchOutcome::Committed
550 ) {
551 return Ok(value);
552 }
553 }
554 Err(ServerError::aborted_retryable(
555 "coordinator lease grant contention",
556 ))
557}
558
559impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H> {
560 pub(super) async fn observe_lease(
561 &self,
562 op: &Operation,
563 p: &Partition,
564 ahead: Option<&Snapshot>,
565 ) -> Result<LeaseObservation, ServerError> {
566 let snap = ahead.ok_or_else(|| internal("D34 lease requires an atomic snapshot"))?;
567 let seed_window = snap.namespace_window.filter(|window| {
568 let key = keys::quota_view(*window);
569 snap.contains(&key) && snap.get(&key).is_none()
570 });
571 let observed_el = snap
572 .get(&keys::epoch_lease())
573 .map(codec::decode_epoch_lease)
574 .transpose()
575 .map_err(meta_error)?;
576 if self.cfg.authority_fence.is_none()
577 && observed_el.is_some_and(|el| el.authority_generation.is_some())
578 {
579 return Err(ServerError::unavailable(
580 "persisted authority lease requires enabled executor",
581 ));
582 }
583 if let Some(lease) = observed_el.filter(|lease| {
584 self.cfg.authority_fence.is_none()
585 || (lease.authority_generation.is_some() && lease.authority_ready == Some(true))
586 }) {
587 let now = ms(self.clock.now_ms());
588 let usable_until = lease.expires_at_ms.checked_sub(self.cfg.lease_margin_ms);
589 if usable_until
590 .and_then(|end| end.checked_sub(now))
591 .is_some_and(|budget| budget >= self.cfg.min_lease_budget_ms)
592 {
593 return Ok(LeaseObservation::Usable(lease));
594 }
595 }
596 Ok(LeaseObservation::Renew(Box::new(
597 self.read_lease(op, p, observed_el, seed_window).await?,
598 )))
599 }
600
601 async fn read_lease(
602 &self,
603 op: &Operation,
604 p: &Partition,
605 observed_el: Option<codec::EpochLease>,
606 seed_window: Option<u64>,
607 ) -> Result<CoordinatorLease, ServerError> {
608 let read = read_lease_rows(
609 &self.meta,
610 &self.meta,
611 self.shards.as_ref(),
612 self.clock.as_ref(),
613 &op.repo,
614 p,
615 observed_el,
616 seed_window,
617 self.cfg.authority_fence.is_some(),
618 )
619 .await?;
620 if self.cfg.authority_fence.is_some()
621 && read.recovery.is_none_or(|mode| {
622 mode.authority_fence != Some(true) || mode.authority_ready != Some(true)
623 })
624 {
625 Box::pin(self.ensure_authority_activation(&op.repo.namespace)).await?;
626 return read_lease_rows(
627 &self.meta,
628 &self.meta,
629 self.shards.as_ref(),
630 self.clock.as_ref(),
631 &op.repo,
632 p,
633 observed_el,
634 seed_window,
635 true,
636 )
637 .await;
638 }
639 Ok(read)
640 }
641
642 pub(super) async fn admit_lease(
643 &self,
644 op: &Operation,
645 p: &Partition,
646 observed: LeaseObservation,
647 skew_ms: i64,
648 ) -> Result<(Creation, LeaseWrite), ServerError> {
649 let LeaseObservation::Renew(read) = observed else {
650 let LeaseObservation::Usable(value) = observed else {
651 unreachable!()
652 };
653 return Ok((
654 Creation::default(),
655 LeaseWrite {
656 value,
657 install: false,
658 },
659 ));
660 };
661 let mut read = *read;
662 let coordinator = self.shards.coordinator(&op.repo.namespace);
663 for _ in 0..LEASE_GRANT_ATTEMPTS {
664 if op
665 .authz
666 .grant
667 .as_ref()
668 .is_some_and(|grant| grant.epoch != read.leased_epoch)
669 {
670 return Err(super::plan::epoch_moved());
671 }
672 if let Some(generation) = op.authz.authority_generation
673 && Some(generation) != read.authority_generation
674 {
675 return Err(crate::authority::moved());
676 }
677 let now = ms(self.clock.now_ms());
678 let created_at_ms = ms(self.clock.now_ms().saturating_add(skew_ms));
679 let grant = grant_batch(
680 &read,
681 &op.repo.name,
682 p,
683 now,
684 created_at_ms,
685 &LeaseParams::from(&self.cfg),
686 )?;
687 let outcome = match self.meta.apply(&coordinator, grant.batch).await {
688 Ok(outcome) => outcome,
689 Err(StoreError::Full) => return Err(self.partition_full(&coordinator, None).await),
690 Err(error) => return Err(meta_error(error)),
691 };
692 match outcome {
693 BatchOutcome::Committed => {
694 let created = if matches!(self.cfg.addressing, Addressing::Multi(_)) {
695 grant.creation
696 } else {
697 Creation::default()
698 };
699 return Ok((
700 created,
701 LeaseWrite {
702 value: grant.value,
703 install: true,
704 },
705 ));
706 }
707 BatchOutcome::PreconditionFailed { .. } => {
708 read = self
709 .read_lease(op, p, read.observed_el, read.quota_seed.map(|(w, _)| w))
710 .await?;
711 }
712 BatchOutcome::DeadlinePassed { .. } => {
713 return Err(internal("lease grant had no deadline"));
714 }
715 }
716 }
717 tracing::warn!(shard = ?p, attempts = LEASE_GRANT_ATTEMPTS, "coordinator lease grant did not settle");
718 Err(ServerError::aborted_retryable(
719 "coordinator lease grant contention",
720 ))
721 }
722}
723
724#[cfg(test)]
725#[path = "lease_model_tests.rs"]
726mod model_tests;