mod admission;
mod advance;
mod auth;
mod authority;
mod begin;
pub mod clearance;
mod coordinator;
mod download;
mod durable_outcome;
mod epoch;
#[cfg(feature = "test-faults")]
pub(crate) mod faults;
mod gate;
mod hooks;
#[cfg(feature = "http-objects")]
mod http;
#[cfg(feature = "http-objects")]
mod http_admission;
#[cfg(feature = "http-objects")]
mod http_tokens;
#[cfg(feature = "http-objects")]
mod object_reader;
#[cfg(feature = "http-objects")]
pub use object_reader::{
IssuedUrl, OBJECT_READER_BATCH, OBJECT_READER_CALLS, ObjectMetadata, ObjectReader, ReaderView,
};
mod implicit;
mod info;
#[cfg(feature = "remote-hooks")]
pub mod inspection;
mod lease;
pub mod list;
mod list_repos;
pub use list_repos::{RepoEntry, RepoPage};
mod outcome;
mod parts;
mod plan;
#[cfg(feature = "published-view")]
pub mod published;
mod purge;
mod ref_policy;
mod reservation;
mod revocation;
mod scanner_retrieval;
mod shard;
mod staging;
#[cfg(test)]
mod tests;
mod upload;
mod watermark;
use core::future::Future;
use core::time::Duration;
use std::sync::Arc;
pub use mkit_attest::grant::Visibility as RepoVisibility;
use mkit_attest::grant::{Visibility, verify_visibility_statement};
use mkit_core::hash::{Hash, to_hex, to_hex_bytes};
use mkit_core::protocol::{AdvanceOutcome, PackKey};
use mkit_core::repo_identity::{Namespace, RepositoryIdentity};
use mkit_core::write_auth::MAX_CLOCK_LEAD_MS;
use tracing::Instrument;
use crate::download::DOWNLOAD_CHUNK_MAX;
use crate::error::{AbortCause, InvalidHeader, ServerError};
use crate::op::{
AuthzFacts, CallerView, GrantRef, OpKind, Operation, Procedure, RefUpdate, VerifiedAuth,
};
use crate::policy::ff::FastForward;
use crate::policy::{
AuthorizerRole, GrantConfig, NamespacePolicy, WritePolicy, grants, read as read_policy,
};
use crate::principal::Principal;
use crate::quota::{
self, DEFAULT_WRITE_QUOTA, NamespaceCharge, NamespaceDecision, NamespaceView, QuotaCharge,
QuotaLimits, QuotaScope, ViewStatus,
};
use crate::refs::{self, strip_listed_prefix, validate_ref_name};
use crate::replay::{
BeginUploadResult, ReplayDecision, ReplayRecord, ReplayState, StoredResult, UpdateRefResult,
classify,
};
use crate::repo::{Addressing, RepoId};
use crate::rt::Clock;
use crate::storage_error::{StorageOp, describe_and_map};
use crate::store::tickets::TicketCaps;
use crate::store::{
Batch, BatchOutcome, Key, KeyClasses, MAX_BATCH_OPS, MultipartBlobStore, NamespaceStore,
Partition, Precondition, StoreError, Value, codec, keys, read,
};
use crate::telemetry::{Metrics, Redactor};
use crate::upload::{UploadLimits, token::TicketKeys};
use crate::url_token::{MintedToken, UrlTarget};
use begin::BeginWrite;
#[cfg(feature = "remote-hooks")]
pub(crate) use admission::validate_decision;
pub use auth::{AuthMode, Authenticated, HeaderValues, RequestMeta};
pub use download::{DownloadChunk, DownloadStream};
pub use durable_outcome::{DeliveryError, Outcome, OutcomeKind};
#[cfg(feature = "test-faults")]
pub use faults::{
BUMP_EPOCH_HEADER, CLOCK_SKEW_HEADER, FAULT_HEADER, FailOnce, FaultHooks, FaultPoint,
LEASE_RECOVERED_HEADER, RELAY_DELAY_MS_HEADER, RUN_TIMERS_HEADER, TIMER_MS_HEADER,
TestDirectives,
};
pub use hooks::{
ADMISSION_EXPOSE_HEADERS, Admission, AdmissionDecision, AdmissionInput, Authorizer, Challenge,
Choice, CredentialHeader, DefaultAdmission, HookSet, Hooks, NoOutcomes, NoPreReceive,
NoReceipts, OpenAuthorizer, OutcomeSink, PreReceive, ReceiptSigner,
};
#[cfg(feature = "ssh")]
pub(crate) use implicit::IMPLICIT_PACKMAP_UNKNOWN;
pub(crate) use implicit::PendingPack;
pub use info::ServerInfo;
pub use lease::{LeaseParams, renew_for_relay};
use outcome::Outcome as RequestOutcome;
pub use parts::PartUploadSession;
use plan::{
ImplicitConsume, MAX_REPLAN, PRUNE_LIMIT, Plan, PlanClock, Planned, Snapshot, WriteKind,
WriteRequest, plan_write, prune_sampled,
};
pub use revocation::{MAX_EPOCH_STEP, RevokeBudget, RevokeProgress};
pub use shard::{D34Shards, ShardMap, SinglePartition};
pub use upload::{UploadMode, UploadSession};
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct ResponseMeta {
headers: Vec<(String, String)>,
external_ref: Option<String>,
}
impl ResponseMeta {
#[must_use]
pub fn headers(&self) -> &[(String, String)] {
&self.headers
}
#[must_use]
pub fn external_ref(&self) -> Option<&str> {
self.external_ref.as_deref()
}
}
macro_rules! fault {
($pipe:expr, $point:ident, $op:expr, $a:expr) => {
#[cfg(feature = "test-faults")]
$pipe
.fault($crate::pipeline::FaultPoint::$point, $op, $a)
.await?
};
}
pub(crate) use fault;
pub const MAX_APPLY_WINDOW: Duration = Duration::from_secs(10);
const _: () = assert!(MAX_CLOCK_LEAD_MS.unsigned_abs() < read::REPLAY_PRUNE_GRACE_MS);
pub const DEFAULT_LIST_PAGE_LIMIT: u32 = 1000;
pub const METRIC_PARTITION_FULL: &str = "mkit_server_partition_full_total";
pub const METRIC_HEADER_DROPPED: &str = "mkit_server_error_header_dropped_total";
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
#[non_exhaustive]
pub enum Sharding {
#[default]
Single,
D34,
}
fn require_relay_source_lease(
sharding: Sharding,
has_lease: bool,
batch: &Batch,
) -> Result<(), ServerError> {
if sharding == Sharding::D34
&& !has_lease
&& batch.writes.iter().any(|write| {
matches!(write, crate::store::Write::Put(key, _) if matches!(keys::parse(key), Some(keys::ParsedKey::Relay(_))))
})
{
return Err(internal("D34 relay batch lacks source epoch lease"));
}
Ok(())
}
#[derive(Debug)]
struct ReadAuth {
facts: AuthzFacts,
epoch: Option<u64>,
}
#[derive(Debug, Clone)]
#[non_exhaustive]
pub struct PipelineConfig {
pub addressing: Addressing,
pub sharding: Sharding,
pub inspection_mode: bool,
pub auth: AuthMode,
pub grants: Option<GrantConfig>,
pub authority_fence: Option<crate::authority::AuthorityFence>,
pub write_policy: WritePolicy,
pub default_repo_visibility: RepoVisibility,
pub authorizer_role: AuthorizerRole,
pub list_repos_authority_full: bool,
pub upload_limits: UploadLimits,
pub single_upload_max_bytes: Option<u64>,
pub part_size: u64,
pub max_parts: u32,
pub max_list_refs_page_size: u32,
pub begin_upload_threshold_bytes: u64,
pub ticket_keys: Option<TicketKeys>,
pub scanner_retrieval: Option<Arc<crate::scanner_retrieval::RetrievalConfig>>,
pub url_tokens: Option<crate::url_token::UrlTokenConfig>,
pub admin_keys: Vec<[u8; 32]>,
pub receipt_publication: Option<crate::takedown::PublicationConfig>,
pub purge: Option<crate::purge::PurgeConfig>,
pub ticket_ttl_ms: u64,
pub ticket_caps: TicketCaps,
pub download_chunk_max: usize,
pub write_quota: Option<QuotaLimits>,
pub list_page_limit: u32,
pub max_apply_window: Duration,
pub epoch_lease_ms: u64,
pub lease_margin_ms: u64,
pub min_lease_budget_ms: u64,
pub redactor: Redactor,
pub admission_credential_headers: Vec<String>,
pub outbox_backlog_cap: Option<OutboxBacklogCap>,
pub indexed: Option<crate::indexed::IndexedConfig>,
pub takedown_denial: bool,
pub ref_policy: Option<crate::policy::RefPolicy>,
#[cfg(feature = "http-objects")]
pub http_objects: Option<crate::http_objects::HttpObjectsConfig>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct OutboxBacklogCap {
pub rows: u64,
pub bytes: u64,
}
impl PipelineConfig {
pub(crate) fn indexed_mode(&self) -> bool {
self.indexed.is_some()
}
#[must_use]
pub fn new(addressing: Addressing, auth: AuthMode, upload_limits: UploadLimits) -> Self {
let write_quota = matches!(auth, AuthMode::AuthV2(_)).then_some(DEFAULT_WRITE_QUOTA);
let write_policy = match &addressing {
Addressing::Single { .. } => WritePolicy::Open,
Addressing::Multi(_) => WritePolicy::Owner,
};
Self {
write_policy,
default_repo_visibility: RepoVisibility::Public,
authorizer_role: AuthorizerRole::Check,
list_repos_authority_full: false,
addressing,
sharding: Sharding::Single,
inspection_mode: false,
auth,
grants: None,
authority_fence: None,
upload_limits,
single_upload_max_bytes: None,
part_size: mkit_core::upload_parts::MIN_PART_SIZE,
max_parts: 10_000,
max_list_refs_page_size: DEFAULT_LIST_PAGE_LIMIT,
begin_upload_threshold_bytes: u64::MAX,
ticket_keys: None,
scanner_retrieval: None,
url_tokens: None,
admin_keys: Vec::new(),
receipt_publication: None,
purge: None,
ticket_ttl_ms: 86_400_000,
ticket_caps: TicketCaps {
per_ref: 1024,
per_signer: 64,
},
download_chunk_max: DOWNLOAD_CHUNK_MAX,
write_quota,
list_page_limit: DEFAULT_LIST_PAGE_LIMIT,
max_apply_window: MAX_APPLY_WINDOW,
epoch_lease_ms: 30_000,
lease_margin_ms: 5_000,
min_lease_budget_ms: 1_000,
redactor: Redactor::default(),
admission_credential_headers: Vec::new(),
outbox_backlog_cap: Some(OutboxBacklogCap {
rows: 100_000,
bytes: 64 * 1024 * 1024,
}),
indexed: None,
takedown_denial: false,
ref_policy: None,
#[cfg(feature = "http-objects")]
http_objects: None,
}
}
#[must_use]
pub fn advertised_namespace_policy(&self) -> &'static str {
match &self.addressing {
Addressing::Single { .. } => "single-repository",
Addressing::Multi(multi) => match &multi.namespace_policy {
NamespacePolicy::Allowlist(_) => "allowlist",
NamespacePolicy::Any { .. } => "any",
},
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
#[non_exhaustive]
pub struct PipelineCapabilities {
pub atomic_advance: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct HealthStatus {
pub blobs: bool,
pub meta: bool,
}
impl HealthStatus {
#[must_use]
pub fn is_healthy(&self) -> bool {
self.blobs && self.meta
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RefEntry {
pub name: String,
pub id: Hash,
}
#[derive(Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum VisibilityRequest {
Envelope(Visibility),
Statement(String),
}
impl core::fmt::Debug for VisibilityRequest {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
match self {
Self::Envelope(visibility) => f.debug_tuple("Envelope").field(visibility).finish(),
Self::Statement(statement) => f
.debug_tuple("Statement")
.field(&format_args!("<{} bytes>", statement.len()))
.finish(),
}
}
}
pub struct Pipeline<B, N, H = Hooks> {
publication_policy: Option<Arc<dyn clearance::PublicationPolicy>>,
#[cfg(feature = "remote-hooks")]
inspectors: Vec<Arc<dyn inspection::ContentInspector>>,
#[cfg(feature = "remote-hooks")]
inspect_limit: usize,
#[cfg(feature = "published-view")]
published: Option<Arc<dyn published::PublishedSource>>,
blobs: B,
meta: N,
hooks: H,
shards: Arc<dyn ShardMap>,
cfg: PipelineConfig,
clock: Arc<dyn Clock>,
metrics: Arc<dyn Metrics>,
#[cfg(feature = "test-faults")]
faults: Option<Arc<dyn faults::DynFaultHooks>>,
#[cfg(feature = "test-faults")]
test_timer_gate: Option<Arc<tokio::sync::Mutex<()>>>,
gate: Option<Arc<gate::WriteGate>>,
#[cfg(feature = "http-objects")]
http_seams: Option<crate::http_objects::HttpSeams>,
}
struct WriteInputs<'a> {
denial_ids: &'a std::collections::BTreeSet<Hash>,
denial_packs: &'a [Hash],
pending: Option<&'a reservation::PendingGuard>,
implicit: Option<&'a [PendingPack]>,
external_bases: &'a std::collections::BTreeSet<Hash>,
inspected: Option<&'a mut crate::indexed::inspection::InspectionSet>,
}
impl<B, N, H> core::fmt::Debug for Pipeline<B, N, H> {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct("Pipeline")
.field("cfg", &self.cfg)
.finish_non_exhaustive()
}
}
fn internal(detail: &'static str) -> ServerError {
ServerError::internal("ref store request failed", detail)
}
fn store_error(op: StorageOp, err: StoreError) -> ServerError {
let op = match err {
StoreError::Corrupt(_) => StorageOp::MetaDecode,
_ => op,
};
let (line, err) = describe_and_map(op, err);
tracing::warn!(detail = %line, "storage failure");
err
}
fn meta_error(err: StoreError) -> ServerError {
store_error(StorageOp::MetaCall, err)
}
fn method(procedure: Procedure) -> &'static str {
let path = procedure.connect_path();
path.rsplit('/').next().unwrap_or(path)
}
fn replay_answer(decision: ReplayDecision) -> Result<Option<StoredResult>, ServerError> {
match decision {
ReplayDecision::New => Ok(None),
ReplayDecision::Return(result) => Ok(Some(result)),
ReplayDecision::FingerprintMismatch => Err(ServerError::invalid_argument(
"nonce reused for a different operation",
)),
ReplayDecision::RetryLater | ReplayDecision::Resume => Err(ServerError::aborted_retryable(
"operation already in flight; retry",
)),
}
}
fn ms(ms: i64) -> u64 {
u64::try_from(ms).unwrap_or(0)
}
#[must_use]
pub fn repo_is_private(stored: Option<&codec::RepoVisibilityV1>, default: RepoVisibility) -> bool {
stored.map_or(default == RepoVisibility::Private, |row| {
row.visibility == codec::StoredVisibility::Private
})
}
fn stored_visibility(visibility: Visibility) -> codec::StoredVisibility {
match visibility {
Visibility::Public => codec::StoredVisibility::Public,
Visibility::Private => codec::StoredVisibility::Private,
}
}
fn validate_upload_ticket_config<H: HookSet>(
cfg: &PipelineConfig,
hooks: &H,
) -> Result<(), ServerError> {
if cfg.begin_upload_threshold_bytes != u64::MAX
&& !matches!(cfg.auth, AuthMode::TransportIdentity)
&& (!matches!(cfg.auth, AuthMode::AuthV2(_)) || cfg.ticket_keys.is_none())
{
return Err(ServerError::invalid_argument(
"a ticket threshold requires auth v2 and upload ticket keys",
));
}
if !matches!(cfg.auth, AuthMode::TransportIdentity)
&& !hooks.admission().is_default()
&& (!matches!(cfg.auth, AuthMode::AuthV2(_)) || cfg.ticket_keys.is_none())
{
return Err(ServerError::invalid_argument(
"admission requires auth v2 and upload ticket keys",
));
}
if matches!(cfg.addressing, Addressing::Multi(_))
&& matches!(cfg.auth, AuthMode::AuthV2(_))
&& cfg.ticket_keys.is_none()
{
return Err(ServerError::invalid_argument(
"multi-repository auth v2 deployments require upload ticket keys",
));
}
Ok(())
}
impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H> {
#[allow(clippy::too_many_lines)] pub fn new(
blobs: B,
meta: N,
hooks: H,
mut cfg: PipelineConfig,
clock: Arc<dyn Clock>,
metrics: Arc<dyn Metrics>,
) -> Result<Self, ServerError> {
if let Some(publication) = &cfg.receipt_publication {
if cfg.indexed.is_none() {
return Err(ServerError::invalid_argument(
"takedown requires indexed mode",
));
}
cfg.admin_keys.extend_from_slice(publication.public_keys());
}
crate::scanner_retrieval::service::validate_config(&cfg)?;
if let Some(purge) = &cfg.purge {
purge.validate().map_err(meta_error)?;
}
if cfg.authority_fence.is_some()
&& (cfg.authorizer_role != AuthorizerRole::Authority
|| hooks.authorizer().is_open()
|| !meta.capabilities().atomic_multi_key
|| !matches!(cfg.auth, AuthMode::AuthV2(_))
|| !matches!(cfg.addressing, Addressing::Multi(_)))
{
return Err(ServerError::invalid_argument(
"authority fencing requires Multi, auth v2, an Authority hook and transactional storage",
));
}
#[allow(clippy::collapsible_if)]
if let Some(fence) = &cfg.authority_fence {
if fence.public_keys().any(|key| {
cfg.ticket_keys
.as_ref()
.is_some_and(|tickets| tickets.contains_ed25519_public(&key))
}) {
return Err(ServerError::invalid_argument(
"authority keys must differ from ticket keys",
));
}
#[cfg(feature = "http-objects")]
if fence.public_keys().any(|key| {
cfg.url_tokens
.as_ref()
.is_some_and(|tokens| tokens.keys().public_keys().any(|public| public == key))
}) {
return Err(ServerError::invalid_argument(
"authority keys must differ from URL-token keys",
));
}
}
if cfg.takedown_denial && cfg.indexed.is_none() {
return Err(ServerError::invalid_argument(
"takedown denial requires indexed mode",
));
}
if let Some(indexed) = &cfg.indexed {
if !matches!(cfg.auth, AuthMode::AuthV2(_))
|| cfg.ticket_keys.is_none()
|| !matches!(cfg.addressing, Addressing::Multi(_))
{
return Err(ServerError::invalid_argument(
"indexed mode requires auth v2, ticket keys, zero ticket threshold, and Multi addressing",
));
}
if indexed.max_delta_chain_depth == 0
|| indexed.max_delta_chain_depth > u32::from(u16::MAX)
|| indexed.max_pack_bytes == 0
|| indexed.max_pack_bytes > indexed.decode_budget
|| indexed.relay_lag_bound_ms == 0
|| indexed.extract_min_bytes == 0
|| indexed.max_ancestry_commits == 0
|| indexed.max_ancestry_commits > crate::indexed::MAX_ANCESTRY_COMMITS_LIMIT
|| indexed
.max_extract_bytes
.is_some_and(|max| max < indexed.max_pack_bytes)
{
return Err(ServerError::invalid_argument("invalid indexed limits"));
}
cfg.upload_limits.max_total_bytes = cfg
.upload_limits
.max_total_bytes
.min(indexed.max_pack_bytes);
}
if let Some(policy) = &cfg.ref_policy {
policy.validate_for_indexed(cfg.indexed.is_some())?;
}
#[cfg(feature = "http-objects")]
if let Some(http) = &cfg.http_objects {
let Some(indexed) = &cfg.indexed else {
return Err(ServerError::invalid_argument(
"HTTP object serving requires indexed mode",
));
};
http.validate(indexed.extract_min_bytes)?;
if http.admit_reads && hooks.admission().is_default() {
return Err(ServerError::invalid_argument(
"paid HTTP reads require a real Admission hook",
));
}
}
cfg.validate_server_info_limits()?;
let mut credential_names = std::collections::BTreeSet::new();
if cfg.admission_credential_headers.iter().any(|name| {
!admission::valid_extra_name(name)
|| ["payment-authorization", "payment-signature"]
.contains(&name.to_ascii_lowercase().as_str())
|| !credential_names.insert(name.to_ascii_lowercase())
}) {
return Err(ServerError::invalid_argument(
"invalid admission credential header name",
));
}
if cfg.admission_credential_headers.len() + 3 > admission::MAX_CREDENTIAL_HEADERS {
return Err(ServerError::invalid_argument(
"too many admission credential headers",
));
}
cfg.redactor.add_names(&cfg.admission_credential_headers);
if cfg.max_parts > B::MAX_PARTS {
return Err(ServerError::invalid_argument(
"max_parts exceeds storage backend capacity",
));
}
validate_upload_ticket_config(&cfg, &hooks)?;
if cfg.ticket_ttl_ms == 0
|| cfg.ticket_ttl_ms >= 604_800_000
|| cfg.ticket_caps.per_ref == 0
|| cfg.ticket_caps.per_signer == 0
{
return Err(ServerError::invalid_argument(
"invalid upload ticket lifetime or caps",
));
}
let policy_refusal = match (&cfg.addressing, cfg.write_policy) {
(Addressing::Multi(_), WritePolicy::Open) => {
Some("write_policy open is single-repository only (SPEC-TRANSPORT-CONNECT §7.5)")
}
(Addressing::Single { repo }, WritePolicy::Owner)
if Namespace::parse(repo.namespace.as_str()).is_err() =>
{
Some("write_policy owner needs multi-repository addressing")
}
(Addressing::Multi(multi), _)
if matches!(
multi.namespace_policy,
NamespacePolicy::Any {
unsafe_without_admission: false
}
) && hooks.admission().is_default() =>
{
Some(
"namespace_policy any needs a non-default admission step, or the explicit unsafe override (D27)",
)
}
_ => None,
};
if let Some(message) = policy_refusal {
return Err(ServerError::invalid_argument(message));
}
if let Some(grants) = &cfg.grants {
let AuthMode::AuthV2(auth) = &cfg.auth else {
return Err(ServerError::invalid_argument(
"write grants require auth v2",
));
};
if !matches!(cfg.addressing, Addressing::Multi(_))
|| cfg.write_policy != WritePolicy::Owner
{
return Err(ServerError::invalid_argument(
"write grants require Multi addressing and owner write policy",
));
}
if grants.audience() != auth.audience() {
return Err(ServerError::invalid_argument(
"write grant audience must match auth v2 audience",
));
}
}
if let Some(tokens) = &cfg.url_tokens {
if !matches!(cfg.auth, AuthMode::AuthV2(_)) {
return Err(ServerError::invalid_argument("URL tokens require auth v2"));
}
if let Some(tickets) = &cfg.ticket_keys
&& tokens
.keys()
.public_keys()
.any(|public| tickets.contains_ed25519_public(&public))
{
return Err(ServerError::invalid_argument(
"the URL token key must differ from the upload ticket keys",
));
}
}
if cfg.authorizer_role == AuthorizerRole::Authority && hooks.authorizer().is_open() {
return Err(ServerError::invalid_argument(
"an authority authorizer must be a real authority source",
));
}
let caps = meta.capabilities();
let full = caps.atomic_multi_key && caps.key_classes == KeyClasses::All;
let refused = if matches!(cfg.auth, AuthMode::AuthV2(_)) && !full {
"auth v2 needs every key class and atomic multi-key batches"
} else if (cfg.sharding == Sharding::D34 || matches!(cfg.addressing, Addressing::Multi(_)))
&& !full
{
"sharded or multi-repository routing needs every key class and atomic multi-key batches"
} else if !caps.atomic_multi_key && caps.implicit_layout_version.is_none() {
"a store without atomic batches must report its layout version"
} else if caps
.implicit_layout_version
.is_some_and(|v| v != keys::LAYOUT_VERSION)
{
"the store's layout version is not this server's"
} else if cfg.lease_margin_ms == 0
|| cfg.epoch_lease_ms <= cfg.lease_margin_ms.saturating_add(cfg.min_lease_budget_ms)
{
"epoch lease must exceed its positive margin plus minimum budget"
} else if cfg.list_page_limit == 0 || cfg.max_apply_window.is_zero() {
"list page limit and apply window must be positive"
} else {
""
};
if !refused.is_empty() {
return Err(ServerError::invalid_argument(refused));
}
let shards: Arc<dyn ShardMap> = match cfg.sharding {
Sharding::Single => Arc::new(SinglePartition),
Sharding::D34 => Arc::new(D34Shards),
};
#[cfg(feature = "http-objects")]
let http_seams = cfg.http_objects.as_ref().map(|http| {
let mut seams = crate::http_objects::HttpSeams::new(http);
if let Some(tokens) = &cfg.url_tokens {
seams.tokens = Arc::new(tokens.clone());
}
seams
});
Ok(Self {
#[cfg(feature = "published-view")]
published: None,
publication_policy: None,
#[cfg(feature = "remote-hooks")]
inspectors: Vec::new(),
#[cfg(feature = "remote-hooks")]
inspect_limit: inspection::MAX_OBJECTS,
blobs,
meta,
hooks,
shards,
cfg,
clock,
metrics,
#[cfg(feature = "test-faults")]
faults: None,
#[cfg(feature = "test-faults")]
test_timer_gate: None,
gate: None,
#[cfg(feature = "http-objects")]
http_seams,
})
}
#[must_use]
pub fn inspection_max_objects(&self) -> Option<u32> {
#[cfg(feature = "remote-hooks")]
if !self.inspectors.is_empty() {
return u32::try_from(self.inspect_limit).ok();
}
None
}
#[cfg(feature = "remote-hooks")]
pub fn with_inspectors(
mut self,
inspectors: Vec<Arc<dyn inspection::ContentInspector>>,
batch_max: usize,
) -> Result<Self, ServerError> {
if inspectors.is_empty() {
return Ok(self);
}
let mut names = std::collections::BTreeSet::new();
if inspectors.len() > inspection::MAX_INSPECTORS
|| !(1..=inspection::MAX_OBJECTS).contains(&batch_max)
|| inspectors.iter().any(|i| {
i.id().is_empty()
|| !names.insert(i.id())
|| i.phase() != inspection::InspectorPhase::Sync
|| i.on_unavailable() != inspection::OnUnavailable::FailClosed
})
{
return Err(ServerError::invalid_argument(
"invalid launch inspection configuration",
));
}
if self.cfg.indexed.is_none()
|| self.cfg.write_policy == WritePolicy::Open
|| self.cfg.ticket_keys.is_none()
|| self.cfg.begin_upload_threshold_bytes != 0
|| !self.meta.capabilities().atomic_multi_key
{
return Err(ServerError::invalid_argument(
"inspection requires indexed mode, restricted writes and ticketed uploads with threshold zero",
));
}
self.inspectors = inspectors;
self.inspect_limit = batch_max;
Ok(self)
}
#[cfg(feature = "test-faults")]
#[must_use]
pub fn with_faults(mut self, hooks: impl FaultHooks + 'static) -> Self {
self.faults = Some(Arc::new(hooks));
self
}
#[cfg(feature = "test-faults")]
#[must_use]
pub fn with_test_timer_gate(mut self, gate: Arc<tokio::sync::Mutex<()>>) -> Self {
self.test_timer_gate = Some(gate);
self
}
#[cfg(feature = "test-faults")]
async fn fault(
&self,
point: FaultPoint,
op: &Operation,
a: &Authenticated,
) -> Result<(), ServerError> {
match &self.faults {
Some(hooks) => hooks.at_boxed(point, op, a.test_directives()).await,
None => Ok(()),
}
}
#[must_use]
pub fn with_write_gate(mut self) -> Self {
self.gate = Some(Arc::new(gate::WriteGate::new()));
self
}
pub fn with_auth(&self, auth: AuthMode) -> Result<Self, ServerError>
where
B: Clone,
N: Clone,
H: Clone,
{
let mut cfg = self.cfg.clone();
cfg.auth = auth;
cfg.url_tokens = None;
if !matches!(cfg.auth, AuthMode::AuthV2(_)) {
cfg.grants = None;
}
let mut sibling = Self::new(
self.blobs.clone(),
self.meta.clone(),
self.hooks.clone(),
cfg,
Arc::clone(&self.clock),
Arc::clone(&self.metrics),
)?;
sibling.shards = Arc::clone(&self.shards);
sibling.gate.clone_from(&self.gate);
#[cfg(feature = "remote-hooks")]
{
sibling.inspectors.clone_from(&self.inspectors);
sibling.inspect_limit = self.inspect_limit;
}
#[cfg(feature = "test-faults")]
{
sibling.faults.clone_from(&self.faults);
sibling.test_timer_gate.clone_from(&self.test_timer_gate);
}
#[cfg(feature = "http-objects")]
sibling.http_seams.clone_from(&self.http_seams);
#[cfg(feature = "published-view")]
sibling.published.clone_from(&self.published);
Ok(sibling)
}
pub fn authenticate(&self, meta: &RequestMeta<'_>) -> Result<Authenticated, ServerError> {
tracing::debug!(stage = "authenticate", procedure = method(meta.procedure));
let result = self.authenticate_inner(meta);
if let Err(err) = &result {
self.outcome_for(meta.procedure, "none", "-")
.record(Err(err));
}
result
}
fn authenticate_inner(&self, meta: &RequestMeta<'_>) -> Result<Authenticated, ServerError> {
let signed = auth::signed_request(&self.cfg.auth, meta);
if !signed && (meta.header)("x-write-grant").is_some() {
return Err(ServerError::unauthenticated(
"write grant requires auth v2 authorization",
));
}
let repo = self
.cfg
.addressing
.resolve((meta.header)("x-repository").as_deref(), signed)?;
let expected_repository = match (&self.cfg.addressing, &self.cfg.auth) {
(Addressing::Single { .. }, AuthMode::AuthV2(cfg)) => cfg.repository(),
_ => &repo.identity,
}
.to_owned();
#[cfg(feature = "test-faults")]
let directives = TestDirectives::from_headers(meta.header)?;
#[cfg(feature = "test-faults")]
let skew = directives.clock_skew_ms;
#[cfg(not(feature = "test-faults"))]
let skew = 0;
let now = self.clock.now_ms().saturating_add(skew);
let mut a = auth::authenticate(&self.cfg.auth, meta, now, repo, &expected_repository)?;
if a.principal
.ed25519()
.is_some_and(|key| self.cfg.admin_keys.contains(key))
{
return Err(ServerError::unauthenticated(
"admin key cannot authenticate client calls",
));
}
if a.principal.ed25519().is_some_and(|key| {
self.cfg
.scanner_retrieval
.as_ref()
.is_some_and(|config| config.scanner_keys().any(|public| &public == key))
}) {
return Err(ServerError::unauthenticated(
"scanner key cannot authenticate client calls",
));
}
if matches!(self.cfg.auth, AuthMode::AuthV2(_))
&& a.auth.is_some()
&& meta.procedure.is_write()
&& meta.procedure != Procedure::SetRepoVisibility
{
a.credential_capture =
admission::capture_credentials(meta, &self.cfg.admission_credential_headers);
}
a.ref_hint =
(meta.header)("x-mkit-ref").filter(|name| name.len() <= refs::MAX_REF_NAME_BYTES);
a.business_skew_ms = skew;
#[cfg(feature = "test-faults")]
a.set_test_directives(directives);
Ok(a)
}
pub async fn list_refs(
&self,
a: &Authenticated,
prefix: &str,
) -> Result<Vec<RefEntry>, ServerError> {
let mut out = Vec::new();
let mut token = None;
loop {
let page = self
.list_refs_page(a, prefix, Some(self.cfg.list_page_limit), token.as_deref())
.await?;
out.extend(page.refs);
match page.next {
Some(next) => token = Some(next),
None => return Ok(out),
}
}
}
#[allow(clippy::too_many_lines)] pub(crate) async fn list_refs_page(
&self,
a: &Authenticated,
prefix: &str,
page_size: Option<u32>,
token: Option<&[u8]>,
) -> Result<list::ListPage, ServerError> {
let kind = OpKind::ListRefs {
prefix: prefix.to_owned(),
};
self.observe(a, async {
if prefix.trim_end_matches('/').len() > refs::MAX_REF_NAME_BYTES {
return Err(ServerError::invalid_argument(refs::REF_NAME_TOO_LONG));
}
if !refs::validate_ref_prefix(prefix) {
return Err(ServerError::invalid_argument(
"prefix is invalid (SPEC-REFS §3)",
));
}
let op = self.identify(a, kind)?;
let authorized = self.authorize_read(&op).await?;
let scan = refs::list_scan_prefix(prefix);
let last = token
.map(|bytes| {
list::decode_token(&op.repo, &scan, bytes)
.ok_or_else(|| ServerError::invalid_argument("invalid page token"))
})
.transpose()?;
#[cfg(feature = "test-faults")]
{
if a.test_directives().lease_recovered {
self.mark_lease_table_recovered(&op.repo.namespace).await?;
}
if let Some(epoch) = a.test_directives().bump_epoch {
self.test_bump_epoch(&op.repo.namespace, epoch).await?;
}
}
#[cfg(feature = "test-faults")]
faults::run_timers(
a.test_directives(),
&self.meta,
&self.blobs,
self.shards.as_ref(),
&op.repo,
self.clock.as_ref(),
ms(self.clock.now_ms().saturating_add(a.business_skew_ms)),
self.test_timer_gate.as_deref(),
)
.await?;
let requested = page_size.unwrap_or(0);
let limit = if requested == 0 {
self.cfg.max_list_refs_page_size
} else {
requested.min(self.cfg.max_list_refs_page_size)
};
let partitions = self.shards.ref_index_partitions(&op.repo);
#[cfg(feature = "published-view")]
if authorized.facts.caller_view != CallerView::Writer
&& self
.published
.as_ref()
.is_some_and(|s| s.inspection_configured() && !s.uses_published_values())
{
return Err(ServerError::unavailable("published view unavailable"));
}
#[cfg(feature = "published-view")]
let source = self.published.as_deref().filter(|_| {
self.visibility_gates_reads()
&& op.auth.is_none()
&& matches!(op.principal, Principal::Anonymous)
&& (self.publication_policy.is_none()
|| self
.published
.as_ref()
.is_some_and(|s| s.uses_published_values()))
});
let view = crate::store::view::ViewStore {
store: &self.meta,
repo: &op.repo,
writer: authorized.facts.caller_view == CallerView::Writer,
policy: self.publication_policy.as_deref(),
};
let result = if partitions.len() == 1 {
let bucket = list::RefBucket {
store: &view,
partition: &partitions[0],
};
list::page(
&[bucket],
&op.repo,
&scan,
last.as_deref(),
limit,
list::MAX_RESPONSE_BYTES,
)
.await
} else {
#[cfg(feature = "published-view")]
let buckets = partitions
.iter()
.map(|partition| published::ReaderBucket {
store: &view,
partition,
source,
now_ms: ms(self.clock.now_ms()),
})
.collect::<Vec<_>>();
#[cfg(not(feature = "published-view"))]
let buckets = partitions
.iter()
.map(|partition| list::IndexBucket {
store: &view,
partition,
})
.collect::<Vec<_>>();
list::page(
&buckets,
&op.repo,
&scan,
last.as_deref(),
limit,
list::MAX_RESPONSE_BYTES,
)
.await
};
let mut page = result.map_err(|err| {
tracing::warn!(detail = %err, "ref listing scan failed");
ServerError::unavailable("ref listing unavailable")
})?;
for entry in &mut page.refs {
entry.name = strip_listed_prefix(&entry.name, prefix)
.ok_or_else(|| ServerError::unavailable("ref listing unavailable"))?
.to_owned();
}
Ok(page)
})
.await
}
#[cfg(feature = "published-view")]
#[must_use]
pub fn with_published_source(mut self, source: Arc<dyn published::PublishedSource>) -> Self {
self.published = Some(source);
self
}
pub fn with_publication_policy(
mut self,
policy: Arc<dyn clearance::PublicationPolicy>,
) -> Result<Self, ServerError> {
if self.cfg.indexed.is_none() || self.cfg.write_policy == WritePolicy::Open {
return Err(ServerError::invalid_argument(
"inspection requires indexed mode and restricted writes",
));
}
self.publication_policy = Some(policy);
Ok(self)
}
pub async fn read_ref(
&self,
a: &Authenticated,
name: &str,
) -> Result<Option<Hash>, ServerError> {
let kind = OpKind::ReadRef {
name: name.to_owned(),
};
self.observe(a, async {
check_ref_name(name)?;
let op = self.identify(a, kind)?;
let authorized = self.authorize_read(&op).await?;
#[cfg(feature = "published-view")]
if let Some(source) = &self.published {
if source.inspection_configured()
&& !source.uses_published_values()
&& authorized.facts.caller_view != CallerView::Writer
{
return Err(ServerError::unavailable("published view unavailable"));
}
if self.visibility_gates_reads()
&& op.auth.is_none()
&& matches!(op.principal, Principal::Anonymous)
&& source.read_ref_enabled()
&& (self.publication_policy.is_none() || source.uses_published_values())
{
let partition = self.shards.ref_index(&op.repo, name);
if let Some(rows) = source
.bucket(&op.repo, &partition, ms(self.clock.now_ms()))
.await
.map_err(meta_error)?
{
if self.meta.capabilities().atomic_multi_key
&& !source.uses_published_values()
{
return Err(ServerError::unavailable("published view unavailable"));
}
return Ok(rows.into_iter().find(|(n, _)| n == name).map(|(_, id)| id));
}
}
}
let p = self.shards.ref_shard(&op.repo, name);
let view = crate::store::view::ViewStore {
store: &self.meta,
repo: &op.repo,
writer: authorized.facts.caller_view == CallerView::Writer,
policy: self.publication_policy.as_deref(),
};
read::read_ref(&view, &p, &op.repo.name, name)
.await
.map_err(meta_error)
})
.await
}
pub async fn update_ref(
&self,
a: &Authenticated,
upd: RefUpdate,
) -> Result<UpdateRefResult, ServerError> {
self.update_ref_with_meta(a, upd)
.await
.map(|(result, _)| result)
}
pub async fn update_ref_with_meta(
&self,
a: &Authenticated,
upd: RefUpdate,
) -> Result<(UpdateRefResult, ResponseMeta), ServerError> {
self.observe(a, async {
check_ref_name(&upd.name)?;
if upd.new.is_none()
&& !matches!(upd.condition, mkit_core::refs::RefWriteCondition::Match(_))
{
return Err(ServerError::invalid_argument(
"delete requires MATCH and an empty new_id",
));
}
match self.write(a, OpKind::UpdateRef(upd)).await? {
(StoredResult::UpdateRef(result), meta) => Ok((result, meta)),
(other, _) => Err(stored_mismatch(&other)),
}
})
.await
}
pub async fn advance_refs(
&self,
a: &Authenticated,
head: RefUpdate,
packmap: RefUpdate,
) -> Result<AdvanceOutcome, ServerError> {
self.advance_refs_with_meta(a, head, packmap)
.await
.map(|(result, _)| result)
}
pub async fn advance_refs_with_meta(
&self,
a: &Authenticated,
head: RefUpdate,
packmap: RefUpdate,
) -> Result<(AdvanceOutcome, ResponseMeta), ServerError> {
self.advance_refs_with_tickets_with_meta(a, head, packmap, Vec::new())
.await
}
pub async fn advance_refs_with_tickets(
&self,
a: &Authenticated,
head: RefUpdate,
packmap: RefUpdate,
tickets: Vec<Hash>,
) -> Result<AdvanceOutcome, ServerError> {
self.advance_refs_with_tickets_with_meta(a, head, packmap, tickets)
.await
.map(|(result, _)| result)
}
pub async fn advance_refs_with_tickets_with_meta(
&self,
a: &Authenticated,
head: RefUpdate,
packmap: RefUpdate,
tickets: Vec<Hash>,
) -> Result<(AdvanceOutcome, ResponseMeta), ServerError> {
self.observe(a, async {
check_ref_name(&head.name)?;
check_ref_name(&packmap.name)?;
if head.new.is_none() != packmap.new.is_none()
|| (head.new.is_none()
&& (!matches!(head.condition, mkit_core::refs::RefWriteCondition::Match(_))
|| !matches!(
packmap.condition,
mkit_core::refs::RefWriteCondition::Match(_)
)))
{
return Err(ServerError::invalid_argument(
"delete requires MATCH and an empty new_id",
));
}
if head.new.is_none() && !tickets.is_empty() {
return Err(ServerError::invalid_argument("delete consumes no tickets"));
}
if tickets.len() > crate::store::outbox::MAX_TICKETS_PER_ADVANCE {
return Err(ServerError::invalid_argument(
"too many tickets in one advance",
));
}
let mut distinct = std::collections::BTreeSet::new();
if tickets.iter().any(|id| !distinct.insert(id)) {
return Err(ServerError::invalid_argument("duplicate ticket id"));
}
if !tickets.is_empty()
|| self.cfg.sharding == Sharding::D34
|| self.cfg.ref_policy.is_some()
|| self.cfg.indexed.is_some()
{
let head_branch = head.name.strip_prefix("refs/heads/");
let packmap_branch = packmap
.name
.strip_prefix(mkit_core::refs::PACKMAP_REF_PREFIX);
if head_branch.is_none() || head_branch != packmap_branch {
if a.write_grant.is_some()
&& self.cfg.grants.is_some()
&& matches!(self.cfg.addressing, Addressing::Multi(_))
{
return Err(ServerError::permission_denied(
"write grant rejected: ref scope",
));
}
return Err(ServerError::invalid_argument(if tickets.is_empty() {
"AdvanceRefs pairs refs/heads/<x> with refs/mkit/packmap/<x> on this server"
} else {
"ticketed advance requires a branch head and its packmap"
}));
}
}
match self
.write(
a,
OpKind::AdvanceRefs {
head,
packmap,
tickets,
},
)
.await?
{
(StoredResult::AdvanceRefs(outcome), meta) => Ok((outcome, meta)),
(other, _) => Err(stored_mismatch(&other)),
}
})
.await
}
pub async fn pack_exists(&self, a: &Authenticated, key: PackKey) -> Result<bool, ServerError> {
self.observe(a, async {
let op = self.identify(a, OpKind::PackExists { key })?;
let authorized = self.authorize_read(&op).await?;
if !self
.pack_is_member(a, &key, authorized.facts.caller_view)
.await?
{
return Ok(false);
}
let head = self.blobs.head(&key.into()).await;
Ok(head
.map_err(|e| store_error(StorageOp::BlobHead, e))?
.is_some())
})
.await
}
pub async fn issue_object_url(
&self,
a: &Authenticated,
target: UrlTarget,
ttl_seconds: u32,
) -> Result<MintedToken, ServerError> {
self.observe(a, async {
if self.cfg.url_tokens.is_none() {
return Err(ServerError::unimplemented("URL tokens not configured"));
}
let op = self.identify(
a,
OpKind::IssueObjectUrl {
target: target.clone(),
ttl_seconds,
},
)?;
self.issue_url(
&op,
&a.repo().identity,
&target,
ttl_seconds,
a.business_now_ms,
)
.await
})
.await
}
async fn issue_url(
&self,
op: &Operation,
repository: &str,
target: &UrlTarget,
ttl_seconds: u32,
now_ms: i64,
) -> Result<MintedToken, ServerError> {
let Some(tokens) = &self.cfg.url_tokens else {
return Err(ServerError::unimplemented("URL tokens not configured"));
};
let read = self.authorize_read(op).await?;
let epoch = match read.epoch {
Some(epoch) => epoch,
None => self.stored_grant_epoch(&op.repo.namespace).await?,
};
let AuthMode::AuthV2(auth) = &self.cfg.auth else {
return Err(internal("URL tokens without auth v2"));
};
tokens.mint(
auth.audience(),
repository,
target,
epoch,
now_ms,
ttl_seconds,
)
}
pub async fn set_repo_visibility(
&self,
a: &Authenticated,
req: VisibilityRequest,
) -> Result<(), ServerError> {
self.observe(a, async {
if !self.visibility_applies() {
return Err(ServerError::failed_precondition(
"repository visibility is not supported by this deployment",
));
}
match (&req, a.auth.is_some()) {
(VisibilityRequest::Statement(_), true) => {
return Err(ServerError::invalid_argument(
"signed_statement is not allowed on a signed request",
));
}
(VisibilityRequest::Envelope(_), false) => {
return Err(ServerError::unauthenticated(
"visibility requires auth v2 authorization",
));
}
_ => {}
}
let repo = &a.repo().repo;
let namespace = Namespace::parse(repo.namespace.as_str())
.map_err(|_| internal("invalid resolved Multi namespace"))?;
if let Addressing::Multi(multi) = &self.cfg.addressing
&& let NamespacePolicy::Allowlist(allowed) = &multi.namespace_policy
&& !allowed.contains(&namespace)
{
return Err(ServerError::permission_denied("namespace not served"));
}
let p = self.shards.coordinator(&repo.namespace);
let _gate = match &self.gate {
Some(gate) => Some(gate.enter(&p).await),
None => None,
};
match req {
VisibilityRequest::Envelope(visibility) => {
self.visibility_envelope(a, repo, &p, visibility).await
}
VisibilityRequest::Statement(statement) => {
self.visibility_statement(a, repo, &p, &statement).await
}
}
})
.await
}
#[allow(clippy::too_many_lines)] async fn visibility_envelope(
&self,
a: &Authenticated,
repo: &RepoId,
p: &Partition,
visibility: Visibility,
) -> Result<(), ServerError> {
if a.write_grant.is_some() {
return Err(ServerError::permission_denied(
"a grant never authorizes SetRepoVisibility",
));
}
let op = self.identify(a, OpKind::SetRepoVisibility { visibility })?;
let auth = a.auth.as_ref().ok_or_else(|| {
ServerError::unauthenticated("visibility requires auth v2 authorization")
})?;
let replay_key = keys::replay(&auth.replay_scope);
let rv_key = keys::repo_visibility(&repo.name);
let rows = self
.meta
.get_many(
p,
&[
replay_key.clone(),
rv_key.clone(),
keys::authority_generation(),
keys::lease_recovery(),
],
)
.await
.map_err(meta_error)?;
let mut rows = rows.into_iter();
let record = rows
.next()
.flatten()
.map(|v| codec::decode_replay_record(&v))
.transpose()
.map_err(meta_error)?;
let mut stored = rows.next().flatten();
match classify(record.as_ref(), &auth.fingerprint) {
ReplayDecision::New => {}
ReplayDecision::Return(StoredResult::RepoVisibility) => return Ok(()),
ReplayDecision::Return(other) => return Err(stored_mismatch(&other)),
ReplayDecision::FingerprintMismatch => {
return Err(ServerError::invalid_argument(
"nonce reused for a different operation",
));
}
ReplayDecision::Resume | ReplayDecision::RetryLater => {
return Err(ServerError::aborted_retryable(
"operation already in flight; retry",
));
}
}
Box::pin(self.ensure_authority_activation(&repo.namespace)).await?;
let facts = self.authorize_visibility_envelope(&op).await?;
let mut replans = 0;
loop {
let (mut batch, prune) = self
.plan_visibility(p, auth, repo, visibility, stored.as_ref())
.await?;
let fence_rows = self
.meta
.get_many(p, &[keys::authority_generation(), keys::lease_recovery()])
.await
.map_err(meta_error)?;
let [generation_value, mode_value] = fence_rows.as_slice() else {
return Err(internal("visibility fence row count"));
};
let mode = mode_value
.as_ref()
.map(codec::decode_lease_recovery)
.transpose()
.map_err(meta_error)?;
if facts.authority_generation.is_none()
&& (generation_value.is_some()
|| mode.is_some_and(|m| m.authority_fence == Some(true)))
{
return Err(ServerError::unavailable(
"persisted authority fence requires enabled executor",
));
}
batch = batch.require(lease::observed_guard(
keys::lease_recovery(),
mode_value.as_ref(),
));
if let Some(generation) = facts.authority_generation {
let current = generation_value
.as_ref()
.map(codec::decode_u64)
.transpose()
.map_err(meta_error)?
.unwrap_or(0);
if current != generation {
return Err(crate::authority::moved());
}
batch = batch.require(lease::observed_guard(
keys::authority_generation(),
generation_value.as_ref(),
));
}
match self.meta.apply(p, batch).await {
Ok(BatchOutcome::Committed) => {
self.invalidate_local_cache(repo).await;
return Ok(());
}
Ok(BatchOutcome::DeadlinePassed { .. }) => {
return Err(ServerError::unavailable("commit deadline passed; retry"));
}
Ok(BatchOutcome::PreconditionFailed { index, .. }) => {
if index == 1 {
let value = self.meta.get(p, &replay_key).await.map_err(meta_error)?;
let record = value
.as_ref()
.map(codec::decode_replay_record)
.transpose()
.map_err(meta_error)?;
return match classify(record.as_ref(), &auth.fingerprint) {
ReplayDecision::Return(StoredResult::RepoVisibility) => Ok(()),
ReplayDecision::Return(other) => Err(stored_mismatch(&other)),
ReplayDecision::FingerprintMismatch => {
Err(ServerError::invalid_argument(
"nonce reused for a different operation",
))
}
_ => Err(ServerError::aborted_retryable(
"operation already in flight; retry",
)),
};
}
replans += 1;
if replans > MAX_REPLAN {
return Err(ServerError::aborted_retryable("write contention; retry"));
}
stored = self.meta.get(p, &rv_key).await.map_err(meta_error)?;
}
Err(StoreError::Full) => {
return Err(self.partition_full(p, prune).await);
}
Err(e) => return Err(meta_error(e)),
}
}
}
async fn authorize_visibility_envelope(
&self,
op: &Operation,
) -> Result<AuthzFacts, ServerError> {
let owner = matches!(
Namespace::parse(op.repo.namespace.as_str()),
Ok(Namespace::Ed25519(key)) if op.principal.ed25519() == Some(&key)
);
if self.cfg.authorizer_role == AuthorizerRole::Check && !owner {
return Err(ServerError::permission_denied(
"SetRepoVisibility not permitted",
));
}
let mut authorized = op.clone();
authorized.authz = AuthzFacts {
authority_generation: None,
grant: None,
owner,
caller_view: if owner {
CallerView::Writer
} else {
CallerView::Reader
},
};
let returned = self
.hooks
.authorizer()
.authorize(&authorized)
.await
.map_err(ServerError::strip_admission_shape)?;
self.merge_authority_facts(&mut authorized.authz, &returned)?;
Ok(authorized.authz)
}
async fn plan_visibility(
&self,
p: &Partition,
auth: &VerifiedAuth,
repo: &RepoId,
visibility: Visibility,
stored: Option<&Value>,
) -> Result<(Batch, Option<Batch>), ServerError> {
let now = ms(self.clock.now_ms());
let window = u64::try_from(self.cfg.max_apply_window.as_millis()).unwrap_or(u64::MAX);
let deadline = now
.saturating_add(window)
.min(ms(auth.expires_at_ms).saturating_add(MAX_CLOCK_LEAD_MS.unsigned_abs()));
let replay_key = keys::replay(&auth.replay_scope);
let rv_key = keys::repo_visibility(&repo.name);
let row = stored
.map(codec::decode_repo_visibility)
.transpose()
.map_err(meta_error)?;
let mut batch = Batch::new()
.require(Precondition::NotAfter(deadline))
.require(Precondition::Absent(replay_key.clone()))
.require(match stored {
Some(value) => Precondition::Equals(rv_key.clone(), value.clone()),
None => Precondition::Absent(rv_key.clone()),
})
.put(
rv_key.clone(),
codec::encode_repo_visibility(&codec::RepoVisibilityV1 {
visibility: stored_visibility(visibility),
last_created_ms: row.as_ref().map_or(0, |r| r.last_created_ms).max(now),
last_statement_id: row.and_then(|r| r.last_statement_id),
changed_ms: Some(now),
}),
)
.put(
replay_key.clone(),
codec::encode_replay_record(&ReplayRecord {
fingerprint: auth.fingerprint,
expires_at_ms: auth.expires_at_ms,
state: ReplayState::Committed(StoredResult::RepoVisibility),
}),
)
.put(
keys::replay_expiry(ms(auth.expires_at_ms), &auth.replay_scope),
Value::default(),
);
self.plan_listing_visibility(p, repo, &mut batch).await?;
let expired = read::expired_replay_keys(&self.meta, p, now, 32)
.await
.map_err(meta_error)?;
let purge = self
.plan_repository_purge(
p,
repo,
crate::purge::Trigger::VisibilityChange,
&mkit_core::hash::to_hex(&auth.replay_scope),
now,
)
.await?;
batch.preconditions.extend(purge.preconditions);
batch.writes.extend(purge.writes);
let mut prune = Batch::new().require(Precondition::NotAfter(deadline));
for (index, target) in &expired {
if batch.preconditions.len() + batch.writes.len() + 2 > MAX_BATCH_OPS {
break;
}
batch = batch.delete(index.clone()).delete(target.clone());
prune = prune.delete(index.clone()).delete(target.clone());
}
Ok((batch, (!prune.writes.is_empty()).then_some(prune)))
}
async fn visibility_statement(
&self,
a: &Authenticated,
repo: &RepoId,
p: &Partition,
statement: &str,
) -> Result<(), ServerError> {
let rejected = |e: mkit_attest::grant::GrantError| {
ServerError::permission_denied(format!("visibility statement rejected: {}", e.reason()))
};
if statement.len() > mkit_attest::grant::MAX_GRANT_HEADER_BYTES {
return Err(ServerError::permission_denied(
"visibility statement rejected: too long",
));
}
let grants = self.cfg.grants.as_ref().ok_or_else(|| {
ServerError::permission_denied(
"visibility statement rejected: no owner schemes configured",
)
})?;
let identity = RepositoryIdentity::parse(&a.repo().identity)
.map_err(|_| internal("invalid resolved repository identity"))?;
let verified =
verify_visibility_statement(grants.verifier(), statement, &identity, a.business_now_ms)
.map_err(rejected)?;
if let Some(namespace) = identity.namespace() {
self.require_client_owner_key(namespace)?;
}
let created = u64::try_from(verified.statement().created_ms)
.map_err(|_| rejected(mkit_attest::grant::GrantError::DecimalOutOfRange))?;
let id = to_hex(verified.id());
let rv_key = keys::repo_visibility(&repo.name);
let window = u64::try_from(self.cfg.max_apply_window.as_millis()).unwrap_or(u64::MAX);
let mut replans = 0;
loop {
let stored = self.meta.get(p, &rv_key).await.map_err(meta_error)?;
let row = stored
.as_ref()
.map(codec::decode_repo_visibility)
.transpose()
.map_err(meta_error)?;
match &row {
Some(r)
if r.last_created_ms == created
&& r.last_statement_id.as_deref() == Some(id.as_str())
&& r.visibility == stored_visibility(verified.statement().visibility) =>
{
return Ok(());
}
Some(r) if created <= r.last_created_ms => {
return Err(ServerError::permission_denied(
"visibility statement rejected: not newer than the stored statement",
));
}
_ => {}
}
let deadline = ms(self.clock.now_ms()).saturating_add(window);
let rv_guard = match &stored {
Some(value) => Precondition::Equals(rv_key.clone(), value.clone()),
None => Precondition::Absent(rv_key.clone()),
};
let mut batch = Batch::new()
.require(Precondition::NotAfter(deadline))
.require(rv_guard)
.put(
rv_key.clone(),
codec::encode_repo_visibility(&codec::RepoVisibilityV1 {
visibility: stored_visibility(verified.statement().visibility),
last_created_ms: created,
last_statement_id: Some(id.clone()),
changed_ms: Some(ms(self.clock.now_ms())),
}),
);
self.plan_listing_visibility(p, repo, &mut batch).await?;
let purge = self
.plan_repository_purge(
p,
repo,
crate::purge::Trigger::VisibilityChange,
&id,
ms(self.clock.now_ms()),
)
.await?;
batch.preconditions.extend(purge.preconditions);
batch.writes.extend(purge.writes);
match self.meta.apply(p, batch).await {
Ok(BatchOutcome::Committed) => {
self.invalidate_local_cache(repo).await;
return Ok(());
}
Ok(BatchOutcome::DeadlinePassed { .. }) => {
return Err(ServerError::unavailable("commit deadline passed; retry"));
}
Ok(BatchOutcome::PreconditionFailed { .. }) => {
replans += 1;
if replans > MAX_REPLAN {
return Err(ServerError::aborted_retryable("write contention; retry"));
}
}
Err(StoreError::Full) => return Err(self.partition_full(p, None).await),
Err(e) => return Err(meta_error(e)),
}
}
}
pub async fn open_upload(
&self,
a: &Authenticated,
pack_id: Option<&[u8]>,
total_bytes: Option<u64>,
) -> Result<UploadSession<'_, B, N, H>, ServerError> {
UploadSession::begin(self, a, pack_id, total_bytes).await
}
pub async fn open_ticketed_upload(
&self,
a: &Authenticated,
pack_id: Option<&[u8]>,
total_bytes: Option<u64>,
token: &[u8],
) -> Result<UploadSession<'_, B, N, H>, ServerError> {
UploadSession::begin_ticketed(self, a, pack_id, total_bytes, token).await
}
pub async fn download(
&self,
a: &Authenticated,
key: PackKey,
) -> Result<DownloadStream, ServerError> {
let mut outcome = self.outcome(a);
let opened = async {
let op = self.identify(a, OpKind::DownloadPack { key })?;
let authorized = self.authorize_read(&op).await?;
if !self
.pack_is_member(a, &key, authorized.facts.caller_view)
.await?
{
return Err(ServerError::not_found("pack not found"));
}
let body = self.blobs.get(&key.into(), None).await;
match body.map_err(|e| store_error(StorageOp::BlobGet, e))? {
Some(body) => Ok(body),
None => Err(ServerError::not_found("pack not found")),
}
}
.instrument(outcome.span.clone())
.await;
match opened {
Ok(body) => {
let max = self.cfg.download_chunk_max;
Ok(DownloadStream::new(body, max, Some(outcome)))
}
Err(err) => {
outcome.record(Err(&err));
Err(err)
}
}
}
pub async fn health(&self) -> HealthStatus {
HealthStatus {
blobs: self.blobs.probe().await.is_ok(),
meta: self.meta.probe().await.is_ok(),
}
}
#[must_use]
pub fn auth_mode(&self) -> &AuthMode {
&self.cfg.auth
}
#[cfg(feature = "ssh")]
pub(crate) fn upload_limits(&self) -> UploadLimits {
self.cfg.upload_limits
}
#[cfg(all(test, feature = "ssh"))]
pub(crate) fn meta_store(&self) -> &N {
&self.meta
}
pub fn capabilities(&self) -> PipelineCapabilities {
PipelineCapabilities {
atomic_advance: self.meta.capabilities().atomic_multi_key,
}
}
#[must_use]
pub fn with_header(&self, err: ServerError, name: &str, value: &str) -> ServerError {
match err.clone().try_with_header(name, value) {
Ok(err) => err,
Err(reason) => {
let label = match reason {
InvalidHeader::Name => "name",
InvalidHeader::Reserved => "reserved",
InvalidHeader::Value => "value",
};
tracing::warn!(header = ?name, %reason, "dropped an error response header");
self.metrics
.incr(METRIC_HEADER_DROPPED, &[("reason", label)], 1);
err
}
}
}
async fn observe<T>(
&self,
a: &Authenticated,
fut: impl Future<Output = Result<T, ServerError>>,
) -> Result<T, ServerError> {
let mut outcome = self.outcome(a);
let result = fut.instrument(outcome.span.clone()).await;
outcome.record(result.as_ref().map(|_| ()));
result
}
fn outcome(&self, a: &Authenticated) -> RequestOutcome {
self.outcome_for(a.procedure(), a.principal.kind(), &a.repo().identity)
}
fn outcome_for(
&self,
procedure: Procedure,
principal: &'static str,
repo: &str,
) -> RequestOutcome {
let procedure = method(procedure);
let span = tracing::info_span!("mkit.server.rpc", procedure, repo, principal);
let (metrics, clock) = (self.metrics.clone(), self.clock.clone());
RequestOutcome::new(span, procedure, metrics, clock, self.cfg.redactor.clone())
}
fn write<'a>(
&'a self,
a: &'a Authenticated,
kind: OpKind,
) -> crate::rt::BoxFuture<'a, Result<(StoredResult, ResponseMeta), ServerError>> {
self.write_with(a, kind, None)
}
fn write_with<'a>(
&'a self,
a: &'a Authenticated,
kind: OpKind,
implicit: Option<&'a [PendingPack]>,
) -> crate::rt::BoxFuture<'a, Result<(StoredResult, ResponseMeta), ServerError>> {
Box::pin(self.write_inner(a, kind, implicit))
}
#[allow(clippy::too_many_lines)] async fn write_inner(
&self,
a: &Authenticated,
kind: OpKind,
implicit: Option<&[PendingPack]>,
) -> Result<(StoredResult, ResponseMeta), ServerError> {
let mut op = self.identify(a, kind)?;
fault!(self, AfterAuthenticate, &op, a);
let (kind, mut refs, p) = self.ref_writes(&op)?;
let mut ahead = self.read_ahead(&op, &p, &refs, a.business_skew_ms).await?;
if let Some(stored) = Self::replay_lookup(&op, ahead.as_ref())? {
return Ok((stored, ResponseMeta::default()));
}
let lease = if self.cfg.sharding == Sharding::D34 {
let observed = self.observe_lease(&op, &p, ahead.as_ref()).await?;
if let Some((window, total)) = observed.quota_seed()
&& let Some(snapshot) = ahead.as_mut()
{
snapshot.namespace_seed = Some((
window,
codec::encode_namespace_view(NamespaceView {
total,
pushed: crate::quota::NamespaceUsage::default(),
observed_at_ms: ms(self.clock.now_ms()),
}),
));
}
op.creation = observed.creation(&self.cfg.addressing);
op.leased_epoch = Some(observed.epoch());
Some(observed)
} else {
op.creation = self.creation_facts(&op, ahead.as_ref()).await?;
None
};
if self.cfg.sharding == Sharding::Single {
let key = keys::grant_epoch();
if let Some(snapshot) = ahead.as_ref().filter(|snapshot| snapshot.contains(&key)) {
op.observed_epoch = Some(
snapshot
.get(&key)
.map(codec::decode_u64)
.transpose()
.map_err(meta_error)?
.unwrap_or(0),
);
}
}
let (authz, fast_forward) = self.authorize(&op).await?;
op.authz = authz;
fault!(self, AfterAuthorize, &op, a);
if let Err(error) = self.check_ref_signers(&op) {
return self
.store_policy_denial(&op, a, &p, ahead, error, false)
.await
.map(|stored| (stored, ResponseMeta::default()));
}
let ticketed =
matches!(&op.kind, OpKind::AdvanceRefs { tickets, .. } if !tickets.is_empty());
let mut staged = crate::indexed::verify::StagedCommits::default();
let mut ticket_ms = None;
if ticketed {
let snap = ahead
.as_mut()
.ok_or_else(|| internal("ticket advance requires atomic metadata"))?;
if let Some(stored) = self.ticket_decision(&op, a, &p, snap).await? {
return Ok((stored, ResponseMeta::default()));
}
if let Some(indexed) = self.cfg.indexed
&& let OpKind::AdvanceRefs { head, tickets, .. } = &op.kind
{
#[cfg(feature = "test-faults")]
if a.test_directives().fault.as_deref() == Some("indexed-pending") {
return Err(crate::indexed::pending(5_000));
}
let rows = tickets
.iter()
.map(|id| {
snap.get(&keys::ticket(id))
.ok_or_else(|| {
ServerError::failed_precondition("invalid or expired upload ticket")
})
.and_then(|raw| codec::decode_ticket(raw).map_err(meta_error))
})
.collect::<Result<Vec<_>, _>>()?;
let tip = head
.new
.ok_or_else(|| ServerError::invalid_argument("delete consumes no tickets"))?;
ticket_ms = rows.iter().map(|ticket| ticket.created_at_ms).min();
staged = if let Some(limit) = self.inspection_max_objects() {
if indexed.verification == crate::indexed::VerificationMode::Scheduled {
crate::indexed::scheduled::check_inspected(
&self.blobs,
&self.meta,
self.shards.as_ref(),
&op.repo,
&p,
&rows,
tickets,
tip,
indexed,
self.clock.as_ref(),
self.metrics.as_ref(),
limit as usize,
)
.await?
} else {
crate::indexed::verify::verify_ticketed_inspected(
&self.blobs,
&self.meta,
self.shards.as_ref(),
&op.repo,
&p,
&rows,
tickets,
tip,
indexed,
self.clock.as_ref(),
self.metrics.as_ref(),
limit as usize,
)
.await?
}
} else if indexed.verification == crate::indexed::VerificationMode::Scheduled {
crate::indexed::scheduled::check(
&self.blobs,
&self.meta,
self.shards.as_ref(),
&op.repo,
&p,
&rows,
tickets,
tip,
indexed,
self.clock.as_ref(),
self.metrics.as_ref(),
)
.await?
} else {
crate::indexed::verify::verify_ticketed(
&self.blobs,
&self.meta,
self.shards.as_ref(),
&op.repo,
&p,
&rows,
tickets,
tip,
indexed,
self.clock.as_ref(),
self.metrics.as_ref(),
)
.await?
};
}
} else if let Err(error) = self.check_ticketless_head(&op).await {
return self
.store_policy_denial(&op, a, &p, ahead, error, false)
.await
.map(|stored| (stored, ResponseMeta::default()));
}
if staged.inspection.is_none() {
staged.inspection = self
.inspection_max_objects()
.map(|limit| crate::indexed::inspection::InspectionSet::new(limit as usize));
}
if let Some(pending) = implicit {
let upd = refs
.first_mut()
.ok_or_else(|| internal("implicit consumption needs an UpdateRef"))?;
self.check_implicit_packmap(&op, pending, ahead.as_ref(), upd)
.await?;
}
let existing = self.begin_decision(&op, a, ahead.as_mut()).await?;
if existing.is_none() && !ticketed {
self.check_outbox_backpressure(&p, ahead.as_ref()).await?;
}
let allowance = if existing.is_some() || ticketed || implicit.is_some_and(|p| !p.is_empty())
{
Allowance::default()
} else {
let credentials = admission::validate_credentials(&a.credential_capture)?;
let mut input = AdmissionInput::new(&op);
input.credential_headers = &credentials;
if let OpKind::BeginUpload { key, bytes, .. } = &op.kind {
input.declared_bytes = *bytes;
input.pack_id = Some(*key);
input.new_to_repo_bytes = Some(*bytes);
}
self.admit(input).await?
};
let pending = match allowance.reservation.as_deref() {
Some(rid) => Some(self.record_pending(a, &p, rid).await?),
None => None,
};
let write_result = async {
let mut begin = self.begin_write(&op, a, existing, allowance.reservation.clone())?;
self.precheck_namespace(&p, &allowance.charges, &mut ahead, a.business_skew_ms)?;
let mut opened_session = None;
if let Some(BeginWrite::Open(open)) = &mut begin
&& open.spec.bytes > open.spec.part_size
{
let key = PackKey(open.spec.pack_id).into();
let ticket_id = crate::store::tickets::ticket_id(&open.spec.reservation_id);
let session = self
.blobs
.begin_multipart_for_ticket(key, open.spec.bytes, open.spec.part_size, ticket_id)
.await
.map_err(|e| {
if open.reserved() {
tracing::warn!(error = %e, "reserved multipart session creation failed");
ServerError::unavailable("multipart session creation failed; retry")
} else {
store_error(StorageOp::MultipartSession, e)
}
})?;
if session.is_empty() || session.len() > u16::MAX as usize {
if let Err(err) = self.blobs.abort(key, &session).await {
tracing::warn!(error = %err, "failed to abort invalid multipart session");
}
return Err(ServerError::internal(
"object storage request failed",
"multipart store returned an invalid session identifier",
));
}
opened_session = Some((key, session.clone(), ticket_id));
open.spec.upload_session = Some(session);
}
let write_result = async {
let lease = if let Some(observed) = lease {
let (created, lease) = {
let _grant_gate = match (&self.gate, &observed) {
(Some(gate), lease::LeaseObservation::Renew(_)) => {
Some(gate.enter(&p).await)
}
_ => None,
};
self.admit_lease(&op, &p, observed, a.business_skew_ms)
.await?
};
op.created = created;
op.leased_epoch = Some(lease.value.epoch);
if lease.install {
fault!(self, AfterLeaseGrant, &op, a);
}
Some(lease)
} else {
op.created = self.commit_creation(&op, a.business_skew_ms).await?;
None
};
if let Err(error) = self
.check_fast_forward(
&op,
&mut refs,
ahead.as_ref(),
fast_forward.as_ref(),
(&staged, ticket_ms),
)
.await
{
return self
.store_policy_denial(&op, a, &p, ahead, error, pending.is_some())
.await;
}
self.pre_receive(&op).await?;
let write = (kind, refs.as_slice(), allowance.charges.as_slice());
self.plan_and_apply(
&op,
a,
&p,
write,
ahead,
(lease, begin.as_ref()),
WriteInputs {
denial_ids: &staged.denial_ids,
denial_packs: &staged.denial_packs,
pending: pending.as_ref(),
implicit,
external_bases: &staged.external_bases,
inspected: staged.inspection.as_mut(),
},
)
.await
}
.await;
if let Some((key, session, fresh_id)) = opened_session {
self.cleanup_opened_session(&p, &write_result, key, &session, fresh_id)
.await;
}
write_result
}
.await;
if let (Some(pending), Err(err)) = (&pending, &write_result) {
let (reason, detail) = reservation::abort_reason(err);
self.resolve_pending(&p, pending, reason, detail).await;
}
write_result.map(|result| {
let committed = matches!(
&result,
StoredResult::UpdateRef(UpdateRefResult::Committed)
| StoredResult::AdvanceRefs(AdvanceOutcome::Committed)
| StoredResult::BeginUpload(BeginUploadResult::Ticket { .. })
);
let meta = if committed
&& (!allowance.response_headers.is_empty() || allowance.external_ref.is_some())
{
let mut headers = allowance.response_headers;
if !headers.is_empty() {
headers.push(("Cache-Control".into(), "private".into()));
}
ResponseMeta {
headers,
external_ref: allowance.external_ref,
}
} else {
ResponseMeta::default()
};
(result, meta)
})
}
async fn cleanup_opened_session(
&self,
partition: &Partition,
result: &Result<StoredResult, ServerError>,
key: crate::store::BlobKey,
session: &[u8],
fresh_id: Hash,
) {
if session == fresh_id {
return;
}
let committed_fresh = match result {
Ok(StoredResult::BeginUpload(BeginUploadResult::Ticket { id, token, .. }))
if *id == fresh_id =>
{
self.cfg
.ticket_keys
.as_ref()
.and_then(|keys| keys.verify(token, 0).ok())
.is_some_and(|claims| claims.upload_session == session)
}
_ => false,
};
let stored_fresh = if result.is_err() {
match self.meta.get(partition, &keys::ticket(&fresh_id)).await {
Ok(Some(raw)) => codec::decode_ticket(&raw)
.ok()
.is_none_or(|ticket| ticket.upload_session.as_deref() == Some(session)),
Ok(None) => false,
Err(err) => {
tracing::warn!(error = %err, "could not confirm multipart ticket after failed write");
true
}
}
} else {
false
};
if !committed_fresh
&& !stored_fresh
&& let Err(err) = self.blobs.abort(key, session).await
{
tracing::warn!(error = %err, "failed to abort unused multipart session");
}
}
fn identify(&self, a: &Authenticated, kind: OpKind) -> Result<Operation, ServerError> {
if a.procedure() != kind.procedure() {
return Err(ServerError::unauthenticated(
"credentials were checked for another procedure",
));
}
if matches!(self.cfg.auth, AuthMode::AuthV2(_))
&& kind.procedure().is_write()
&& kind.procedure() != Procedure::SetRepoVisibility
&& a.auth.is_none()
{
return Err(ServerError::unauthenticated(
"missing auth v2 authorization",
));
}
let repo = a.repo().repo.clone();
let principal = a.principal.clone();
let mut op = Operation::new(repo, principal, a.auth.clone(), kind);
op.write_grant.clone_from(&a.write_grant);
op.business_now_ms = Some(a.business_now_ms);
Ok(op)
}
async fn require_repository(&self, repo: &crate::repo::RepoId) -> Result<(), ServerError> {
if matches!(self.cfg.addressing, Addressing::Multi(_)) {
let p = self.shards.coordinator(&repo.namespace);
let value = self
.meta
.get(&p, &keys::repo_record(&repo.name))
.await
.map_err(meta_error)?;
match value {
Some(value) => {
codec::decode_repo_record(&value).map_err(meta_error)?;
}
None => return Err(ServerError::repository_not_found()),
}
}
Ok(())
}
fn visibility_applies(&self) -> bool {
matches!(self.cfg.addressing, Addressing::Multi(_))
&& self.cfg.write_policy == WritePolicy::Owner
&& matches!(self.cfg.auth, AuthMode::AuthV2(_))
}
fn visibility_gates_reads(&self) -> bool {
matches!(self.cfg.addressing, Addressing::Multi(_))
&& self.cfg.write_policy == WritePolicy::Owner
}
async fn authorize_read(&self, op: &Operation) -> Result<ReadAuth, ServerError> {
tracing::debug!(stage = "authorize");
if !self.visibility_gates_reads() {
return self.authorize_read_ungated(op).await;
}
let admin_owner = Namespace::parse(op.repo.namespace.as_str())
.is_ok_and(|namespace| self.owner_key_is_admin(&namespace));
let check = op
.write_grant
.as_ref()
.filter(|_| !admin_owner)
.and_then(|header| {
self.cfg
.grants
.as_ref()
.and_then(|g| read_policy::check_grant(g, header.expose(), op))
});
let signed = op.auth.is_some();
let owner = op.write_grant.is_none()
&& matches!(Namespace::parse(op.repo.namespace.as_str()),
Ok(Namespace::Ed25519(key)) if op.principal.ed25519() == Some(&key));
let Some((private, epoch)) = self.read_repo_state(&op.repo).await? else {
if signed && op.write_grant.is_none() {
let provisional = AuthzFacts {
authority_generation: None,
grant: None,
owner,
caller_view: if owner {
CallerView::Writer
} else {
CallerView::Reader
},
};
let _ = self.read_hook(op, true, true, &provisional).await;
}
return Err(ServerError::repository_not_found());
};
let grant = check.filter(|c| c.epoch == epoch);
let grant_ref = grant.map(|c| GrantRef {
id: c.id,
epoch,
presence_requirement: None,
});
let authority = self.cfg.authorizer_role == AuthorizerRole::Authority;
if private && !signed {
return Err(ServerError::repository_not_found());
}
let provisional = AuthzFacts {
authority_generation: None,
grant: grant_ref.clone(),
owner,
caller_view: if owner || grant.is_some_and(|g| g.write) {
CallerView::Writer
} else {
CallerView::Reader
},
};
let hook = self.read_hook(op, private, signed, &provisional).await?;
let caller = read_policy::Caller {
#[cfg(feature = "http-objects")]
http_token_authorized: false,
signed,
owner,
grant: grant.map(|g| read_policy::GrantEval {
read: g.read,
write: g.write,
}),
hook,
authority,
};
match read_policy::decide(op.procedure(), private, caller) {
read_policy::ReadDecision::NotFound => Err(ServerError::repository_not_found()),
read_policy::ReadDecision::Allow(caller_view) => Ok(ReadAuth {
facts: AuthzFacts {
authority_generation: None,
grant: grant_ref,
owner,
caller_view,
},
epoch: Some(epoch),
}),
}
}
async fn read_repo_state(
&self,
repo: &crate::repo::RepoId,
) -> Result<Option<(bool, u64)>, ServerError> {
let coordinator = self.shards.coordinator(&repo.namespace);
let rows = self
.meta
.get_many(
&coordinator,
&[
keys::repo_record(&repo.name),
keys::repo_visibility(&repo.name),
keys::grant_epoch(),
],
)
.await
.map_err(|e| {
tracing::warn!(detail = %e, "repository state read failed");
ServerError::unavailable("repository state unavailable")
})?;
let mut rows = rows.into_iter();
let Some(record) = rows.next().flatten() else {
return Ok(None);
};
codec::decode_repo_record(&record).map_err(meta_error)?;
let stored = rows
.next()
.flatten()
.map(|v| codec::decode_repo_visibility(&v))
.transpose()
.map_err(meta_error)?;
let epoch = rows
.next()
.flatten()
.map(|v| codec::decode_u64(&v))
.transpose()
.map_err(meta_error)?
.unwrap_or(0);
Ok(Some((
repo_is_private(stored.as_ref(), self.cfg.default_repo_visibility),
epoch,
)))
}
async fn read_hook(
&self,
op: &Operation,
private: bool,
signed: bool,
provisional: &AuthzFacts,
) -> Result<read_policy::HookEval, ServerError> {
let authorized = || {
let mut authorized = op.clone();
authorized.authz = provisional.clone();
authorized
};
let writer = |facts: &AuthzFacts| facts.caller_view == CallerView::Writer;
if private {
if op.write_grant.is_some() {
return Ok(read_policy::HookEval::NotConsulted);
}
return Ok(
match self.hooks.authorizer().authorize(&authorized()).await {
Ok(facts) => read_policy::HookEval::Allow {
writer_view: writer(&facts),
},
Err(_) => read_policy::HookEval::Deny,
},
);
}
if self.cfg.authorizer_role == AuthorizerRole::Authority {
if !signed || provisional.caller_view == CallerView::Writer {
return Ok(read_policy::HookEval::NotConsulted);
}
return Ok(
match self.hooks.authorizer().authorize(&authorized()).await {
Ok(facts) => read_policy::HookEval::Allow {
writer_view: writer(&facts),
},
Err(_) => read_policy::HookEval::NotConsulted,
},
);
}
self.hooks
.authorizer()
.authorize(op)
.await
.map_err(ServerError::strip_admission_shape)?;
Ok(read_policy::HookEval::NotConsulted)
}
async fn authorize_read_ungated(&self, op: &Operation) -> Result<ReadAuth, ServerError> {
let facts = self
.hooks
.authorizer()
.authorize(op)
.await
.map_err(ServerError::strip_admission_shape)?;
self.require_repository(&op.repo).await?;
let caller_view = match &op.principal {
Principal::Anonymous => CallerView::Anonymous,
Principal::BearerHolder => CallerView::Reader,
principal => {
if self.cfg.write_policy == WritePolicy::Open
|| matches!(Namespace::parse(op.repo.namespace.as_str()),
Ok(Namespace::Ed25519(key)) if principal.ed25519() == Some(&key))
{
CallerView::Writer
} else {
CallerView::Reader
}
}
};
Ok(ReadAuth {
facts: AuthzFacts {
caller_view,
..facts
},
epoch: None,
})
}
async fn pack_is_member(
&self,
a: &Authenticated,
key: &PackKey,
caller: CallerView,
) -> Result<bool, ServerError> {
if matches!(self.cfg.addressing, Addressing::Single { .. }) {
return Ok(true);
}
let view = crate::store::view::ViewStore {
store: &self.meta,
repo: &a.repo().repo,
writer: caller == CallerView::Writer,
policy: self.publication_policy.as_deref(),
};
let member = read::is_member(
&view,
self.shards.as_ref(),
&a.repo().repo,
&key.0,
a.ref_hint.as_deref(),
)
.await
.map_err(meta_error)?;
if member && self.cfg.indexed.is_some() {
match crate::takedown::denial::require_clear(&self.meta, &key.0).await {
Ok(()) => {}
Err(e) if e.public_message() == "object blocked" => return Ok(false),
Err(e) => return Err(e),
}
}
if member && self.cfg.takedown_denial && self.cfg.indexed.is_some() {
return match crate::takedown::denial::require_pack_clear(
&self.meta,
self.shards.as_ref(),
&a.repo().repo,
&key.0,
)
.await
{
Ok(()) => Ok(true),
Err(e) if e.public_message() == "object blocked" => Ok(false),
Err(e) => Err(e),
};
}
Ok(member)
}
fn ref_writes(
&self,
op: &Operation,
) -> Result<(WriteKind, Vec<RefUpdate>, Partition), ServerError> {
let (kind, refs) = match &op.kind {
OpKind::UpdateRef(u) => (WriteKind::UpdateRef, vec![u.clone()]),
OpKind::AdvanceRefs { head, packmap, .. } => {
(WriteKind::AdvanceRefs, vec![packmap.clone(), head.clone()])
}
OpKind::BeginUpload { ref_name, .. } => {
return Ok((
WriteKind::BeginUpload,
vec![],
self.shards.ref_shard(&op.repo, ref_name),
));
}
_ => return Err(internal("not a unary write")),
};
let p = self.shards.ref_shard(&op.repo, &refs[0].name);
if refs
.iter()
.any(|r| self.shards.ref_shard(&op.repo, &r.name) != p)
{
return Err(internal("shard map splits a head from its packmap"));
}
Ok((kind, refs, p))
}
async fn read_ahead(
&self,
op: &Operation,
p: &Partition,
refs: &[RefUpdate],
business_skew_ms: i64,
) -> Result<Option<Snapshot>, ServerError> {
let caps = self.meta.capabilities();
if !caps.atomic_multi_key {
return Ok(None);
}
let namespace_window = if matches!(self.cfg.addressing, Addressing::Multi(_))
&& self.hooks.admission().is_default()
{
self.cfg.write_quota.map(|limits| {
quota::namespace_window(
self.clock.now_ms().saturating_add(business_skew_ms),
limits.window_ms,
)
})
} else {
None
};
let mut wanted: Vec<Key> = refs
.iter()
.map(|r| keys::ref_key(&op.repo.name, &r.name))
.collect();
if let Some(update) = refs.first() {
let name = crate::store::publication::sequence_ref(&update.name);
wanted.push(keys::publication(&op.repo.name, &name));
wanted.push(keys::ref_key(&op.repo.name, &name));
if let Some(packmap) = mkit_attest::grant::head_packmap(&name) {
wanted.push(keys::ref_key(&op.repo.name, &packmap));
}
wanted.push(keys::outbox_sequence());
}
if self.cfg.sharding == Sharding::Single && caps.key_classes == KeyClasses::All {
wanted.extend([keys::authority_generation(), keys::lease_recovery()]);
}
if caps.implicit_layout_version.is_none() {
wanted.push(keys::layout_version());
}
wanted.push(keys::outcome_backlog());
if matches!(self.cfg.addressing, Addressing::Multi(_)) {
wanted.push(keys::repo_known(&op.repo.name));
}
if self.cfg.sharding == Sharding::D34 {
wanted.push(keys::epoch_lease());
if !refs.is_empty() {
wanted.push(keys::outbox_sequence());
}
}
if let Some(auth) = &op.auth {
wanted.push(keys::replay(&auth.replay_scope));
if self.cfg.sharding == Sharding::Single {
wanted.push(keys::grant_epoch());
}
if self.cfg.write_quota.is_some() {
let scope = QuotaScope::for_signer(&op.repo.namespace, &auth.signer);
wanted.push(keys::quota(&scope));
}
if let (Some(window), Some(limits)) = (namespace_window, self.cfg.write_quota) {
let charge = NamespaceCharge {
limits,
window,
bytes: 0,
rollup: matches!(p, Partition::Ref { .. }),
};
wanted.push(quota::counter_key(charge, window));
if charge.rollup {
wanted.push(keys::quota_view(window));
}
let next = window.saturating_add(1);
wanted.push(quota::counter_key(charge, next));
if charge.rollup {
wanted.push(keys::quota_view(next));
}
}
}
if let OpKind::BeginUpload { ref_name, key, .. } = &op.kind {
let signer = op
.auth
.as_ref()
.ok_or_else(|| internal("missing ticket signer"))?
.signer;
wanted.extend(begin::decision_keys(
&op.repo.name,
ref_name,
&key.0,
&signer,
)?);
}
let mut snap = Snapshot::default();
snap.namespace_window = namespace_window;
self.fill(p, &mut snap, wanted).await?;
Ok(Some(snap))
}
fn namespace_charge(
&self,
p: &Partition,
charges: &[QuotaCharge],
window: Option<u64>,
) -> Result<Option<NamespaceCharge>, ServerError> {
if !matches!(self.cfg.addressing, Addressing::Multi(_))
|| !self.hooks.admission().is_default()
{
return Ok(None);
}
let Some(charge) = charges.first() else {
return Ok(None);
};
let window = window.ok_or_else(|| internal("namespace quota read-ahead missing"))?;
Ok(Some(NamespaceCharge {
limits: charge.limits,
window,
bytes: charge.bytes,
rollup: matches!(p, Partition::Ref { .. }),
}))
}
fn precheck_namespace(
&self,
p: &Partition,
charges: &[QuotaCharge],
ahead: &mut Option<Snapshot>,
business_skew_ms: i64,
) -> Result<(), ServerError> {
let Some(mut charge) = self.namespace_charge(
p,
charges,
ahead
.as_ref()
.and_then(|snapshot| snapshot.namespace_window),
)?
else {
return Ok(());
};
let now = self.clock.now_ms().saturating_add(business_skew_ms);
let window = quota::namespace_window(now, charge.limits.window_ms);
if window != charge.window && window != charge.window.saturating_add(1) {
return Err(
ServerError::aborted_retryable("namespace quota window advanced; retry")
.with_abort_cause(AbortCause::QuotaWindow),
);
}
charge.window = window;
let snap = ahead
.as_mut()
.ok_or_else(|| internal("namespace quota needs atomic reads"))?;
let stored_view = charge
.rollup
.then(|| snap.get(&keys::quota_view(window)))
.flatten();
let seeded_view = snap
.namespace_seed
.as_ref()
.filter(|(seed_window, _)| *seed_window == window)
.map(|(_, value)| value);
let decision = quota::check_namespace(
snap.get("a::counter_key(charge, window)),
stored_view.or(seeded_view),
now,
charge,
)?;
match decision {
NamespaceDecision::Exhausted => {
return Err(ServerError::resource_exhausted(
"namespace write op/byte quota exceeded for this window; try again later",
));
}
NamespaceDecision::Allowed {
view: ViewStatus::Missing,
..
} => self.metrics.incr(
crate::telemetry::METRIC_NAMESPACE_QUOTA_VIEW_FALLBACK,
&[("state", "missing")],
1,
),
NamespaceDecision::Allowed {
view: ViewStatus::Stale,
..
} => self.metrics.incr(
crate::telemetry::METRIC_NAMESPACE_QUOTA_VIEW_FALLBACK,
&[("state", "stale")],
1,
),
NamespaceDecision::Allowed { .. } => {}
}
snap.namespace_window = Some(window);
Ok(())
}
async fn fill(
&self,
p: &Partition,
snap: &mut Snapshot,
mut wanted: Vec<Key>,
) -> Result<(), ServerError> {
wanted.retain(|k| !snap.contains(k));
wanted.sort();
wanted.dedup();
if wanted.is_empty() {
return Ok(());
}
let values = self.meta.get_many(p, &wanted).await.map_err(meta_error)?;
for (key, value) in wanted.into_iter().zip(values) {
snap.insert(key, value);
}
Ok(())
}
fn replay_lookup(
op: &Operation,
ahead: Option<&Snapshot>,
) -> Result<Option<StoredResult>, ServerError> {
let (Some(auth), Some(snap)) = (&op.auth, ahead) else {
return Ok(None);
};
tracing::debug!(stage = "replay_lookup");
let stored = snap.get(&keys::replay(&auth.replay_scope));
let record = stored.map(codec::decode_replay_record).transpose();
let record = record.map_err(meta_error)?;
replay_answer(classify(record.as_ref(), &auth.fingerprint))
}
fn merge_authority_facts(
&self,
built_in: &mut AuthzFacts,
returned: &AuthzFacts,
) -> Result<(), ServerError> {
if self.cfg.authorizer_role == AuthorizerRole::Authority
&& self.cfg.authority_fence.is_some()
{
built_in.authority_generation =
Some(returned.authority_generation.ok_or_else(|| {
ServerError::unavailable("Authority allowance missing generation")
})?);
}
Ok(())
}
async fn authorize(
&self,
op: &Operation,
) -> Result<(AuthzFacts, Option<FastForward>), ServerError> {
tracing::debug!(stage = "authorize");
if !op.procedure().is_write() {
return Ok((self.authorize_read(op).await?.facts, None));
}
if self.cfg.authority_fence.is_some() && self.cfg.sharding == Sharding::Single {
Box::pin(self.ensure_authority_activation(&op.repo.namespace)).await?;
}
match &self.cfg.addressing {
Addressing::Multi(multi) => self.owner_rule(op, Some(&multi.namespace_policy)).await,
Addressing::Single { .. } if self.cfg.write_policy == WritePolicy::Owner => {
self.owner_rule(op, None).await
}
Addressing::Single { .. } => self
.hooks
.authorizer()
.authorize(op)
.await
.map(|mut facts| {
facts.authority_generation = None;
(facts, None)
})
.map_err(ServerError::strip_admission_shape),
}
}
fn owner_key_is_admin(&self, namespace: &Namespace) -> bool {
matches!(namespace, Namespace::Ed25519(key) if self.cfg.admin_keys.contains(key))
}
fn require_client_owner_key(&self, namespace: &Namespace) -> Result<(), ServerError> {
if self.owner_key_is_admin(namespace) {
return Err(ServerError::permission_denied(
"admin key cannot authorize client calls",
));
}
Ok(())
}
async fn owner_rule(
&self,
op: &Operation,
policy: Option<&NamespacePolicy>,
) -> Result<(AuthzFacts, Option<FastForward>), ServerError> {
let namespace = Namespace::parse(op.repo.namespace.as_str())
.map_err(|_| internal("invalid resolved owner-policy namespace"))?;
if op.write_grant.is_some() {
self.require_client_owner_key(&namespace)?;
}
if let Some(NamespacePolicy::Allowlist(allowed)) = policy
&& !allowed.contains(&namespace)
{
return Err(ServerError::permission_denied("write not permitted"));
}
if op.write_grant.is_some()
&& matches!(namespace, Namespace::Address(_))
&& !matches!(policy, Some(NamespacePolicy::Allowlist(_)))
{
return Err(ServerError::permission_denied("write not permitted"));
}
let mut fast_forward = None;
let grant = if let Some(header) = &op.write_grant {
let cfg = self.cfg.grants.as_ref().ok_or_else(|| {
grants::rejected(mkit_attest::grant::GrantError::SchemeNotAdvertised)
})?;
let verified = cfg.verify(header.expose(), op)?;
let scope =
crate::policy::ref_scopes::authorize(&verified, &op.kind, self.cfg.indexed_mode())?;
fast_forward = scope.fast_forward;
let presence_requirement = scope.presence;
let observed = match self.cfg.sharding {
Sharding::Single => op.observed_epoch,
Sharding::D34 => op.leased_epoch,
};
if observed != Some(verified.epoch()) {
return Err(plan::epoch_moved());
}
tracing::debug!(grant_id = %to_hex_bytes(verified.id()), "write grant accepted");
Some(crate::op::GrantRef {
id: *verified.id(),
epoch: verified.epoch(),
presence_requirement,
})
} else {
None
};
let owner = grant.is_none()
&& matches!(&namespace, Namespace::Ed25519(key)
if op.principal.ed25519() == Some(key));
if self.cfg.authorizer_role == AuthorizerRole::Check && !owner && grant.is_none() {
return Err(ServerError::permission_denied("write not permitted"));
}
let mut facts = AuthzFacts {
authority_generation: None,
grant,
owner,
caller_view: CallerView::Writer,
};
let mut authorized = op.clone();
authorized.authz = facts.clone();
let returned = self
.hooks
.authorizer()
.authorize(&authorized)
.await
.map_err(ServerError::strip_admission_shape)?;
self.merge_authority_facts(&mut facts, &returned)?;
Ok((facts, fast_forward))
}
async fn pre_receive(&self, op: &Operation) -> Result<(), ServerError> {
tracing::debug!(stage = "pre_receive");
self.hooks
.pre_receive()
.check(op, None)
.await
.map_err(ServerError::strip_admission_shape)
}
async fn plan_and_apply(
&self,
op: &Operation,
a: &Authenticated,
p: &Partition,
(kind, refs, charges): (WriteKind, &[RefUpdate], &[QuotaCharge]),
ahead: Option<Snapshot>,
(lease, begin): (Option<lease::LeaseWrite>, Option<&BeginWrite>),
inputs: WriteInputs<'_>,
) -> Result<StoredResult, ServerError> {
let WriteInputs {
denial_ids,
denial_packs,
pending,
implicit,
external_bases,
inspected,
} = inputs;
let caps = self.meta.capabilities();
let replay = upload::replay_guard(op);
let implicit_ids = implicit
.filter(|pending| {
!pending.is_empty() && matches!(self.cfg.addressing, Addressing::Multi(_))
})
.map(implicit::implicit_packs);
let advance = self.publication_ticket_write(op, a, p)?;
let mut req = WriteRequest {
denial_ids: Some(denial_ids),
denial_packs,
authority_store: plan::AuthorityStore::from_capabilities(self.meta.capabilities()),
authority_generation: op.authz.authority_generation,
repo: &op.repo.name,
kind,
refs,
ref_index: (self.cfg.sharding == Sharding::D34 && !refs.is_empty()).then_some((
&op.repo,
p,
self.shards.as_ref(),
)),
replay,
charges,
namespace_charge: self.namespace_charge(
p,
charges,
ahead.as_ref().and_then(|s| s.namespace_window),
)?,
grant: op.authz.grant.clone(),
lease,
layout_version: caps.implicit_layout_version.is_none(),
mark_repo_known: matches!(self.cfg.addressing, Addressing::Multi(_))
&& ahead
.as_ref()
.is_none_or(|snap| snap.get(&keys::repo_known(&op.repo.name)).is_none()),
rejection: None,
publication: (caps.atomic_multi_key && !refs.is_empty()).then_some(
clearance::PublicationWrite {
repo: &op.repo,
source: p,
shards: self.shards.as_ref(),
prepared: None,
},
),
pending,
begin,
advance,
implicit: implicit_ids.as_deref().map(|packs| ImplicitConsume {
packs,
repo_id: &op.repo,
source: p,
shards: self.shards.as_ref(),
}),
};
let mut ahead = ahead;
let prepared = match Box::pin(self.prepare_publication(
op,
&a.repo().identity,
p,
&req,
&mut ahead,
implicit_ids.as_deref(),
external_bases,
inspected,
))
.await
{
Ok(prepared) => prepared,
Err(error)
if self.inspection_max_objects().is_some()
&& error.code() == crate::Code::PermissionDenied =>
{
return self
.store_policy_denial(op, a, p, ahead, error, pending.is_some())
.await;
}
Err(error) => return Err(error),
};
if let Some(publication) = &mut req.publication {
publication.prepared = prepared.as_ref();
}
if caps.atomic_multi_key || replay.is_some() || !charges.is_empty() {
return self.apply_atomic(op, a, p, &req, ahead).await;
}
self.apply_sequential(op, a, p, req).await
}
async fn apply_sequential(
&self,
op: &Operation,
a: &Authenticated,
p: &Partition,
mut req: WriteRequest<'_>,
) -> Result<StoredResult, ServerError> {
let kind = req.kind;
let refs = req.refs;
req.kind = WriteKind::UpdateRef;
for (i, update) in refs.iter().enumerate() {
req.refs = core::slice::from_ref(update);
let result = self.apply_loop(op, a, p, &req, None).await?;
if let StoredResult::UpdateRef(UpdateRefResult::Conflict { .. }) = result {
if kind == WriteKind::UpdateRef {
return Ok(result);
}
return Ok(StoredResult::AdvanceRefs(if i == 0 {
AdvanceOutcome::PackmapConflict
} else {
AdvanceOutcome::HeadConflict
}));
}
}
Ok(match kind {
WriteKind::UpdateRef => StoredResult::UpdateRef(UpdateRefResult::Committed),
_ => StoredResult::AdvanceRefs(AdvanceOutcome::Committed),
})
}
fn publication_ticket_write<'a>(
&'a self,
op: &'a Operation,
a: &'a Authenticated,
p: &'a Partition,
) -> Result<Option<advance::AdvanceWrite<'a>>, ServerError> {
Ok(match &op.kind {
OpKind::AdvanceRefs { head, tickets, .. } if !tickets.is_empty() => {
Some(advance::AdvanceWrite {
ids: tickets,
signer: op
.auth
.as_ref()
.ok_or_else(|| {
ServerError::failed_precondition("invalid or expired upload ticket")
})?
.signer,
head_ref: &head.name,
repo_id: &op.repo,
repository: &a.repo().identity,
source: p,
shards: self.shards.as_ref(),
})
}
_ => None,
})
}
#[allow(clippy::too_many_arguments, clippy::too_many_lines)] async fn prepare_publication(
&self,
op: &Operation,
repository: &str,
p: &Partition,
req: &WriteRequest<'_>,
ahead: &mut Option<Snapshot>,
implicit_ids: Option<&[Hash]>,
external_bases: &std::collections::BTreeSet<Hash>,
inspected: Option<&mut crate::indexed::inspection::InspectionSet>,
) -> Result<Option<crate::store::publication::Advance>, ServerError> {
#[cfg(not(feature = "remote-hooks"))]
let _ = repository;
let policy = self.publication_policy.as_deref().or_else(|| {
self.cfg
.takedown_denial
.then_some(&clearance::Immediate as &dyn clearance::PublicationPolicy)
});
#[cfg(feature = "remote-hooks")]
let policy = policy.or_else(|| {
(!self.inspectors.is_empty())
.then_some(&inspection::Immediate as &dyn clearance::PublicationPolicy)
});
if let Some(policy) = policy
&& !req.refs.is_empty()
&& req.refs.iter().all(|update| update.new.is_some())
{
let snapshot = ahead.get_or_insert_with(Snapshot::default);
self.fill(p, snapshot, req.read_keys()).await?;
let pair = clearance::resulting_pair(&op.repo.name, req.refs, snapshot)?;
let mut prepared = policy.prepare(op, &pair).await?;
if prepared.value != pair {
return Err(internal("publication policy changed the resulting pair"));
}
prepared.generation =
crate::store::publication::Publication::decode(snapshot.get(&keys::publication(
&op.repo.name,
&crate::store::publication::sequence_ref(&req.refs[0].name),
)))
.map_err(meta_error)?
.generation;
prepared.additions = if let Some(advance) = &req.advance {
advance
.ids
.iter()
.map(|id| {
snapshot
.get(&keys::ticket(id))
.ok_or_else(|| internal("publication ticket missing"))
.and_then(|raw| {
codec::decode_ticket(raw)
.map(|t| t.pack_id)
.map_err(meta_error)
})
})
.collect::<Result<Vec<_>, _>>()?
} else {
implicit_ids.map(<[Hash]>::to_vec).unwrap_or_default()
};
#[cfg(feature = "remote-hooks")]
let mut inspected = inspected;
#[cfg(feature = "remote-hooks")]
let verification_set = inspected.as_deref_mut();
#[cfg(not(feature = "remote-hooks"))]
let verification_set = inspected;
self.verify_publication(
op,
p,
&mut prepared,
policy,
mkit_attest::grant::head_packmap(&crate::store::publication::sequence_ref(
&req.refs[0].name,
))
.is_some(),
external_bases,
verification_set,
req.advance
.as_ref()
.and_then(|a| {
a.ids
.iter()
.filter_map(|id| {
snapshot
.get(&keys::ticket(id))
.and_then(|raw| codec::decode_ticket(raw).ok())
.filter(|t| Some(t.pack_id) == pair.packmap)
.map(|t| t.created_at_ms)
})
.min()
})
.unwrap_or_else(|| {
op.auth
.as_ref()
.map_or(ms(self.clock.now_ms()), |a| ms(a.created_at_ms))
}),
)
.await?;
#[cfg(feature = "remote-hooks")]
if let Some(set) = inspected {
let assignment = if self.cfg.scanner_retrieval.is_some() {
Some(Self::retrieval_assignment(
op,
req.advance.as_ref(),
repository,
snapshot,
set,
)?)
} else {
None
};
self.inspect_advance(
op,
&prepared.value,
set.clone().finalize(),
assignment.as_ref(),
)
.await?;
}
Ok(Some(prepared))
} else {
Ok(None)
}
}
#[allow(clippy::too_many_arguments, clippy::too_many_lines)] async fn verify_publication(
&self,
op: &Operation,
p: &Partition,
prepared: &mut crate::store::publication::Advance,
policy: &dyn clearance::PublicationPolicy,
branch: bool,
external_bases: &std::collections::BTreeSet<Hash>,
mut inspected: Option<&mut crate::indexed::inspection::InspectionSet>,
created: u64,
) -> Result<(), ServerError> {
let indexed = self
.cfg
.indexed
.ok_or_else(|| internal("publication requires indexed mode"))?;
let resumed = self.cfg.takedown_denial
&& self.publication_policy.is_none()
&& inspected.is_none()
&& Box::pin(crate::indexed::publication::resume::prepare(
&self.meta,
p,
self.shards.as_ref(),
&op.repo,
prepared,
indexed,
ms(self.clock.now_ms()),
self.metrics.as_ref(),
created,
))
.await?;
let mut indexed = indexed;
if self.cfg.takedown_denial && !resumed {
indexed.decode_budget = indexed.decode_budget.min(8 << 20);
}
let inspection_budget = crate::indexed::budget::SliceBudget::new(256);
let inspection_blobs =
crate::indexed::budget::Budgeted::new(&self.blobs, &inspection_budget);
let inspection_meta = crate::indexed::budget::Budgeted::new(&self.meta, &inspection_budget);
if let Some(set) = inspected.as_deref_mut() {
crate::indexed::publication::verify_inspected(
&inspection_blobs,
&inspection_meta,
self.shards.as_ref(),
&op.repo,
&prepared.value.clone(),
branch,
prepared,
policy,
indexed,
self.metrics.as_ref(),
set,
)
.await?;
} else if !resumed {
crate::indexed::publication::verify(
&self.blobs,
&self.meta,
self.shards.as_ref(),
&op.repo,
&prepared.value.clone(),
branch,
prepared,
policy,
indexed,
self.metrics.as_ref(),
)
.await?;
}
prepared.external_bases = prepared
.external_bases
.iter()
.copied()
.chain(external_bases.iter().copied())
.collect::<std::collections::BTreeSet<_>>()
.into_iter()
.collect();
if prepared.external_bases.len() > crate::store::publication::MAX_ADVANCE_ITEMS {
return Err(ServerError::invalid_argument("object index limit exceeded"));
}
if prepared.state.publishable() {
let visible = if inspected.is_some() {
crate::timers::publication_recheck::dependencies(
&inspection_meta,
&inspection_meta,
p,
self.shards.as_ref(),
&op.repo,
prepared,
)
.await
.map_err(|error| {
if crate::indexed::budget::is_exhausted(&error) {
ServerError::invalid_argument("object index limit exceeded")
} else {
meta_error(error)
}
})?
} else {
crate::timers::publication_recheck::dependencies(
&self.meta,
&self.meta,
p,
self.shards.as_ref(),
&op.repo,
prepared,
)
.await
.map_err(meta_error)?
};
if !visible {
prepared.state = crate::store::publication::Clearance::Pending;
}
}
Ok(())
}
async fn apply_atomic(
&self,
op: &Operation,
a: &Authenticated,
p: &Partition,
req: &WriteRequest<'_>,
ahead: Option<Snapshot>,
) -> Result<StoredResult, ServerError> {
if !self.meta.capabilities().atomic_multi_key {
return Err(internal("replay and quota need atomic multi-key batches"));
}
self.apply_loop(op, a, p, req, ahead).await
}
#[allow(clippy::too_many_lines)]
async fn apply_loop(
&self,
op: &Operation,
a: &Authenticated,
p: &Partition,
req: &WriteRequest<'_>,
mut ahead: Option<Snapshot>,
) -> Result<StoredResult, ServerError> {
let mut req = req.for_store(self.meta.capabilities());
let _gate = match &self.gate {
Some(gate) => Some(gate.enter(p).await),
None => None,
};
let skew_ms = a.business_skew_ms;
let (mut replans, mut deadline_missed, mut prune_ok) = (0, false, true);
let mut first_attempt = true;
let denial_budget = crate::indexed::budget::SliceBudget::new(9000);
loop {
let denial_packs: Vec<_> = req
.denial_packs
.iter()
.chain(
req.publication
.iter()
.filter_map(|p| p.prepared)
.flat_map(|a| a.dependencies.iter().chain(&a.external_bases)),
)
.copied()
.collect();
let prove = self.cfg.takedown_denial
&& req
.denial_ids
.is_some_and(|ids| !ids.is_empty() || !denial_packs.is_empty());
let proof_plan_time = prove.then(|| ms(self.clock.now_ms()));
if let Some(ids) = req.denial_ids.filter(|_| prove) {
crate::takedown::denial::require_repo_clear_budgeted(
&self.meta,
self.shards.as_ref(),
&op.repo,
ids,
&denial_packs,
&denial_budget,
)
.await?;
}
let base = ahead.take().unwrap_or_default();
let base = if prove {
let mut fresh = Snapshot::default();
fresh.namespace_window = base.namespace_window;
fresh.namespace_seed = base.namespace_seed;
fresh
} else {
base
};
let clock = self.plan_clock(skew_ms, &req);
let snap = self.read_snapshot(p, &req, &clock, base, prune_ok).await?;
if let Some(lease) = req.lease
&& (prove
|| !first_attempt
|| lease
.value
.expires_at_ms
.saturating_sub(self.cfg.lease_margin_ms)
< ms(self.clock.now_ms()).saturating_add(self.cfg.min_lease_budget_ms))
{
let observed = self.observe_lease(op, p, Some(&snap)).await?;
let (_, renewed) = self.admit_lease(op, p, observed, skew_ms).await?;
req.lease = Some(renewed);
}
first_attempt = false;
let mut clock = self.plan_clock(skew_ms, &req);
if let Some(plan_time) = proof_plan_time {
clock.plan_time_ms = plan_time;
}
let plan = match plan_write(&req, &snap, &clock)? {
Planned::Done(result) => return Ok(result),
Planned::Apply(plan) => plan,
};
let Plan {
batch,
on_commit,
replay_index,
epoch_index,
pending_index,
prune,
prune_from,
} = plan;
require_relay_source_lease(self.cfg.sharding, req.lease.is_some(), &batch)?;
#[cfg(feature = "test-faults")]
let batch = faults::delay_relay_batch(
batch,
a.test_directives(),
op,
ms(clock.business_now_ms),
);
if req.kind != WriteKind::UploadReserve {
#[cfg(feature = "test-faults")]
{
let mut attempt = op.clone();
attempt.leased_epoch = req.lease.map(|l| l.value.epoch);
fault!(self, BeforeFinalApply, &attempt, a);
}
}
tracing::debug!(stage = "apply", replans);
match self.meta.apply(p, batch).await {
Ok(BatchOutcome::Committed) => {
#[cfg(feature = "test-faults")]
self.schedule_test_ref_timer(op, a, &on_commit, ms(clock.business_now_ms))
.await?;
return Ok(on_commit);
}
Ok(BatchOutcome::DeadlinePassed { backend_now }) => {
tracing::info!(
backend_now,
deadline = clock.deadline(),
"commit deadline passed"
);
let now = self.clock.now_ms().saturating_add(skew_ms);
let backend = i64::try_from(backend_now).unwrap_or(i64::MAX);
let valid = req
.replay
.is_none_or(|r| now.max(backend) <= r.expires_at_ms);
if deadline_missed || !valid {
return Err(ServerError::unavailable("commit deadline passed; retry"));
}
deadline_missed = true;
}
Ok(BatchOutcome::PreconditionFailed { index, observed }) => {
if Some(index) == pending_index {
return Err(ServerError::unavailable(
"admission reservation changed; retry",
));
}
if Some(index) == replay_index {
return replay_raced(&req, observed.as_ref());
}
if Some(index) == epoch_index {
return Err(plan::epoch_moved());
}
if index >= prune_from && prune_ok {
prune_ok = false;
continue;
}
replans += 1;
if replans > MAX_REPLAN {
return Err(ServerError::aborted_retryable("write contention; retry")
.with_abort_cause(AbortCause::Contention));
}
}
Err(StoreError::Full) => return Err(self.partition_full(p, prune).await),
Err(e) => return Err(meta_error(e)),
}
}
}
#[cfg(feature = "test-faults")]
async fn schedule_test_ref_timer(
&self,
op: &Operation,
a: &Authenticated,
on_commit: &StoredResult,
now_ms: u64,
) -> Result<(), ServerError> {
if let OpKind::UpdateRef(upd) = &op.kind
&& matches!(
on_commit,
StoredResult::UpdateRef(UpdateRefResult::Committed)
)
{
let timer_partition = self.shards.ref_shard(&op.repo, &upd.name);
faults::schedule_timer(
a.test_directives(),
&self.meta,
&timer_partition,
&op.repo.name,
&upd.name,
now_ms,
)
.await?;
}
Ok(())
}
fn plan_clock(&self, skew_ms: i64, req: &WriteRequest<'_>) -> PlanClock {
let now = self.clock.now_ms();
let lead = MAX_CLOCK_LEAD_MS.unsigned_abs();
PlanClock {
plan_time_ms: ms(now),
business_now_ms: now.saturating_add(skew_ms),
max_apply_window_ms: u64::try_from(self.cfg.max_apply_window.as_millis())
.unwrap_or(u64::MAX),
deadline_cap: match (
req.replay.map(|r| ms(r.expires_at_ms).saturating_add(lead)),
req.lease.map(|l| {
l.value
.expires_at_ms
.saturating_sub(self.cfg.lease_margin_ms)
}),
) {
(Some(replay), Some(lease)) => Some(replay.min(lease)),
(replay, lease) => replay.or(lease),
},
}
}
async fn read_snapshot(
&self,
p: &Partition,
req: &WriteRequest<'_>,
clock: &PlanClock,
mut snap: Snapshot,
prune: bool,
) -> Result<Snapshot, ServerError> {
let now = clock.plan_time_ms;
if prune && prune_sampled(req, now) {
if req.replay.is_some() {
snap.expired_replays = read::expired_replay_keys(&self.meta, p, now, PRUNE_LIMIT)
.await
.map_err(meta_error)?;
}
if let Some(window) = req.charges.iter().map(|c| c.limits.window_ms).max() {
snap.stale_quotas =
read::stale_quota_keys(&self.meta, p, now, ms(window), PRUNE_LIMIT)
.await
.map_err(meta_error)?;
}
}
let mut wanted = req.read_keys();
wanted.extend(snap.stale_quotas.iter().map(|(_, quota)| quota.clone()));
self.fill(p, &mut snap, wanted).await?;
if let Some(advance) = &req.advance {
let detail = advance::detail_keys(&snap, advance)?;
self.fill(p, &mut snap, detail).await?;
}
if let Some(BeginWrite::Open(open)) = req.begin {
begin::read_indexed(&self.meta, p, &open.spec, &mut snap).await?;
let reservation = crate::store::tickets::keys(&open.spec).reservation;
if snap.get(&reservation).is_some()
&& let Some(replay) = req.replay
{
let key = keys::replay(&replay.scope);
let value = self.meta.get(p, &key).await.map_err(meta_error)?;
snap.insert(key, value);
}
}
Ok(snap)
}
async fn apply_meta(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, ServerError> {
match self.meta.apply(p, batch).await {
Ok(outcome) => Ok(outcome),
Err(StoreError::Full) => Err(self.partition_full(p, None).await),
Err(error) => Err(meta_error(error)),
}
}
async fn partition_full(&self, p: &Partition, prune: Option<Batch>) -> ServerError {
self.metrics
.incr(METRIC_PARTITION_FULL, &[("kind", p.kind())], 1);
let name = p.encode().map(|b| to_hex_bytes(&b)).unwrap_or_default();
tracing::error!(partition = %name, kind = p.kind(), "storage partition full");
if let Some(prune) = prune
&& let Err(e) = self.meta.apply(p, prune).await
{
tracing::warn!(error = %e, "prune on a full partition failed");
}
ServerError::unavailable("storage partition full")
}
}
#[derive(Default)]
struct Allowance {
charges: Vec<QuotaCharge>,
reservation: Option<String>,
response_headers: Vec<(String, String)>,
external_ref: Option<String>,
}
fn check_ref_name(name: &str) -> Result<(), ServerError> {
if name.len() > refs::MAX_REF_NAME_BYTES {
Err(ServerError::invalid_argument(refs::REF_NAME_TOO_LONG))
} else if !validate_ref_name(name) {
Err(ServerError::invalid_argument(
"ref name is invalid (SPEC-REFS §3)",
))
} else if refs::is_served_ref_name(name) {
Ok(())
} else {
Err(ServerError::invalid_argument(refs::REF_NAME_OUTSIDE_REFS))
}
}
fn stored_mismatch(result: &StoredResult) -> ServerError {
match result {
StoredResult::Rejected(rejection) => {
ServerError::new(rejection.code(), rejection.message().to_owned())
}
_ => internal("stored result is for another procedure"),
}
}
fn replay_raced(
req: &WriteRequest<'_>,
observed: Option<&Value>,
) -> Result<StoredResult, ServerError> {
if req.pending.is_some() {
return Err(
ServerError::aborted_retryable("operation already in flight; retry")
.with_abort_cause(AbortCause::ReplayRace),
);
}
let (Some(replay), Some(value)) = (req.replay, observed) else {
return Err(ServerError::aborted_retryable(
"operation already in flight; retry",
));
};
let record = codec::decode_replay_record(value).map_err(meta_error)?;
replay_answer(classify(Some(&record), &replay.fingerprint))?
.ok_or_else(|| ServerError::aborted_retryable("operation already in flight; retry"))
}