1mod admission;
23mod advance;
24mod auth;
25mod authority;
26mod begin;
27pub mod clearance;
28mod coordinator;
29mod download;
30mod durable_outcome;
31mod epoch;
32#[cfg(feature = "test-faults")]
33pub(crate) mod faults;
34mod gate;
35mod hooks;
36#[cfg(feature = "http-objects")]
37mod http;
38#[cfg(feature = "http-objects")]
39mod http_admission;
40#[cfg(feature = "http-objects")]
41mod http_tokens;
42#[cfg(feature = "http-objects")]
43mod object_reader;
44#[cfg(feature = "http-objects")]
45pub use object_reader::{
46 IssuedUrl, OBJECT_READER_BATCH, OBJECT_READER_CALLS, ObjectMetadata, ObjectReader, ReaderView,
47};
48mod implicit;
49mod info;
50#[cfg(feature = "remote-hooks")]
51pub mod inspection;
52mod lease;
53pub mod list;
54mod list_repos;
55pub use list_repos::{RepoEntry, RepoPage};
56mod outcome;
57mod parts;
58mod plan;
59#[cfg(feature = "published-view")]
60pub mod published;
61mod purge;
62mod ref_policy;
63mod reservation;
64mod revocation;
65mod scanner_retrieval;
66mod shard;
67mod staging;
68#[cfg(test)]
69mod tests;
70mod upload;
71mod watermark;
72
73use core::future::Future;
74use core::time::Duration;
75use std::sync::Arc;
76
77pub use mkit_attest::grant::Visibility as RepoVisibility;
78use mkit_attest::grant::{Visibility, verify_visibility_statement};
79use mkit_core::hash::{Hash, to_hex, to_hex_bytes};
80use mkit_core::protocol::{AdvanceOutcome, PackKey};
81use mkit_core::repo_identity::{Namespace, RepositoryIdentity};
82use mkit_core::write_auth::MAX_CLOCK_LEAD_MS;
83use tracing::Instrument;
84
85use crate::download::DOWNLOAD_CHUNK_MAX;
86use crate::error::{AbortCause, InvalidHeader, ServerError};
87use crate::op::{
88 AuthzFacts, CallerView, GrantRef, OpKind, Operation, Procedure, RefUpdate, VerifiedAuth,
89};
90use crate::policy::ff::FastForward;
91use crate::policy::{
92 AuthorizerRole, GrantConfig, NamespacePolicy, WritePolicy, grants, read as read_policy,
93};
94use crate::principal::Principal;
95use crate::quota::{
96 self, DEFAULT_WRITE_QUOTA, NamespaceCharge, NamespaceDecision, NamespaceView, QuotaCharge,
97 QuotaLimits, QuotaScope, ViewStatus,
98};
99use crate::refs::{self, strip_listed_prefix, validate_ref_name};
100use crate::replay::{
101 BeginUploadResult, ReplayDecision, ReplayRecord, ReplayState, StoredResult, UpdateRefResult,
102 classify,
103};
104use crate::repo::{Addressing, RepoId};
105use crate::rt::Clock;
106use crate::storage_error::{StorageOp, describe_and_map};
107use crate::store::tickets::TicketCaps;
108use crate::store::{
109 Batch, BatchOutcome, Key, KeyClasses, MAX_BATCH_OPS, MultipartBlobStore, NamespaceStore,
110 Partition, Precondition, StoreError, Value, codec, keys, read,
111};
112use crate::telemetry::{Metrics, Redactor};
113use crate::upload::{UploadLimits, token::TicketKeys};
114use crate::url_token::{MintedToken, UrlTarget};
115use begin::BeginWrite;
116
117#[cfg(feature = "remote-hooks")]
118pub(crate) use admission::validate_decision;
119pub use auth::{AuthMode, Authenticated, HeaderValues, RequestMeta};
120pub use download::{DownloadChunk, DownloadStream};
121pub use durable_outcome::{DeliveryError, Outcome, OutcomeKind};
122#[cfg(feature = "test-faults")]
123pub use faults::{
124 BUMP_EPOCH_HEADER, CLOCK_SKEW_HEADER, FAULT_HEADER, FailOnce, FaultHooks, FaultPoint,
125 LEASE_RECOVERED_HEADER, RELAY_DELAY_MS_HEADER, RUN_TIMERS_HEADER, TIMER_MS_HEADER,
126 TestDirectives,
127};
128pub use hooks::{
129 ADMISSION_EXPOSE_HEADERS, Admission, AdmissionDecision, AdmissionInput, Authorizer, Challenge,
130 Choice, CredentialHeader, DefaultAdmission, HookSet, Hooks, NoOutcomes, NoPreReceive,
131 NoReceipts, OpenAuthorizer, OutcomeSink, PreReceive, ReceiptSigner,
132};
133#[cfg(feature = "ssh")]
134pub(crate) use implicit::IMPLICIT_PACKMAP_UNKNOWN;
135pub(crate) use implicit::PendingPack;
136pub use info::ServerInfo;
137pub use lease::{LeaseParams, renew_for_relay};
138use outcome::Outcome as RequestOutcome;
139pub use parts::PartUploadSession;
140use plan::{
141 ImplicitConsume, MAX_REPLAN, PRUNE_LIMIT, Plan, PlanClock, Planned, Snapshot, WriteKind,
142 WriteRequest, plan_write, prune_sampled,
143};
144pub use revocation::{MAX_EPOCH_STEP, RevokeBudget, RevokeProgress};
145pub use shard::{D34Shards, ShardMap, SinglePartition};
146pub use upload::{UploadMode, UploadSession};
147
148#[derive(Debug, Clone, Default, PartialEq, Eq)]
151pub struct ResponseMeta {
152 headers: Vec<(String, String)>,
153 external_ref: Option<String>,
154}
155
156impl ResponseMeta {
157 #[must_use]
159 pub fn headers(&self) -> &[(String, String)] {
160 &self.headers
161 }
162 #[must_use]
164 pub fn external_ref(&self) -> Option<&str> {
165 self.external_ref.as_deref()
166 }
167}
168
169macro_rules! fault {
172 ($pipe:expr, $point:ident, $op:expr, $a:expr) => {
173 #[cfg(feature = "test-faults")]
174 $pipe
175 .fault($crate::pipeline::FaultPoint::$point, $op, $a)
176 .await?
177 };
178}
179pub(crate) use fault;
180
181pub const MAX_APPLY_WINDOW: Duration = Duration::from_secs(10);
187
188const _: () = assert!(MAX_CLOCK_LEAD_MS.unsigned_abs() < read::REPLAY_PRUNE_GRACE_MS);
192
193pub const DEFAULT_LIST_PAGE_LIMIT: u32 = 1000;
195
196pub const METRIC_PARTITION_FULL: &str = "mkit_server_partition_full_total";
199
200pub const METRIC_HEADER_DROPPED: &str = "mkit_server_error_header_dropped_total";
203
204#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
206#[non_exhaustive]
207pub enum Sharding {
208 #[default]
210 Single,
211 D34,
213}
214
215fn require_relay_source_lease(
216 sharding: Sharding,
217 has_lease: bool,
218 batch: &Batch,
219) -> Result<(), ServerError> {
220 if sharding == Sharding::D34
221 && !has_lease
222 && batch.writes.iter().any(|write| {
223 matches!(write, crate::store::Write::Put(key, _) if matches!(keys::parse(key), Some(keys::ParsedKey::Relay(_))))
224 })
225 {
226 return Err(internal("D34 relay batch lacks source epoch lease"));
227 }
228 Ok(())
229}
230
231#[derive(Debug)]
235struct ReadAuth {
236 facts: AuthzFacts,
238 epoch: Option<u64>,
242}
243
244#[derive(Debug, Clone)]
246#[non_exhaustive]
247pub struct PipelineConfig {
248 pub addressing: Addressing,
250 pub sharding: Sharding,
252 pub inspection_mode: bool,
254 pub auth: AuthMode,
256 pub grants: Option<GrantConfig>,
258 pub authority_fence: Option<crate::authority::AuthorityFence>,
260 pub write_policy: WritePolicy,
262 pub default_repo_visibility: RepoVisibility,
264 pub authorizer_role: AuthorizerRole,
266 pub list_repos_authority_full: bool,
269 pub upload_limits: UploadLimits,
271 pub single_upload_max_bytes: Option<u64>,
273 pub part_size: u64,
275 pub max_parts: u32,
277 pub max_list_refs_page_size: u32,
280 pub begin_upload_threshold_bytes: u64,
283 pub ticket_keys: Option<TicketKeys>,
285 pub scanner_retrieval: Option<Arc<crate::scanner_retrieval::RetrievalConfig>>,
287 pub url_tokens: Option<crate::url_token::UrlTokenConfig>,
290 pub admin_keys: Vec<[u8; 32]>,
292 pub receipt_publication: Option<crate::takedown::PublicationConfig>,
294 pub purge: Option<crate::purge::PurgeConfig>,
296 pub ticket_ttl_ms: u64,
298 pub ticket_caps: TicketCaps,
300 pub download_chunk_max: usize,
302 pub write_quota: Option<QuotaLimits>,
306 pub list_page_limit: u32,
308 pub max_apply_window: Duration,
310 pub epoch_lease_ms: u64,
312 pub lease_margin_ms: u64,
316 pub min_lease_budget_ms: u64,
318 pub redactor: Redactor,
320 pub admission_credential_headers: Vec<String>,
322 pub outbox_backlog_cap: Option<OutboxBacklogCap>,
324 pub indexed: Option<crate::indexed::IndexedConfig>,
326 pub takedown_denial: bool,
328 pub ref_policy: Option<crate::policy::RefPolicy>,
332 #[cfg(feature = "http-objects")]
335 pub http_objects: Option<crate::http_objects::HttpObjectsConfig>,
336}
337
338#[derive(Debug, Clone, Copy, PartialEq, Eq)]
341pub struct OutboxBacklogCap {
342 pub rows: u64,
344 pub bytes: u64,
346}
347
348impl PipelineConfig {
349 pub(crate) fn indexed_mode(&self) -> bool {
351 self.indexed.is_some()
352 }
353
354 #[must_use]
356 pub fn new(addressing: Addressing, auth: AuthMode, upload_limits: UploadLimits) -> Self {
357 let write_quota = matches!(auth, AuthMode::AuthV2(_)).then_some(DEFAULT_WRITE_QUOTA);
358 let write_policy = match &addressing {
359 Addressing::Single { .. } => WritePolicy::Open,
360 Addressing::Multi(_) => WritePolicy::Owner,
361 };
362 Self {
363 write_policy,
364 default_repo_visibility: RepoVisibility::Public,
365 authorizer_role: AuthorizerRole::Check,
366 list_repos_authority_full: false,
367 addressing,
368 sharding: Sharding::Single,
369 inspection_mode: false,
370 auth,
371 grants: None,
372 authority_fence: None,
373 upload_limits,
374 single_upload_max_bytes: None,
375 part_size: mkit_core::upload_parts::MIN_PART_SIZE,
376 max_parts: 10_000,
377 max_list_refs_page_size: DEFAULT_LIST_PAGE_LIMIT,
378 begin_upload_threshold_bytes: u64::MAX,
379 ticket_keys: None,
380 scanner_retrieval: None,
381 url_tokens: None,
382 admin_keys: Vec::new(),
383 receipt_publication: None,
384 purge: None,
385 ticket_ttl_ms: 86_400_000,
386 ticket_caps: TicketCaps {
387 per_ref: 1024,
388 per_signer: 64,
389 },
390 download_chunk_max: DOWNLOAD_CHUNK_MAX,
391 write_quota,
392 list_page_limit: DEFAULT_LIST_PAGE_LIMIT,
393 max_apply_window: MAX_APPLY_WINDOW,
394 epoch_lease_ms: 30_000,
395 lease_margin_ms: 5_000,
396 min_lease_budget_ms: 1_000,
397 redactor: Redactor::default(),
398 admission_credential_headers: Vec::new(),
399 outbox_backlog_cap: Some(OutboxBacklogCap {
400 rows: 100_000,
401 bytes: 64 * 1024 * 1024,
402 }),
403 indexed: None,
404 takedown_denial: false,
405 ref_policy: None,
406 #[cfg(feature = "http-objects")]
407 http_objects: None,
408 }
409 }
410
411 #[must_use]
413 pub fn advertised_namespace_policy(&self) -> &'static str {
414 match &self.addressing {
415 Addressing::Single { .. } => "single-repository",
416 Addressing::Multi(multi) => match &multi.namespace_policy {
417 NamespacePolicy::Allowlist(_) => "allowlist",
418 NamespacePolicy::Any { .. } => "any",
419 },
420 }
421 }
422}
423
424#[derive(Debug, Clone, Copy, PartialEq, Eq)]
426#[non_exhaustive]
427pub struct PipelineCapabilities {
428 pub atomic_advance: bool,
430}
431
432#[derive(Debug, Clone, Copy, PartialEq, Eq)]
434pub struct HealthStatus {
435 pub blobs: bool,
437 pub meta: bool,
439}
440
441impl HealthStatus {
442 #[must_use]
444 pub fn is_healthy(&self) -> bool {
445 self.blobs && self.meta
446 }
447}
448
449#[derive(Debug, Clone, PartialEq, Eq)]
451pub struct RefEntry {
452 pub name: String,
454 pub id: Hash,
456}
457
458#[derive(Clone, PartialEq, Eq)]
461#[non_exhaustive]
462pub enum VisibilityRequest {
463 Envelope(Visibility),
465 Statement(String),
467}
468
469impl core::fmt::Debug for VisibilityRequest {
470 fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
472 match self {
473 Self::Envelope(visibility) => f.debug_tuple("Envelope").field(visibility).finish(),
474 Self::Statement(statement) => f
475 .debug_tuple("Statement")
476 .field(&format_args!("<{} bytes>", statement.len()))
477 .finish(),
478 }
479 }
480}
481
482pub struct Pipeline<B, N, H = Hooks> {
484 publication_policy: Option<Arc<dyn clearance::PublicationPolicy>>,
485 #[cfg(feature = "remote-hooks")]
486 inspectors: Vec<Arc<dyn inspection::ContentInspector>>,
487 #[cfg(feature = "remote-hooks")]
488 inspect_limit: usize,
489 #[cfg(feature = "published-view")]
490 published: Option<Arc<dyn published::PublishedSource>>,
491 blobs: B,
492 meta: N,
493 hooks: H,
494 shards: Arc<dyn ShardMap>,
495 cfg: PipelineConfig,
496 clock: Arc<dyn Clock>,
497 metrics: Arc<dyn Metrics>,
498 #[cfg(feature = "test-faults")]
499 faults: Option<Arc<dyn faults::DynFaultHooks>>,
500 #[cfg(feature = "test-faults")]
501 test_timer_gate: Option<Arc<tokio::sync::Mutex<()>>>,
502 gate: Option<Arc<gate::WriteGate>>,
503 #[cfg(feature = "http-objects")]
504 http_seams: Option<crate::http_objects::HttpSeams>,
505}
506
507struct WriteInputs<'a> {
508 denial_ids: &'a std::collections::BTreeSet<Hash>,
509 denial_packs: &'a [Hash],
510 pending: Option<&'a reservation::PendingGuard>,
511 implicit: Option<&'a [PendingPack]>,
512 external_bases: &'a std::collections::BTreeSet<Hash>,
513 inspected: Option<&'a mut crate::indexed::inspection::InspectionSet>,
514}
515
516impl<B, N, H> core::fmt::Debug for Pipeline<B, N, H> {
517 fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
518 f.debug_struct("Pipeline")
519 .field("cfg", &self.cfg)
520 .finish_non_exhaustive()
521 }
522}
523
524fn internal(detail: &'static str) -> ServerError {
526 ServerError::internal("ref store request failed", detail)
527}
528
529fn store_error(op: StorageOp, err: StoreError) -> ServerError {
531 let op = match err {
532 StoreError::Corrupt(_) => StorageOp::MetaDecode,
533 _ => op,
534 };
535 let (line, err) = describe_and_map(op, err);
536 tracing::warn!(detail = %line, "storage failure");
537 err
538}
539
540fn meta_error(err: StoreError) -> ServerError {
541 store_error(StorageOp::MetaCall, err)
542}
543
544fn method(procedure: Procedure) -> &'static str {
546 let path = procedure.connect_path();
547 path.rsplit('/').next().unwrap_or(path)
548}
549
550fn replay_answer(decision: ReplayDecision) -> Result<Option<StoredResult>, ServerError> {
552 match decision {
553 ReplayDecision::New => Ok(None),
554 ReplayDecision::Return(result) => Ok(Some(result)),
555 ReplayDecision::FingerprintMismatch => Err(ServerError::invalid_argument(
556 "nonce reused for a different operation",
557 )),
558 ReplayDecision::RetryLater | ReplayDecision::Resume => Err(ServerError::aborted_retryable(
561 "operation already in flight; retry",
562 )),
563 }
564}
565
566fn ms(ms: i64) -> u64 {
567 u64::try_from(ms).unwrap_or(0)
568}
569
570#[must_use]
573pub fn repo_is_private(stored: Option<&codec::RepoVisibilityV1>, default: RepoVisibility) -> bool {
574 stored.map_or(default == RepoVisibility::Private, |row| {
575 row.visibility == codec::StoredVisibility::Private
576 })
577}
578
579fn stored_visibility(visibility: Visibility) -> codec::StoredVisibility {
581 match visibility {
582 Visibility::Public => codec::StoredVisibility::Public,
583 Visibility::Private => codec::StoredVisibility::Private,
584 }
585}
586
587fn validate_upload_ticket_config<H: HookSet>(
588 cfg: &PipelineConfig,
589 hooks: &H,
590) -> Result<(), ServerError> {
591 if cfg.begin_upload_threshold_bytes != u64::MAX
595 && !matches!(cfg.auth, AuthMode::TransportIdentity)
596 && (!matches!(cfg.auth, AuthMode::AuthV2(_)) || cfg.ticket_keys.is_none())
597 {
598 return Err(ServerError::invalid_argument(
599 "a ticket threshold requires auth v2 and upload ticket keys",
600 ));
601 }
602 if !matches!(cfg.auth, AuthMode::TransportIdentity)
603 && !hooks.admission().is_default()
604 && (!matches!(cfg.auth, AuthMode::AuthV2(_)) || cfg.ticket_keys.is_none())
605 {
606 return Err(ServerError::invalid_argument(
607 "admission requires auth v2 and upload ticket keys",
608 ));
609 }
610 if matches!(cfg.addressing, Addressing::Multi(_))
611 && matches!(cfg.auth, AuthMode::AuthV2(_))
612 && cfg.ticket_keys.is_none()
613 {
614 return Err(ServerError::invalid_argument(
615 "multi-repository auth v2 deployments require upload ticket keys",
616 ));
617 }
618 Ok(())
619}
620
621impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H> {
622 #[allow(clippy::too_many_lines)] pub fn new(
637 blobs: B,
638 meta: N,
639 hooks: H,
640 mut cfg: PipelineConfig,
641 clock: Arc<dyn Clock>,
642 metrics: Arc<dyn Metrics>,
643 ) -> Result<Self, ServerError> {
644 if let Some(publication) = &cfg.receipt_publication {
645 if cfg.indexed.is_none() {
646 return Err(ServerError::invalid_argument(
647 "takedown requires indexed mode",
648 ));
649 }
650 cfg.admin_keys.extend_from_slice(publication.public_keys());
651 }
652 crate::scanner_retrieval::service::validate_config(&cfg)?;
653 if let Some(purge) = &cfg.purge {
654 purge.validate().map_err(meta_error)?;
655 }
656 if cfg.authority_fence.is_some()
657 && (cfg.authorizer_role != AuthorizerRole::Authority
658 || hooks.authorizer().is_open()
659 || !meta.capabilities().atomic_multi_key
660 || !matches!(cfg.auth, AuthMode::AuthV2(_))
661 || !matches!(cfg.addressing, Addressing::Multi(_)))
662 {
663 return Err(ServerError::invalid_argument(
664 "authority fencing requires Multi, auth v2, an Authority hook and transactional storage",
665 ));
666 }
667 #[allow(clippy::collapsible_if)]
669 if let Some(fence) = &cfg.authority_fence {
670 if fence.public_keys().any(|key| {
671 cfg.ticket_keys
672 .as_ref()
673 .is_some_and(|tickets| tickets.contains_ed25519_public(&key))
674 }) {
675 return Err(ServerError::invalid_argument(
676 "authority keys must differ from ticket keys",
677 ));
678 }
679 #[cfg(feature = "http-objects")]
680 if fence.public_keys().any(|key| {
681 cfg.url_tokens
682 .as_ref()
683 .is_some_and(|tokens| tokens.keys().public_keys().any(|public| public == key))
684 }) {
685 return Err(ServerError::invalid_argument(
686 "authority keys must differ from URL-token keys",
687 ));
688 }
689 }
690 if cfg.takedown_denial && cfg.indexed.is_none() {
691 return Err(ServerError::invalid_argument(
692 "takedown denial requires indexed mode",
693 ));
694 }
695 if let Some(indexed) = &cfg.indexed {
696 if !matches!(cfg.auth, AuthMode::AuthV2(_))
697 || cfg.ticket_keys.is_none()
698 || !matches!(cfg.addressing, Addressing::Multi(_))
699 {
700 return Err(ServerError::invalid_argument(
701 "indexed mode requires auth v2, ticket keys, zero ticket threshold, and Multi addressing",
702 ));
703 }
704 if indexed.max_delta_chain_depth == 0
705 || indexed.max_delta_chain_depth > u32::from(u16::MAX)
706 || indexed.max_pack_bytes == 0
707 || indexed.max_pack_bytes > indexed.decode_budget
708 || indexed.relay_lag_bound_ms == 0
709 || indexed.extract_min_bytes == 0
710 || indexed.max_ancestry_commits == 0
711 || indexed.max_ancestry_commits > crate::indexed::MAX_ANCESTRY_COMMITS_LIMIT
712 || indexed
713 .max_extract_bytes
714 .is_some_and(|max| max < indexed.max_pack_bytes)
715 {
716 return Err(ServerError::invalid_argument("invalid indexed limits"));
717 }
718 cfg.upload_limits.max_total_bytes = cfg
719 .upload_limits
720 .max_total_bytes
721 .min(indexed.max_pack_bytes);
722 }
723 if let Some(policy) = &cfg.ref_policy {
724 policy.validate_for_indexed(cfg.indexed.is_some())?;
725 }
726 #[cfg(feature = "http-objects")]
727 if let Some(http) = &cfg.http_objects {
728 let Some(indexed) = &cfg.indexed else {
729 return Err(ServerError::invalid_argument(
730 "HTTP object serving requires indexed mode",
731 ));
732 };
733 http.validate(indexed.extract_min_bytes)?;
734 if http.admit_reads && hooks.admission().is_default() {
735 return Err(ServerError::invalid_argument(
736 "paid HTTP reads require a real Admission hook",
737 ));
738 }
739 }
740 cfg.validate_server_info_limits()?;
741 let mut credential_names = std::collections::BTreeSet::new();
742 if cfg.admission_credential_headers.iter().any(|name| {
743 !admission::valid_extra_name(name)
744 || ["payment-authorization", "payment-signature"]
745 .contains(&name.to_ascii_lowercase().as_str())
746 || !credential_names.insert(name.to_ascii_lowercase())
747 }) {
748 return Err(ServerError::invalid_argument(
749 "invalid admission credential header name",
750 ));
751 }
752 if cfg.admission_credential_headers.len() + 3 > admission::MAX_CREDENTIAL_HEADERS {
755 return Err(ServerError::invalid_argument(
756 "too many admission credential headers",
757 ));
758 }
759 cfg.redactor.add_names(&cfg.admission_credential_headers);
760
761 if cfg.max_parts > B::MAX_PARTS {
762 return Err(ServerError::invalid_argument(
763 "max_parts exceeds storage backend capacity",
764 ));
765 }
766
767 validate_upload_ticket_config(&cfg, &hooks)?;
768 if cfg.ticket_ttl_ms == 0
769 || cfg.ticket_ttl_ms >= 604_800_000
770 || cfg.ticket_caps.per_ref == 0
771 || cfg.ticket_caps.per_signer == 0
772 {
773 return Err(ServerError::invalid_argument(
774 "invalid upload ticket lifetime or caps",
775 ));
776 }
777 let policy_refusal = match (&cfg.addressing, cfg.write_policy) {
778 (Addressing::Multi(_), WritePolicy::Open) => {
779 Some("write_policy open is single-repository only (SPEC-TRANSPORT-CONNECT §7.5)")
780 }
781 (Addressing::Single { repo }, WritePolicy::Owner)
785 if Namespace::parse(repo.namespace.as_str()).is_err() =>
786 {
787 Some("write_policy owner needs multi-repository addressing")
788 }
789 (Addressing::Multi(multi), _)
790 if matches!(
791 multi.namespace_policy,
792 NamespacePolicy::Any {
793 unsafe_without_admission: false
794 }
795 ) && hooks.admission().is_default() =>
796 {
797 Some(
798 "namespace_policy any needs a non-default admission step, or the explicit unsafe override (D27)",
799 )
800 }
801 _ => None,
802 };
803 if let Some(message) = policy_refusal {
804 return Err(ServerError::invalid_argument(message));
805 }
806 if let Some(grants) = &cfg.grants {
807 let AuthMode::AuthV2(auth) = &cfg.auth else {
808 return Err(ServerError::invalid_argument(
809 "write grants require auth v2",
810 ));
811 };
812 if !matches!(cfg.addressing, Addressing::Multi(_))
813 || cfg.write_policy != WritePolicy::Owner
814 {
815 return Err(ServerError::invalid_argument(
816 "write grants require Multi addressing and owner write policy",
817 ));
818 }
819 if grants.audience() != auth.audience() {
820 return Err(ServerError::invalid_argument(
821 "write grant audience must match auth v2 audience",
822 ));
823 }
824 }
825 if let Some(tokens) = &cfg.url_tokens {
826 if !matches!(cfg.auth, AuthMode::AuthV2(_)) {
827 return Err(ServerError::invalid_argument("URL tokens require auth v2"));
828 }
829 if let Some(tickets) = &cfg.ticket_keys
832 && tokens
833 .keys()
834 .public_keys()
835 .any(|public| tickets.contains_ed25519_public(&public))
836 {
837 return Err(ServerError::invalid_argument(
838 "the URL token key must differ from the upload ticket keys",
839 ));
840 }
841 }
842 if cfg.authorizer_role == AuthorizerRole::Authority && hooks.authorizer().is_open() {
843 return Err(ServerError::invalid_argument(
844 "an authority authorizer must be a real authority source",
845 ));
846 }
847 let caps = meta.capabilities();
848 let full = caps.atomic_multi_key && caps.key_classes == KeyClasses::All;
849 let refused = if matches!(cfg.auth, AuthMode::AuthV2(_)) && !full {
850 "auth v2 needs every key class and atomic multi-key batches"
851 } else if (cfg.sharding == Sharding::D34 || matches!(cfg.addressing, Addressing::Multi(_)))
852 && !full
853 {
854 "sharded or multi-repository routing needs every key class and atomic multi-key batches"
855 } else if !caps.atomic_multi_key && caps.implicit_layout_version.is_none() {
856 "a store without atomic batches must report its layout version"
857 } else if caps
858 .implicit_layout_version
859 .is_some_and(|v| v != keys::LAYOUT_VERSION)
860 {
861 "the store's layout version is not this server's"
862 } else if cfg.lease_margin_ms == 0
863 || cfg.epoch_lease_ms <= cfg.lease_margin_ms.saturating_add(cfg.min_lease_budget_ms)
864 {
865 "epoch lease must exceed its positive margin plus minimum budget"
866 } else if cfg.list_page_limit == 0 || cfg.max_apply_window.is_zero() {
867 "list page limit and apply window must be positive"
868 } else {
869 ""
870 };
871 if !refused.is_empty() {
872 return Err(ServerError::invalid_argument(refused));
873 }
874 let shards: Arc<dyn ShardMap> = match cfg.sharding {
875 Sharding::Single => Arc::new(SinglePartition),
876 Sharding::D34 => Arc::new(D34Shards),
877 };
878 #[cfg(feature = "http-objects")]
879 let http_seams = cfg.http_objects.as_ref().map(|http| {
880 let mut seams = crate::http_objects::HttpSeams::new(http);
881 if let Some(tokens) = &cfg.url_tokens {
882 seams.tokens = Arc::new(tokens.clone());
883 }
884 seams
885 });
886 Ok(Self {
887 #[cfg(feature = "published-view")]
888 published: None,
889 publication_policy: None,
890 #[cfg(feature = "remote-hooks")]
891 inspectors: Vec::new(),
892 #[cfg(feature = "remote-hooks")]
893 inspect_limit: inspection::MAX_OBJECTS,
894 blobs,
895 meta,
896 hooks,
897 shards,
898 cfg,
899 clock,
900 metrics,
901 #[cfg(feature = "test-faults")]
902 faults: None,
903 #[cfg(feature = "test-faults")]
904 test_timer_gate: None,
905 gate: None,
906 #[cfg(feature = "http-objects")]
907 http_seams,
908 })
909 }
910
911 #[must_use]
913 pub fn inspection_max_objects(&self) -> Option<u32> {
914 #[cfg(feature = "remote-hooks")]
915 if !self.inspectors.is_empty() {
916 return u32::try_from(self.inspect_limit).ok();
917 }
918 None
919 }
920
921 #[cfg(feature = "remote-hooks")]
927 pub fn with_inspectors(
928 mut self,
929 inspectors: Vec<Arc<dyn inspection::ContentInspector>>,
930 batch_max: usize,
931 ) -> Result<Self, ServerError> {
932 if inspectors.is_empty() {
933 return Ok(self);
934 }
935 let mut names = std::collections::BTreeSet::new();
936 if inspectors.len() > inspection::MAX_INSPECTORS
937 || !(1..=inspection::MAX_OBJECTS).contains(&batch_max)
938 || inspectors.iter().any(|i| {
939 i.id().is_empty()
940 || !names.insert(i.id())
941 || i.phase() != inspection::InspectorPhase::Sync
942 || i.on_unavailable() != inspection::OnUnavailable::FailClosed
943 })
944 {
945 return Err(ServerError::invalid_argument(
946 "invalid launch inspection configuration",
947 ));
948 }
949 if self.cfg.indexed.is_none()
950 || self.cfg.write_policy == WritePolicy::Open
951 || self.cfg.ticket_keys.is_none()
952 || self.cfg.begin_upload_threshold_bytes != 0
953 || !self.meta.capabilities().atomic_multi_key
954 {
955 return Err(ServerError::invalid_argument(
956 "inspection requires indexed mode, restricted writes and ticketed uploads with threshold zero",
957 ));
958 }
959 self.inspectors = inspectors;
960 self.inspect_limit = batch_max;
961 Ok(self)
962 }
963
964 #[cfg(feature = "test-faults")]
966 #[must_use]
967 pub fn with_faults(mut self, hooks: impl FaultHooks + 'static) -> Self {
968 self.faults = Some(Arc::new(hooks));
969 self
970 }
971
972 #[cfg(feature = "test-faults")]
975 #[must_use]
976 pub fn with_test_timer_gate(mut self, gate: Arc<tokio::sync::Mutex<()>>) -> Self {
977 self.test_timer_gate = Some(gate);
978 self
979 }
980
981 #[cfg(feature = "test-faults")]
982 async fn fault(
983 &self,
984 point: FaultPoint,
985 op: &Operation,
986 a: &Authenticated,
987 ) -> Result<(), ServerError> {
988 match &self.faults {
989 Some(hooks) => hooks.at_boxed(point, op, a.test_directives()).await,
990 None => Ok(()),
991 }
992 }
993
994 #[must_use]
1004 pub fn with_write_gate(mut self) -> Self {
1005 self.gate = Some(Arc::new(gate::WriteGate::new()));
1006 self
1007 }
1008
1009 pub fn with_auth(&self, auth: AuthMode) -> Result<Self, ServerError>
1019 where
1020 B: Clone,
1021 N: Clone,
1022 H: Clone,
1023 {
1024 let mut cfg = self.cfg.clone();
1025 cfg.auth = auth;
1026 cfg.url_tokens = None;
1028 if !matches!(cfg.auth, AuthMode::AuthV2(_)) {
1033 cfg.grants = None;
1034 }
1035 let mut sibling = Self::new(
1036 self.blobs.clone(),
1037 self.meta.clone(),
1038 self.hooks.clone(),
1039 cfg,
1040 Arc::clone(&self.clock),
1041 Arc::clone(&self.metrics),
1042 )?;
1043 sibling.shards = Arc::clone(&self.shards);
1044 sibling.gate.clone_from(&self.gate);
1045 #[cfg(feature = "remote-hooks")]
1046 {
1047 sibling.inspectors.clone_from(&self.inspectors);
1048 sibling.inspect_limit = self.inspect_limit;
1049 }
1050 #[cfg(feature = "test-faults")]
1051 {
1052 sibling.faults.clone_from(&self.faults);
1053 sibling.test_timer_gate.clone_from(&self.test_timer_gate);
1054 }
1055 #[cfg(feature = "http-objects")]
1056 sibling.http_seams.clone_from(&self.http_seams);
1057 #[cfg(feature = "published-view")]
1058 sibling.published.clone_from(&self.published);
1059 Ok(sibling)
1060 }
1061
1062 pub fn authenticate(&self, meta: &RequestMeta<'_>) -> Result<Authenticated, ServerError> {
1074 tracing::debug!(stage = "authenticate", procedure = method(meta.procedure));
1075 let result = self.authenticate_inner(meta);
1076 if let Err(err) = &result {
1077 self.outcome_for(meta.procedure, "none", "-")
1078 .record(Err(err));
1079 }
1080 result
1081 }
1082
1083 fn authenticate_inner(&self, meta: &RequestMeta<'_>) -> Result<Authenticated, ServerError> {
1084 let signed = auth::signed_request(&self.cfg.auth, meta);
1085 if !signed && (meta.header)("x-write-grant").is_some() {
1088 return Err(ServerError::unauthenticated(
1089 "write grant requires auth v2 authorization",
1090 ));
1091 }
1092 let repo = self
1093 .cfg
1094 .addressing
1095 .resolve((meta.header)("x-repository").as_deref(), signed)?;
1096 let expected_repository = match (&self.cfg.addressing, &self.cfg.auth) {
1097 (Addressing::Single { .. }, AuthMode::AuthV2(cfg)) => cfg.repository(),
1098 _ => &repo.identity,
1099 }
1100 .to_owned();
1101 #[cfg(feature = "test-faults")]
1102 let directives = TestDirectives::from_headers(meta.header)?;
1103 #[cfg(feature = "test-faults")]
1104 let skew = directives.clock_skew_ms;
1105 #[cfg(not(feature = "test-faults"))]
1106 let skew = 0;
1107 let now = self.clock.now_ms().saturating_add(skew);
1108 let mut a = auth::authenticate(&self.cfg.auth, meta, now, repo, &expected_repository)?;
1109 if a.principal
1110 .ed25519()
1111 .is_some_and(|key| self.cfg.admin_keys.contains(key))
1112 {
1113 return Err(ServerError::unauthenticated(
1114 "admin key cannot authenticate client calls",
1115 ));
1116 }
1117 if a.principal.ed25519().is_some_and(|key| {
1118 self.cfg
1119 .scanner_retrieval
1120 .as_ref()
1121 .is_some_and(|config| config.scanner_keys().any(|public| &public == key))
1122 }) {
1123 return Err(ServerError::unauthenticated(
1124 "scanner key cannot authenticate client calls",
1125 ));
1126 }
1127 if matches!(self.cfg.auth, AuthMode::AuthV2(_))
1130 && a.auth.is_some()
1131 && meta.procedure.is_write()
1132 && meta.procedure != Procedure::SetRepoVisibility
1133 {
1134 a.credential_capture =
1135 admission::capture_credentials(meta, &self.cfg.admission_credential_headers);
1136 }
1137 a.ref_hint =
1139 (meta.header)("x-mkit-ref").filter(|name| name.len() <= refs::MAX_REF_NAME_BYTES);
1140 a.business_skew_ms = skew;
1141 #[cfg(feature = "test-faults")]
1142 a.set_test_directives(directives);
1143 Ok(a)
1144 }
1145
1146 pub async fn list_refs(
1156 &self,
1157 a: &Authenticated,
1158 prefix: &str,
1159 ) -> Result<Vec<RefEntry>, ServerError> {
1160 let mut out = Vec::new();
1161 let mut token = None;
1162 loop {
1163 let page = self
1164 .list_refs_page(a, prefix, Some(self.cfg.list_page_limit), token.as_deref())
1165 .await?;
1166 out.extend(page.refs);
1167 match page.next {
1168 Some(next) => token = Some(next),
1169 None => return Ok(out),
1170 }
1171 }
1172 }
1173
1174 #[allow(clippy::too_many_lines)] pub(crate) async fn list_refs_page(
1178 &self,
1179 a: &Authenticated,
1180 prefix: &str,
1181 page_size: Option<u32>,
1182 token: Option<&[u8]>,
1183 ) -> Result<list::ListPage, ServerError> {
1184 let kind = OpKind::ListRefs {
1185 prefix: prefix.to_owned(),
1186 };
1187 self.observe(a, async {
1188 if prefix.trim_end_matches('/').len() > refs::MAX_REF_NAME_BYTES {
1189 return Err(ServerError::invalid_argument(refs::REF_NAME_TOO_LONG));
1190 }
1191 if !refs::validate_ref_prefix(prefix) {
1192 return Err(ServerError::invalid_argument(
1193 "prefix is invalid (SPEC-REFS §3)",
1194 ));
1195 }
1196 let op = self.identify(a, kind)?;
1197 let authorized = self.authorize_read(&op).await?;
1198 let scan = refs::list_scan_prefix(prefix);
1199 let last = token
1200 .map(|bytes| {
1201 list::decode_token(&op.repo, &scan, bytes)
1202 .ok_or_else(|| ServerError::invalid_argument("invalid page token"))
1203 })
1204 .transpose()?;
1205 #[cfg(feature = "test-faults")]
1206 {
1207 if a.test_directives().lease_recovered {
1208 self.mark_lease_table_recovered(&op.repo.namespace).await?;
1209 }
1210 if let Some(epoch) = a.test_directives().bump_epoch {
1211 self.test_bump_epoch(&op.repo.namespace, epoch).await?;
1212 }
1213 }
1214 #[cfg(feature = "test-faults")]
1215 faults::run_timers(
1216 a.test_directives(),
1217 &self.meta,
1218 &self.blobs,
1219 self.shards.as_ref(),
1220 &op.repo,
1221 self.clock.as_ref(),
1222 ms(self.clock.now_ms().saturating_add(a.business_skew_ms)),
1223 self.test_timer_gate.as_deref(),
1224 )
1225 .await?;
1226 let requested = page_size.unwrap_or(0);
1227 let limit = if requested == 0 {
1228 self.cfg.max_list_refs_page_size
1229 } else {
1230 requested.min(self.cfg.max_list_refs_page_size)
1231 };
1232 let partitions = self.shards.ref_index_partitions(&op.repo);
1233 #[cfg(feature = "published-view")]
1234 if authorized.facts.caller_view != CallerView::Writer
1235 && self
1236 .published
1237 .as_ref()
1238 .is_some_and(|s| s.inspection_configured() && !s.uses_published_values())
1239 {
1240 return Err(ServerError::unavailable("published view unavailable"));
1241 }
1242 #[cfg(feature = "published-view")]
1243 let source = self.published.as_deref().filter(|_| {
1244 self.visibility_gates_reads()
1245 && op.auth.is_none()
1246 && matches!(op.principal, Principal::Anonymous)
1247 && (self.publication_policy.is_none()
1248 || self
1249 .published
1250 .as_ref()
1251 .is_some_and(|s| s.uses_published_values()))
1252 });
1253 let view = crate::store::view::ViewStore {
1254 store: &self.meta,
1255 repo: &op.repo,
1256 writer: authorized.facts.caller_view == CallerView::Writer,
1257 policy: self.publication_policy.as_deref(),
1258 };
1259 let result = if partitions.len() == 1 {
1260 let bucket = list::RefBucket {
1261 store: &view,
1262 partition: &partitions[0],
1263 };
1264 list::page(
1265 &[bucket],
1266 &op.repo,
1267 &scan,
1268 last.as_deref(),
1269 limit,
1270 list::MAX_RESPONSE_BYTES,
1271 )
1272 .await
1273 } else {
1274 #[cfg(feature = "published-view")]
1275 let buckets = partitions
1276 .iter()
1277 .map(|partition| published::ReaderBucket {
1278 store: &view,
1279 partition,
1280 source,
1281 now_ms: ms(self.clock.now_ms()),
1282 })
1283 .collect::<Vec<_>>();
1284 #[cfg(not(feature = "published-view"))]
1285 let buckets = partitions
1286 .iter()
1287 .map(|partition| list::IndexBucket {
1288 store: &view,
1289 partition,
1290 })
1291 .collect::<Vec<_>>();
1292 list::page(
1293 &buckets,
1294 &op.repo,
1295 &scan,
1296 last.as_deref(),
1297 limit,
1298 list::MAX_RESPONSE_BYTES,
1299 )
1300 .await
1301 };
1302 let mut page = result.map_err(|err| {
1303 tracing::warn!(detail = %err, "ref listing scan failed");
1304 ServerError::unavailable("ref listing unavailable")
1305 })?;
1306 for entry in &mut page.refs {
1307 entry.name = strip_listed_prefix(&entry.name, prefix)
1308 .ok_or_else(|| ServerError::unavailable("ref listing unavailable"))?
1309 .to_owned();
1310 }
1311 Ok(page)
1312 })
1313 .await
1314 }
1315
1316 #[cfg(feature = "published-view")]
1318 #[must_use]
1319 pub fn with_published_source(mut self, source: Arc<dyn published::PublishedSource>) -> Self {
1320 self.published = Some(source);
1321 self
1322 }
1323
1324 pub fn with_publication_policy(
1327 mut self,
1328 policy: Arc<dyn clearance::PublicationPolicy>,
1329 ) -> Result<Self, ServerError> {
1330 if self.cfg.indexed.is_none() || self.cfg.write_policy == WritePolicy::Open {
1331 return Err(ServerError::invalid_argument(
1332 "inspection requires indexed mode and restricted writes",
1333 ));
1334 }
1335 self.publication_policy = Some(policy);
1336 Ok(self)
1337 }
1338
1339 pub async fn read_ref(
1346 &self,
1347 a: &Authenticated,
1348 name: &str,
1349 ) -> Result<Option<Hash>, ServerError> {
1350 let kind = OpKind::ReadRef {
1351 name: name.to_owned(),
1352 };
1353 self.observe(a, async {
1354 check_ref_name(name)?;
1355 let op = self.identify(a, kind)?;
1356 let authorized = self.authorize_read(&op).await?;
1357 #[cfg(feature = "published-view")]
1358 if let Some(source) = &self.published {
1359 if source.inspection_configured()
1360 && !source.uses_published_values()
1361 && authorized.facts.caller_view != CallerView::Writer
1362 {
1363 return Err(ServerError::unavailable("published view unavailable"));
1364 }
1365 if self.visibility_gates_reads()
1366 && op.auth.is_none()
1367 && matches!(op.principal, Principal::Anonymous)
1368 && source.read_ref_enabled()
1369 && (self.publication_policy.is_none() || source.uses_published_values())
1370 {
1371 let partition = self.shards.ref_index(&op.repo, name);
1372 if let Some(rows) = source
1373 .bucket(&op.repo, &partition, ms(self.clock.now_ms()))
1374 .await
1375 .map_err(meta_error)?
1376 {
1377 if self.meta.capabilities().atomic_multi_key
1378 && !source.uses_published_values()
1379 {
1380 return Err(ServerError::unavailable("published view unavailable"));
1381 }
1382 return Ok(rows.into_iter().find(|(n, _)| n == name).map(|(_, id)| id));
1383 }
1384 }
1385 }
1386 let p = self.shards.ref_shard(&op.repo, name);
1387 let view = crate::store::view::ViewStore {
1388 store: &self.meta,
1389 repo: &op.repo,
1390 writer: authorized.facts.caller_view == CallerView::Writer,
1391 policy: self.publication_policy.as_deref(),
1392 };
1393 read::read_ref(&view, &p, &op.repo.name, name)
1394 .await
1395 .map_err(meta_error)
1396 })
1397 .await
1398 }
1399
1400 pub async fn update_ref(
1405 &self,
1406 a: &Authenticated,
1407 upd: RefUpdate,
1408 ) -> Result<UpdateRefResult, ServerError> {
1409 self.update_ref_with_meta(a, upd)
1410 .await
1411 .map(|(result, _)| result)
1412 }
1413
1414 pub async fn update_ref_with_meta(
1416 &self,
1417 a: &Authenticated,
1418 upd: RefUpdate,
1419 ) -> Result<(UpdateRefResult, ResponseMeta), ServerError> {
1420 self.observe(a, async {
1421 check_ref_name(&upd.name)?;
1422 if upd.new.is_none()
1423 && !matches!(upd.condition, mkit_core::refs::RefWriteCondition::Match(_))
1424 {
1425 return Err(ServerError::invalid_argument(
1426 "delete requires MATCH and an empty new_id",
1427 ));
1428 }
1429 match self.write(a, OpKind::UpdateRef(upd)).await? {
1430 (StoredResult::UpdateRef(result), meta) => Ok((result, meta)),
1431 (other, _) => Err(stored_mismatch(&other)),
1432 }
1433 })
1434 .await
1435 }
1436
1437 pub async fn advance_refs(
1444 &self,
1445 a: &Authenticated,
1446 head: RefUpdate,
1447 packmap: RefUpdate,
1448 ) -> Result<AdvanceOutcome, ServerError> {
1449 self.advance_refs_with_meta(a, head, packmap)
1450 .await
1451 .map(|(result, _)| result)
1452 }
1453
1454 pub async fn advance_refs_with_meta(
1456 &self,
1457 a: &Authenticated,
1458 head: RefUpdate,
1459 packmap: RefUpdate,
1460 ) -> Result<(AdvanceOutcome, ResponseMeta), ServerError> {
1461 self.advance_refs_with_tickets_with_meta(a, head, packmap, Vec::new())
1462 .await
1463 }
1464
1465 pub async fn advance_refs_with_tickets(
1470 &self,
1471 a: &Authenticated,
1472 head: RefUpdate,
1473 packmap: RefUpdate,
1474 tickets: Vec<Hash>,
1475 ) -> Result<AdvanceOutcome, ServerError> {
1476 self.advance_refs_with_tickets_with_meta(a, head, packmap, tickets)
1477 .await
1478 .map(|(result, _)| result)
1479 }
1480
1481 pub async fn advance_refs_with_tickets_with_meta(
1483 &self,
1484 a: &Authenticated,
1485 head: RefUpdate,
1486 packmap: RefUpdate,
1487 tickets: Vec<Hash>,
1488 ) -> Result<(AdvanceOutcome, ResponseMeta), ServerError> {
1489 self.observe(a, async {
1490 check_ref_name(&head.name)?;
1491 check_ref_name(&packmap.name)?;
1492 if head.new.is_none() != packmap.new.is_none()
1493 || (head.new.is_none()
1494 && (!matches!(head.condition, mkit_core::refs::RefWriteCondition::Match(_))
1495 || !matches!(
1496 packmap.condition,
1497 mkit_core::refs::RefWriteCondition::Match(_)
1498 )))
1499 {
1500 return Err(ServerError::invalid_argument(
1501 "delete requires MATCH and an empty new_id",
1502 ));
1503 }
1504 if head.new.is_none() && !tickets.is_empty() {
1505 return Err(ServerError::invalid_argument("delete consumes no tickets"));
1506 }
1507 if tickets.len() > crate::store::outbox::MAX_TICKETS_PER_ADVANCE {
1508 return Err(ServerError::invalid_argument(
1509 "too many tickets in one advance",
1510 ));
1511 }
1512 let mut distinct = std::collections::BTreeSet::new();
1513 if tickets.iter().any(|id| !distinct.insert(id)) {
1514 return Err(ServerError::invalid_argument("duplicate ticket id"));
1515 }
1516 if !tickets.is_empty()
1519 || self.cfg.sharding == Sharding::D34
1520 || self.cfg.ref_policy.is_some()
1521 || self.cfg.indexed.is_some()
1522 {
1523 let head_branch = head.name.strip_prefix("refs/heads/");
1524 let packmap_branch = packmap
1525 .name
1526 .strip_prefix(mkit_core::refs::PACKMAP_REF_PREFIX);
1527 if head_branch.is_none() || head_branch != packmap_branch {
1528 if a.write_grant.is_some()
1529 && self.cfg.grants.is_some()
1530 && matches!(self.cfg.addressing, Addressing::Multi(_))
1531 {
1532 return Err(ServerError::permission_denied(
1533 "write grant rejected: ref scope",
1534 ));
1535 }
1536 return Err(ServerError::invalid_argument(if tickets.is_empty() {
1537 "AdvanceRefs pairs refs/heads/<x> with refs/mkit/packmap/<x> on this server"
1538 } else {
1539 "ticketed advance requires a branch head and its packmap"
1540 }));
1541 }
1542 }
1543 match self
1544 .write(
1545 a,
1546 OpKind::AdvanceRefs {
1547 head,
1548 packmap,
1549 tickets,
1550 },
1551 )
1552 .await?
1553 {
1554 (StoredResult::AdvanceRefs(outcome), meta) => Ok((outcome, meta)),
1555 (other, _) => Err(stored_mismatch(&other)),
1556 }
1557 })
1558 .await
1559 }
1560
1561 pub async fn pack_exists(&self, a: &Authenticated, key: PackKey) -> Result<bool, ServerError> {
1568 self.observe(a, async {
1569 let op = self.identify(a, OpKind::PackExists { key })?;
1570 let authorized = self.authorize_read(&op).await?;
1571 if !self
1572 .pack_is_member(a, &key, authorized.facts.caller_view)
1573 .await?
1574 {
1575 return Ok(false);
1576 }
1577 let head = self.blobs.head(&key.into()).await;
1578 Ok(head
1579 .map_err(|e| store_error(StorageOp::BlobHead, e))?
1580 .is_some())
1581 })
1582 .await
1583 }
1584
1585 pub async fn issue_object_url(
1599 &self,
1600 a: &Authenticated,
1601 target: UrlTarget,
1602 ttl_seconds: u32,
1603 ) -> Result<MintedToken, ServerError> {
1604 self.observe(a, async {
1605 if self.cfg.url_tokens.is_none() {
1606 return Err(ServerError::unimplemented("URL tokens not configured"));
1607 }
1608 let op = self.identify(
1609 a,
1610 OpKind::IssueObjectUrl {
1611 target: target.clone(),
1612 ttl_seconds,
1613 },
1614 )?;
1615 self.issue_url(
1616 &op,
1617 &a.repo().identity,
1618 &target,
1619 ttl_seconds,
1620 a.business_now_ms,
1621 )
1622 .await
1623 })
1624 .await
1625 }
1626
1627 async fn issue_url(
1630 &self,
1631 op: &Operation,
1632 repository: &str,
1633 target: &UrlTarget,
1634 ttl_seconds: u32,
1635 now_ms: i64,
1636 ) -> Result<MintedToken, ServerError> {
1637 let Some(tokens) = &self.cfg.url_tokens else {
1638 return Err(ServerError::unimplemented("URL tokens not configured"));
1639 };
1640 let read = self.authorize_read(op).await?;
1641 let epoch = match read.epoch {
1642 Some(epoch) => epoch,
1643 None => self.stored_grant_epoch(&op.repo.namespace).await?,
1644 };
1645 let AuthMode::AuthV2(auth) = &self.cfg.auth else {
1646 return Err(internal("URL tokens without auth v2"));
1647 };
1648 tokens.mint(
1649 auth.audience(),
1650 repository,
1651 target,
1652 epoch,
1653 now_ms,
1654 ttl_seconds,
1655 )
1656 }
1657
1658 pub async fn set_repo_visibility(
1670 &self,
1671 a: &Authenticated,
1672 req: VisibilityRequest,
1673 ) -> Result<(), ServerError> {
1674 self.observe(a, async {
1675 if !self.visibility_applies() {
1676 return Err(ServerError::failed_precondition(
1677 "repository visibility is not supported by this deployment",
1678 ));
1679 }
1680 match (&req, a.auth.is_some()) {
1681 (VisibilityRequest::Statement(_), true) => {
1682 return Err(ServerError::invalid_argument(
1683 "signed_statement is not allowed on a signed request",
1684 ));
1685 }
1686 (VisibilityRequest::Envelope(_), false) => {
1687 return Err(ServerError::unauthenticated(
1688 "visibility requires auth v2 authorization",
1689 ));
1690 }
1691 _ => {}
1692 }
1693 let repo = &a.repo().repo;
1694 let namespace = Namespace::parse(repo.namespace.as_str())
1695 .map_err(|_| internal("invalid resolved Multi namespace"))?;
1696 if let Addressing::Multi(multi) = &self.cfg.addressing
1697 && let NamespacePolicy::Allowlist(allowed) = &multi.namespace_policy
1698 && !allowed.contains(&namespace)
1699 {
1700 return Err(ServerError::permission_denied("namespace not served"));
1701 }
1702 let p = self.shards.coordinator(&repo.namespace);
1703 let _gate = match &self.gate {
1704 Some(gate) => Some(gate.enter(&p).await),
1705 None => None,
1706 };
1707 match req {
1708 VisibilityRequest::Envelope(visibility) => {
1709 self.visibility_envelope(a, repo, &p, visibility).await
1710 }
1711 VisibilityRequest::Statement(statement) => {
1712 self.visibility_statement(a, repo, &p, &statement).await
1713 }
1714 }
1715 })
1716 .await
1717 }
1718
1719 #[allow(clippy::too_many_lines)] async fn visibility_envelope(
1724 &self,
1725 a: &Authenticated,
1726 repo: &RepoId,
1727 p: &Partition,
1728 visibility: Visibility,
1729 ) -> Result<(), ServerError> {
1730 if a.write_grant.is_some() {
1731 return Err(ServerError::permission_denied(
1732 "a grant never authorizes SetRepoVisibility",
1733 ));
1734 }
1735 let op = self.identify(a, OpKind::SetRepoVisibility { visibility })?;
1736 let auth = a.auth.as_ref().ok_or_else(|| {
1737 ServerError::unauthenticated("visibility requires auth v2 authorization")
1738 })?;
1739 let replay_key = keys::replay(&auth.replay_scope);
1740 let rv_key = keys::repo_visibility(&repo.name);
1741 let rows = self
1742 .meta
1743 .get_many(
1744 p,
1745 &[
1746 replay_key.clone(),
1747 rv_key.clone(),
1748 keys::authority_generation(),
1749 keys::lease_recovery(),
1750 ],
1751 )
1752 .await
1753 .map_err(meta_error)?;
1754 let mut rows = rows.into_iter();
1755 let record = rows
1756 .next()
1757 .flatten()
1758 .map(|v| codec::decode_replay_record(&v))
1759 .transpose()
1760 .map_err(meta_error)?;
1761 let mut stored = rows.next().flatten();
1762 match classify(record.as_ref(), &auth.fingerprint) {
1763 ReplayDecision::New => {}
1764 ReplayDecision::Return(StoredResult::RepoVisibility) => return Ok(()),
1765 ReplayDecision::Return(other) => return Err(stored_mismatch(&other)),
1766 ReplayDecision::FingerprintMismatch => {
1767 return Err(ServerError::invalid_argument(
1768 "nonce reused for a different operation",
1769 ));
1770 }
1771 ReplayDecision::Resume | ReplayDecision::RetryLater => {
1772 return Err(ServerError::aborted_retryable(
1773 "operation already in flight; retry",
1774 ));
1775 }
1776 }
1777 Box::pin(self.ensure_authority_activation(&repo.namespace)).await?;
1778 let facts = self.authorize_visibility_envelope(&op).await?;
1779 let mut replans = 0;
1780 loop {
1781 let (mut batch, prune) = self
1782 .plan_visibility(p, auth, repo, visibility, stored.as_ref())
1783 .await?;
1784 let fence_rows = self
1785 .meta
1786 .get_many(p, &[keys::authority_generation(), keys::lease_recovery()])
1787 .await
1788 .map_err(meta_error)?;
1789 let [generation_value, mode_value] = fence_rows.as_slice() else {
1790 return Err(internal("visibility fence row count"));
1791 };
1792 let mode = mode_value
1793 .as_ref()
1794 .map(codec::decode_lease_recovery)
1795 .transpose()
1796 .map_err(meta_error)?;
1797 if facts.authority_generation.is_none()
1798 && (generation_value.is_some()
1799 || mode.is_some_and(|m| m.authority_fence == Some(true)))
1800 {
1801 return Err(ServerError::unavailable(
1802 "persisted authority fence requires enabled executor",
1803 ));
1804 }
1805 batch = batch.require(lease::observed_guard(
1806 keys::lease_recovery(),
1807 mode_value.as_ref(),
1808 ));
1809 if let Some(generation) = facts.authority_generation {
1810 let current = generation_value
1811 .as_ref()
1812 .map(codec::decode_u64)
1813 .transpose()
1814 .map_err(meta_error)?
1815 .unwrap_or(0);
1816 if current != generation {
1817 return Err(crate::authority::moved());
1818 }
1819 batch = batch.require(lease::observed_guard(
1820 keys::authority_generation(),
1821 generation_value.as_ref(),
1822 ));
1823 }
1824 match self.meta.apply(p, batch).await {
1825 Ok(BatchOutcome::Committed) => {
1826 self.invalidate_local_cache(repo).await;
1827 return Ok(());
1828 }
1829 Ok(BatchOutcome::DeadlinePassed { .. }) => {
1830 return Err(ServerError::unavailable("commit deadline passed; retry"));
1831 }
1832 Ok(BatchOutcome::PreconditionFailed { index, .. }) => {
1833 if index == 1 {
1834 let value = self.meta.get(p, &replay_key).await.map_err(meta_error)?;
1837 let record = value
1838 .as_ref()
1839 .map(codec::decode_replay_record)
1840 .transpose()
1841 .map_err(meta_error)?;
1842 return match classify(record.as_ref(), &auth.fingerprint) {
1843 ReplayDecision::Return(StoredResult::RepoVisibility) => Ok(()),
1844 ReplayDecision::Return(other) => Err(stored_mismatch(&other)),
1845 ReplayDecision::FingerprintMismatch => {
1846 Err(ServerError::invalid_argument(
1847 "nonce reused for a different operation",
1848 ))
1849 }
1850 _ => Err(ServerError::aborted_retryable(
1851 "operation already in flight; retry",
1852 )),
1853 };
1854 }
1855 replans += 1;
1856 if replans > MAX_REPLAN {
1857 return Err(ServerError::aborted_retryable("write contention; retry"));
1858 }
1859 stored = self.meta.get(p, &rv_key).await.map_err(meta_error)?;
1860 }
1861 Err(StoreError::Full) => {
1862 return Err(self.partition_full(p, prune).await);
1863 }
1864 Err(e) => return Err(meta_error(e)),
1865 }
1866 }
1867 }
1868
1869 async fn authorize_visibility_envelope(
1873 &self,
1874 op: &Operation,
1875 ) -> Result<AuthzFacts, ServerError> {
1876 let owner = matches!(
1877 Namespace::parse(op.repo.namespace.as_str()),
1878 Ok(Namespace::Ed25519(key)) if op.principal.ed25519() == Some(&key)
1879 );
1880 if self.cfg.authorizer_role == AuthorizerRole::Check && !owner {
1881 return Err(ServerError::permission_denied(
1882 "SetRepoVisibility not permitted",
1883 ));
1884 }
1885 let mut authorized = op.clone();
1886 authorized.authz = AuthzFacts {
1887 authority_generation: None,
1888 grant: None,
1889 owner,
1890 caller_view: if owner {
1893 CallerView::Writer
1894 } else {
1895 CallerView::Reader
1896 },
1897 };
1898 let returned = self
1899 .hooks
1900 .authorizer()
1901 .authorize(&authorized)
1902 .await
1903 .map_err(ServerError::strip_admission_shape)?;
1904 self.merge_authority_facts(&mut authorized.authz, &returned)?;
1905 Ok(authorized.authz)
1906 }
1907
1908 async fn plan_visibility(
1914 &self,
1915 p: &Partition,
1916 auth: &VerifiedAuth,
1917 repo: &RepoId,
1918 visibility: Visibility,
1919 stored: Option<&Value>,
1920 ) -> Result<(Batch, Option<Batch>), ServerError> {
1921 let now = ms(self.clock.now_ms());
1922 let window = u64::try_from(self.cfg.max_apply_window.as_millis()).unwrap_or(u64::MAX);
1923 let deadline = now
1924 .saturating_add(window)
1925 .min(ms(auth.expires_at_ms).saturating_add(MAX_CLOCK_LEAD_MS.unsigned_abs()));
1926 let replay_key = keys::replay(&auth.replay_scope);
1927 let rv_key = keys::repo_visibility(&repo.name);
1928 let row = stored
1929 .map(codec::decode_repo_visibility)
1930 .transpose()
1931 .map_err(meta_error)?;
1932 let mut batch = Batch::new()
1933 .require(Precondition::NotAfter(deadline))
1934 .require(Precondition::Absent(replay_key.clone()))
1935 .require(match stored {
1936 Some(value) => Precondition::Equals(rv_key.clone(), value.clone()),
1937 None => Precondition::Absent(rv_key.clone()),
1938 })
1939 .put(
1940 rv_key.clone(),
1941 codec::encode_repo_visibility(&codec::RepoVisibilityV1 {
1942 visibility: stored_visibility(visibility),
1943 last_created_ms: row.as_ref().map_or(0, |r| r.last_created_ms).max(now),
1947 last_statement_id: row.and_then(|r| r.last_statement_id),
1948 changed_ms: Some(now),
1949 }),
1950 )
1951 .put(
1952 replay_key.clone(),
1953 codec::encode_replay_record(&ReplayRecord {
1954 fingerprint: auth.fingerprint,
1955 expires_at_ms: auth.expires_at_ms,
1956 state: ReplayState::Committed(StoredResult::RepoVisibility),
1957 }),
1958 )
1959 .put(
1960 keys::replay_expiry(ms(auth.expires_at_ms), &auth.replay_scope),
1961 Value::default(),
1962 );
1963 self.plan_listing_visibility(p, repo, &mut batch).await?;
1964 let expired = read::expired_replay_keys(&self.meta, p, now, 32)
1965 .await
1966 .map_err(meta_error)?;
1967 let purge = self
1968 .plan_repository_purge(
1969 p,
1970 repo,
1971 crate::purge::Trigger::VisibilityChange,
1972 &mkit_core::hash::to_hex(&auth.replay_scope),
1973 now,
1974 )
1975 .await?;
1976 batch.preconditions.extend(purge.preconditions);
1977 batch.writes.extend(purge.writes);
1978 let mut prune = Batch::new().require(Precondition::NotAfter(deadline));
1979 for (index, target) in &expired {
1980 if batch.preconditions.len() + batch.writes.len() + 2 > MAX_BATCH_OPS {
1981 break;
1982 }
1983 batch = batch.delete(index.clone()).delete(target.clone());
1984 prune = prune.delete(index.clone()).delete(target.clone());
1985 }
1986 Ok((batch, (!prune.writes.is_empty()).then_some(prune)))
1987 }
1988
1989 async fn visibility_statement(
1993 &self,
1994 a: &Authenticated,
1995 repo: &RepoId,
1996 p: &Partition,
1997 statement: &str,
1998 ) -> Result<(), ServerError> {
1999 let rejected = |e: mkit_attest::grant::GrantError| {
2000 ServerError::permission_denied(format!("visibility statement rejected: {}", e.reason()))
2001 };
2002 if statement.len() > mkit_attest::grant::MAX_GRANT_HEADER_BYTES {
2003 return Err(ServerError::permission_denied(
2004 "visibility statement rejected: too long",
2005 ));
2006 }
2007 let grants = self.cfg.grants.as_ref().ok_or_else(|| {
2008 ServerError::permission_denied(
2009 "visibility statement rejected: no owner schemes configured",
2010 )
2011 })?;
2012 let identity = RepositoryIdentity::parse(&a.repo().identity)
2013 .map_err(|_| internal("invalid resolved repository identity"))?;
2014 let verified =
2015 verify_visibility_statement(grants.verifier(), statement, &identity, a.business_now_ms)
2016 .map_err(rejected)?;
2017 if let Some(namespace) = identity.namespace() {
2018 self.require_client_owner_key(namespace)?;
2019 }
2020 let created = u64::try_from(verified.statement().created_ms)
2021 .map_err(|_| rejected(mkit_attest::grant::GrantError::DecimalOutOfRange))?;
2022 let id = to_hex(verified.id());
2023 let rv_key = keys::repo_visibility(&repo.name);
2024 let window = u64::try_from(self.cfg.max_apply_window.as_millis()).unwrap_or(u64::MAX);
2025 let mut replans = 0;
2026 loop {
2027 let stored = self.meta.get(p, &rv_key).await.map_err(meta_error)?;
2028 let row = stored
2029 .as_ref()
2030 .map(codec::decode_repo_visibility)
2031 .transpose()
2032 .map_err(meta_error)?;
2033 match &row {
2034 Some(r)
2036 if r.last_created_ms == created
2037 && r.last_statement_id.as_deref() == Some(id.as_str())
2038 && r.visibility == stored_visibility(verified.statement().visibility) =>
2039 {
2040 return Ok(());
2041 }
2042 Some(r) if created <= r.last_created_ms => {
2043 return Err(ServerError::permission_denied(
2044 "visibility statement rejected: not newer than the stored statement",
2045 ));
2046 }
2047 _ => {}
2048 }
2049 let deadline = ms(self.clock.now_ms()).saturating_add(window);
2050 let rv_guard = match &stored {
2051 Some(value) => Precondition::Equals(rv_key.clone(), value.clone()),
2052 None => Precondition::Absent(rv_key.clone()),
2053 };
2054 let mut batch = Batch::new()
2055 .require(Precondition::NotAfter(deadline))
2056 .require(rv_guard)
2057 .put(
2058 rv_key.clone(),
2059 codec::encode_repo_visibility(&codec::RepoVisibilityV1 {
2060 visibility: stored_visibility(verified.statement().visibility),
2061 last_created_ms: created,
2062 last_statement_id: Some(id.clone()),
2063 changed_ms: Some(ms(self.clock.now_ms())),
2064 }),
2065 );
2066 self.plan_listing_visibility(p, repo, &mut batch).await?;
2067 let purge = self
2068 .plan_repository_purge(
2069 p,
2070 repo,
2071 crate::purge::Trigger::VisibilityChange,
2072 &id,
2073 ms(self.clock.now_ms()),
2074 )
2075 .await?;
2076 batch.preconditions.extend(purge.preconditions);
2077 batch.writes.extend(purge.writes);
2078 match self.meta.apply(p, batch).await {
2079 Ok(BatchOutcome::Committed) => {
2080 self.invalidate_local_cache(repo).await;
2081 return Ok(());
2082 }
2083 Ok(BatchOutcome::DeadlinePassed { .. }) => {
2084 return Err(ServerError::unavailable("commit deadline passed; retry"));
2085 }
2086 Ok(BatchOutcome::PreconditionFailed { .. }) => {
2087 replans += 1;
2088 if replans > MAX_REPLAN {
2089 return Err(ServerError::aborted_retryable("write contention; retry"));
2090 }
2091 }
2092 Err(StoreError::Full) => return Err(self.partition_full(p, None).await),
2093 Err(e) => return Err(meta_error(e)),
2094 }
2095 }
2096 }
2097
2098 pub async fn open_upload(
2110 &self,
2111 a: &Authenticated,
2112 pack_id: Option<&[u8]>,
2113 total_bytes: Option<u64>,
2114 ) -> Result<UploadSession<'_, B, N, H>, ServerError> {
2115 UploadSession::begin(self, a, pack_id, total_bytes).await
2116 }
2117
2118 pub async fn open_ticketed_upload(
2124 &self,
2125 a: &Authenticated,
2126 pack_id: Option<&[u8]>,
2127 total_bytes: Option<u64>,
2128 token: &[u8],
2129 ) -> Result<UploadSession<'_, B, N, H>, ServerError> {
2130 UploadSession::begin_ticketed(self, a, pack_id, total_bytes, token).await
2131 }
2132
2133 pub async fn download(
2143 &self,
2144 a: &Authenticated,
2145 key: PackKey,
2146 ) -> Result<DownloadStream, ServerError> {
2147 let mut outcome = self.outcome(a);
2148 let opened = async {
2149 let op = self.identify(a, OpKind::DownloadPack { key })?;
2150 let authorized = self.authorize_read(&op).await?;
2151 if !self
2152 .pack_is_member(a, &key, authorized.facts.caller_view)
2153 .await?
2154 {
2155 return Err(ServerError::not_found("pack not found"));
2156 }
2157 let body = self.blobs.get(&key.into(), None).await;
2158 match body.map_err(|e| store_error(StorageOp::BlobGet, e))? {
2159 Some(body) => Ok(body),
2160 None => Err(ServerError::not_found("pack not found")),
2161 }
2162 }
2163 .instrument(outcome.span.clone())
2164 .await;
2165 match opened {
2166 Ok(body) => {
2167 let max = self.cfg.download_chunk_max;
2168 Ok(DownloadStream::new(body, max, Some(outcome)))
2169 }
2170 Err(err) => {
2171 outcome.record(Err(&err));
2172 Err(err)
2173 }
2174 }
2175 }
2176
2177 pub async fn health(&self) -> HealthStatus {
2179 HealthStatus {
2180 blobs: self.blobs.probe().await.is_ok(),
2181 meta: self.meta.probe().await.is_ok(),
2182 }
2183 }
2184
2185 #[must_use]
2189 pub fn auth_mode(&self) -> &AuthMode {
2190 &self.cfg.auth
2191 }
2192
2193 #[cfg(feature = "ssh")]
2196 pub(crate) fn upload_limits(&self) -> UploadLimits {
2197 self.cfg.upload_limits
2198 }
2199
2200 #[cfg(all(test, feature = "ssh"))]
2202 pub(crate) fn meta_store(&self) -> &N {
2203 &self.meta
2204 }
2205
2206 pub fn capabilities(&self) -> PipelineCapabilities {
2208 PipelineCapabilities {
2209 atomic_advance: self.meta.capabilities().atomic_multi_key,
2210 }
2211 }
2212
2213 #[must_use]
2217 pub fn with_header(&self, err: ServerError, name: &str, value: &str) -> ServerError {
2218 match err.clone().try_with_header(name, value) {
2219 Ok(err) => err,
2220 Err(reason) => {
2221 let label = match reason {
2222 InvalidHeader::Name => "name",
2223 InvalidHeader::Reserved => "reserved",
2224 InvalidHeader::Value => "value",
2225 };
2226 tracing::warn!(header = ?name, %reason, "dropped an error response header");
2227 self.metrics
2228 .incr(METRIC_HEADER_DROPPED, &[("reason", label)], 1);
2229 err
2230 }
2231 }
2232 }
2233
2234 async fn observe<T>(
2236 &self,
2237 a: &Authenticated,
2238 fut: impl Future<Output = Result<T, ServerError>>,
2239 ) -> Result<T, ServerError> {
2240 let mut outcome = self.outcome(a);
2241 let result = fut.instrument(outcome.span.clone()).await;
2242 outcome.record(result.as_ref().map(|_| ()));
2243 result
2244 }
2245
2246 fn outcome(&self, a: &Authenticated) -> RequestOutcome {
2248 self.outcome_for(a.procedure(), a.principal.kind(), &a.repo().identity)
2249 }
2250
2251 fn outcome_for(
2253 &self,
2254 procedure: Procedure,
2255 principal: &'static str,
2256 repo: &str,
2257 ) -> RequestOutcome {
2258 let procedure = method(procedure);
2259 let span = tracing::info_span!("mkit.server.rpc", procedure, repo, principal);
2260 let (metrics, clock) = (self.metrics.clone(), self.clock.clone());
2261 RequestOutcome::new(span, procedure, metrics, clock, self.cfg.redactor.clone())
2262 }
2263
2264 fn write<'a>(
2268 &'a self,
2269 a: &'a Authenticated,
2270 kind: OpKind,
2271 ) -> crate::rt::BoxFuture<'a, Result<(StoredResult, ResponseMeta), ServerError>> {
2272 self.write_with(a, kind, None)
2273 }
2274
2275 fn write_with<'a>(
2280 &'a self,
2281 a: &'a Authenticated,
2282 kind: OpKind,
2283 implicit: Option<&'a [PendingPack]>,
2284 ) -> crate::rt::BoxFuture<'a, Result<(StoredResult, ResponseMeta), ServerError>> {
2285 Box::pin(self.write_inner(a, kind, implicit))
2286 }
2287
2288 #[allow(clippy::too_many_lines)] async fn write_inner(
2290 &self,
2291 a: &Authenticated,
2292 kind: OpKind,
2293 implicit: Option<&[PendingPack]>,
2294 ) -> Result<(StoredResult, ResponseMeta), ServerError> {
2295 let mut op = self.identify(a, kind)?;
2296 fault!(self, AfterAuthenticate, &op, a);
2297 let (kind, mut refs, p) = self.ref_writes(&op)?;
2298 let mut ahead = self.read_ahead(&op, &p, &refs, a.business_skew_ms).await?;
2299 if let Some(stored) = Self::replay_lookup(&op, ahead.as_ref())? {
2300 return Ok((stored, ResponseMeta::default()));
2301 }
2302 let lease = if self.cfg.sharding == Sharding::D34 {
2303 let observed = self.observe_lease(&op, &p, ahead.as_ref()).await?;
2304 if let Some((window, total)) = observed.quota_seed()
2305 && let Some(snapshot) = ahead.as_mut()
2306 {
2307 snapshot.namespace_seed = Some((
2308 window,
2309 codec::encode_namespace_view(NamespaceView {
2310 total,
2311 pushed: crate::quota::NamespaceUsage::default(),
2312 observed_at_ms: ms(self.clock.now_ms()),
2313 }),
2314 ));
2315 }
2316 op.creation = observed.creation(&self.cfg.addressing);
2317 op.leased_epoch = Some(observed.epoch());
2318 Some(observed)
2319 } else {
2320 op.creation = self.creation_facts(&op, ahead.as_ref()).await?;
2321 None
2322 };
2323 if self.cfg.sharding == Sharding::Single {
2324 let key = keys::grant_epoch();
2325 if let Some(snapshot) = ahead.as_ref().filter(|snapshot| snapshot.contains(&key)) {
2326 op.observed_epoch = Some(
2327 snapshot
2328 .get(&key)
2329 .map(codec::decode_u64)
2330 .transpose()
2331 .map_err(meta_error)?
2332 .unwrap_or(0),
2333 );
2334 }
2335 }
2336 let (authz, fast_forward) = self.authorize(&op).await?;
2337 op.authz = authz;
2338 fault!(self, AfterAuthorize, &op, a);
2339 if let Err(error) = self.check_ref_signers(&op) {
2341 return self
2342 .store_policy_denial(&op, a, &p, ahead, error, false)
2343 .await
2344 .map(|stored| (stored, ResponseMeta::default()));
2345 }
2346 let ticketed =
2347 matches!(&op.kind, OpKind::AdvanceRefs { tickets, .. } if !tickets.is_empty());
2348 let mut staged = crate::indexed::verify::StagedCommits::default();
2349 let mut ticket_ms = None;
2350 if ticketed {
2351 let snap = ahead
2352 .as_mut()
2353 .ok_or_else(|| internal("ticket advance requires atomic metadata"))?;
2354 if let Some(stored) = self.ticket_decision(&op, a, &p, snap).await? {
2355 return Ok((stored, ResponseMeta::default()));
2356 }
2357 if let Some(indexed) = self.cfg.indexed
2358 && let OpKind::AdvanceRefs { head, tickets, .. } = &op.kind
2359 {
2360 #[cfg(feature = "test-faults")]
2364 if a.test_directives().fault.as_deref() == Some("indexed-pending") {
2365 return Err(crate::indexed::pending(5_000));
2366 }
2367 let rows = tickets
2368 .iter()
2369 .map(|id| {
2370 snap.get(&keys::ticket(id))
2371 .ok_or_else(|| {
2372 ServerError::failed_precondition("invalid or expired upload ticket")
2373 })
2374 .and_then(|raw| codec::decode_ticket(raw).map_err(meta_error))
2375 })
2376 .collect::<Result<Vec<_>, _>>()?;
2377 let tip = head
2378 .new
2379 .ok_or_else(|| ServerError::invalid_argument("delete consumes no tickets"))?;
2380 ticket_ms = rows.iter().map(|ticket| ticket.created_at_ms).min();
2381 staged = if let Some(limit) = self.inspection_max_objects() {
2382 if indexed.verification == crate::indexed::VerificationMode::Scheduled {
2383 crate::indexed::scheduled::check_inspected(
2384 &self.blobs,
2385 &self.meta,
2386 self.shards.as_ref(),
2387 &op.repo,
2388 &p,
2389 &rows,
2390 tickets,
2391 tip,
2392 indexed,
2393 self.clock.as_ref(),
2394 self.metrics.as_ref(),
2395 limit as usize,
2396 )
2397 .await?
2398 } else {
2399 crate::indexed::verify::verify_ticketed_inspected(
2400 &self.blobs,
2401 &self.meta,
2402 self.shards.as_ref(),
2403 &op.repo,
2404 &p,
2405 &rows,
2406 tickets,
2407 tip,
2408 indexed,
2409 self.clock.as_ref(),
2410 self.metrics.as_ref(),
2411 limit as usize,
2412 )
2413 .await?
2414 }
2415 } else if indexed.verification == crate::indexed::VerificationMode::Scheduled {
2416 crate::indexed::scheduled::check(
2418 &self.blobs,
2419 &self.meta,
2420 self.shards.as_ref(),
2421 &op.repo,
2422 &p,
2423 &rows,
2424 tickets,
2425 tip,
2426 indexed,
2427 self.clock.as_ref(),
2428 self.metrics.as_ref(),
2429 )
2430 .await?
2431 } else {
2432 crate::indexed::verify::verify_ticketed(
2433 &self.blobs,
2434 &self.meta,
2435 self.shards.as_ref(),
2436 &op.repo,
2437 &p,
2438 &rows,
2439 tickets,
2440 tip,
2441 indexed,
2442 self.clock.as_ref(),
2443 self.metrics.as_ref(),
2444 )
2445 .await?
2446 };
2447 }
2448 } else if let Err(error) = self.check_ticketless_head(&op).await {
2449 return self
2450 .store_policy_denial(&op, a, &p, ahead, error, false)
2451 .await
2452 .map(|stored| (stored, ResponseMeta::default()));
2453 }
2454 if staged.inspection.is_none() {
2455 staged.inspection = self
2456 .inspection_max_objects()
2457 .map(|limit| crate::indexed::inspection::InspectionSet::new(limit as usize));
2458 }
2459 if let Some(pending) = implicit {
2460 let upd = refs
2461 .first_mut()
2462 .ok_or_else(|| internal("implicit consumption needs an UpdateRef"))?;
2463 self.check_implicit_packmap(&op, pending, ahead.as_ref(), upd)
2466 .await?;
2467 }
2468 let existing = self.begin_decision(&op, a, ahead.as_mut()).await?;
2469 if existing.is_none() && !ticketed {
2470 self.check_outbox_backpressure(&p, ahead.as_ref()).await?;
2471 }
2472 let allowance = if existing.is_some() || ticketed || implicit.is_some_and(|p| !p.is_empty())
2476 {
2477 Allowance::default()
2478 } else {
2479 let credentials = admission::validate_credentials(&a.credential_capture)?;
2480 let mut input = AdmissionInput::new(&op);
2481 input.credential_headers = &credentials;
2482 if let OpKind::BeginUpload { key, bytes, .. } = &op.kind {
2483 input.declared_bytes = *bytes;
2484 input.pack_id = Some(*key);
2485 input.new_to_repo_bytes = Some(*bytes);
2486 }
2487 self.admit(input).await?
2488 };
2489 let pending = match allowance.reservation.as_deref() {
2490 Some(rid) => Some(self.record_pending(a, &p, rid).await?),
2491 None => None,
2492 };
2493 let write_result = async {
2494 let mut begin = self.begin_write(&op, a, existing, allowance.reservation.clone())?;
2495 self.precheck_namespace(&p, &allowance.charges, &mut ahead, a.business_skew_ms)?;
2496 let mut opened_session = None;
2497 if let Some(BeginWrite::Open(open)) = &mut begin
2498 && open.spec.bytes > open.spec.part_size
2499 {
2500 let key = PackKey(open.spec.pack_id).into();
2501 let ticket_id = crate::store::tickets::ticket_id(&open.spec.reservation_id);
2502 let session = self
2503 .blobs
2504 .begin_multipart_for_ticket(key, open.spec.bytes, open.spec.part_size, ticket_id)
2505 .await
2506 .map_err(|e| {
2507 if open.reserved() {
2508 tracing::warn!(error = %e, "reserved multipart session creation failed");
2509 ServerError::unavailable("multipart session creation failed; retry")
2510 } else {
2511 store_error(StorageOp::MultipartSession, e)
2512 }
2513 })?;
2514 if session.is_empty() || session.len() > u16::MAX as usize {
2515 if let Err(err) = self.blobs.abort(key, &session).await {
2516 tracing::warn!(error = %err, "failed to abort invalid multipart session");
2517 }
2518 return Err(ServerError::internal(
2519 "object storage request failed",
2520 "multipart store returned an invalid session identifier",
2521 ));
2522 }
2523 opened_session = Some((key, session.clone(), ticket_id));
2524 open.spec.upload_session = Some(session);
2525 }
2526 let write_result = async {
2527 let lease = if let Some(observed) = lease {
2528 let (created, lease) = {
2529 let _grant_gate = match (&self.gate, &observed) {
2533 (Some(gate), lease::LeaseObservation::Renew(_)) => {
2534 Some(gate.enter(&p).await)
2535 }
2536 _ => None,
2537 };
2538 self.admit_lease(&op, &p, observed, a.business_skew_ms)
2539 .await?
2540 };
2541 op.created = created;
2542 op.leased_epoch = Some(lease.value.epoch);
2543 if lease.install {
2544 fault!(self, AfterLeaseGrant, &op, a);
2545 }
2546 Some(lease)
2547 } else {
2548 op.created = self.commit_creation(&op, a.business_skew_ms).await?;
2549 None
2550 };
2551 if let Err(error) = self
2552 .check_fast_forward(
2553 &op,
2554 &mut refs,
2555 ahead.as_ref(),
2556 fast_forward.as_ref(),
2557 (&staged, ticket_ms),
2558 )
2559 .await
2560 {
2561 return self
2562 .store_policy_denial(&op, a, &p, ahead, error, pending.is_some())
2563 .await;
2564 }
2565 self.pre_receive(&op).await?;
2566 let write = (kind, refs.as_slice(), allowance.charges.as_slice());
2567 self.plan_and_apply(
2568 &op,
2569 a,
2570 &p,
2571 write,
2572 ahead,
2573 (lease, begin.as_ref()),
2574 WriteInputs {
2575 denial_ids: &staged.denial_ids,
2576 denial_packs: &staged.denial_packs,
2577 pending: pending.as_ref(),
2578 implicit,
2579 external_bases: &staged.external_bases,
2580 inspected: staged.inspection.as_mut(),
2581 },
2582 )
2583 .await
2584 }
2585 .await;
2586 if let Some((key, session, fresh_id)) = opened_session {
2587 self.cleanup_opened_session(&p, &write_result, key, &session, fresh_id)
2588 .await;
2589 }
2590 write_result
2591 }
2592 .await;
2593 if let (Some(pending), Err(err)) = (&pending, &write_result) {
2594 let (reason, detail) = reservation::abort_reason(err);
2595 self.resolve_pending(&p, pending, reason, detail).await;
2596 }
2597 write_result.map(|result| {
2598 let committed = matches!(
2599 &result,
2600 StoredResult::UpdateRef(UpdateRefResult::Committed)
2601 | StoredResult::AdvanceRefs(AdvanceOutcome::Committed)
2602 | StoredResult::BeginUpload(BeginUploadResult::Ticket { .. })
2603 );
2604 let meta = if committed
2605 && (!allowance.response_headers.is_empty() || allowance.external_ref.is_some())
2606 {
2607 let mut headers = allowance.response_headers;
2608 if !headers.is_empty() {
2609 headers.push(("Cache-Control".into(), "private".into()));
2610 }
2611 ResponseMeta {
2612 headers,
2613 external_ref: allowance.external_ref,
2614 }
2615 } else {
2616 ResponseMeta::default()
2617 };
2618 (result, meta)
2619 })
2620 }
2621
2622 async fn cleanup_opened_session(
2623 &self,
2624 partition: &Partition,
2625 result: &Result<StoredResult, ServerError>,
2626 key: crate::store::BlobKey,
2627 session: &[u8],
2628 fresh_id: Hash,
2629 ) {
2630 if session == fresh_id {
2634 return;
2635 }
2636 let committed_fresh = match result {
2639 Ok(StoredResult::BeginUpload(BeginUploadResult::Ticket { id, token, .. }))
2640 if *id == fresh_id =>
2641 {
2642 self.cfg
2643 .ticket_keys
2644 .as_ref()
2645 .and_then(|keys| keys.verify(token, 0).ok())
2646 .is_some_and(|claims| claims.upload_session == session)
2647 }
2648 _ => false,
2649 };
2650 let stored_fresh = if result.is_err() {
2653 match self.meta.get(partition, &keys::ticket(&fresh_id)).await {
2654 Ok(Some(raw)) => codec::decode_ticket(&raw)
2655 .ok()
2656 .is_none_or(|ticket| ticket.upload_session.as_deref() == Some(session)),
2657 Ok(None) => false,
2658 Err(err) => {
2659 tracing::warn!(error = %err, "could not confirm multipart ticket after failed write");
2660 true
2661 }
2662 }
2663 } else {
2664 false
2665 };
2666 if !committed_fresh
2667 && !stored_fresh
2668 && let Err(err) = self.blobs.abort(key, session).await
2669 {
2670 tracing::warn!(error = %err, "failed to abort unused multipart session");
2671 }
2672 }
2673
2674 fn identify(&self, a: &Authenticated, kind: OpKind) -> Result<Operation, ServerError> {
2677 if a.procedure() != kind.procedure() {
2678 return Err(ServerError::unauthenticated(
2679 "credentials were checked for another procedure",
2680 ));
2681 }
2682 if matches!(self.cfg.auth, AuthMode::AuthV2(_))
2683 && kind.procedure().is_write()
2684 && kind.procedure() != Procedure::SetRepoVisibility
2685 && a.auth.is_none()
2686 {
2687 return Err(ServerError::unauthenticated(
2688 "missing auth v2 authorization",
2689 ));
2690 }
2691 let repo = a.repo().repo.clone();
2692 let principal = a.principal.clone();
2693 let mut op = Operation::new(repo, principal, a.auth.clone(), kind);
2694 op.write_grant.clone_from(&a.write_grant);
2695 op.business_now_ms = Some(a.business_now_ms);
2696 Ok(op)
2697 }
2698
2699 async fn require_repository(&self, repo: &crate::repo::RepoId) -> Result<(), ServerError> {
2701 if matches!(self.cfg.addressing, Addressing::Multi(_)) {
2702 let p = self.shards.coordinator(&repo.namespace);
2703 let value = self
2704 .meta
2705 .get(&p, &keys::repo_record(&repo.name))
2706 .await
2707 .map_err(meta_error)?;
2708 match value {
2709 Some(value) => {
2710 codec::decode_repo_record(&value).map_err(meta_error)?;
2711 }
2712 None => return Err(ServerError::repository_not_found()),
2713 }
2714 }
2715 Ok(())
2716 }
2717
2718 fn visibility_applies(&self) -> bool {
2721 matches!(self.cfg.addressing, Addressing::Multi(_))
2722 && self.cfg.write_policy == WritePolicy::Owner
2723 && matches!(self.cfg.auth, AuthMode::AuthV2(_))
2724 }
2725
2726 fn visibility_gates_reads(&self) -> bool {
2731 matches!(self.cfg.addressing, Addressing::Multi(_))
2732 && self.cfg.write_policy == WritePolicy::Owner
2733 }
2734
2735 async fn authorize_read(&self, op: &Operation) -> Result<ReadAuth, ServerError> {
2742 tracing::debug!(stage = "authorize");
2743 if !self.visibility_gates_reads() {
2744 return self.authorize_read_ungated(op).await;
2745 }
2746 let admin_owner = Namespace::parse(op.repo.namespace.as_str())
2749 .is_ok_and(|namespace| self.owner_key_is_admin(&namespace));
2750 let check = op
2751 .write_grant
2752 .as_ref()
2753 .filter(|_| !admin_owner)
2754 .and_then(|header| {
2755 self.cfg
2756 .grants
2757 .as_ref()
2758 .and_then(|g| read_policy::check_grant(g, header.expose(), op))
2759 });
2760 let signed = op.auth.is_some();
2761 let owner = op.write_grant.is_none()
2762 && matches!(Namespace::parse(op.repo.namespace.as_str()),
2763 Ok(Namespace::Ed25519(key)) if op.principal.ed25519() == Some(&key));
2764 let Some((private, epoch)) = self.read_repo_state(&op.repo).await? else {
2765 if signed && op.write_grant.is_none() {
2768 let provisional = AuthzFacts {
2769 authority_generation: None,
2770 grant: None,
2771 owner,
2772 caller_view: if owner {
2773 CallerView::Writer
2774 } else {
2775 CallerView::Reader
2776 },
2777 };
2778 let _ = self.read_hook(op, true, true, &provisional).await;
2779 }
2780 return Err(ServerError::repository_not_found());
2781 };
2782 let grant = check.filter(|c| c.epoch == epoch);
2784 let grant_ref = grant.map(|c| GrantRef {
2785 id: c.id,
2786 epoch,
2787 presence_requirement: None,
2789 });
2790 let authority = self.cfg.authorizer_role == AuthorizerRole::Authority;
2791 if private && !signed {
2792 return Err(ServerError::repository_not_found());
2793 }
2794 let provisional = AuthzFacts {
2795 authority_generation: None,
2796 grant: grant_ref.clone(),
2797 owner,
2798 caller_view: if owner || grant.is_some_and(|g| g.write) {
2799 CallerView::Writer
2800 } else {
2801 CallerView::Reader
2802 },
2803 };
2804 let hook = self.read_hook(op, private, signed, &provisional).await?;
2805 let caller = read_policy::Caller {
2806 #[cfg(feature = "http-objects")]
2807 http_token_authorized: false,
2808 signed,
2809 owner,
2810 grant: grant.map(|g| read_policy::GrantEval {
2811 read: g.read,
2812 write: g.write,
2813 }),
2814 hook,
2815 authority,
2816 };
2817 match read_policy::decide(op.procedure(), private, caller) {
2818 read_policy::ReadDecision::NotFound => Err(ServerError::repository_not_found()),
2819 read_policy::ReadDecision::Allow(caller_view) => Ok(ReadAuth {
2820 facts: AuthzFacts {
2821 authority_generation: None,
2822 grant: grant_ref,
2823 owner,
2824 caller_view,
2825 },
2826 epoch: Some(epoch),
2827 }),
2828 }
2829 }
2830
2831 async fn read_repo_state(
2835 &self,
2836 repo: &crate::repo::RepoId,
2837 ) -> Result<Option<(bool, u64)>, ServerError> {
2838 let coordinator = self.shards.coordinator(&repo.namespace);
2839 let rows = self
2840 .meta
2841 .get_many(
2842 &coordinator,
2843 &[
2844 keys::repo_record(&repo.name),
2845 keys::repo_visibility(&repo.name),
2846 keys::grant_epoch(),
2847 ],
2848 )
2849 .await
2850 .map_err(|e| {
2851 tracing::warn!(detail = %e, "repository state read failed");
2852 ServerError::unavailable("repository state unavailable")
2853 })?;
2854 let mut rows = rows.into_iter();
2855 let Some(record) = rows.next().flatten() else {
2856 return Ok(None);
2857 };
2858 codec::decode_repo_record(&record).map_err(meta_error)?;
2859 let stored = rows
2860 .next()
2861 .flatten()
2862 .map(|v| codec::decode_repo_visibility(&v))
2863 .transpose()
2864 .map_err(meta_error)?;
2865 let epoch = rows
2866 .next()
2867 .flatten()
2868 .map(|v| codec::decode_u64(&v))
2869 .transpose()
2870 .map_err(meta_error)?
2871 .unwrap_or(0);
2872 Ok(Some((
2873 repo_is_private(stored.as_ref(), self.cfg.default_repo_visibility),
2874 epoch,
2875 )))
2876 }
2877
2878 async fn read_hook(
2884 &self,
2885 op: &Operation,
2886 private: bool,
2887 signed: bool,
2888 provisional: &AuthzFacts,
2889 ) -> Result<read_policy::HookEval, ServerError> {
2890 let authorized = || {
2891 let mut authorized = op.clone();
2892 authorized.authz = provisional.clone();
2893 authorized
2894 };
2895 let writer = |facts: &AuthzFacts| facts.caller_view == CallerView::Writer;
2896 if private {
2897 if op.write_grant.is_some() {
2898 return Ok(read_policy::HookEval::NotConsulted);
2901 }
2902 return Ok(
2903 match self.hooks.authorizer().authorize(&authorized()).await {
2904 Ok(facts) => read_policy::HookEval::Allow {
2905 writer_view: writer(&facts),
2906 },
2907 Err(_) => read_policy::HookEval::Deny,
2910 },
2911 );
2912 }
2913 if self.cfg.authorizer_role == AuthorizerRole::Authority {
2914 if !signed || provisional.caller_view == CallerView::Writer {
2917 return Ok(read_policy::HookEval::NotConsulted);
2918 }
2919 return Ok(
2920 match self.hooks.authorizer().authorize(&authorized()).await {
2921 Ok(facts) => read_policy::HookEval::Allow {
2922 writer_view: writer(&facts),
2923 },
2924 Err(_) => read_policy::HookEval::NotConsulted,
2925 },
2926 );
2927 }
2928 self.hooks
2931 .authorizer()
2932 .authorize(op)
2933 .await
2934 .map_err(ServerError::strip_admission_shape)?;
2935 Ok(read_policy::HookEval::NotConsulted)
2936 }
2937
2938 async fn authorize_read_ungated(&self, op: &Operation) -> Result<ReadAuth, ServerError> {
2942 let facts = self
2943 .hooks
2944 .authorizer()
2945 .authorize(op)
2946 .await
2947 .map_err(ServerError::strip_admission_shape)?;
2948 self.require_repository(&op.repo).await?;
2949 let caller_view = match &op.principal {
2950 Principal::Anonymous => CallerView::Anonymous,
2951 Principal::BearerHolder => CallerView::Reader,
2952 principal => {
2953 if self.cfg.write_policy == WritePolicy::Open
2954 || matches!(Namespace::parse(op.repo.namespace.as_str()),
2955 Ok(Namespace::Ed25519(key)) if principal.ed25519() == Some(&key))
2956 {
2957 CallerView::Writer
2958 } else {
2959 CallerView::Reader
2960 }
2961 }
2962 };
2963 Ok(ReadAuth {
2964 facts: AuthzFacts {
2965 caller_view,
2966 ..facts
2967 },
2968 epoch: None,
2969 })
2970 }
2971
2972 async fn pack_is_member(
2973 &self,
2974 a: &Authenticated,
2975 key: &PackKey,
2976 caller: CallerView,
2977 ) -> Result<bool, ServerError> {
2978 if matches!(self.cfg.addressing, Addressing::Single { .. }) {
2979 return Ok(true);
2980 }
2981 let view = crate::store::view::ViewStore {
2982 store: &self.meta,
2983 repo: &a.repo().repo,
2984 writer: caller == CallerView::Writer,
2985 policy: self.publication_policy.as_deref(),
2986 };
2987 let member = read::is_member(
2988 &view,
2989 self.shards.as_ref(),
2990 &a.repo().repo,
2991 &key.0,
2992 a.ref_hint.as_deref(),
2993 )
2994 .await
2995 .map_err(meta_error)?;
2996 if member && self.cfg.indexed.is_some() {
2997 match crate::takedown::denial::require_clear(&self.meta, &key.0).await {
2998 Ok(()) => {}
2999 Err(e) if e.public_message() == "object blocked" => return Ok(false),
3000 Err(e) => return Err(e),
3001 }
3002 }
3003 if member && self.cfg.takedown_denial && self.cfg.indexed.is_some() {
3004 return match crate::takedown::denial::require_pack_clear(
3005 &self.meta,
3006 self.shards.as_ref(),
3007 &a.repo().repo,
3008 &key.0,
3009 )
3010 .await
3011 {
3012 Ok(()) => Ok(true),
3013 Err(e) if e.public_message() == "object blocked" => Ok(false),
3014 Err(e) => Err(e),
3015 };
3016 }
3017 Ok(member)
3018 }
3019
3020 fn ref_writes(
3023 &self,
3024 op: &Operation,
3025 ) -> Result<(WriteKind, Vec<RefUpdate>, Partition), ServerError> {
3026 let (kind, refs) = match &op.kind {
3027 OpKind::UpdateRef(u) => (WriteKind::UpdateRef, vec![u.clone()]),
3028 OpKind::AdvanceRefs { head, packmap, .. } => {
3029 (WriteKind::AdvanceRefs, vec![packmap.clone(), head.clone()])
3030 }
3031 OpKind::BeginUpload { ref_name, .. } => {
3032 return Ok((
3033 WriteKind::BeginUpload,
3034 vec![],
3035 self.shards.ref_shard(&op.repo, ref_name),
3036 ));
3037 }
3038 _ => return Err(internal("not a unary write")),
3039 };
3040 let p = self.shards.ref_shard(&op.repo, &refs[0].name);
3041 if refs
3042 .iter()
3043 .any(|r| self.shards.ref_shard(&op.repo, &r.name) != p)
3044 {
3045 return Err(internal("shard map splits a head from its packmap"));
3046 }
3047 Ok((kind, refs, p))
3048 }
3049
3050 async fn read_ahead(
3055 &self,
3056 op: &Operation,
3057 p: &Partition,
3058 refs: &[RefUpdate],
3059 business_skew_ms: i64,
3060 ) -> Result<Option<Snapshot>, ServerError> {
3061 let caps = self.meta.capabilities();
3062 if !caps.atomic_multi_key {
3063 return Ok(None);
3064 }
3065 let namespace_window = if matches!(self.cfg.addressing, Addressing::Multi(_))
3066 && self.hooks.admission().is_default()
3067 {
3068 self.cfg.write_quota.map(|limits| {
3069 quota::namespace_window(
3070 self.clock.now_ms().saturating_add(business_skew_ms),
3071 limits.window_ms,
3072 )
3073 })
3074 } else {
3075 None
3076 };
3077 let mut wanted: Vec<Key> = refs
3078 .iter()
3079 .map(|r| keys::ref_key(&op.repo.name, &r.name))
3080 .collect();
3081 if let Some(update) = refs.first() {
3082 let name = crate::store::publication::sequence_ref(&update.name);
3083 wanted.push(keys::publication(&op.repo.name, &name));
3084 wanted.push(keys::ref_key(&op.repo.name, &name));
3085 if let Some(packmap) = mkit_attest::grant::head_packmap(&name) {
3086 wanted.push(keys::ref_key(&op.repo.name, &packmap));
3087 }
3088 wanted.push(keys::outbox_sequence());
3089 }
3090 if self.cfg.sharding == Sharding::Single && caps.key_classes == KeyClasses::All {
3091 wanted.extend([keys::authority_generation(), keys::lease_recovery()]);
3092 }
3093 if caps.implicit_layout_version.is_none() {
3094 wanted.push(keys::layout_version());
3095 }
3096 wanted.push(keys::outcome_backlog());
3097 if matches!(self.cfg.addressing, Addressing::Multi(_)) {
3098 wanted.push(keys::repo_known(&op.repo.name));
3099 }
3100 if self.cfg.sharding == Sharding::D34 {
3101 wanted.push(keys::epoch_lease());
3102 if !refs.is_empty() {
3103 wanted.push(keys::outbox_sequence());
3104 }
3105 }
3106 if let Some(auth) = &op.auth {
3107 wanted.push(keys::replay(&auth.replay_scope));
3108 if self.cfg.sharding == Sharding::Single {
3109 wanted.push(keys::grant_epoch());
3110 }
3111 if self.cfg.write_quota.is_some() {
3112 let scope = QuotaScope::for_signer(&op.repo.namespace, &auth.signer);
3113 wanted.push(keys::quota(&scope));
3114 }
3115 if let (Some(window), Some(limits)) = (namespace_window, self.cfg.write_quota) {
3116 let charge = NamespaceCharge {
3117 limits,
3118 window,
3119 bytes: 0,
3120 rollup: matches!(p, Partition::Ref { .. }),
3121 };
3122 wanted.push(quota::counter_key(charge, window));
3123 if charge.rollup {
3124 wanted.push(keys::quota_view(window));
3125 }
3126 let next = window.saturating_add(1);
3129 wanted.push(quota::counter_key(charge, next));
3130 if charge.rollup {
3131 wanted.push(keys::quota_view(next));
3132 }
3133 }
3134 }
3135 if let OpKind::BeginUpload { ref_name, key, .. } = &op.kind {
3136 let signer = op
3137 .auth
3138 .as_ref()
3139 .ok_or_else(|| internal("missing ticket signer"))?
3140 .signer;
3141 wanted.extend(begin::decision_keys(
3142 &op.repo.name,
3143 ref_name,
3144 &key.0,
3145 &signer,
3146 )?);
3147 }
3148 let mut snap = Snapshot::default();
3149 snap.namespace_window = namespace_window;
3150 self.fill(p, &mut snap, wanted).await?;
3151 Ok(Some(snap))
3152 }
3153
3154 fn namespace_charge(
3155 &self,
3156 p: &Partition,
3157 charges: &[QuotaCharge],
3158 window: Option<u64>,
3159 ) -> Result<Option<NamespaceCharge>, ServerError> {
3160 if !matches!(self.cfg.addressing, Addressing::Multi(_))
3161 || !self.hooks.admission().is_default()
3162 {
3163 return Ok(None);
3164 }
3165 let Some(charge) = charges.first() else {
3166 return Ok(None);
3167 };
3168 let window = window.ok_or_else(|| internal("namespace quota read-ahead missing"))?;
3169 Ok(Some(NamespaceCharge {
3170 limits: charge.limits,
3171 window,
3172 bytes: charge.bytes,
3173 rollup: matches!(p, Partition::Ref { .. }),
3174 }))
3175 }
3176
3177 fn precheck_namespace(
3180 &self,
3181 p: &Partition,
3182 charges: &[QuotaCharge],
3183 ahead: &mut Option<Snapshot>,
3184 business_skew_ms: i64,
3185 ) -> Result<(), ServerError> {
3186 let Some(mut charge) = self.namespace_charge(
3187 p,
3188 charges,
3189 ahead
3190 .as_ref()
3191 .and_then(|snapshot| snapshot.namespace_window),
3192 )?
3193 else {
3194 return Ok(());
3195 };
3196 let now = self.clock.now_ms().saturating_add(business_skew_ms);
3197 let window = quota::namespace_window(now, charge.limits.window_ms);
3198 if window != charge.window && window != charge.window.saturating_add(1) {
3199 return Err(
3200 ServerError::aborted_retryable("namespace quota window advanced; retry")
3201 .with_abort_cause(AbortCause::QuotaWindow),
3202 );
3203 }
3204 charge.window = window;
3205 let snap = ahead
3206 .as_mut()
3207 .ok_or_else(|| internal("namespace quota needs atomic reads"))?;
3208 let stored_view = charge
3209 .rollup
3210 .then(|| snap.get(&keys::quota_view(window)))
3211 .flatten();
3212 let seeded_view = snap
3213 .namespace_seed
3214 .as_ref()
3215 .filter(|(seed_window, _)| *seed_window == window)
3216 .map(|(_, value)| value);
3217 let decision = quota::check_namespace(
3218 snap.get("a::counter_key(charge, window)),
3219 stored_view.or(seeded_view),
3220 now,
3221 charge,
3222 )?;
3223 match decision {
3224 NamespaceDecision::Exhausted => {
3225 return Err(ServerError::resource_exhausted(
3226 "namespace write op/byte quota exceeded for this window; try again later",
3227 ));
3228 }
3229 NamespaceDecision::Allowed {
3230 view: ViewStatus::Missing,
3231 ..
3232 } => self.metrics.incr(
3233 crate::telemetry::METRIC_NAMESPACE_QUOTA_VIEW_FALLBACK,
3234 &[("state", "missing")],
3235 1,
3236 ),
3237 NamespaceDecision::Allowed {
3238 view: ViewStatus::Stale,
3239 ..
3240 } => self.metrics.incr(
3241 crate::telemetry::METRIC_NAMESPACE_QUOTA_VIEW_FALLBACK,
3242 &[("state", "stale")],
3243 1,
3244 ),
3245 NamespaceDecision::Allowed { .. } => {}
3246 }
3247 snap.namespace_window = Some(window);
3248 Ok(())
3249 }
3250
3251 async fn fill(
3253 &self,
3254 p: &Partition,
3255 snap: &mut Snapshot,
3256 mut wanted: Vec<Key>,
3257 ) -> Result<(), ServerError> {
3258 wanted.retain(|k| !snap.contains(k));
3259 wanted.sort();
3260 wanted.dedup();
3261 if wanted.is_empty() {
3262 return Ok(());
3263 }
3264 let values = self.meta.get_many(p, &wanted).await.map_err(meta_error)?;
3265 for (key, value) in wanted.into_iter().zip(values) {
3266 snap.insert(key, value);
3267 }
3268 Ok(())
3269 }
3270
3271 fn replay_lookup(
3275 op: &Operation,
3276 ahead: Option<&Snapshot>,
3277 ) -> Result<Option<StoredResult>, ServerError> {
3278 let (Some(auth), Some(snap)) = (&op.auth, ahead) else {
3279 return Ok(None);
3280 };
3281 tracing::debug!(stage = "replay_lookup");
3282 let stored = snap.get(&keys::replay(&auth.replay_scope));
3283 let record = stored.map(codec::decode_replay_record).transpose();
3284 let record = record.map_err(meta_error)?;
3285 replay_answer(classify(record.as_ref(), &auth.fingerprint))
3286 }
3287
3288 fn merge_authority_facts(
3289 &self,
3290 built_in: &mut AuthzFacts,
3291 returned: &AuthzFacts,
3292 ) -> Result<(), ServerError> {
3293 if self.cfg.authorizer_role == AuthorizerRole::Authority
3294 && self.cfg.authority_fence.is_some()
3295 {
3296 built_in.authority_generation =
3297 Some(returned.authority_generation.ok_or_else(|| {
3298 ServerError::unavailable("Authority allowance missing generation")
3299 })?);
3300 }
3301 Ok(())
3302 }
3303
3304 async fn authorize(
3306 &self,
3307 op: &Operation,
3308 ) -> Result<(AuthzFacts, Option<FastForward>), ServerError> {
3309 tracing::debug!(stage = "authorize");
3310 if !op.procedure().is_write() {
3311 return Ok((self.authorize_read(op).await?.facts, None));
3312 }
3313 if self.cfg.authority_fence.is_some() && self.cfg.sharding == Sharding::Single {
3314 Box::pin(self.ensure_authority_activation(&op.repo.namespace)).await?;
3315 }
3316 match &self.cfg.addressing {
3317 Addressing::Multi(multi) => self.owner_rule(op, Some(&multi.namespace_policy)).await,
3318 Addressing::Single { .. } if self.cfg.write_policy == WritePolicy::Owner => {
3321 self.owner_rule(op, None).await
3322 }
3323 Addressing::Single { .. } => self
3324 .hooks
3325 .authorizer()
3326 .authorize(op)
3327 .await
3328 .map(|mut facts| {
3329 facts.authority_generation = None;
3332 (facts, None)
3333 })
3334 .map_err(ServerError::strip_admission_shape),
3335 }
3336 }
3337
3338 fn owner_key_is_admin(&self, namespace: &Namespace) -> bool {
3339 matches!(namespace, Namespace::Ed25519(key) if self.cfg.admin_keys.contains(key))
3340 }
3341
3342 fn require_client_owner_key(&self, namespace: &Namespace) -> Result<(), ServerError> {
3343 if self.owner_key_is_admin(namespace) {
3344 return Err(ServerError::permission_denied(
3345 "admin key cannot authorize client calls",
3346 ));
3347 }
3348 Ok(())
3349 }
3350
3351 async fn owner_rule(
3359 &self,
3360 op: &Operation,
3361 policy: Option<&NamespacePolicy>,
3362 ) -> Result<(AuthzFacts, Option<FastForward>), ServerError> {
3363 let namespace = Namespace::parse(op.repo.namespace.as_str())
3364 .map_err(|_| internal("invalid resolved owner-policy namespace"))?;
3365 if op.write_grant.is_some() {
3366 self.require_client_owner_key(&namespace)?;
3367 }
3368 if let Some(NamespacePolicy::Allowlist(allowed)) = policy
3369 && !allowed.contains(&namespace)
3370 {
3371 return Err(ServerError::permission_denied("write not permitted"));
3372 }
3373 if op.write_grant.is_some()
3374 && matches!(namespace, Namespace::Address(_))
3375 && !matches!(policy, Some(NamespacePolicy::Allowlist(_)))
3376 {
3377 return Err(ServerError::permission_denied("write not permitted"));
3378 }
3379 let mut fast_forward = None;
3380 let grant = if let Some(header) = &op.write_grant {
3381 let cfg = self.cfg.grants.as_ref().ok_or_else(|| {
3382 grants::rejected(mkit_attest::grant::GrantError::SchemeNotAdvertised)
3383 })?;
3384 let verified = cfg.verify(header.expose(), op)?;
3385 let scope =
3387 crate::policy::ref_scopes::authorize(&verified, &op.kind, self.cfg.indexed_mode())?;
3388 fast_forward = scope.fast_forward;
3389 let presence_requirement = scope.presence;
3390 let observed = match self.cfg.sharding {
3391 Sharding::Single => op.observed_epoch,
3392 Sharding::D34 => op.leased_epoch,
3393 };
3394 if observed != Some(verified.epoch()) {
3395 return Err(plan::epoch_moved());
3396 }
3397 tracing::debug!(grant_id = %to_hex_bytes(verified.id()), "write grant accepted");
3398 Some(crate::op::GrantRef {
3399 id: *verified.id(),
3400 epoch: verified.epoch(),
3401 presence_requirement,
3402 })
3403 } else {
3404 None
3405 };
3406 let owner = grant.is_none()
3407 && matches!(&namespace, Namespace::Ed25519(key)
3408 if op.principal.ed25519() == Some(key));
3409 if self.cfg.authorizer_role == AuthorizerRole::Check && !owner && grant.is_none() {
3412 return Err(ServerError::permission_denied("write not permitted"));
3413 }
3414 let mut facts = AuthzFacts {
3415 authority_generation: None,
3416 grant,
3417 owner,
3418 caller_view: CallerView::Writer,
3419 };
3420 let mut authorized = op.clone();
3422 authorized.authz = facts.clone();
3423 let returned = self
3424 .hooks
3425 .authorizer()
3426 .authorize(&authorized)
3427 .await
3428 .map_err(ServerError::strip_admission_shape)?;
3429 self.merge_authority_facts(&mut facts, &returned)?;
3430 Ok((facts, fast_forward))
3431 }
3432
3433 async fn pre_receive(&self, op: &Operation) -> Result<(), ServerError> {
3435 tracing::debug!(stage = "pre_receive");
3436 self.hooks
3437 .pre_receive()
3438 .check(op, None)
3439 .await
3440 .map_err(ServerError::strip_admission_shape)
3441 }
3442
3443 async fn plan_and_apply(
3446 &self,
3447 op: &Operation,
3448 a: &Authenticated,
3449 p: &Partition,
3450 (kind, refs, charges): (WriteKind, &[RefUpdate], &[QuotaCharge]),
3451 ahead: Option<Snapshot>,
3452 (lease, begin): (Option<lease::LeaseWrite>, Option<&BeginWrite>),
3453 inputs: WriteInputs<'_>,
3454 ) -> Result<StoredResult, ServerError> {
3455 let WriteInputs {
3456 denial_ids,
3457 denial_packs,
3458 pending,
3459 implicit,
3460 external_bases,
3461 inspected,
3462 } = inputs;
3463 let caps = self.meta.capabilities();
3464 let replay = upload::replay_guard(op);
3465 let implicit_ids = implicit
3469 .filter(|pending| {
3470 !pending.is_empty() && matches!(self.cfg.addressing, Addressing::Multi(_))
3471 })
3472 .map(implicit::implicit_packs);
3473 let advance = self.publication_ticket_write(op, a, p)?;
3474 let mut req = WriteRequest {
3475 denial_ids: Some(denial_ids),
3476 denial_packs,
3477 authority_store: plan::AuthorityStore::from_capabilities(self.meta.capabilities()),
3478 authority_generation: op.authz.authority_generation,
3479 repo: &op.repo.name,
3480 kind,
3481 refs,
3482 ref_index: (self.cfg.sharding == Sharding::D34 && !refs.is_empty()).then_some((
3483 &op.repo,
3484 p,
3485 self.shards.as_ref(),
3486 )),
3487 replay,
3488 charges,
3489 namespace_charge: self.namespace_charge(
3490 p,
3491 charges,
3492 ahead.as_ref().and_then(|s| s.namespace_window),
3493 )?,
3494 grant: op.authz.grant.clone(),
3495 lease,
3496 layout_version: caps.implicit_layout_version.is_none(),
3497 mark_repo_known: matches!(self.cfg.addressing, Addressing::Multi(_))
3498 && ahead
3499 .as_ref()
3500 .is_none_or(|snap| snap.get(&keys::repo_known(&op.repo.name)).is_none()),
3501 rejection: None,
3502 publication: (caps.atomic_multi_key && !refs.is_empty()).then_some(
3503 clearance::PublicationWrite {
3504 repo: &op.repo,
3505 source: p,
3506 shards: self.shards.as_ref(),
3507 prepared: None,
3508 },
3509 ),
3510 pending,
3511 begin,
3512 advance,
3513 implicit: implicit_ids.as_deref().map(|packs| ImplicitConsume {
3514 packs,
3515 repo_id: &op.repo,
3516 source: p,
3517 shards: self.shards.as_ref(),
3518 }),
3519 };
3520 let mut ahead = ahead;
3521 let prepared = match Box::pin(self.prepare_publication(
3522 op,
3523 &a.repo().identity,
3524 p,
3525 &req,
3526 &mut ahead,
3527 implicit_ids.as_deref(),
3528 external_bases,
3529 inspected,
3530 ))
3531 .await
3532 {
3533 Ok(prepared) => prepared,
3534 Err(error)
3535 if self.inspection_max_objects().is_some()
3536 && error.code() == crate::Code::PermissionDenied =>
3537 {
3538 return self
3539 .store_policy_denial(op, a, p, ahead, error, pending.is_some())
3540 .await;
3541 }
3542 Err(error) => return Err(error),
3543 };
3544 if let Some(publication) = &mut req.publication {
3545 publication.prepared = prepared.as_ref();
3546 }
3547 if caps.atomic_multi_key || replay.is_some() || !charges.is_empty() {
3548 return self.apply_atomic(op, a, p, &req, ahead).await;
3549 }
3550 self.apply_sequential(op, a, p, req).await
3551 }
3552
3553 async fn apply_sequential(
3554 &self,
3555 op: &Operation,
3556 a: &Authenticated,
3557 p: &Partition,
3558 mut req: WriteRequest<'_>,
3559 ) -> Result<StoredResult, ServerError> {
3560 let kind = req.kind;
3562 let refs = req.refs;
3563 req.kind = WriteKind::UpdateRef;
3564 for (i, update) in refs.iter().enumerate() {
3565 req.refs = core::slice::from_ref(update);
3566 let result = self.apply_loop(op, a, p, &req, None).await?;
3567 if let StoredResult::UpdateRef(UpdateRefResult::Conflict { .. }) = result {
3568 if kind == WriteKind::UpdateRef {
3569 return Ok(result);
3570 }
3571 return Ok(StoredResult::AdvanceRefs(if i == 0 {
3572 AdvanceOutcome::PackmapConflict
3573 } else {
3574 AdvanceOutcome::HeadConflict
3575 }));
3576 }
3577 }
3578 Ok(match kind {
3579 WriteKind::UpdateRef => StoredResult::UpdateRef(UpdateRefResult::Committed),
3580 _ => StoredResult::AdvanceRefs(AdvanceOutcome::Committed),
3581 })
3582 }
3583
3584 fn publication_ticket_write<'a>(
3585 &'a self,
3586 op: &'a Operation,
3587 a: &'a Authenticated,
3588 p: &'a Partition,
3589 ) -> Result<Option<advance::AdvanceWrite<'a>>, ServerError> {
3590 Ok(match &op.kind {
3591 OpKind::AdvanceRefs { head, tickets, .. } if !tickets.is_empty() => {
3592 Some(advance::AdvanceWrite {
3593 ids: tickets,
3594 signer: op
3595 .auth
3596 .as_ref()
3597 .ok_or_else(|| {
3598 ServerError::failed_precondition("invalid or expired upload ticket")
3599 })?
3600 .signer,
3601 head_ref: &head.name,
3602 repo_id: &op.repo,
3603 repository: &a.repo().identity,
3604 source: p,
3605 shards: self.shards.as_ref(),
3606 })
3607 }
3608 _ => None,
3609 })
3610 }
3611
3612 #[allow(clippy::too_many_arguments, clippy::too_many_lines)] async fn prepare_publication(
3615 &self,
3616 op: &Operation,
3617 repository: &str,
3618 p: &Partition,
3619 req: &WriteRequest<'_>,
3620 ahead: &mut Option<Snapshot>,
3621 implicit_ids: Option<&[Hash]>,
3622 external_bases: &std::collections::BTreeSet<Hash>,
3623 inspected: Option<&mut crate::indexed::inspection::InspectionSet>,
3624 ) -> Result<Option<crate::store::publication::Advance>, ServerError> {
3625 #[cfg(not(feature = "remote-hooks"))]
3626 let _ = repository;
3627 let policy = self.publication_policy.as_deref().or_else(|| {
3630 self.cfg
3631 .takedown_denial
3632 .then_some(&clearance::Immediate as &dyn clearance::PublicationPolicy)
3633 });
3634 #[cfg(feature = "remote-hooks")]
3635 let policy = policy.or_else(|| {
3636 (!self.inspectors.is_empty())
3637 .then_some(&inspection::Immediate as &dyn clearance::PublicationPolicy)
3638 });
3639 if let Some(policy) = policy
3640 && !req.refs.is_empty()
3641 && req.refs.iter().all(|update| update.new.is_some())
3642 {
3643 let snapshot = ahead.get_or_insert_with(Snapshot::default);
3644 self.fill(p, snapshot, req.read_keys()).await?;
3645 let pair = clearance::resulting_pair(&op.repo.name, req.refs, snapshot)?;
3646 let mut prepared = policy.prepare(op, &pair).await?;
3647 if prepared.value != pair {
3648 return Err(internal("publication policy changed the resulting pair"));
3649 }
3650 prepared.generation =
3653 crate::store::publication::Publication::decode(snapshot.get(&keys::publication(
3654 &op.repo.name,
3655 &crate::store::publication::sequence_ref(&req.refs[0].name),
3656 )))
3657 .map_err(meta_error)?
3658 .generation;
3659 prepared.additions = if let Some(advance) = &req.advance {
3660 advance
3661 .ids
3662 .iter()
3663 .map(|id| {
3664 snapshot
3665 .get(&keys::ticket(id))
3666 .ok_or_else(|| internal("publication ticket missing"))
3667 .and_then(|raw| {
3668 codec::decode_ticket(raw)
3669 .map(|t| t.pack_id)
3670 .map_err(meta_error)
3671 })
3672 })
3673 .collect::<Result<Vec<_>, _>>()?
3674 } else {
3675 implicit_ids.map(<[Hash]>::to_vec).unwrap_or_default()
3676 };
3677 #[cfg(feature = "remote-hooks")]
3678 let mut inspected = inspected;
3679 #[cfg(feature = "remote-hooks")]
3680 let verification_set = inspected.as_deref_mut();
3681 #[cfg(not(feature = "remote-hooks"))]
3682 let verification_set = inspected;
3683 self.verify_publication(
3684 op,
3685 p,
3686 &mut prepared,
3687 policy,
3688 mkit_attest::grant::head_packmap(&crate::store::publication::sequence_ref(
3689 &req.refs[0].name,
3690 ))
3691 .is_some(),
3692 external_bases,
3693 verification_set,
3694 req.advance
3695 .as_ref()
3696 .and_then(|a| {
3697 a.ids
3698 .iter()
3699 .filter_map(|id| {
3700 snapshot
3701 .get(&keys::ticket(id))
3702 .and_then(|raw| codec::decode_ticket(raw).ok())
3703 .filter(|t| Some(t.pack_id) == pair.packmap)
3704 .map(|t| t.created_at_ms)
3705 })
3706 .min()
3707 })
3708 .unwrap_or_else(|| {
3709 op.auth
3710 .as_ref()
3711 .map_or(ms(self.clock.now_ms()), |a| ms(a.created_at_ms))
3712 }),
3713 )
3714 .await?;
3715 #[cfg(feature = "remote-hooks")]
3716 if let Some(set) = inspected {
3717 let assignment = if self.cfg.scanner_retrieval.is_some() {
3718 Some(Self::retrieval_assignment(
3719 op,
3720 req.advance.as_ref(),
3721 repository,
3722 snapshot,
3723 set,
3724 )?)
3725 } else {
3726 None
3727 };
3728 self.inspect_advance(
3729 op,
3730 &prepared.value,
3731 set.clone().finalize(),
3732 assignment.as_ref(),
3733 )
3734 .await?;
3735 }
3736 Ok(Some(prepared))
3737 } else {
3738 Ok(None)
3739 }
3740 }
3741
3742 #[allow(clippy::too_many_arguments, clippy::too_many_lines)] async fn verify_publication(
3744 &self,
3745 op: &Operation,
3746 p: &Partition,
3747 prepared: &mut crate::store::publication::Advance,
3748 policy: &dyn clearance::PublicationPolicy,
3749 branch: bool,
3750 external_bases: &std::collections::BTreeSet<Hash>,
3751 mut inspected: Option<&mut crate::indexed::inspection::InspectionSet>,
3752 created: u64,
3753 ) -> Result<(), ServerError> {
3754 let indexed = self
3755 .cfg
3756 .indexed
3757 .ok_or_else(|| internal("publication requires indexed mode"))?;
3758 let resumed = self.cfg.takedown_denial
3759 && self.publication_policy.is_none()
3760 && inspected.is_none()
3761 && Box::pin(crate::indexed::publication::resume::prepare(
3762 &self.meta,
3763 p,
3764 self.shards.as_ref(),
3765 &op.repo,
3766 prepared,
3767 indexed,
3768 ms(self.clock.now_ms()),
3769 self.metrics.as_ref(),
3770 created,
3771 ))
3772 .await?;
3773 let mut indexed = indexed;
3776 if self.cfg.takedown_denial && !resumed {
3777 indexed.decode_budget = indexed.decode_budget.min(8 << 20);
3778 }
3779 let inspection_budget = crate::indexed::budget::SliceBudget::new(256);
3781 let inspection_blobs =
3782 crate::indexed::budget::Budgeted::new(&self.blobs, &inspection_budget);
3783 let inspection_meta = crate::indexed::budget::Budgeted::new(&self.meta, &inspection_budget);
3784 if let Some(set) = inspected.as_deref_mut() {
3785 crate::indexed::publication::verify_inspected(
3786 &inspection_blobs,
3787 &inspection_meta,
3788 self.shards.as_ref(),
3789 &op.repo,
3790 &prepared.value.clone(),
3791 branch,
3792 prepared,
3793 policy,
3794 indexed,
3795 self.metrics.as_ref(),
3796 set,
3797 )
3798 .await?;
3799 } else if !resumed {
3800 crate::indexed::publication::verify(
3801 &self.blobs,
3802 &self.meta,
3803 self.shards.as_ref(),
3804 &op.repo,
3805 &prepared.value.clone(),
3806 branch,
3807 prepared,
3808 policy,
3809 indexed,
3810 self.metrics.as_ref(),
3811 )
3812 .await?;
3813 }
3814 prepared.external_bases = prepared
3815 .external_bases
3816 .iter()
3817 .copied()
3818 .chain(external_bases.iter().copied())
3819 .collect::<std::collections::BTreeSet<_>>()
3820 .into_iter()
3821 .collect();
3822 if prepared.external_bases.len() > crate::store::publication::MAX_ADVANCE_ITEMS {
3823 return Err(ServerError::invalid_argument("object index limit exceeded"));
3824 }
3825 if prepared.state.publishable() {
3826 let visible = if inspected.is_some() {
3827 crate::timers::publication_recheck::dependencies(
3828 &inspection_meta,
3829 &inspection_meta,
3830 p,
3831 self.shards.as_ref(),
3832 &op.repo,
3833 prepared,
3834 )
3835 .await
3836 .map_err(|error| {
3837 if crate::indexed::budget::is_exhausted(&error) {
3838 ServerError::invalid_argument("object index limit exceeded")
3839 } else {
3840 meta_error(error)
3841 }
3842 })?
3843 } else {
3844 crate::timers::publication_recheck::dependencies(
3845 &self.meta,
3846 &self.meta,
3847 p,
3848 self.shards.as_ref(),
3849 &op.repo,
3850 prepared,
3851 )
3852 .await
3853 .map_err(meta_error)?
3854 };
3855 if !visible {
3856 prepared.state = crate::store::publication::Clearance::Pending;
3857 }
3858 }
3859 Ok(())
3860 }
3861
3862 async fn apply_atomic(
3865 &self,
3866 op: &Operation,
3867 a: &Authenticated,
3868 p: &Partition,
3869 req: &WriteRequest<'_>,
3870 ahead: Option<Snapshot>,
3871 ) -> Result<StoredResult, ServerError> {
3872 if !self.meta.capabilities().atomic_multi_key {
3873 return Err(internal("replay and quota need atomic multi-key batches"));
3874 }
3875 self.apply_loop(op, a, p, req, ahead).await
3876 }
3877
3878 #[allow(clippy::too_many_lines)]
3888 async fn apply_loop(
3889 &self,
3890 op: &Operation,
3891 a: &Authenticated,
3892 p: &Partition,
3893 req: &WriteRequest<'_>,
3894 mut ahead: Option<Snapshot>,
3895 ) -> Result<StoredResult, ServerError> {
3896 let mut req = req.for_store(self.meta.capabilities());
3897 let _gate = match &self.gate {
3899 Some(gate) => Some(gate.enter(p).await),
3900 None => None,
3901 };
3902 let skew_ms = a.business_skew_ms;
3903 let (mut replans, mut deadline_missed, mut prune_ok) = (0, false, true);
3904 let mut first_attempt = true;
3905 let denial_budget = crate::indexed::budget::SliceBudget::new(9000);
3906 loop {
3907 let denial_packs: Vec<_> = req
3908 .denial_packs
3909 .iter()
3910 .chain(
3911 req.publication
3912 .iter()
3913 .filter_map(|p| p.prepared)
3914 .flat_map(|a| a.dependencies.iter().chain(&a.external_bases)),
3915 )
3916 .copied()
3917 .collect();
3918 let prove = self.cfg.takedown_denial
3919 && req
3920 .denial_ids
3921 .is_some_and(|ids| !ids.is_empty() || !denial_packs.is_empty());
3922 let proof_plan_time = prove.then(|| ms(self.clock.now_ms()));
3925 if let Some(ids) = req.denial_ids.filter(|_| prove) {
3926 crate::takedown::denial::require_repo_clear_budgeted(
3927 &self.meta,
3928 self.shards.as_ref(),
3929 &op.repo,
3930 ids,
3931 &denial_packs,
3932 &denial_budget,
3933 )
3934 .await?;
3935 }
3936 let base = ahead.take().unwrap_or_default();
3940 let base = if prove {
3941 let mut fresh = Snapshot::default();
3942 fresh.namespace_window = base.namespace_window;
3943 fresh.namespace_seed = base.namespace_seed;
3944 fresh
3945 } else {
3946 base
3947 };
3948 let clock = self.plan_clock(skew_ms, &req);
3949 let snap = self.read_snapshot(p, &req, &clock, base, prune_ok).await?;
3950 if let Some(lease) = req.lease
3951 && (prove
3952 || !first_attempt
3953 || lease
3954 .value
3955 .expires_at_ms
3956 .saturating_sub(self.cfg.lease_margin_ms)
3957 < ms(self.clock.now_ms()).saturating_add(self.cfg.min_lease_budget_ms))
3958 {
3959 let observed = self.observe_lease(op, p, Some(&snap)).await?;
3960 let (_, renewed) = self.admit_lease(op, p, observed, skew_ms).await?;
3961 req.lease = Some(renewed);
3962 }
3963 first_attempt = false;
3964 let mut clock = self.plan_clock(skew_ms, &req);
3965 if let Some(plan_time) = proof_plan_time {
3966 clock.plan_time_ms = plan_time;
3969 }
3970 let plan = match plan_write(&req, &snap, &clock)? {
3971 Planned::Done(result) => return Ok(result),
3972 Planned::Apply(plan) => plan,
3973 };
3974 let Plan {
3975 batch,
3976 on_commit,
3977 replay_index,
3978 epoch_index,
3979 pending_index,
3980 prune,
3981 prune_from,
3982 } = plan;
3983 require_relay_source_lease(self.cfg.sharding, req.lease.is_some(), &batch)?;
3984 #[cfg(feature = "test-faults")]
3985 let batch = faults::delay_relay_batch(
3986 batch,
3987 a.test_directives(),
3988 op,
3989 ms(clock.business_now_ms),
3990 );
3991 if req.kind != WriteKind::UploadReserve {
3992 #[cfg(feature = "test-faults")]
3993 {
3994 let mut attempt = op.clone();
3995 attempt.leased_epoch = req.lease.map(|l| l.value.epoch);
3996 fault!(self, BeforeFinalApply, &attempt, a);
3997 }
3998 }
3999 tracing::debug!(stage = "apply", replans);
4000 match self.meta.apply(p, batch).await {
4001 Ok(BatchOutcome::Committed) => {
4002 #[cfg(feature = "test-faults")]
4003 self.schedule_test_ref_timer(op, a, &on_commit, ms(clock.business_now_ms))
4004 .await?;
4005 return Ok(on_commit);
4006 }
4007 Ok(BatchOutcome::DeadlinePassed { backend_now }) => {
4008 tracing::info!(
4009 backend_now,
4010 deadline = clock.deadline(),
4011 "commit deadline passed"
4012 );
4013 let now = self.clock.now_ms().saturating_add(skew_ms);
4015 let backend = i64::try_from(backend_now).unwrap_or(i64::MAX);
4016 let valid = req
4017 .replay
4018 .is_none_or(|r| now.max(backend) <= r.expires_at_ms);
4019 if deadline_missed || !valid {
4020 return Err(ServerError::unavailable("commit deadline passed; retry"));
4021 }
4022 deadline_missed = true;
4023 }
4024 Ok(BatchOutcome::PreconditionFailed { index, observed }) => {
4025 if Some(index) == pending_index {
4026 return Err(ServerError::unavailable(
4027 "admission reservation changed; retry",
4028 ));
4029 }
4030 if Some(index) == replay_index {
4031 return replay_raced(&req, observed.as_ref());
4032 }
4033 if Some(index) == epoch_index {
4034 return Err(plan::epoch_moved());
4035 }
4036 if index >= prune_from && prune_ok {
4037 prune_ok = false;
4038 continue;
4039 }
4040 replans += 1;
4041 if replans > MAX_REPLAN {
4042 return Err(ServerError::aborted_retryable("write contention; retry")
4043 .with_abort_cause(AbortCause::Contention));
4044 }
4045 }
4046 Err(StoreError::Full) => return Err(self.partition_full(p, prune).await),
4047 Err(e) => return Err(meta_error(e)),
4048 }
4049 }
4050 }
4051
4052 #[cfg(feature = "test-faults")]
4053 async fn schedule_test_ref_timer(
4054 &self,
4055 op: &Operation,
4056 a: &Authenticated,
4057 on_commit: &StoredResult,
4058 now_ms: u64,
4059 ) -> Result<(), ServerError> {
4060 if let OpKind::UpdateRef(upd) = &op.kind
4061 && matches!(
4062 on_commit,
4063 StoredResult::UpdateRef(UpdateRefResult::Committed)
4064 )
4065 {
4066 let timer_partition = self.shards.ref_shard(&op.repo, &upd.name);
4067 faults::schedule_timer(
4068 a.test_directives(),
4069 &self.meta,
4070 &timer_partition,
4071 &op.repo.name,
4072 &upd.name,
4073 now_ms,
4074 )
4075 .await?;
4076 }
4077 Ok(())
4078 }
4079
4080 fn plan_clock(&self, skew_ms: i64, req: &WriteRequest<'_>) -> PlanClock {
4084 let now = self.clock.now_ms();
4085 let lead = MAX_CLOCK_LEAD_MS.unsigned_abs();
4086 PlanClock {
4087 plan_time_ms: ms(now),
4088 business_now_ms: now.saturating_add(skew_ms),
4089 max_apply_window_ms: u64::try_from(self.cfg.max_apply_window.as_millis())
4090 .unwrap_or(u64::MAX),
4091 deadline_cap: match (
4092 req.replay.map(|r| ms(r.expires_at_ms).saturating_add(lead)),
4093 req.lease.map(|l| {
4094 l.value
4095 .expires_at_ms
4096 .saturating_sub(self.cfg.lease_margin_ms)
4097 }),
4098 ) {
4099 (Some(replay), Some(lease)) => Some(replay.min(lease)),
4100 (replay, lease) => replay.or(lease),
4101 },
4102 }
4103 }
4104
4105 async fn read_snapshot(
4110 &self,
4111 p: &Partition,
4112 req: &WriteRequest<'_>,
4113 clock: &PlanClock,
4114 mut snap: Snapshot,
4115 prune: bool,
4116 ) -> Result<Snapshot, ServerError> {
4117 let now = clock.plan_time_ms;
4118 if prune && prune_sampled(req, now) {
4119 if req.replay.is_some() {
4120 snap.expired_replays = read::expired_replay_keys(&self.meta, p, now, PRUNE_LIMIT)
4121 .await
4122 .map_err(meta_error)?;
4123 }
4124 if let Some(window) = req.charges.iter().map(|c| c.limits.window_ms).max() {
4125 snap.stale_quotas =
4126 read::stale_quota_keys(&self.meta, p, now, ms(window), PRUNE_LIMIT)
4127 .await
4128 .map_err(meta_error)?;
4129 }
4130 }
4131 let mut wanted = req.read_keys();
4132 wanted.extend(snap.stale_quotas.iter().map(|(_, quota)| quota.clone()));
4133 self.fill(p, &mut snap, wanted).await?;
4134 if let Some(advance) = &req.advance {
4135 let detail = advance::detail_keys(&snap, advance)?;
4136 self.fill(p, &mut snap, detail).await?;
4137 }
4138 if let Some(BeginWrite::Open(open)) = req.begin {
4139 begin::read_indexed(&self.meta, p, &open.spec, &mut snap).await?;
4140 let reservation = crate::store::tickets::keys(&open.spec).reservation;
4141 if snap.get(&reservation).is_some()
4142 && let Some(replay) = req.replay
4143 {
4144 let key = keys::replay(&replay.scope);
4145 let value = self.meta.get(p, &key).await.map_err(meta_error)?;
4146 snap.insert(key, value);
4147 }
4148 }
4149 Ok(snap)
4150 }
4151
4152 async fn apply_meta(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, ServerError> {
4154 match self.meta.apply(p, batch).await {
4155 Ok(outcome) => Ok(outcome),
4156 Err(StoreError::Full) => Err(self.partition_full(p, None).await),
4157 Err(error) => Err(meta_error(error)),
4158 }
4159 }
4160
4161 async fn partition_full(&self, p: &Partition, prune: Option<Batch>) -> ServerError {
4164 self.metrics
4165 .incr(METRIC_PARTITION_FULL, &[("kind", p.kind())], 1);
4166 let name = p.encode().map(|b| to_hex_bytes(&b)).unwrap_or_default();
4167 tracing::error!(partition = %name, kind = p.kind(), "storage partition full");
4168 if let Some(prune) = prune
4169 && let Err(e) = self.meta.apply(p, prune).await
4170 {
4171 tracing::warn!(error = %e, "prune on a full partition failed");
4172 }
4173 ServerError::unavailable("storage partition full")
4174 }
4175}
4176
4177#[derive(Default)]
4179struct Allowance {
4180 charges: Vec<QuotaCharge>,
4181 reservation: Option<String>,
4182 response_headers: Vec<(String, String)>,
4183 external_ref: Option<String>,
4184}
4185
4186fn check_ref_name(name: &str) -> Result<(), ServerError> {
4190 if name.len() > refs::MAX_REF_NAME_BYTES {
4191 Err(ServerError::invalid_argument(refs::REF_NAME_TOO_LONG))
4192 } else if !validate_ref_name(name) {
4193 Err(ServerError::invalid_argument(
4194 "ref name is invalid (SPEC-REFS §3)",
4195 ))
4196 } else if refs::is_served_ref_name(name) {
4197 Ok(())
4198 } else {
4199 Err(ServerError::invalid_argument(refs::REF_NAME_OUTSIDE_REFS))
4200 }
4201}
4202
4203fn stored_mismatch(result: &StoredResult) -> ServerError {
4206 match result {
4207 StoredResult::Rejected(rejection) => {
4208 ServerError::new(rejection.code(), rejection.message().to_owned())
4209 }
4210 _ => internal("stored result is for another procedure"),
4211 }
4212}
4213
4214fn replay_raced(
4217 req: &WriteRequest<'_>,
4218 observed: Option<&Value>,
4219) -> Result<StoredResult, ServerError> {
4220 if req.pending.is_some() {
4224 return Err(
4225 ServerError::aborted_retryable("operation already in flight; retry")
4226 .with_abort_cause(AbortCause::ReplayRace),
4227 );
4228 }
4229 let (Some(replay), Some(value)) = (req.replay, observed) else {
4230 return Err(ServerError::aborted_retryable(
4231 "operation already in flight; retry",
4232 ));
4233 };
4234 let record = codec::decode_replay_record(value).map_err(meta_error)?;
4235 replay_answer(classify(Some(&record), &replay.fingerprint))?
4236 .ok_or_else(|| ServerError::aborted_retryable("operation already in flight; retry"))
4237}