Skip to main content

Pipeline

Struct Pipeline 

Source
pub struct Pipeline<B, N, H = Hooks> { /* private fields */ }
Expand description

The request pipeline over blobs B, metadata N and hooks H.

Implementations§

Source§

impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H>

Source

pub async fn get_authority_generation( &self, namespace: &str, ) -> Result<u64, ServerError>

Read the optional namespace authority generation outside auth-v2.

§Errors

Disabled fencing, invalid namespace or storage failure.

Source

pub async fn set_authority_generation( &self, wire: &str, ) -> Result<u64, ServerError>

Verify a deployment statement, install its target and finish the lease barrier. Completion retries are idempotent; each call scans at most one four-shard slice.

§Errors

Disabled configuration, rejected statement or pending completion (Retry-After: 1).

Source§

impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H>

Source

pub async fn begin_upload( &self, a: &Authenticated, ref_name: &str, pack_id: &[u8], bytes: u64, ) -> Result<BeginUploadResult, ServerError>

Open a stateless authenticated upload ticket in the target ref shard.

§Errors

Invalid geometry/ref, unsupported auth or missing keys, policy or cap refusal, and the usual replay, quota, lease and storage errors.

Source

pub async fn begin_upload_with_meta( &self, a: &Authenticated, ref_name: &str, pack_id: &[u8], bytes: u64, ) -> Result<(BeginUploadResult, ResponseMeta), ServerError>

Begin a resumable upload and return success-only admission headers.

Source§

impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H>

Source

pub async fn get_grant_epoch(&self, namespace: &str) -> Result<u64, ServerError>

Read a namespace epoch without authentication or repository resolution.

§Errors

invalid_argument for a bad namespace, unimplemented under Single, or a mapped storage error.

Source

pub async fn set_grant_epoch( &self, signed_statement: &str, ) -> Result<u64, ServerError>

Verify an owner epoch statement, CAS the coordinator epoch, then fence every leased shard before reporting completion.

§Errors

permission_denied for a rejected statement, unimplemented under Single, or unavailable while completion is pending.

Source§

impl<B: MultipartBlobStore, N: NamespaceStore + Clone + 'static, H: HookSet> Pipeline<B, N, H>

Source

pub fn http_objects_enabled(&self) -> bool

Whether this pipeline can mount HTTP object routes. Opaque pipelines cannot.

Source

pub fn url_token_config(&self) -> Option<&UrlTokenConfig>

Public-key publication configuration; never exposes seeds.

Source

pub fn with_http_seams(self, edit: impl FnOnce(HttpSeams) -> HttpSeams) -> Self

Replace the inert seams of an HTTP-objects pipeline: admission (WP-4.13), proofs (WP-4.14b), tokens (WP-4.15), takedown (WP-5.9a) or a maintained reachable set (WP-5.3a). A pipeline built without PipelineConfig::http_objects is unchanged.

Source

pub async fn serve_http_object( &self, req: &HttpObjectRequest<'_>, ) -> HttpObjectResponse

Serve one HTTP object request (GET, HEAD or OPTIONS) and never fail: every error is mapped to its §3 response. The principal is always anonymous, whatever credentials the request carries (§7); the raw path, query and any token are never logged.

Source

pub async fn serve_http_object_with_runtime( &self, req: &HttpObjectRequest<'_>, runtime: HttpReadRuntime, ) -> HttpObjectResponse

Adapter-provided retained settlement runtime for this request.

Source

pub async fn serve_http_object_with_proofs( &self, req: &HttpObjectRequest<'_>, runtime: HttpReadRuntime, proofs: Arc<dyn ProofServer>, ) -> HttpObjectResponse

Adapter-provided proof builder and retained settlement runtime.

Source§

impl<B: MultipartBlobStore, N: NamespaceStore + Clone + 'static, H: HookSet> Pipeline<B, N, H>

Source

pub async fn object_reader<'a>( &'a self, repo: RepoId, view: ReaderView<'a>, ) -> Result<ObjectReader<'a, B, N, H>, ServerError>

Construct a repository reader using verified owner authority.

§Errors

Missing indexed/HTTP configuration, or invalid/non-owner credentials.

Source§

impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H>

Source

pub fn server_info(&self) -> ServerInfo

Configured deployment information. Never resolves a repository or reads a store; the atomic flag comes from store capabilities only.

Source§

impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H>

Source

pub async fn list_repos_page( &self, a: &Authenticated, namespace: &str, name_prefix: &str, page_size: Option<u32>, token: Option<&[u8]>, ) -> Result<RepoPage, ServerError>

List a namespace without visiting private rows for the public view. Signed envelopes use the ordinary X-Repository syntax; its namespace must match namespace. The name is only a namespace selector for this RPC. Grants never extend listing rights. Check and Authority hooks concern the entire namespace.

§Errors

Invalid namespace/prefix/size/token, mismatched authentication, or unavailable storage/hooks.

Source§

impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H>

Source

pub async fn open_part( &self, a: &Authenticated, token: &[u8], index: u32, ) -> Result<PartUploadSession<'_, B, N, H>, ServerError>

Validate a part header and ticket before opening any part sink.

§Errors

Invalid or mismatched ticket, commitment or part geometry; storage errors while opening the part.

Source

pub async fn complete_upload( &self, a: &Authenticated, token: &[u8], receipts: &[Vec<u8>], ) -> Result<(), ServerError>

Authenticate all receipts and the merged pack root before any store call; completion writes no metadata rows.

§Errors

Invalid ticket, receipt, root, length or storage session.

Source§

impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H>

Source

pub async fn plan_repository_purge( &self, partition: &Partition, repo: &RepoId, trigger: Trigger, operation_id: &str, now_ms: u64, ) -> Result<Batch, ServerError>

Plan an audited safety purge in the same batch as a serving stop. operation_id is stable across retries of the automatic action. Callers merge this batch into the authoritative state mutation.

Source

pub async fn invalidate_local_cache(&self, repo: &RepoId)

Immediately invalidate local entries after committing a serving stop and its durable purge work. Timer 11 retries any failed cache deletion.

Source§

impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H>

Source

pub async fn bump_epoch( &self, ns: &NamespaceKey, new_epoch: u64, ) -> Result<(), ServerError>

Advance the coordinator epoch, serialized with every epoch lease grant.

§Errors

invalid_argument unless the increment is in 1..=1024; unavailable on repeated contention, or a mapped storage failure.

Source

pub async fn mark_lease_table_recovered( &self, ns: &NamespaceKey, ) -> Result<(), ServerError>

Declare lease-table recovery using the real pipeline clock. Any procedure that restores or rebuilds a coordinator partition MUST call this before serving writes for that namespace. The durable marker prevents completion for epoch_lease + margin after recovery.

§Errors

A mapped storage failure, or internal for an impossible batch outcome.

Source

pub async fn revoke_step( &self, ns: &NamespaceKey, budget: &RevokeBudget, ) -> Result<RevokeProgress, ServerError>

Push and acknowledge at most four live shards in one bounded slice. Outside a declared recovery hold-off, ls.acked_epoch = n only if the shard’s el durably holds epoch >= n, or all older-epoch writes are already past their commit deadline. Recovery fences surviving copies until then. Renewal alone never raises a live row’s acknowledgement.

§Errors

A mapped storage or codec failure. Contention leaves progress pending.

Source§

impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H>

Source

pub fn scanner_retrieval_enabled(&self) -> bool

Whether adapters may mount the private route. Default false.

Source

pub async fn retrieve_scanner_pack( &self, body: &[u8], headers: &Headers, ) -> Result<RetrievalResponse, ServerError>

Verify both credentials, strongly check existing bound ticket state and global denial, then read one bounded raw pack range. No public membership authorization or pack decoding is involved.

§Errors

Uniform not_found, including storage errors and exhausted budgets.

Source§

impl<B: BlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H>

Source

pub async fn namespace_relay_watermark( &self, ns: &NamespaceKey, ) -> Result<u64, ServerError>

Namespace lower bound on commit time of undelivered relay rows for native GC and takedown. Consumers compare it against T + MAX_APPLY_WINDOW + margin. Workers use the resumable store step. A new shard’s stale-low first report or recovery can lower the value; it is unavailable until the recovered lease table is reconciled.

§Errors

Unavailable during lease-table recovery or on an unreadable store.

Source

pub async fn active_shards( &self, ns: &NamespaceKey, cursor: Option<&Cursor>, limit: u32, ) -> Result<ActiveShardsPage, ServerError>

Enumerate coordinator ref-shard rows for a safety re-check. Expired rows retained for relay backlog are included.

§Errors

Unavailable during lease-table recovery or on an unreadable store.

Source§

impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H>

Source

pub fn new( blobs: B, meta: N, hooks: H, cfg: PipelineConfig, clock: Arc<dyn Clock>, metrics: Arc<dyn Metrics>, ) -> Result<Self, ServerError>

A pipeline over blobs and meta, routed by cfg.sharding.

§Errors

invalid_argument for a configuration the store cannot serve: auth v2 needs every key class and atomic multi-key batches (so FsLayoutStore never runs auth v2); a store without atomic batches must report an implicit layout version; a store’s layout version must be this binary’s; the page limit and apply window must be positive. Resumable upload and advertised page limits must be valid and the part capacity must reach the upload byte cap. Namespace/write policy combinations must be compatible; any requires non-default admission or its explicit unsafe override, and an authority authorizer must not be the open default.

Source

pub fn inspection_max_objects(&self) -> Option<u32>

Total launch inspected-set cap, absent when inspection is disabled.

Source

pub fn with_inspectors( self, inspectors: Vec<Arc<dyn ContentInspector>>, batch_max: usize, ) -> Result<Self, ServerError>

Configure the launch’s synchronous, fail-closed inspectors.

§Errors

Invalid launch settings or a deployment without indexed, ticketed, restricted writes and atomic metadata.

Source

pub fn with_write_gate(self) -> Self

Run the writes to one partition one at a time in this process: each write’s read-plan-apply loop waits for the partition’s gate, as a Durable Object’s input gate serializes them on Workers. For a single-process server over a single-writer store (native SQLite, which commits one batch at a time anyway): concurrent writes that share a key, such as one signer’s quota window, then never exhaust the re-plan bound and fail aborted. D34 lease-grant batches also take this gate after admission. Reads are not gated. Several processes on one store still race, through the optimistic loop.

Source

pub fn with_auth(&self, auth: AuthMode) -> Result<Self, ServerError>
where B: Clone, N: Clone, H: Clone,

A second pipeline over the same stores, hooks, shard map, clock, metrics, test fault hooks and write gate, authenticating with auth: how one server hosts bindings with different identity sources on one root (an enc listener’s TransportIdentity beside an HTTP listener’s bearer token or auth v2) while its writes to a partition still pass one gate. Every other setting is self’s.

§Errors

As Self::new for auth over these stores.

Source

pub fn authenticate( &self, meta: &RequestMeta<'_>, ) -> Result<Authenticated, ServerError>

Stages 0a and 1: verify credentials and map the identity. Pure and synchronous; writes no state. The result is bound to meta.procedure. Under test-faults it also reads the request’s test directives: the clock skew shifts business time, including the auth v2 validity window, for this request only.

§Errors

unauthenticated for missing or invalid credentials; invalid_argument for a malformed test directive. A rejection is recorded like any failed request (procedure, code, latency), with principal none: no entry point runs after it to record it.

Source

pub async fn list_refs( &self, a: &Authenticated, prefix: &str, ) -> Result<Vec<RefEntry>, ServerError>

Every ref of the repository under prefix at a path-component boundary, with the prefix and its / stripped (SPEC-REFS §4, see refs::list_scan_prefix), read page by page.

§Errors

not_found for a nonexistent Multi repository; invalid_argument for an invalid prefix or one over refs::MAX_REF_NAME_BYTES; the authorizer’s error; unavailable for a bucket scan failure.

Source

pub fn with_published_source(self, source: Arc<dyn PublishedSource>) -> Self

Attach an explicit published reader source (snapshot opt-in).

Source

pub fn with_publication_policy( self, policy: Arc<dyn PublicationPolicy>, ) -> Result<Self, ServerError>

Install the inspection preparation and immediate hold gate (post-launch WP-5.5c). Inspection requires indexed mode and an owner/authority write policy.

Source

pub async fn read_ref( &self, a: &Authenticated, name: &str, ) -> Result<Option<Hash>, ServerError>

One ref’s id, if it exists.

§Errors

not_found for a nonexistent Multi repository; invalid_argument for an invalid name; the authorizer’s error; internal for a storage failure.

Source

pub async fn update_ref( &self, a: &Authenticated, upd: RefUpdate, ) -> Result<UpdateRefResult, ServerError>

Compare-and-swap one ref. A conflict is a result, not an error.

§Errors

See the stage functions; a stored rejection comes back as its error.

Source

pub async fn update_ref_with_meta( &self, a: &Authenticated, upd: RefUpdate, ) -> Result<(UpdateRefResult, ResponseMeta), ServerError>

Compare-and-swap a ref and return success-only admission headers.

Source

pub async fn advance_refs( &self, a: &Authenticated, head: RefUpdate, packmap: RefUpdate, ) -> Result<AdvanceOutcome, ServerError>

Advance a branch head and its packmap together: one batch on an atomic store, else packmap then head (Transport::advance_refs’s default order).

§Errors

See the stage functions; a stored rejection comes back as its error.

Source

pub async fn advance_refs_with_meta( &self, a: &Authenticated, head: RefUpdate, packmap: RefUpdate, ) -> Result<(AdvanceOutcome, ResponseMeta), ServerError>

Advance refs and return success-only admission headers.

Source

pub async fn advance_refs_with_tickets( &self, a: &Authenticated, head: RefUpdate, packmap: RefUpdate, tickets: Vec<Hash>, ) -> Result<AdvanceOutcome, ServerError>

Advance both refs while consuming the named upload tickets.

§Errors

Invalid ticket bindings, incomplete uploads, authorization or storage errors.

Source

pub async fn advance_refs_with_tickets_with_meta( &self, a: &Authenticated, head: RefUpdate, packmap: RefUpdate, tickets: Vec<Hash>, ) -> Result<(AdvanceOutcome, ResponseMeta), ServerError>

Advance refs with tickets and return success-only admission headers.

Source

pub async fn pack_exists( &self, a: &Authenticated, key: PackKey, ) -> Result<bool, ServerError>

Whether the pack is present in a Single deployment’s blob store. Multi deployments require repository membership before serving packs.

§Errors

not_found for a nonexistent Multi repository; the authorizer’s error; internal for a storage failure.

Source

pub async fn issue_object_url( &self, a: &Authenticated, target: UrlTarget, ttl_seconds: u32, ) -> Result<MintedToken, ServerError>

IssueObjectUrl (SPEC-WRITE-GRANTS §9.4): mint a token binding the deployment’s audience, the request’s repository identity and target at the stored grant epoch. The target is never resolved — serving decides what it names, so minting reveals nothing about the repository’s contents. The token is a credential: it is never logged.

§Errors

unimplemented when no URL-token key is configured, before any repository access; unauthenticated for an unsigned request (stage 0 already rejects it under auth v2); the read errors of Self::authorize_read, including the uniform not_found on a private repository.

Source

pub async fn set_repo_visibility( &self, a: &Authenticated, req: VisibilityRequest, ) -> Result<(), ServerError>

SetRepoVisibility (SPEC-WRITE-GRANTS §9.1): the envelope mode is a signed, replay-protected write of the rv row; the statement mode verifies an unsigned owner-signed statement and keeps the newest created. Only deployments where Self::visibility_applies holds have visibility at all.

§Errors

unauthenticated for an unsigned envelope request, invalid_argument for a statement on a signed request, failed_precondition when the deployment has no repository visibility, permission_denied for a grant, a non-owner, an unserved namespace or a rejected statement.

Source

pub async fn open_upload( &self, a: &Authenticated, pack_id: Option<&[u8]>, total_bytes: Option<u64>, ) -> Result<UploadSession<'_, B, N, H>, ServerError>

Stages 0–3 of an UploadPack whose header declared pack_id and total_bytes (None when absent): framing, the signed pack: commitment, the replay lookup, authorization, admission and, for a new signed operation, the reservation. Nothing is read from the stream before it returns.

§Errors

failed_precondition when the advertised threshold requires a ticket; the header’s crate::upload::UploadError; unauthenticated when the header differs from the signed commitment; a stored or in-flight replay answer; a hook’s error; the reservation’s error.

Source

pub async fn open_ticketed_upload( &self, a: &Authenticated, pack_id: Option<&[u8]>, total_bytes: Option<u64>, token: &[u8], ) -> Result<UploadSession<'_, B, N, H>, ServerError>

Open a stateless ticketed UploadPack stream after checking the header, signed commitment and token, before reading any body bytes.

§Errors

Framing, authentication, binding and blob-store failures.

Source

pub async fn download( &self, a: &Authenticated, key: PackKey, ) -> Result<DownloadStream, ServerError>

A pack’s bytes as chunks of at most download_chunk_max bytes.

§Errors

not_found for a missing repository or non-member pack, before any chunk; the authorizer’s error; internal for a storage failure.

The request is recorded ok when the last chunk is yielded, with its error at the first failure, and as canceled when the stream is dropped before either.

Source

pub async fn health(&self) -> HealthStatus

Probe both stores.

Source

pub fn auth_mode(&self) -> &AuthMode

How this pipeline authenticates (the ssh session requires TransportIdentity; an HTTP adapter may pre-check a bearer token before it spends resources on the request).

Source

pub fn capabilities(&self) -> PipelineCapabilities

What this pipeline offers.

Source

pub fn with_header( &self, err: ServerError, name: &str, value: &str, ) -> ServerError

Add a response header to err, counting a dropped one in METRIC_HEADER_DROPPED (what ServerError::with_header cannot do without a metrics sink). The value is never logged.

Trait Implementations§

Source§

impl<B, N, H> Debug for Pipeline<B, N, H>

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more

Auto Trait Implementations§

§

impl<B, N, H = Hooks> !RefUnwindSafe for Pipeline<B, N, H>

§

impl<B, N, H = Hooks> !UnwindSafe for Pipeline<B, N, H>

§

impl<B, N, H> Freeze for Pipeline<B, N, H>
where B: Freeze, N: Freeze, H: Freeze,

§

impl<B, N, H> Send for Pipeline<B, N, H>
where B: Send, N: Send, H: Send,

§

impl<B, N, H> Sync for Pipeline<B, N, H>
where B: Sync, N: Sync, H: Sync,

§

impl<B, N, H> Unpin for Pipeline<B, N, H>
where B: Unpin, N: Unpin, H: Unpin,

§

impl<B, N, H> UnsafeUnpin for Pipeline<B, N, H>

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self> ⓘ

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> MaybeSend for T
where T: Send + ?Sized,

Source§

impl<T> MaybeSync for T
where T: Sync + ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more