Skip to main content

mkit_server/pipeline/
authority.rs

1//! Optional authority fence RPCs; independent of the grant epoch.
2
3use super::{
4    AuthMode, HookSet, Pipeline, meta_error, ms,
5    revocation::{RevokeBudget, RevokeProgress},
6};
7use crate::{
8    authority::FenceKind,
9    error::ServerError,
10    policy::NamespacePolicy,
11    repo::{Addressing, NamespaceKey},
12    store::{MultipartBlobStore, NamespaceStore, codec, keys},
13};
14use mkit_attest::grant::EpochTransition;
15use mkit_core::repo_identity::Namespace;
16
17impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H> {
18    pub(super) async fn stored_authority_generation(
19        &self,
20        ns: &NamespaceKey,
21    ) -> Result<u64, ServerError> {
22        self.meta
23            .get(&self.shards.coordinator(ns), &keys::authority_generation())
24            .await
25            .map_err(meta_error)?
26            .as_ref()
27            .map(codec::decode_u64)
28            .transpose()
29            .map_err(meta_error)
30            .map(|generation| generation.unwrap_or(0))
31    }
32
33    /// Establish durable mode without creating an accounting namespace, then
34    /// complete the shared initial barrier before issuing any fenced grant.
35    pub(super) async fn ensure_authority_activation(
36        &self,
37        ns: &NamespaceKey,
38    ) -> Result<Option<u64>, ServerError> {
39        let p = self.shards.coordinator(ns);
40        let wanted = [keys::authority_generation(), keys::lease_recovery()];
41        for _ in 0..1 {
42            let rows = self.meta.get_many(&p, &wanted).await.map_err(meta_error)?;
43            let [generation, raw_mode] = rows.as_slice() else {
44                return Err(super::internal("authority activation row count"));
45            };
46            let current = generation
47                .as_ref()
48                .map(codec::decode_u64)
49                .transpose()
50                .map_err(meta_error)?
51                .unwrap_or(0);
52            let mut mode = raw_mode
53                .as_ref()
54                .map(codec::decode_lease_recovery)
55                .transpose()
56                .map_err(meta_error)?
57                .unwrap_or(codec::LeaseRecovery {
58                    resumed_at_ms: 0,
59                    authority_fence: None,
60                    authority_ready: None,
61                    activation_only: Some(true),
62                });
63            if self.cfg.authority_fence.is_none() {
64                if generation.is_some() || mode.authority_fence == Some(true) {
65                    return Err(ServerError::unavailable(
66                        "persisted authority fence requires enabled executor",
67                    ));
68                }
69                return Ok(None);
70            }
71            if mode.authority_fence == Some(true) && generation.is_none() {
72                return Err(ServerError::unavailable(
73                    "authority generation missing from fenced namespace",
74                ));
75            }
76            if mode.authority_fence != Some(true) {
77                mode.authority_fence = Some(true);
78                mode.authority_ready = Some(false);
79                let batch = crate::store::Batch::new()
80                    .require(super::lease::observed_guard(
81                        wanted[0].clone(),
82                        generation.as_ref(),
83                    ))
84                    .require(super::lease::observed_guard(
85                        wanted[1].clone(),
86                        raw_mode.as_ref(),
87                    ))
88                    .put(wanted[0].clone(), codec::encode_u64(current))
89                    .put(wanted[1].clone(), codec::encode_lease_recovery(&mode));
90                match self.apply_meta(&p, batch).await? {
91                    crate::store::BatchOutcome::Committed => {}
92                    crate::store::BatchOutcome::PreconditionFailed { .. } => continue,
93                    crate::store::BatchOutcome::DeadlinePassed { .. } => {
94                        return Err(super::internal("activation had no deadline"));
95                    }
96                }
97                // Use the installed raw value below; the barrier itself rereads
98                // generation and recovery, and the ready CAS guards both.
99            }
100            if mode.authority_ready == Some(true) {
101                return Ok(None);
102            }
103            if self
104                .revoke_fence_step(ns, &RevokeBudget::default(), FenceKind::Authority)
105                .await?
106                != RevokeProgress::Complete
107            {
108                return Err(
109                    ServerError::unavailable("authority activation pending; retry")
110                        .with_header("Retry-After", "1"),
111                );
112            }
113            let prior = codec::encode_lease_recovery(&mode);
114            mode.authority_ready = Some(true);
115            let batch = crate::store::Batch::new()
116                .require(crate::store::Precondition::Equals(
117                    wanted[0].clone(),
118                    codec::encode_u64(current),
119                ))
120                .require(crate::store::Precondition::Equals(wanted[1].clone(), prior))
121                .put(wanted[1].clone(), codec::encode_lease_recovery(&mode));
122            match self.apply_meta(&p, batch).await? {
123                crate::store::BatchOutcome::Committed => return Ok(Some(current)),
124                crate::store::BatchOutcome::PreconditionFailed { .. } => {}
125                crate::store::BatchOutcome::DeadlinePassed { .. } => {
126                    return Err(super::internal("activation ready had no deadline"));
127                }
128            }
129        }
130        Err(
131            ServerError::unavailable("authority activation contention; retry")
132                .with_header("Retry-After", "1"),
133        )
134    }
135
136    // Keep streaming-session futures small while the authoritative read is pending.
137    pub(super) fn check_ticket_generation<'a>(
138        &'a self,
139        ns: &'a NamespaceKey,
140        generation: Option<u64>,
141    ) -> crate::BoxFuture<'a, Result<(), ServerError>> {
142        Box::pin(async move {
143            if self.cfg.authority_fence.is_none() {
144                if generation.is_some() {
145                    return Err(ServerError::unavailable(
146                        "fenced ticket requires enabled executor",
147                    ));
148                }
149                Box::pin(self.ensure_authority_activation(ns)).await?;
150            } else {
151                // A generation-bearing authenticated ticket is issued only after
152                // activation. Do not repeat the activation scan on every chunk.
153                let current = self
154                    .meta
155                    .get(&self.shards.coordinator(ns), &keys::authority_generation())
156                    .await
157                    .map_err(meta_error)?;
158                let current = current
159                    .as_ref()
160                    .map(codec::decode_u64)
161                    .transpose()
162                    .map_err(meta_error)?
163                    .ok_or_else(|| {
164                        ServerError::unavailable(
165                            "authority generation missing from fenced namespace",
166                        )
167                    })?;
168                if generation != Some(current) {
169                    return Err(crate::authority::moved());
170                }
171            }
172            Ok(())
173        })
174    }
175
176    /// Read the optional namespace authority generation outside auth-v2.
177    /// # Errors
178    /// Disabled fencing, invalid namespace or storage failure.
179    pub async fn get_authority_generation(&self, namespace: &str) -> Result<u64, ServerError> {
180        if self.cfg.authority_fence.is_none() {
181            return Err(ServerError::unimplemented("authority fencing is disabled"));
182        }
183        let ns = Namespace::parse(namespace)
184            .map_err(|_| ServerError::invalid_argument("invalid namespace"))?;
185        self.authority_namespace(&ns).await?;
186        self.stored_authority_generation(&NamespaceKey::from_namespace(&ns))
187            .await
188    }
189
190    async fn authority_namespace(&self, ns: &Namespace) -> Result<(), ServerError> {
191        let Addressing::Multi(multi) = &self.cfg.addressing else {
192            return Err(ServerError::unimplemented(
193                "authority fencing requires multi addressing",
194            ));
195        };
196        match &multi.namespace_policy {
197            NamespacePolicy::Allowlist(allowed) if allowed.contains(ns) => Ok(()),
198            NamespacePolicy::Any { .. } => {
199                let key = NamespaceKey::from_namespace(ns);
200                if self
201                    .meta
202                    .get(&self.shards.coordinator(&key), &keys::namespace_record())
203                    .await
204                    .map_err(meta_error)?
205                    .is_some()
206                {
207                    Ok(())
208                } else {
209                    Err(ServerError::permission_denied("namespace not served"))
210                }
211            }
212            _ => Err(ServerError::permission_denied("namespace not served")),
213        }
214    }
215
216    /// Verify a deployment statement, install its target and finish the lease barrier.
217    /// Completion retries are idempotent; each call scans at most one four-shard slice.
218    /// # Errors
219    /// Disabled configuration, rejected statement or pending completion (`Retry-After: 1`).
220    pub async fn set_authority_generation(&self, wire: &str) -> Result<u64, ServerError> {
221        let fence = self
222            .cfg
223            .authority_fence
224            .as_ref()
225            .ok_or_else(|| ServerError::unimplemented("authority fencing is disabled"))?;
226        let AuthMode::AuthV2(auth) = &self.cfg.auth else {
227            return Err(ServerError::unimplemented(
228                "authority fencing requires auth-v2",
229            ));
230        };
231        let statement = fence.verify(wire, auth.audience(), self.clock.now_ms())?;
232        self.authority_namespace(&statement.namespace).await?;
233        let ns = NamespaceKey::from_namespace(&statement.namespace);
234        if self
235            .transition_fence(&ns, statement.generation, FenceKind::Authority)
236            .await?
237            == EpochTransition::Reject
238        {
239            return Err(ServerError::permission_denied(
240                "authority generation step rejected",
241            ));
242        }
243        if let Some(completed) = Box::pin(self.ensure_authority_activation(&ns)).await? {
244            if self.stored_authority_generation(&ns).await? == completed {
245                return Ok(completed);
246            }
247            return Err(
248                ServerError::unavailable("authority revocation pending; retry")
249                    .with_header("Retry-After", "1"),
250            );
251        }
252        let start = ms(self.clock.now_ms());
253        for _ in 0..1 {
254            let elapsed = ms(self.clock.now_ms()).saturating_sub(start);
255            if elapsed >= 5_000 {
256                break;
257            }
258            let before = self.stored_authority_generation(&ns).await?;
259            let budget = RevokeBudget::new((5_000 - elapsed).min(1_000));
260            if self
261                .revoke_fence_step(&ns, &budget, FenceKind::Authority)
262                .await?
263                == RevokeProgress::Complete
264            {
265                let after = self.stored_authority_generation(&ns).await?;
266                if before == after {
267                    return Ok(after);
268                }
269            }
270        }
271        Err(
272            ServerError::unavailable("authority revocation pending; retry")
273                .with_header("Retry-After", "1"),
274        )
275    }
276}