mkit_server/pipeline/
epoch.rs1use 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
17const 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 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 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 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 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 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 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}