use mkit_attest::grant::{EpochTransition, GrantError, MAX_GRANT_HEADER_BYTES};
use mkit_core::hash::to_hex_bytes;
use mkit_core::repo_identity::Namespace;
use crate::error::ServerError;
use crate::policy::NamespacePolicy;
use crate::repo::{Addressing, NamespaceKey};
use crate::store::{MultipartBlobStore, NamespaceStore, codec, keys};
use super::{
HookSet, Pipeline, meta_error, ms,
revocation::{RevokeBudget, RevokeProgress},
};
const MAX_REVOKE_SLICES: usize = 5;
const MAX_REVOKE_ELAPSED_MS: u64 = 5_000;
fn rejected(reason: &str) -> ServerError {
ServerError::permission_denied(format!("epoch statement rejected: {reason}"))
}
impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H> {
pub(super) async fn stored_grant_epoch(&self, key: &NamespaceKey) -> Result<u64, ServerError> {
self.meta
.get(&self.shards.coordinator(key), &keys::grant_epoch())
.await
.map_err(meta_error)?
.as_ref()
.map(codec::decode_u64)
.transpose()
.map_err(meta_error)
.map(|epoch| epoch.unwrap_or(0))
}
pub async fn get_grant_epoch(&self, namespace: &str) -> Result<u64, ServerError> {
let namespace = Namespace::parse(namespace)
.map_err(|_| ServerError::invalid_argument("invalid namespace"))?;
let Addressing::Multi(multi) = &self.cfg.addressing else {
return Err(ServerError::unimplemented(
"grant epochs require multi-repository addressing",
));
};
let key = NamespaceKey::from_namespace(&namespace);
let coordinator = self.shards.coordinator(&key);
match &multi.namespace_policy {
NamespacePolicy::Allowlist(allowed) if !allowed.contains(&namespace) => return Ok(0),
NamespacePolicy::Any { .. } => {
if matches!(namespace, Namespace::Address(_)) {
return Ok(0);
}
if self
.meta
.get(&coordinator, &keys::namespace_record())
.await
.map_err(meta_error)?
.is_none()
{
return Ok(0);
}
}
_ => {}
}
self.stored_grant_epoch(&key).await
}
pub async fn set_grant_epoch(&self, signed_statement: &str) -> Result<u64, ServerError> {
let Addressing::Multi(multi) = &self.cfg.addressing else {
return Err(ServerError::unimplemented(
"grant epochs require multi-repository addressing",
));
};
if signed_statement.len() > MAX_GRANT_HEADER_BYTES {
return Err(rejected(GrantError::HeaderTooLong.reason()));
}
let config = self
.cfg
.grants
.as_ref()
.ok_or_else(|| rejected(GrantError::SchemeNotAdvertised.reason()))?;
let verified = config
.verify_epoch(signed_statement, self.clock.now_ms())
.map_err(|error| rejected(error.reason()))?;
let statement = verified.statement();
let namespace = &statement.namespace;
if self.owner_key_is_admin(namespace) {
return Err(rejected("admin key cannot authorize client calls"));
}
let key = NamespaceKey::from_namespace(namespace);
match &multi.namespace_policy {
NamespacePolicy::Allowlist(allowed) if !allowed.contains(namespace) => {
return Err(rejected("namespace not served"));
}
NamespacePolicy::Any { .. } => {
if matches!(namespace, Namespace::Address(_)) {
return Err(rejected("namespace not served"));
}
let coordinator = self.shards.coordinator(&key);
if self
.meta
.get(&coordinator, &keys::namespace_record())
.await
.map_err(meta_error)?
.is_none()
{
return Err(rejected("namespace not served"));
}
}
_ => {}
}
tracing::debug!(statement_id = %to_hex_bytes(verified.id()), "epoch statement accepted");
if self.transition_epoch(&key, statement.new_epoch).await? == EpochTransition::Reject {
return Err(rejected("epoch step"));
}
let start = ms(self.clock.now_ms());
for _ in 0..MAX_REVOKE_SLICES {
let elapsed = ms(self.clock.now_ms()).saturating_sub(start);
if elapsed >= MAX_REVOKE_ELAPSED_MS {
break;
}
let budget = RevokeBudget::new((MAX_REVOKE_ELAPSED_MS - elapsed).min(1_000));
let before = self.stored_grant_epoch(&key).await?;
if self.revoke_step(&key, &budget).await? == RevokeProgress::Complete {
let after = self.stored_grant_epoch(&key).await?;
if before == after {
return Ok(after);
}
}
}
Err(ServerError::unavailable("epoch revocation pending; retry")
.with_header("Retry-After", "1"))
}
}