Skip to main content

mkit_server/pipeline/
revocation.rs

1//! Epoch changes serialize with grants in the coordinator. Pushes and
2//! acknowledgements are separate guarded transactions, always in that order.
3
4use 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
15/// Largest allowed epoch increment (SPEC-WRITE-GRANTS ยง1.1).
16pub const MAX_EPOCH_STEP: u64 = mkit_attest::grant::MAX_EPOCH_STEP;
17const PAGE_SIZE: u32 = 4;
18// A frozen clock and acknowledged/expired rows must not permit an unbounded
19// walk. Each slice reads at most 32 rows in eight pages, besides four bounded
20// push loops (one attempt, at most five store calls each).
21const MAX_SCAN_PAGES: u32 = 8;
22// Bound contention even with a frozen injected clock.
23const MAX_PUSH_ATTEMPTS: u32 = 1;
24
25/// One revocation slice processes at most four shards, within this time budget.
26#[derive(Debug, Clone, Copy)]
27#[non_exhaustive]
28pub struct RevokeBudget {
29    /// Milliseconds of elapsed pipeline time allowed between store calls.
30    pub max_elapsed_ms: u64,
31}
32
33impl RevokeBudget {
34    /// A positive elapsed-time budget; each slice visits at most four shards.
35    #[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/// Completion is reported only after every lease is acknowledged or expired,
54/// and any declared recovery holdoff has passed.
55#[derive(Debug, Clone, Copy, PartialEq, Eq)]
56#[non_exhaustive]
57pub enum RevokeProgress {
58    /// No old-epoch write can still commit.
59    Complete,
60    /// Outstanding shards counted so far. At a budget stop this is a lower
61    /// bound (at least one); recovery holdoff alone also returns one.
62    Pending {
63        /// Conservative lower bound on remaining work.
64        remaining: u64,
65    },
66}
67
68/// An epoch-bound scan checkpoint is only an optimization. Completed prefix
69/// rows cannot become outstanding again at the same epoch: live renewal
70/// preserves their ack, expired renewal acknowledges that epoch immediately,
71/// and any newly inserted row grants the coordinator's current epoch.
72#[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    /// Advance the coordinator epoch, serialized with every epoch lease grant.
103    ///
104    /// # Errors
105    /// `invalid_argument` unless the increment is in `1..=1024`;
106    /// `unavailable` on repeated contention, or a mapped storage failure.
107    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    /// CAS the epoch. A lost CAS re-reads and re-classifies, including a
117    /// concurrent winner that installed the same epoch (a valid retry).
118    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    /// Declare lease-table recovery using the real pipeline clock.
160    /// Any procedure that restores or rebuilds a coordinator partition MUST
161    /// call this before serving writes for that namespace. The durable marker
162    /// prevents completion for `epoch_lease + margin` after recovery.
163    ///
164    /// # Errors
165    /// A mapped storage failure, or `internal` for an impossible batch outcome.
166    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    /// Push and acknowledge at most four live shards in one bounded slice.
295    /// Outside a declared recovery hold-off, `ls.acked_epoch = n` only if the
296    /// shard's el durably holds epoch >= n, or all older-epoch writes are already
297    /// past their commit deadline. Recovery fences surviving copies until then.
298    /// Renewal alone never raises a live row's acknowledgement.
299    ///
300    /// # Errors
301    /// A mapped storage or codec failure. Contention leaves progress pending.
302    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        // Detect a newer epoch or recovery declaration during the slice.
383        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)] // The guarded push-before-ack sequence preserves both independently updated generations.
412    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                // A stale slice must never undo a newer push or shorten a
454                // concurrent renewal. Retry that coordinator observation.
455                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    /// Test-only bump driver used by the `ListRefs` directive. No HTTP route.
535    ///
536    /// # Errors
537    /// The bump/step error, or `unavailable` after ten seconds or 10,000 steps.
538    /// The step cap also terminates when an injected clock stays frozen.
539    #[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            // Yield even when the memory store completes synchronously.
555            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}