Skip to main content

mkit_server/pipeline/
mod.rs

1//! The transport-neutral request pipeline (PRD §5.4): the unary RPCs of
2//! `mkit.transport.v1` over the storage contract.
3//!
4//! A binding calls [`Pipeline::authenticate`] (stages 0 and 1) from its
5//! interceptor, then one entry point. A signed write runs one function per
6//! stage, in order: `identify` (stage 1), `replay_lookup` (the stage 0
7//! lookup), `authorize` (2), `admit` (3), `pre_receive` (5) and
8//! `plan_and_apply` (4 and 6, one batch). Admission receipts are attached
9//! only to a committed response; durable outcomes are queued for delivery.
10//! The streaming procedures are [`Pipeline::open_upload`]
11//! ([`UploadSession`]) and [`Pipeline::download`] ([`DownloadStream`]).
12//!
13//! With the `test-faults` feature, `Pipeline::with_faults` installs
14//! `FaultHooks` and [`Pipeline::authenticate`] reads per-request
15//! `TestDirectives`; without it none of that exists in the binary.
16//!
17//! Behavior equals the servers it replaces: `mkit serve --http` (`Bearer`/`Open`:
18//! no replay or quota, packmap-then-head on a non-atomic store),
19//! `vcs-worker` (`AuthV2`: replay ledger, per-signer quota, atomic
20//! advance) and `mkit serve` over ssh (`TransportIdentity`).
21
22mod 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/// Success-only headers returned by admission. They are best-effort and are
149/// never stored in the signed replay ledger or returned on a replay.
150#[derive(Debug, Clone, Default, PartialEq, Eq)]
151pub struct ResponseMeta {
152    headers: Vec<(String, String)>,
153    external_ref: Option<String>,
154}
155
156impl ResponseMeta {
157    /// Headers in admission order, followed by the private cache directive.
158    #[must_use]
159    pub fn headers(&self) -> &[(String, String)] {
160        &self.headers
161    }
162    /// Hook reference carried for the later storage receipt integration.
163    #[must_use]
164    pub fn external_ref(&self) -> Option<&str> {
165        self.external_ref.as_deref()
166    }
167}
168
169/// Call the installed fault hooks at a fault point, returning early on
170/// their error. Compiled out without `test-faults`.
171macro_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
181/// Default bound from planning a batch to its commit (00-plan P-21,
182/// SPEC-WRITE-GRANTS §5.5). It MUST exceed the clock skew between the
183/// planner and the storage backend (Worker isolate vs Durable Object on
184/// Workers; zero natively) plus the longest synchronous span before the
185/// commit, by a wide margin: a batch that misses it commits nothing.
186pub const MAX_APPLY_WINDOW: Duration = Duration::from_secs(10);
187
188// A signed write's deadline is capped at `expires_at + MAX_CLOCK_LEAD_MS`
189// (see `plan_clock`). That stays below the replay prune grace, so a stalled
190// duplicate can never commit after its record could have been pruned.
191const _: () = assert!(MAX_CLOCK_LEAD_MS.unsigned_abs() < read::REPLAY_PRUNE_GRACE_MS);
192
193/// Default `ListRefs` scan page.
194pub const DEFAULT_LIST_PAGE_LIMIT: u32 = 1000;
195
196/// Counter: a write failed because its partition is full (00-plan P-24).
197/// Label `kind`; alert on any increase.
198pub const METRIC_PARTITION_FULL: &str = "mkit_server_partition_full_total";
199
200/// Counter: [`Pipeline::with_header`] dropped an invalid or reserved
201/// header; label `reason` (`name`, `reserved`, `value`).
202pub const METRIC_HEADER_DROPPED: &str = "mkit_server_error_header_dropped_total";
203
204/// Deployment routing for repository state.
205#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
206#[non_exhaustive]
207pub enum Sharding {
208    /// All rows of a namespace share one partition.
209    #[default]
210    Single,
211    /// Coordinator and per-branch ref shards (D34).
212    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/// What `authorize_read` established: the caller's facts (including its
232/// §10.1 view) and the stored grant epoch it checked a presented grant
233/// against.
234#[derive(Debug)]
235struct ReadAuth {
236    /// Facts threaded into `op.authz` for the hooks.
237    facts: AuthzFacts,
238    /// The coordinator's `e` row, `0` when absent; `None` when
239    /// visibility does not apply and no epoch was read. `IssueObjectUrl`
240    /// reuses it.
241    epoch: Option<u64>,
242}
243
244/// A deployment's pipeline settings. Start from [`PipelineConfig::new`].
245#[derive(Debug, Clone)]
246#[non_exhaustive]
247pub struct PipelineConfig {
248    /// How requests map to a repository (M0: `Single`).
249    pub addressing: Addressing,
250    /// How metadata partitions are routed.
251    pub sharding: Sharding,
252    /// Durable inspection mode, default-off and reserved for asynchronous inspection wiring.
253    pub inspection_mode: bool,
254    /// How requests authenticate.
255    pub auth: AuthMode,
256    /// Owner-signed write grant verifier for Multi/Owner deployments.
257    pub grants: Option<GrantConfig>,
258    /// Independent deployment-authority fence, optional and default-off.
259    pub authority_fence: Option<crate::authority::AuthorityFence>,
260    /// Write authorization policy; Open for Single, Owner for Multi.
261    pub write_policy: WritePolicy,
262    /// Visibility of repositories without an explicit stored setting (default public).
263    pub default_repo_visibility: RepoVisibility,
264    /// Role of the authorizer hook, defaulting to an additional check.
265    pub authorizer_role: AuthorizerRole,
266    /// Opt in to namespace-wide `ListRepos` authority grants, requiring returned writer view.
267    /// Default false: non-owner authority callers receive only the public listing.
268    pub list_repos_authority_full: bool,
269    /// Upload caps, supplied by the binding (used by M0-05b).
270    pub upload_limits: UploadLimits,
271    /// Optional tighter cap for legacy single-part `UploadPack` requests.
272    pub single_upload_max_bytes: Option<u64>,
273    /// Resumable upload part size: a power of two in 8–32 MiB.
274    pub part_size: u64,
275    /// Largest number of parts, sufficient to reach the upload byte cap.
276    pub max_parts: u32,
277    /// Largest requested `ListRefs` page, in 1..=10,000. The default 1000
278    /// refs with at most 512-byte names fit STC §7.9's 2 MiB page bound.
279    pub max_list_refs_page_size: u32,
280    /// Packs smaller than this may skip `BeginUpload`; `u64::MAX` means never
281    /// required. Advertised as zero with Multi addressing or admission.
282    pub begin_upload_threshold_bytes: u64,
283    /// Accepted deployment upload MAC keys; first key signs.
284    pub ticket_keys: Option<TicketKeys>,
285    /// Default-off private raw-pack scanner retrieval, with dedicated keys.
286    pub scanner_retrieval: Option<Arc<crate::scanner_retrieval::RetrievalConfig>>,
287    /// URL-token key set and lifetime for `IssueObjectUrl`
288    /// (SPEC-WRITE-GRANTS §9.4); `None` answers `unimplemented`.
289    pub url_tokens: Option<crate::url_token::UrlTokenConfig>,
290    /// Dedicated role keys, forbidden for client and owner authorization.
291    pub admin_keys: Vec<[u8; 32]>,
292    /// Receipt role key publication required by enabled takedown, without issuing receipts.
293    pub receipt_publication: Option<crate::takedown::PublicationConfig>,
294    /// Durable invalidation; absent keeps launch purge machinery inert.
295    pub purge: Option<crate::purge::PurgeConfig>,
296    /// Ticket lifetime, positive and strictly below seven days.
297    pub ticket_ttl_ms: u64,
298    /// Open-ticket bounds in each ref shard.
299    pub ticket_caps: TicketCaps,
300    /// Largest download chunk (used by M0-05b).
301    pub download_chunk_max: usize,
302    /// The default write quota: `Some(DEFAULT_WRITE_QUOTA)` for auth v2
303    /// deployments (`vcs-worker` parity). Under D34 it counts per ref shard;
304    /// Multi-addressing deployments also use it as the default namespace cap.
305    pub write_quota: Option<QuotaLimits>,
306    /// `ListRefs` scan page size, at least 1.
307    pub list_page_limit: u32,
308    /// Commit deadline window; see [`MAX_APPLY_WINDOW`].
309    pub max_apply_window: Duration,
310    /// Coordinator epoch lease duration, in milliseconds.
311    pub epoch_lease_ms: u64,
312    /// Safety margin in milliseconds. It must exceed the maximum clock skew
313    /// between every pipeline instance (grant, renewal and revoke), the sweep
314    /// driver, and every storage backend. The constructor cannot verify this.
315    pub lease_margin_ms: u64,
316    /// Minimum useful lease budget before renewing, in milliseconds.
317    pub min_lease_budget_ms: u64,
318    /// Extra header names never to log.
319    pub redactor: Redactor,
320    /// Extra request credential names passed to admission, in addition to payment defaults.
321    pub admission_credential_headers: Vec<String>,
322    /// Soft per-shard cap on undelivered outcomes, events and purges.
323    pub outbox_backlog_cap: Option<OutboxBacklogCap>,
324    /// Indexed ingestion and pre-receive verification, off by default.
325    pub indexed: Option<crate::indexed::IndexedConfig>,
326    /// Global takedown proofs; the Worker launch enables them with preservation.
327    pub takedown_denial: bool,
328    /// Per-ref allowed signers and fast-forward-only rules (SPEC-SERVER
329    /// §9.7). Native embedders and the Worker launch can configure them. A
330    /// fast-forward-only rule needs `indexed`.
331    pub ref_policy: Option<crate::policy::RefPolicy>,
332    /// HTTP object serving (SPEC-HTTP-OBJECTS), off by default and
333    /// explicitly configured by the adapter. Requires [`Self::indexed`].
334    #[cfg(feature = "http-objects")]
335    pub http_objects: Option<crate::http_objects::HttpObjectsConfig>,
336}
337
338/// A soft, unguarded backlog threshold; concurrent admissions may overshoot
339/// by at most their in-flight terminal rows and encoded bytes.
340#[derive(Debug, Clone, Copy, PartialEq, Eq)]
341pub struct OutboxBacklogCap {
342    /// Terminal/event/purge rows.
343    pub rows: u64,
344    /// Key plus encoded-value bytes.
345    pub bytes: u64,
346}
347
348impl PipelineConfig {
349    /// One seam for advertised and enforced indexed mode.
350    pub(crate) fn indexed_mode(&self) -> bool {
351        self.indexed.is_some()
352    }
353
354    /// Defaults for `auth`: the default write quota only for auth v2.
355    #[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    /// The namespace policy advertised by `GetServerInfo` (STC §2.1).
412    #[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/// What the pipeline offers bindings and `GetServerInfo`.
425#[derive(Debug, Clone, Copy, PartialEq, Eq)]
426#[non_exhaustive]
427pub struct PipelineCapabilities {
428    /// `AdvanceRefs` commits head and packmap in one batch.
429    pub atomic_advance: bool,
430}
431
432/// [`Pipeline::health`]: whether each store answered its probe.
433#[derive(Debug, Clone, Copy, PartialEq, Eq)]
434pub struct HealthStatus {
435    /// The blob store.
436    pub blobs: bool,
437    /// The metadata store.
438    pub meta: bool,
439}
440
441impl HealthStatus {
442    /// Both stores are healthy.
443    #[must_use]
444    pub fn is_healthy(&self) -> bool {
445        self.blobs && self.meta
446    }
447}
448
449/// One `ListRefs` entry.
450#[derive(Debug, Clone, PartialEq, Eq)]
451pub struct RefEntry {
452    /// The name with the requested prefix stripped (SPEC-REFS §4).
453    pub name: String,
454    /// The object id.
455    pub id: Hash,
456}
457
458/// A `SetRepoVisibility` request: the signed envelope's choice or the
459/// unsigned owner-signed statement (SPEC-WRITE-GRANTS §9.1).
460#[derive(Clone, PartialEq, Eq)]
461#[non_exhaustive]
462pub enum VisibilityRequest {
463    /// Envelope mode: the signed request's `visibility`.
464    Envelope(Visibility),
465    /// `signed_statement` mode: the encoded `scheme:statement:blob` header.
466    Statement(String),
467}
468
469impl core::fmt::Debug for VisibilityRequest {
470    /// Never shows the raw statement.
471    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
482/// The request pipeline over blobs `B`, metadata `N` and hooks `H`.
483pub 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
524/// A fixed-message `internal` whose detail only the server logs.
525fn internal(detail: &'static str) -> ServerError {
526    ServerError::internal("ref store request failed", detail)
527}
528
529/// Map a storage failure to a redacted `internal` and log its detail.
530fn 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
544/// The Connect method name, e.g. `UpdateRef`: the `procedure` label.
545fn method(procedure: Procedure) -> &'static str {
546    let path = procedure.connect_path();
547    path.rsplit('/').next().unwrap_or(path)
548}
549
550/// Stage 0 lookup's answer for a replay decision: `None` continues.
551fn 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        // `Resume` is for `UploadPack` only (M0-05b); a unary write that
559        // finds its record in flight waits like any other.
560        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/// Resolve visibility for reads and snapshot publication. Explicit rows always
571/// override the deployment default; callers strongly read the row, never cache it.
572#[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
579/// The `rv` codec's visibility for a statement/envelope [`Visibility`].
580fn 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    // Transport identity carries no tickets: the ssh and enc transports
592    // earn pack membership implicitly (WP-1.15's session pending set), so
593    // neither the Multi refusal nor the ticket threshold applies to it.
594    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    /// A pipeline over `blobs` and `meta`, routed by `cfg.sharding`.
623    ///
624    /// # Errors
625    /// `invalid_argument` for a configuration the store cannot serve: auth
626    /// v2 needs every key class and atomic multi-key batches (so
627    /// `FsLayoutStore` never runs auth v2); a store without atomic batches
628    /// must report an implicit layout version; a store's layout version
629    /// must be this binary's; the page limit and apply window must be
630    /// positive. Resumable upload and advertised page limits must be valid
631    /// and the part capacity must reach the upload byte cap. Namespace/write
632    /// policy combinations must be compatible;
633    /// `any` requires non-default admission or its explicit unsafe override,
634    /// and an authority authorizer must not be the open default.
635    #[allow(clippy::too_many_lines)] // Startup rejects incompatible storage, auth and grant combinations together.
636    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        // Collapsible only when `http-objects` is off.
668        #[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        // Defaults (Payment-Authorization, PAYMENT-SIGNATURE, Authorization)
753        // plus extras may never exceed the eight-header bound.
754        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            // Owner on Single is the ssh root mode's policy: it needs a
782            // self-certifying namespace (§7.4) to check principals against.
783            // A bare `root` Single has none and stays refused.
784            (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            // §9.4: the URL-token key is dedicated; a shared ticket secret
830            // would let ticket MACs stand in for URL-token signatures.
831            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    /// Total launch inspected-set cap, absent when inspection is disabled.
912    #[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    /// Configure the launch's synchronous, fail-closed inspectors.
922    ///
923    /// # Errors
924    /// Invalid launch settings or a deployment without indexed, ticketed,
925    /// restricted writes and atomic metadata.
926    #[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    /// Install test fault hooks (feature `test-faults` only).
965    #[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    /// Exclude an adapter's autonomous timer tick while a test directive drains.
973    /// The adapter must use this same gate around its own ticks.
974    #[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    /// Run the writes to one partition one at a time in this process: each
995    /// write's read-plan-apply loop waits for the partition's gate, as a
996    /// Durable Object's input gate serializes them on Workers. For a
997    /// single-process server over a single-writer store (native `SQLite`,
998    /// which commits one batch at a time anyway): concurrent writes that
999    /// share a key, such as one signer's quota window, then never exhaust
1000    /// the re-plan bound and fail `aborted`. D34 lease-grant batches also
1001    /// take this gate after admission. Reads are not gated. Several
1002    /// processes on one store still race, through the optimistic loop.
1003    #[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    /// A second pipeline over the same stores, hooks, shard map, clock,
1010    /// metrics, test fault hooks and write gate, authenticating with
1011    /// `auth`: how one server hosts bindings with different identity
1012    /// sources on one root (an enc listener's `TransportIdentity` beside
1013    /// an HTTP listener's bearer token or auth v2) while its writes to a
1014    /// partition still pass one gate. Every other setting is `self`'s.
1015    ///
1016    /// # Errors
1017    /// As [`Self::new`] for `auth` over these stores.
1018    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        // Sibling enc/ssh pipelines do not mint object URL tokens.
1027        cfg.url_tokens = None;
1028        // Grants need auth v2 (`Self::new` refuses them otherwise), and a
1029        // transport-identity write has no header-grant path: `GrantConfig::verify`
1030        // needs the signed operation. WP-2.12 (registered grants over ssh/enc)
1031        // replaces this with a sibling exception in `Self::new`.
1032        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    /// Stages 0a and 1: verify credentials and map the identity. Pure and
1063    /// synchronous; writes no state. The result is bound to
1064    /// `meta.procedure`. Under `test-faults` it also reads the request's
1065    /// test directives: the clock skew shifts business time, including the
1066    /// auth v2 validity window, for this request only.
1067    ///
1068    /// # Errors
1069    /// `unauthenticated` for missing or invalid credentials;
1070    /// `invalid_argument` for a malformed test directive. A rejection is
1071    /// recorded like any failed request (procedure, code, latency), with
1072    /// principal `none`: no entry point runs after it to record it.
1073    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        // SPEC-WRITE-GRANTS §4.2: a grant header without auth v2 fails on
1086        // any procedure of every deployment.
1087        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        // Credentials are captured for admission, which only signed writes
1128        // reach: signed reads and `SetRepoVisibility` never run it (§9.1).
1129        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        // Header adapters supply UTF-8 Strings; undecodable bytes are absent.
1138        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    /// Every ref of the repository under `prefix` at a path-component
1147    /// boundary, with the prefix and its `/` stripped (SPEC-REFS §4, see
1148    /// [`refs::list_scan_prefix`]), read page by page.
1149    ///
1150    /// # Errors
1151    /// `not_found` for a nonexistent Multi repository;
1152    /// `invalid_argument` for an invalid prefix or one over
1153    /// [`refs::MAX_REF_NAME_BYTES`]; the authorizer's error; `unavailable`
1154    /// for a bucket scan failure.
1155    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    /// One bounded `ListRefs` page. The token is opaque core bytes; the
1175    /// Connect binding encodes it as unpadded base64url.
1176    #[allow(clippy::too_many_lines)] // One authorized listing with shared token and response bounds.
1177    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    /// Attach an explicit published reader source (snapshot opt-in).
1317    #[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    /// Install the inspection preparation and immediate hold gate (post-launch WP-5.5c).
1325    /// Inspection requires indexed mode and an owner/authority write policy.
1326    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    /// One ref's id, if it exists.
1340    ///
1341    /// # Errors
1342    /// `not_found` for a nonexistent Multi repository;
1343    /// `invalid_argument` for an invalid name; the authorizer's error;
1344    /// `internal` for a storage failure.
1345    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    /// Compare-and-swap one ref. A conflict is a result, not an error.
1401    ///
1402    /// # Errors
1403    /// See the stage functions; a stored rejection comes back as its error.
1404    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    /// Compare-and-swap a ref and return success-only admission headers.
1415    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    /// Advance a branch head and its packmap together: one batch on an
1438    /// atomic store, else packmap then head (`Transport::advance_refs`'s
1439    /// default order).
1440    ///
1441    /// # Errors
1442    /// See the stage functions; a stored rejection comes back as its error.
1443    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    /// Advance refs and return success-only admission headers.
1455    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    /// Advance both refs while consuming the named upload tickets.
1466    ///
1467    /// # Errors
1468    /// Invalid ticket bindings, incomplete uploads, authorization or storage errors.
1469    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    /// Advance refs with tickets and return success-only admission headers.
1482    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            // A ref policy or indexed mode checks the head only (and its
1517            // packmap through it), so the pair must be a real pair.
1518            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    /// Whether the pack is present in a Single deployment's blob store.
1562    /// Multi deployments require repository membership before serving packs.
1563    ///
1564    /// # Errors
1565    /// `not_found` for a nonexistent Multi repository; the authorizer's
1566    /// error; `internal` for a storage failure.
1567    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    /// `IssueObjectUrl` (SPEC-WRITE-GRANTS §9.4): mint a token binding the
1586    /// deployment's audience, the request's repository identity and
1587    /// `target` at the stored grant epoch. The target is never resolved —
1588    /// serving decides what it names, so minting reveals nothing about
1589    /// the repository's contents. The token is a credential: it is never
1590    /// logged.
1591    ///
1592    /// # Errors
1593    /// `unimplemented` when no URL-token key is configured, before any
1594    /// repository access; `unauthenticated` for an unsigned request
1595    /// (stage 0 already rejects it under auth v2); the read errors of
1596    /// `Self::authorize_read`, including the uniform `not_found` on a
1597    /// private repository.
1598    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    // Shared below the wire envelope boundary: embedders transfer only verified
1628    // reader authority, while the RPC still requires its own signed envelope.
1629    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    /// `SetRepoVisibility` (SPEC-WRITE-GRANTS §9.1): the envelope mode is a
1659    /// signed, replay-protected write of the `rv` row; the statement mode
1660    /// verifies an unsigned owner-signed statement and keeps the newest
1661    /// `created`. Only deployments where `Self::visibility_applies`
1662    /// holds have visibility at all.
1663    ///
1664    /// # Errors
1665    /// `unauthenticated` for an unsigned envelope request, `invalid_argument`
1666    /// for a statement on a signed request, `failed_precondition` when the
1667    /// deployment has no repository visibility, `permission_denied` for a
1668    /// grant, a non-owner, an unserved namespace or a rejected statement.
1669    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    /// Envelope mode: the stage-0 replay lookup, owner-or-hook
1720    /// authorization mirroring the write branch (never a grant), then the
1721    /// guarded `rv` write with its replay record.
1722    #[allow(clippy::too_many_lines)] // Replay, authorization and both guarded visibility retries share one lifecycle.
1723    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                        // The replay guard failed: a request with this
1835                        // nonce landed first; answer it as stage 0 does.
1836                        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    /// Envelope-mode authorization: the owner, or (`authority` role) the
1870    /// hook, which is shown a reader until it decides. Hook errors are
1871    /// stripped of any admission shape.
1872    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            // The hook decides whether a non-owner may write; until it
1891            // does the caller is only a reader.
1892            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    /// One envelope attempt: deadline, replay `Absent`, the `rv` guard and
1909    /// its new row, the replay record and expiry index, then up to 32
1910    /// expired replay `(index, record)` deletes capped at `MAX_BATCH_OPS`.
1911    /// Returns the commit batch plus the delete-only batch a full
1912    /// partition can still apply.
1913    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                    // An envelope write is newer than any statement made
1944                    // before it: an older unsubmitted statement must not
1945                    // undo it (§9.1).
1946                    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    /// Statement mode: the owner-signed statement is its own
1990    /// authorization; only a strictly newer `created` replaces the stored
1991    /// one. No hook, no replay record, no admission.
1992    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                // The stored statement is this one: already applied.
2035                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    /// Stages 0–3 of an `UploadPack` whose header declared `pack_id` and
2099    /// `total_bytes` (`None` when absent): framing, the signed `pack:`
2100    /// commitment, the replay lookup, authorization, admission and, for a
2101    /// new signed operation, the reservation. Nothing is read from the
2102    /// stream before it returns.
2103    ///
2104    /// # Errors
2105    /// `failed_precondition` when the advertised threshold requires a ticket;
2106    /// the header's [`crate::upload::UploadError`]; `unauthenticated` when
2107    /// the header differs from the signed commitment; a stored or
2108    /// in-flight replay answer; a hook's error; the reservation's error.
2109    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    /// Open a stateless ticketed `UploadPack` stream after checking the header,
2119    /// signed commitment and token, before reading any body bytes.
2120    ///
2121    /// # Errors
2122    /// Framing, authentication, binding and blob-store failures.
2123    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    /// A pack's bytes as chunks of at most `download_chunk_max` bytes.
2134    ///
2135    /// # Errors
2136    /// `not_found` for a missing repository or non-member pack, before any chunk; the authorizer's
2137    /// error; `internal` for a storage failure.
2138    ///
2139    /// The request is recorded `ok` when the `last` chunk is yielded, with
2140    /// its error at the first failure, and as `canceled` when the stream is
2141    /// dropped before either.
2142    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    /// Probe both stores.
2178    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    /// How this pipeline authenticates (the ssh session requires
2186    /// `TransportIdentity`; an HTTP adapter may pre-check a bearer token
2187    /// before it spends resources on the request).
2188    #[must_use]
2189    pub fn auth_mode(&self) -> &AuthMode {
2190        &self.cfg.auth
2191    }
2192
2193    /// The upload caps `begin_upload` applies: a binding that validates
2194    /// framing itself uses the same ones.
2195    #[cfg(feature = "ssh")]
2196    pub(crate) fn upload_limits(&self) -> UploadLimits {
2197        self.cfg.upload_limits
2198    }
2199
2200    /// The metadata store, for the ssh tests' state checks.
2201    #[cfg(all(test, feature = "ssh"))]
2202    pub(crate) fn meta_store(&self) -> &N {
2203        &self.meta
2204    }
2205
2206    /// What this pipeline offers.
2207    pub fn capabilities(&self) -> PipelineCapabilities {
2208        PipelineCapabilities {
2209            atomic_advance: self.meta.capabilities().atomic_multi_key,
2210        }
2211    }
2212
2213    /// Add a response header to `err`, counting a dropped one in
2214    /// [`METRIC_HEADER_DROPPED`] (what [`ServerError::with_header`] cannot
2215    /// do without a metrics sink). The value is never logged.
2216    #[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    /// The span, request metrics and outcome log around one entry point.
2235    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    /// The span and recorder of one request.
2247    fn outcome(&self, a: &Authenticated) -> RequestOutcome {
2248        self.outcome_for(a.procedure(), a.principal.kind(), &a.repo().identity)
2249    }
2250
2251    /// The span and recorder of one request to `procedure` as `principal`.
2252    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    /// A signed or unsigned unary write, stage by stage. In steady state
2265    /// a signed write costs two backend calls: one `get_many` before any
2266    /// hook runs (the replay record and the snapshot) and one `apply`.
2267    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    /// [`Self::write`], optionally consuming the session's pending packs
2276    /// as implicit tickets (WP-1.15): the B10 packmap check runs after
2277    /// authorization, admission is skipped exactly when `pending` is
2278    /// non-empty, and under Multi the plan adds membership and relay rows.
2279    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)] // Stage order and multipart session cleanup share this entry point.
2289    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        // §9.7: the signer rule needs no verified content, so it runs first.
2340        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                // The wire conformance fixture exercises the pending response
2361                // after ticket proof validation, before any verification state
2362                // or replay row is written. Release builds omit this seam.
2363                #[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                    // Kind-7 slices verify; the advance only checks their result.
2417                    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            // The B10 check may rewrite `Any` to the exact value it
2464            // observed, so the planned batch guards it.
2465            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        // An implicit consuming write skips admission exactly when it has
2473        // pending packs to consume; an empty pending set runs admission
2474        // like any other UpdateRef (the B10 check still applied).
2475        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                        // The native gate serializes same-shard lease grants too.
2530                        // Read observations remain pre-admission; a waiter rebuilds
2531                        // a stale coordinator observation within the usual three tries.
2532                        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        // A ticket-derived storage id can be shared by two attempts in the
2631        // same replay scope. The loser must leave it for the winning ticket;
2632        // if neither commits, the session contains only meta until the sweep.
2633        if session == fresh_id {
2634            return;
2635        }
2636        // A raced Existing ticket can have the same reservation-derived id
2637        // with a different session. Compare the authenticated answer.
2638        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        // An apply can commit and then lose its acknowledgement. Check the
2651        // row before reclaiming; an unavailable or corrupt read is ambiguous.
2652        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    /// Stage 1: the typed operation, for the procedure `a` was
2675    /// authenticated for only.
2676    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    /// Multi reads require a repository registered in the namespace coordinator.
2700    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    /// Whether repository visibility (`rv` rows) gates reads: a Multi,
2719    /// owner-policy, auth v2 deployment (SPEC-WRITE-GRANTS §9.1).
2720    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    /// Whether a stored `rv = private` gates reads: any Multi
2727    /// owner-policy deployment, whatever its auth mode. A pipeline that
2728    /// cannot verify a signer (a non-auth-v2 sibling) reads every private
2729    /// repository as `not_found`.
2730    fn visibility_gates_reads(&self) -> bool {
2731        matches!(self.cfg.addressing, Addressing::Multi(_))
2732            && self.cfg.write_policy == WritePolicy::Owner
2733    }
2734
2735    /// Stage 2 for a read. Under [`Self::visibility_gates_reads`], one
2736    /// coordinator `get_many` reads `rr`, `rv` and `e`; a private
2737    /// repository denies unauthorized reads with the same `not_found` as a
2738    /// missing one (SPEC-WRITE-GRANTS §9.3). The returned facts carry the
2739    /// caller's view; `epoch` is the stored grant epoch, for
2740    /// `IssueObjectUrl`'s reuse.
2741    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        // §7 steps 1–10 are stateless: run them before the coordinator
2747        // read, even when the repository turns out to be missing.
2748        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            // A missing repository pays the hook round trip a private one
2766            // would, and discards it, so latency does not tell them apart.
2767            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        // §7 step 11: the grant's epoch must equal the stored epoch.
2783        let grant = check.filter(|c| c.epoch == epoch);
2784        let grant_ref = grant.map(|c| GrantRef {
2785            id: c.id,
2786            epoch,
2787            // Ref-scope presence constraints govern writes only.
2788            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    /// One coordinator `get_many` for the `rr`, `rv` and `e` rows of a
2832    /// read: `Some((private, epoch))`. A missing `rr` is `None`
2833    /// (the caller answers the uniform `not_found`); a store failure is `unavailable`, never public.
2834    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    /// Consult the authorizer for a read on an existing repository under
2879    /// `visibility_gates_reads`: on a private repository its verdict can
2880    /// authorize the read; on a public one it only classifies the caller
2881    /// (`Authority`) or checks the read as before (`Check`, error
2882    /// propagates). `provisional` are the facts the hook sees.
2883    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                // A presented grant is the only path to a private read
2899                // (§6, §9.3): the hook never authorizes it.
2900                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                    // Any hook error, `unavailable` included, denies so the
2908                    // read cannot reveal that the repository exists (§9.3).
2909                    Err(_) => read_policy::HookEval::Deny,
2910                },
2911            );
2912        }
2913        if self.cfg.authorizer_role == AuthorizerRole::Authority {
2914            // Classification only: an unsigned caller or an established
2915            // writer needs no hook; a hook error keeps the reader view.
2916            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        // Check role: the hook sees the read as before; it cannot confer
2929        // a view.
2930        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    /// The read path when [`Self::visibility_gates_reads`] does not hold: the
2939    /// hook decides, then the `rr` row gates existence (Multi only). The
2940    /// caller's view comes from its principal and the write policy.
2941    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    /// A unary write's ref writes in decision order (packmap first) and
3021    /// the ref shard they commit in, which a head and its packmap share.
3022    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    /// One `get_many` before any hook, on an atomic store: the write's
3051    /// refs and layout version, and for a signed write its replay record,
3052    /// the grant epoch and the default quota key it will most likely be
3053    /// charged. Keys admission adds later are read after it.
3054    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                // Admission may run after the current window ends. Fetch the
3127                // next window in the same read-ahead call, before any lease.
3128                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    /// Reject a namespace cap before a lease or multipart session is made.
3178    /// The planner repeats the guarded check when it commits the write.
3179    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(&quota::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    /// Read every key of `wanted` that `snap` lacks, in one `get_many`.
3252    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    /// Stage 0 lookup: a signed write's replay record, from the read-ahead
3272    /// and before any hook. A committed record's result is returned, an
3273    /// in-flight one is `aborted`; a new operation continues.
3274    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    /// Stage 2: the facts it returns become `op.authz` before admission.
3305    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            // Owner on a namespaced Single is the ssh root mode's rule;
3319            // Open Single is the authorizer's alone, unchanged.
3320            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                    // The authority fence is valid only for Multi addressing.
3330                    // A hook cannot activate it on an unfenced Single deployment.
3331                    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    /// SPEC-TRANSPORT-CONNECT §7.5 rule 1: an allowlisted namespace whose
3352    /// principal owns it (`ed25519-<key>` ↔ the Ed25519 principal), and the
3353    /// M2 write grants that qualify it. `policy` is the deployment's
3354    /// `NamespacePolicy` on Multi and `None` on an Owner-policy Single —
3355    /// the `None` counts as "not an allowlist" for the 0x-plus-grant rule,
3356    /// which denies fail closed. The authorizer hook sees the established
3357    /// facts on success; under `Check` a non-owner never reaches it.
3358    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            // Step 8 (ref scope) before step 11, which reads state and comes last.
3386            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        // Over ssh/enc (no grant header) a 0x namespace has no owner, so
3410        // its writes stay denied until WP-2.12.
3411        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        // Both Authorize and Admit see the established owner/grant facts (§6.2).
3421        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    /// Stage 5 (no pack on a unary write).
3434    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    /// Stages 4 and 6: plan and apply, as one batch on an atomic store or
3444    /// as sequential single-ref batches on a non-atomic one.
3445    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        // Single's packs live in the repo directory itself; only Multi
3466        // plans `m` rows and relay for the consumed set, and only when
3467        // the set is non-empty (L2).
3468        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        // `Transport::advance_refs`'s default: packmap first, then head.
3561        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    /// Prepare inspection against the complete resulting pair before any write.
3613    #[allow(clippy::too_many_arguments, clippy::too_many_lines)] // Carry authenticated wire identity and ticket lag context.
3614    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        // Deletions establish an immediate boundary without consulting inspection
3628        // or verifying the surviving pair; older membership obligations remain retained.
3629        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            // Membership dependencies belong to the server, not to an
3651            // inspector's verdict. An Inspect Pass alone cannot publish.
3652            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)] // Carry consuming-ticket age for repository lag classification.
3743    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        // The canonical fallback retains delta bases; only the metadata-only
3774        // continuation can use the whole-job allowance without that residency.
3775        let mut indexed = indexed;
3776        if self.cfg.takedown_denial && !resumed {
3777            indexed.decode_budget = indexed.decode_budget.min(8 << 20);
3778        }
3779        // Pair verification and dependency visibility share one allocation.
3780        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    /// [`Self::apply_loop`] on a store with atomic multi-key batches, which
3863    /// replay records and quota need.
3864    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    /// The bounded optimistic loop: read, plan, apply. The first attempt
3879    /// retains read-ahead only when no global denial proof is required.
3880    /// Each attempt fixes its plan time before proof; fresh rows, a usable
3881    /// lease and current business time follow it. A guard another
3882    /// writer broke re-plans up to [`MAX_REPLAN`] times, then `aborted`; a
3883    /// lost prune race retries once without the prune, uncounted. A missed
3884    /// deadline re-plans once while the envelope is still valid at the
3885    /// failed commit, then `unavailable` (SPEC-WRITE-GRANTS §5.5).
3886    // Keep fresh denial, authorization and guarded apply together on every retry.
3887    #[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        // Held until the loop ends (see `with_write_gate`).
3898        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            // Every gating denial read is at/after this attempt's plan time
3923            // (SPEC-SERVER §14.2). Slow proof cannot borrow a new commit window.
3924            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            // A complete proof can outlive both the commit window and the
3937            // initial lease. Only fixed admission quota facts survive it;
3938            // every mutable source row is read again before planning.
3939            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                // Keep fresh business time and current lease/replay caps; only
3967                // the denial attempt's original plan time remains immutable.
3968                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                    // Validity at the failed commit, on both clocks.
4014                    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    /// The deadline uses the injected clock unshifted; business time adds
4081    /// the request's skew. A signed write's deadline is also capped at
4082    /// `expires_at + MAX_CLOCK_LEAD_MS`, below the replay prune grace.
4083    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    /// Everything [`plan_write`] reads that `base` lacks, in one
4106    /// `get_many`. A sampled write ([`prune_sampled`]) first scans a
4107    /// bounded page of prune candidates, on the real clock, never the
4108    /// business clock.
4109    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    /// Apply a batch, typing a full partition as retryable `unavailable`.
4153    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    /// A full partition: count it, retry the prune alone (deletes still
4162    /// work) and fail closed with a retryable `unavailable`.
4163    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/// The allowance stays local until the guarded apply.
4178#[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
4186/// A ref name the pipeline reads or writes: at most
4187/// [`refs::MAX_REF_NAME_BYTES`], the SPEC-REFS §3 grammar, and under
4188/// `refs/` ([`refs::is_served_ref_name`], R-86), each refused by name.
4189fn 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
4203/// A result of another procedure: a stored rejection is its error; any
4204/// other kind is corruption (the fingerprint covers the procedure).
4205fn 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
4214/// The replay guard failed: another request with this nonce committed
4215/// first. Classify its record like stage 0 does.
4216fn replay_raced(
4217    req: &WriteRequest<'_>,
4218    observed: Option<&Value>,
4219) -> Result<StoredResult, ServerError> {
4220    // Receipts and a Committed outcome belong to this request only (brief
4221    // B7): a same-nonce loser holding its own reservation aborts it, whatever
4222    // the winner stored.
4223    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}