Skip to main content

mkit_server/pipeline/
epoch.rs

1//! Unsigned namespace epoch RPCs (SPEC-WRITE-GRANTS §§5.2–5.4).
2
3use mkit_attest::grant::{EpochTransition, GrantError, MAX_GRANT_HEADER_BYTES};
4use mkit_core::hash::to_hex_bytes;
5use mkit_core::repo_identity::Namespace;
6
7use crate::error::ServerError;
8use crate::policy::NamespacePolicy;
9use crate::repo::{Addressing, NamespaceKey};
10use crate::store::{MultipartBlobStore, NamespaceStore, codec, keys};
11
12use super::{
13    HookSet, Pipeline, meta_error, ms,
14    revocation::{RevokeBudget, RevokeProgress},
15};
16
17/// At most five revocation slices and five seconds per `SetGrantEpoch` request.
18/// The slice count also terminates when the injected clock is frozen.
19const MAX_REVOKE_SLICES: usize = 5;
20const MAX_REVOKE_ELAPSED_MS: u64 = 5_000;
21
22fn rejected(reason: &str) -> ServerError {
23    ServerError::permission_denied(format!("epoch statement rejected: {reason}"))
24}
25
26impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H> {
27    pub(super) async fn stored_grant_epoch(&self, key: &NamespaceKey) -> Result<u64, ServerError> {
28        self.meta
29            .get(&self.shards.coordinator(key), &keys::grant_epoch())
30            .await
31            .map_err(meta_error)?
32            .as_ref()
33            .map(codec::decode_u64)
34            .transpose()
35            .map_err(meta_error)
36            .map(|epoch| epoch.unwrap_or(0))
37    }
38
39    /// Read a namespace epoch without authentication or repository resolution.
40    ///
41    /// # Errors
42    /// `invalid_argument` for a bad namespace, `unimplemented` under Single,
43    /// or a mapped storage error.
44    pub async fn get_grant_epoch(&self, namespace: &str) -> Result<u64, ServerError> {
45        let namespace = Namespace::parse(namespace)
46            .map_err(|_| ServerError::invalid_argument("invalid namespace"))?;
47        let Addressing::Multi(multi) = &self.cfg.addressing else {
48            return Err(ServerError::unimplemented(
49                "grant epochs require multi-repository addressing",
50            ));
51        };
52        let key = NamespaceKey::from_namespace(&namespace);
53        let coordinator = self.shards.coordinator(&key);
54        match &multi.namespace_policy {
55            NamespacePolicy::Allowlist(allowed) if !allowed.contains(&namespace) => return Ok(0),
56            NamespacePolicy::Any { .. } => {
57                // A 0x namespace is served only through an allowlist. Under
58                // Any, even an existing record cannot make it grant-writable.
59                if matches!(namespace, Namespace::Address(_)) {
60                    return Ok(0);
61                }
62                if self
63                    .meta
64                    .get(&coordinator, &keys::namespace_record())
65                    .await
66                    .map_err(meta_error)?
67                    .is_none()
68                {
69                    return Ok(0);
70                }
71            }
72            _ => {}
73        }
74        self.stored_grant_epoch(&key).await
75    }
76
77    /// Verify an owner epoch statement, CAS the coordinator epoch, then
78    /// fence every leased shard before reporting completion.
79    ///
80    /// # Errors
81    /// `permission_denied` for a rejected statement, `unimplemented` under
82    /// Single, or `unavailable` while completion is pending.
83    pub async fn set_grant_epoch(&self, signed_statement: &str) -> Result<u64, ServerError> {
84        let Addressing::Multi(multi) = &self.cfg.addressing else {
85            return Err(ServerError::unimplemented(
86                "grant epochs require multi-repository addressing",
87            ));
88        };
89        // Check the byte count before the header's base64 decoding or any
90        // expensive signature work (§5.2 check 1).
91        if signed_statement.len() > MAX_GRANT_HEADER_BYTES {
92            return Err(rejected(GrantError::HeaderTooLong.reason()));
93        }
94        let config = self
95            .cfg
96            .grants
97            .as_ref()
98            .ok_or_else(|| rejected(GrantError::SchemeNotAdvertised.reason()))?;
99        let verified = config
100            .verify_epoch(signed_statement, self.clock.now_ms())
101            .map_err(|error| rejected(error.reason()))?;
102        let statement = verified.statement();
103        let namespace = &statement.namespace;
104        if self.owner_key_is_admin(namespace) {
105            return Err(rejected("admin key cannot authorize client calls"));
106        }
107        let key = NamespaceKey::from_namespace(namespace);
108
109        // §5.2 check 6. Allowlisted namespaces can raise an epoch before a
110        // repository exists. Under Any, only an admitted write can create
111        // the namespace record that makes a later epoch change acceptable.
112        match &multi.namespace_policy {
113            NamespacePolicy::Allowlist(allowed) if !allowed.contains(namespace) => {
114                return Err(rejected("namespace not served"));
115            }
116            NamespacePolicy::Any { .. } => {
117                if matches!(namespace, Namespace::Address(_)) {
118                    return Err(rejected("namespace not served"));
119                }
120                let coordinator = self.shards.coordinator(&key);
121                if self
122                    .meta
123                    .get(&coordinator, &keys::namespace_record())
124                    .await
125                    .map_err(meta_error)?
126                    .is_none()
127                {
128                    return Err(rejected("namespace not served"));
129                }
130            }
131            _ => {}
132        }
133
134        tracing::debug!(statement_id = %to_hex_bytes(verified.id()), "epoch statement accepted");
135        if self.transition_epoch(&key, statement.new_epoch).await? == EpochTransition::Reject {
136            return Err(rejected("epoch step"));
137        }
138
139        let start = ms(self.clock.now_ms());
140        for _ in 0..MAX_REVOKE_SLICES {
141            let elapsed = ms(self.clock.now_ms()).saturating_sub(start);
142            if elapsed >= MAX_REVOKE_ELAPSED_MS {
143                break;
144            }
145            let budget = RevokeBudget::new((MAX_REVOKE_ELAPSED_MS - elapsed).min(1_000));
146            // A concurrent statement may raise e between our CAS and the
147            // scan. Only answer with an epoch that this scan actually fenced.
148            let before = self.stored_grant_epoch(&key).await?;
149            if self.revoke_step(&key, &budget).await? == RevokeProgress::Complete {
150                let after = self.stored_grant_epoch(&key).await?;
151                if before == after {
152                    return Ok(after);
153                }
154            }
155        }
156        Err(ServerError::unavailable("epoch revocation pending; retry")
157            .with_header("Retry-After", "1"))
158    }
159}