1use crate::authority::FenceKind;
5use crate::error::ServerError;
6use crate::repo::NamespaceKey;
7use crate::store::{
8 Batch, BatchOutcome, MultipartBlobStore, NamespaceStore, Partition, Precondition, Value, codec,
9 keys, restore::mark_lease_table_recovered,
10};
11use mkit_attest::grant::{EpochTransition, epoch_transition};
12
13use super::{HookSet, Pipeline, internal, lease::observed_guard, meta_error, ms};
14
15pub const MAX_EPOCH_STEP: u64 = mkit_attest::grant::MAX_EPOCH_STEP;
17const PAGE_SIZE: u32 = 4;
18const MAX_SCAN_PAGES: u32 = 8;
22const MAX_PUSH_ATTEMPTS: u32 = 1;
24
25#[derive(Debug, Clone, Copy)]
27#[non_exhaustive]
28pub struct RevokeBudget {
29 pub max_elapsed_ms: u64,
31}
32
33impl RevokeBudget {
34 #[must_use]
36 pub const fn new(max_elapsed_ms: u64) -> Self {
37 Self {
38 max_elapsed_ms: if max_elapsed_ms == 0 {
39 1
40 } else {
41 max_elapsed_ms
42 },
43 }
44 }
45}
46
47impl Default for RevokeBudget {
48 fn default() -> Self {
49 Self::new(1000)
50 }
51}
52
53#[derive(Debug, Clone, Copy, PartialEq, Eq)]
56#[non_exhaustive]
57pub enum RevokeProgress {
58 Complete,
60 Pending {
63 remaining: u64,
65 },
66}
67
68#[derive(serde::Serialize, serde::Deserialize)]
73#[serde(deny_unknown_fields)]
74struct RevokeCheckpoint {
75 generation: u64,
76 recovery: Option<codec::LeaseRecovery>,
77 cursor: Vec<u8>,
78}
79
80struct CoordinatorState {
81 kind: FenceKind,
82 epoch_value: Option<Value>,
83 epoch: u64,
84 config_version: u64,
85 recovery: Option<codec::LeaseRecovery>,
86}
87
88fn generation(row: &codec::LeasedShard, kind: FenceKind) -> u64 {
89 match kind {
90 FenceKind::Grant => row.epoch,
91 FenceKind::Authority => row.authority_generation.unwrap_or(0),
92 }
93}
94fn acked(row: &codec::LeasedShard, kind: FenceKind) -> Option<u64> {
95 match kind {
96 FenceKind::Grant => Some(row.acked_epoch),
97 FenceKind::Authority => row.acked_authority_generation,
98 }
99}
100
101impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H> {
102 pub async fn bump_epoch(&self, ns: &NamespaceKey, new_epoch: u64) -> Result<(), ServerError> {
108 match self.transition_epoch(ns, new_epoch).await? {
109 EpochTransition::Advance => Ok(()),
110 EpochTransition::Retry | EpochTransition::Reject => Err(ServerError::invalid_argument(
111 "epoch increment must be between 1 and 1024",
112 )),
113 }
114 }
115
116 pub(super) async fn transition_epoch(
119 &self,
120 ns: &NamespaceKey,
121 new_epoch: u64,
122 ) -> Result<EpochTransition, ServerError> {
123 self.transition_fence(ns, new_epoch, FenceKind::Grant).await
124 }
125
126 pub(super) async fn transition_fence(
127 &self,
128 ns: &NamespaceKey,
129 new_epoch: u64,
130 kind: FenceKind,
131 ) -> Result<EpochTransition, ServerError> {
132 let p = self.shards.coordinator(ns);
133 for _ in 0..super::coordinator::CREATION_ATTEMPTS {
134 let current = self.meta.get(&p, &kind.key()).await.map_err(meta_error)?;
135 let epoch = current
136 .as_ref()
137 .map(codec::decode_u64)
138 .transpose()
139 .map_err(meta_error)?
140 .unwrap_or(0);
141 let transition = epoch_transition(epoch, new_epoch);
142 if transition != EpochTransition::Advance {
143 return Ok(transition);
144 }
145 let batch = Batch::new()
146 .require(observed_guard(kind.key(), current.as_ref()))
147 .put(kind.key(), codec::encode_u64(new_epoch));
148 match self.apply_meta(&p, batch).await? {
149 BatchOutcome::Committed => return Ok(EpochTransition::Advance),
150 BatchOutcome::PreconditionFailed { .. } => {}
151 BatchOutcome::DeadlinePassed { .. } => {
152 return Err(internal("epoch bump had no deadline"));
153 }
154 }
155 }
156 Err(ServerError::unavailable("epoch contention; retry").with_header("Retry-After", "1"))
157 }
158
159 pub async fn mark_lease_table_recovered(&self, ns: &NamespaceKey) -> Result<(), ServerError> {
167 mark_lease_table_recovered(
168 &self.meta,
169 &self.shards.coordinator(ns),
170 ms(self.clock.now_ms()),
171 )
172 .await
173 .map_err(meta_error)
174 }
175
176 async fn coordinator_state(
177 &self,
178 p: &Partition,
179 kind: FenceKind,
180 ) -> Result<CoordinatorState, ServerError> {
181 let rows = self
182 .meta
183 .get_many(
184 p,
185 &[kind.key(), keys::namespace_record(), keys::lease_recovery()],
186 )
187 .await
188 .map_err(meta_error)?;
189 let [epoch, nr, lr] = rows.as_slice() else {
190 return Err(internal("revocation get_many returned the wrong row count"));
191 };
192 Ok(CoordinatorState {
193 kind,
194 epoch: epoch
195 .as_ref()
196 .map(codec::decode_u64)
197 .transpose()
198 .map_err(meta_error)?
199 .unwrap_or(0),
200 epoch_value: epoch.clone(),
201 config_version: nr
202 .as_ref()
203 .map(codec::decode_namespace_record)
204 .transpose()
205 .map_err(meta_error)?
206 .map_or(1, |n| n.config_version),
207 recovery: lr
208 .as_ref()
209 .map(codec::decode_lease_recovery)
210 .transpose()
211 .map_err(meta_error)?,
212 })
213 }
214
215 fn recovery_pending(&self, state: &CoordinatorState) -> bool {
216 state
217 .recovery
218 .and_then(codec::LeaseRecovery::recovery_time)
219 .is_some_and(|resumed| {
220 ms(self.clock.now_ms())
221 < resumed
222 .saturating_add(self.cfg.epoch_lease_ms)
223 .saturating_add(self.cfg.lease_margin_ms)
224 })
225 }
226
227 fn revoke_budget_passed(&self, start: u64, budget: RevokeBudget) -> bool {
228 ms(self.clock.now_ms()).saturating_sub(start) >= budget.max_elapsed_ms.max(1)
229 }
230
231 async fn read_revoke_cursor(
232 &self,
233 p: &Partition,
234 state: &CoordinatorState,
235 ) -> Result<(Option<Value>, Option<crate::store::Cursor>), ServerError> {
236 let raw = self
237 .meta
238 .get(p, &keys::revoke_cursor(state.kind == FenceKind::Authority))
239 .await
240 .map_err(meta_error)?;
241 let cursor = if let Some(raw) = raw.as_ref() {
242 let bytes = raw.as_bytes();
243 if bytes.len() > 32_768 || bytes.first() != Some(&1) {
244 return Err(internal("invalid revocation checkpoint"));
245 }
246 let checkpoint: RevokeCheckpoint = serde_json::from_slice(&bytes[1..])
247 .map_err(|_| internal("invalid revocation checkpoint"))?;
248 if checkpoint.cursor.len() > crate::store::MAX_KEY_BYTES {
249 return Err(internal("oversized revocation cursor"));
250 }
251 (checkpoint.generation == state.epoch && checkpoint.recovery == state.recovery)
252 .then(|| crate::store::Cursor::new(checkpoint.cursor))
253 } else {
254 None
255 };
256 Ok((raw, cursor))
257 }
258
259 async fn save_revoke_cursor(
260 &self,
261 p: &Partition,
262 state: &CoordinatorState,
263 prior: Option<&Value>,
264 cursor: Option<crate::store::Cursor>,
265 ) -> Result<bool, ServerError> {
266 let key = keys::revoke_cursor(state.kind == FenceKind::Authority);
267 let recovery = state.recovery.as_ref().map(codec::encode_lease_recovery);
268 let batch = Batch::new()
269 .require(observed_guard(key.clone(), prior))
270 .require(observed_guard(state.kind.key(), state.epoch_value.as_ref()))
271 .require(observed_guard(keys::lease_recovery(), recovery.as_ref()));
272 let batch = if let Some(cursor) = cursor {
273 let checkpoint = RevokeCheckpoint {
274 generation: state.epoch,
275 recovery: state.recovery,
276 cursor: cursor.as_bytes().to_vec(),
277 };
278 let mut bytes = vec![1];
279 bytes.extend(
280 serde_json::to_vec(&checkpoint)
281 .map_err(|_| internal("cannot encode revocation checkpoint"))?,
282 );
283 batch.put(key, Value::new(bytes))
284 } else {
285 batch.delete(key)
286 };
287 match self.apply_meta(p, batch).await? {
288 BatchOutcome::Committed => Ok(true),
289 BatchOutcome::PreconditionFailed { .. } => Ok(false),
290 BatchOutcome::DeadlinePassed { .. } => Err(internal("checkpoint had no deadline")),
291 }
292 }
293
294 pub async fn revoke_step(
303 &self,
304 ns: &NamespaceKey,
305 budget: &RevokeBudget,
306 ) -> Result<RevokeProgress, ServerError> {
307 self.revoke_fence_step(ns, budget, FenceKind::Grant).await
308 }
309
310 pub(super) async fn revoke_fence_step(
311 &self,
312 ns: &NamespaceKey,
313 budget: &RevokeBudget,
314 kind: FenceKind,
315 ) -> Result<RevokeProgress, ServerError> {
316 let start = ms(self.clock.now_ms());
317 let coordinator = self.shards.coordinator(ns);
318 let state = self.coordinator_state(&coordinator, kind).await?;
319 let (first, end) = keys::class_range(keys::TAG_LEASED_SHARD);
320 let (cursor_value, mut cursor) = self.read_revoke_cursor(&coordinator, &state).await?;
321 let mut checkpoint = cursor.clone();
322 let (mut visited, mut remaining, mut prefix_complete) = (0, 0_u64, true);
323 let mut pages = 0;
324 loop {
325 if pages == MAX_SCAN_PAGES || self.revoke_budget_passed(start, *budget) {
326 self.save_revoke_cursor(&coordinator, &state, cursor_value.as_ref(), checkpoint)
327 .await?;
328 return Ok(RevokeProgress::Pending {
329 remaining: remaining.max(1),
330 });
331 }
332 let page = self
333 .meta
334 .scan(&coordinator, &first, &end, cursor.as_ref(), PAGE_SIZE)
335 .await
336 .map_err(meta_error)?;
337 pages += 1;
338 if page.entries.len() > PAGE_SIZE as usize {
339 return Err(internal("revocation scan exceeded its row bound"));
340 }
341 for (key, value) in page.entries {
342 if self.revoke_budget_passed(start, *budget) {
343 self.save_revoke_cursor(
344 &coordinator,
345 &state,
346 cursor_value.as_ref(),
347 checkpoint,
348 )
349 .await?;
350 return Ok(RevokeProgress::Pending {
351 remaining: remaining.max(1),
352 });
353 }
354 let row = codec::decode_leased_shard(&value).map_err(meta_error)?;
355 if row.expires_at_ms <= ms(self.clock.now_ms())
356 || acked(&row, state.kind) == Some(state.epoch)
357 {
358 continue;
359 }
360 if acked(&row, state.kind).is_some_and(|n| n > state.epoch) || visited == 4 {
361 remaining += 1;
362 prefix_complete = false;
363 continue;
364 }
365 visited += 1;
366 if !self
367 .push_and_ack(&coordinator, (&key, value), &state, start, budget)
368 .await?
369 {
370 remaining += 1;
371 prefix_complete = false;
372 }
373 }
374 if prefix_complete {
375 checkpoint = page.next.clone();
376 }
377 cursor = page.next;
378 if cursor.is_none() {
379 break;
380 }
381 }
382 let latest = self.coordinator_state(&coordinator, kind).await?;
384 if remaining == 0 && latest.epoch == state.epoch && !self.recovery_pending(&latest) {
385 if self
386 .save_revoke_cursor(&coordinator, &state, cursor_value.as_ref(), None)
387 .await?
388 {
389 Ok(RevokeProgress::Complete)
390 } else {
391 Ok(RevokeProgress::Pending { remaining: 1 })
392 }
393 } else {
394 self.save_revoke_cursor(
395 &coordinator,
396 &state,
397 cursor_value.as_ref(),
398 if latest.epoch == state.epoch && latest.recovery == state.recovery {
399 checkpoint
400 } else {
401 None
402 },
403 )
404 .await?;
405 Ok(RevokeProgress::Pending {
406 remaining: remaining.max(1),
407 })
408 }
409 }
410
411 #[allow(clippy::too_many_lines)] async fn push_and_ack(
413 &self,
414 coordinator: &Partition,
415 (key, mut value): (&crate::store::Key, Value),
416 state: &CoordinatorState,
417 start: u64,
418 budget: &RevokeBudget,
419 ) -> Result<bool, ServerError> {
420 let Some(keys::ParsedKey::LeasedShard { repo, shard_ref }) = keys::parse(key) else {
421 return Err(internal("invalid leased-shard key"));
422 };
423 let Partition::Coordinator(ns) = coordinator else {
424 return Err(internal("revoke outside coordinator"));
425 };
426 let p = Partition::Ref {
427 ns: ns.clone(),
428 repo,
429 shard_ref,
430 };
431 for _ in 0..MAX_PUSH_ATTEMPTS {
432 if self.revoke_budget_passed(start, *budget) {
433 return Ok(false);
434 }
435 let mut row = codec::decode_leased_shard(&value).map_err(meta_error)?;
436 if row.expires_at_ms <= ms(self.clock.now_ms())
437 || acked(&row, state.kind) == Some(state.epoch)
438 {
439 return Ok(true);
440 }
441 if generation(&row, state.kind) > state.epoch
442 || acked(&row, state.kind).is_some_and(|n| n > state.epoch)
443 {
444 return Ok(false);
445 }
446 let old = self
447 .meta
448 .get(&p, &keys::epoch_lease())
449 .await
450 .map_err(meta_error)?;
451 if let Some(old) = old.as_ref() {
452 let old = codec::decode_epoch_lease(old).map_err(meta_error)?;
453 if match state.kind {
456 FenceKind::Grant => old.epoch,
457 FenceKind::Authority => old.authority_generation.unwrap_or(0),
458 } > state.epoch
459 {
460 return Ok(false);
461 }
462 if old.expires_at_ms > row.expires_at_ms {
463 let Some(latest) = self.meta.get(coordinator, key).await.map_err(meta_error)?
464 else {
465 return Ok(true);
466 };
467 value = latest;
468 continue;
469 }
470 }
471 if self.revoke_budget_passed(start, *budget) {
472 return Ok(false);
473 }
474 let prior = old
475 .as_ref()
476 .map(codec::decode_epoch_lease)
477 .transpose()
478 .map_err(meta_error)?;
479 let lease = codec::EpochLease {
480 authority_ready: prior.and_then(|el| el.authority_ready),
481 epoch: if state.kind == FenceKind::Grant {
482 state.epoch
483 } else {
484 prior.map_or(row.epoch, |el| el.epoch)
485 },
486 authority_generation: if state.kind == FenceKind::Authority {
487 Some(state.epoch)
488 } else {
489 prior
490 .and_then(|el| el.authority_generation)
491 .or(row.authority_generation)
492 },
493 expires_at_ms: row.expires_at_ms,
494 config_version: state.config_version,
495 };
496 let push = Batch::new()
497 .require(observed_guard(keys::epoch_lease(), old.as_ref()))
498 .put(keys::epoch_lease(), codec::encode_epoch_lease(&lease));
499 match self.apply_meta(&p, push).await? {
500 BatchOutcome::PreconditionFailed { .. } => continue,
501 BatchOutcome::DeadlinePassed { .. } => {
502 return Err(internal("epoch push had no deadline"));
503 }
504 BatchOutcome::Committed => {}
505 }
506 if self.revoke_budget_passed(start, *budget) {
507 return Ok(false);
508 }
509 match state.kind {
510 FenceKind::Grant => row.acked_epoch = state.epoch,
511 FenceKind::Authority => row.acked_authority_generation = Some(state.epoch),
512 }
513 let ack = Batch::new()
514 .require(Precondition::Equals(key.clone(), value))
515 .require(observed_guard(state.kind.key(), state.epoch_value.as_ref()))
516 .put(key.clone(), codec::encode_leased_shard(&row));
517 match self.apply_meta(coordinator, ack).await? {
518 BatchOutcome::Committed => return Ok(true),
519 BatchOutcome::DeadlinePassed { .. } => {
520 return Err(internal("epoch acknowledgement had no deadline"));
521 }
522 BatchOutcome::PreconditionFailed { .. } => {
523 let Some(latest) = self.meta.get(coordinator, key).await.map_err(meta_error)?
524 else {
525 return Ok(true);
526 };
527 value = latest;
528 }
529 }
530 }
531 Ok(false)
532 }
533
534 #[cfg(feature = "test-faults")]
540 pub async fn test_bump_epoch(
541 &self,
542 ns: &NamespaceKey,
543 new_epoch: u64,
544 ) -> Result<(), ServerError> {
545 self.bump_epoch(ns, new_epoch).await?;
546 let start = ms(self.clock.now_ms());
547 for _ in 0..10_000 {
548 if self.revoke_step(ns, &RevokeBudget::default()).await? == RevokeProgress::Complete {
549 return Ok(());
550 }
551 if ms(self.clock.now_ms()).saturating_sub(start) >= 10_000 {
552 return Err(ServerError::unavailable("epoch revocation pending; retry"));
553 }
554 let mut yielded = false;
556 core::future::poll_fn(|cx| {
557 if yielded {
558 core::task::Poll::Ready(())
559 } else {
560 yielded = true;
561 cx.waker().wake_by_ref();
562 core::task::Poll::Pending
563 }
564 })
565 .await;
566 }
567 Err(ServerError::unavailable("epoch revocation pending; retry"))
568 }
569}
570
571#[cfg(all(test, feature = "memory", feature = "test-faults"))]
572mod tests {
573 use super::*;
574 use crate::pipeline::{AuthMode, Hooks, PipelineConfig, Sharding};
575 use crate::upload::UploadLimits;
576 use crate::{
577 Addressing, Clock, Code, ManualClock, MemoryBlobStore, MemoryKv, NoopMetrics, RepoId,
578 RepoName,
579 };
580 use std::sync::Arc;
581
582 #[tokio::test]
583 async fn test_bump_epoch_terminates_with_a_frozen_clock_during_recovery() {
584 let clock = Arc::new(ManualClock::new(100_000));
585 let namespace = NamespaceKey::deployment_default();
586 let repo = RepoId {
587 namespace: namespace.clone(),
588 name: RepoName::new("room").expect("valid test repository"),
589 };
590 let mut cfg = PipelineConfig::new(
591 Addressing::Single { repo },
592 AuthMode::Open,
593 UploadLimits {
594 max_total_bytes: 64,
595 max_chunks: 16,
596 },
597 );
598 cfg.sharding = Sharding::D34;
599 let pipe = Pipeline::new(
600 MemoryBlobStore::default(),
601 MemoryKv::with_clock(clock.clone()),
602 Hooks::new(),
603 cfg,
604 clock.clone(),
605 Arc::new(NoopMetrics),
606 )
607 .expect("valid pipeline configuration");
608 pipe.mark_lease_table_recovered(&namespace)
609 .await
610 .expect("recovery marker commits");
611 let error = tokio::time::timeout(
612 std::time::Duration::from_secs(30),
613 pipe.test_bump_epoch(&namespace, 1),
614 )
615 .await
616 .expect("step cap must terminate despite frozen pipeline clock")
617 .expect_err("recovery holdoff cannot complete at frozen time");
618 assert_eq!(error.code(), Code::Unavailable);
619 assert_eq!(clock.now_ms(), 100_000);
620 assert_eq!(
621 pipe.revoke_step(&namespace, &RevokeBudget::default())
622 .await
623 .expect("revoke step succeeds"),
624 RevokeProgress::Pending { remaining: 1 },
625 );
626 }
627}