Skip to main content

mkit_server/pipeline/
upload.rs

1//! `UploadPack` streams into a verifying blob sink with bounded memory.
2//!
3//! [`Pipeline::open_ticketed_upload`] checks framing, the signed `pack:`
4//! commitment and the stateless ticket before reading chunks. It skips the
5//! replay ledger, authorization, admission, quota and `pre_receive`. After
6//! the full pack verifies and commits, it writes a content-addressed upload
7//! marker in a non-pack blob namespace. It writes no metadata shard; bounded
8//! authoritative mode/generation reads guard staging and completion.
9//!
10//! [`Pipeline::open_upload`] remains for un-ticketed single-repository
11//! uploads below the advertised threshold, and for transport identity
12//! uploads. A fresh signed operation reserves a replay row and charges
13//! quota before streaming. The legacy STC ยง7.1 MAY (R-85) lets an in-flight
14//! retry with the same nonce re-read and commit the stream without charging
15//! again. A committed replay verifies the stream without writing a blob.
16//! `pre_receive` runs after the pack commit, followed by the replay commit.
17//! Rejection leaves an unreferenced pack for GC. If the envelope expires
18//! during the stream, the pack may succeed without the final replay write;
19//! the in-flight row is left for the pruner.
20
21use core::fmt;
22
23use bytes::Bytes;
24use mkit_core::hash::Hasher;
25use mkit_core::write_auth::MAX_CLOCK_LEAD_MS;
26use tracing::Instrument;
27
28use super::hooks::{AdmissionInput, PreReceive};
29use super::outcome::Outcome;
30use super::plan::{ReplayGuard, Snapshot, WriteKind, WriteRequest};
31use super::{AuthMode, Authenticated, HookSet, Pipeline, internal, meta_error, ms, store_error};
32use crate::auth_v2::check_pack_commitment;
33use crate::error::{Code, ServerError};
34use crate::op::{OpKind, Operation};
35use crate::replay::{ReplayDecision, StoredRejection, StoredResult, classify};
36use crate::storage_error::StorageOp;
37use crate::store::{
38    BlobStore, MultipartBlobStore, NamespaceStore, PackSink, Partition, StoreError, codec, keys,
39};
40use crate::telemetry::METRIC_UPLOAD_BYTES;
41use crate::upload::marker::write_upload_marker;
42use crate::upload::ticket_auth::verify_ticket;
43use crate::upload::{UploadError, UploadValidator};
44use mkit_core::hash::Hash;
45
46/// How an upload relates to its replay record.
47#[derive(Debug, Clone, Copy, PartialEq, Eq)]
48pub enum UploadMode {
49    /// A new operation, reserved and charged by `open_upload` (also every
50    /// unsigned upload).
51    Fresh,
52    /// The same operation is in flight: re-stream and re-commit, with no
53    /// second charge (legacy M0, see the module docs).
54    Resume,
55    /// The operation already committed: verify the stream without writing
56    /// it, then return OK without an apply.
57    Replay,
58    /// A stateless ticket authorizes the stream; bounded authority-mode reads
59    /// guard physical staging and completion, without metadata writes.
60    Ticketed,
61}
62
63/// Where a session's bytes go: the blob sink, or (on a replay) only a
64/// hasher.
65enum Target<S> {
66    Sink(S),
67    Verify(Box<Hasher>),
68}
69
70/// One `UploadPack` stream in progress. It borrows the pipeline: a
71/// binding that owns an `Arc<Pipeline>` keeps the `Arc` alive across the
72/// stream. Dropping it without [`Self::finish`] or [`Self::abort`] records
73/// the request as `canceled` and leaves nothing visible; an in-flight
74/// record stays resumable until it expires.
75pub struct UploadSession<'p, B: BlobStore, N, H> {
76    pipe: &'p Pipeline<B, N, H>,
77    a: Authenticated,
78    op: Operation,
79    p: Partition,
80    mode: UploadMode,
81    validator: UploadValidator,
82    target: Option<Target<B::Sink>>,
83    failed: Option<ServerError>,
84    outcome: Outcome,
85    ticket_id: Option<Hash>,
86    staging: super::staging::StagingBuffer,
87}
88
89impl<B: BlobStore, N, H> fmt::Debug for UploadSession<'_, B, N, H> {
90    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
91        f.debug_struct("UploadSession")
92            .field("mode", &self.mode)
93            .field("validator", &self.validator)
94            .finish_non_exhaustive()
95    }
96}
97
98/// What [`UploadSession::open`] established.
99struct Opened<S> {
100    op: Operation,
101    p: Partition,
102    mode: UploadMode,
103    validator: UploadValidator,
104    target: Target<S>,
105    ticket_id: Option<Hash>,
106}
107
108/// The replay record a signed operation commits.
109pub(super) fn replay_guard(op: &Operation) -> Option<ReplayGuard> {
110    op.auth.as_ref().map(|auth| ReplayGuard {
111        scope: auth.replay_scope,
112        fingerprint: auth.fingerprint,
113        expires_at_ms: auth.expires_at_ms,
114    })
115}
116
117/// Stage 0 for an upload: its mode from the replay record. A stored
118/// rejection is answered here, before any chunk.
119fn upload_mode(op: &Operation, ahead: Option<&Snapshot>) -> Result<UploadMode, ServerError> {
120    let (Some(auth), Some(snap)) = (&op.auth, ahead) else {
121        return Ok(UploadMode::Fresh);
122    };
123    let stored = snap.get(&keys::replay(&auth.replay_scope));
124    let record = stored.map(codec::decode_replay_record).transpose();
125    match classify(record.map_err(meta_error)?.as_ref(), &auth.fingerprint) {
126        ReplayDecision::New => Ok(UploadMode::Fresh),
127        ReplayDecision::Resume => Ok(UploadMode::Resume),
128        ReplayDecision::Return(StoredResult::UploadPack) => Ok(UploadMode::Replay),
129        ReplayDecision::Return(other) => Err(super::stored_mismatch(&other)),
130        ReplayDecision::FingerprintMismatch => Err(ServerError::invalid_argument(
131            "nonce reused for a different operation",
132        )),
133        ReplayDecision::RetryLater => Err(ServerError::aborted_retryable(
134            "operation already in flight; retry",
135        )),
136    }
137}
138
139/// A `pre_receive` error a retry may get back forever: a final code and no
140/// typed detail (so never an admission challenge).
141fn storable(err: &ServerError) -> Option<StoredRejection> {
142    if err.details().is_empty() {
143        StoredRejection::new(err.code(), err.public_message())
144    } else {
145        None
146    }
147}
148
149impl<'p, B: MultipartBlobStore, N: NamespaceStore, H: HookSet> UploadSession<'p, B, N, H> {
150    /// See [`Pipeline::open_upload`].
151    pub(super) async fn begin(
152        pipe: &'p Pipeline<B, N, H>,
153        a: &Authenticated,
154        pack_id: Option<&[u8]>,
155        total_bytes: Option<u64>,
156    ) -> Result<Self, ServerError> {
157        let mut outcome = pipe.outcome(a);
158        let opened = Self::open(pipe, a, pack_id, total_bytes)
159            .instrument(outcome.span.clone())
160            .await;
161        match opened {
162            Ok(o) => Ok(Self {
163                pipe,
164                a: a.clone(),
165                op: o.op,
166                p: o.p,
167                mode: o.mode,
168                validator: o.validator,
169                target: Some(o.target),
170                failed: None,
171                outcome,
172                ticket_id: o.ticket_id,
173                staging: super::staging::StagingBuffer::default(),
174            }),
175            Err(err) => {
176                outcome.record(Err(&err));
177                Err(err)
178            }
179        }
180    }
181
182    #[allow(clippy::too_many_lines)] // Upload target selection and replay validation share one lifecycle.
183    async fn open(
184        pipe: &'p Pipeline<B, N, H>,
185        a: &Authenticated,
186        pack_id: Option<&[u8]>,
187        total_bytes: Option<u64>,
188    ) -> Result<Opened<B::Sink>, ServerError> {
189        let mut limits = pipe.cfg.upload_limits;
190        if let Some(cap) = pipe.cfg.single_upload_max_bytes {
191            limits.max_total_bytes = limits.max_total_bytes.min(cap);
192        }
193        let validator = UploadValidator::new(pack_id, total_bytes, limits)?;
194        if !matches!(pipe.cfg.auth, AuthMode::TransportIdentity)
195            && validator.declared() >= pipe.effective_threshold()
196        {
197            return Err(ServerError::new(
198                Code::FailedPrecondition,
199                "upload requires a ticket from BeginUpload",
200            ));
201        }
202        let (key, declared) = (validator.key(), validator.declared());
203        let kind = OpKind::UploadPack {
204            key,
205            declared_len: declared,
206        };
207        let mut op = pipe.identify(a, kind)?;
208        if let Some(auth) = &op.auth {
209            check_pack_commitment(auth, &key.0, declared)
210                .map_err(|e| ServerError::unauthenticated(e.to_string()))?;
211        }
212        super::fault!(pipe, AfterAuthenticate, &op, a);
213        let p = pipe.shards.coordinator(&op.repo.namespace);
214        let ahead = pipe.read_ahead(&op, &p, &[], a.business_skew_ms).await?;
215        let mode = upload_mode(&op, ahead.as_ref())?;
216        tracing::debug!(stage = "replay_lookup", ?mode);
217        if mode != UploadMode::Replay {
218            op.authz = pipe.authorize(&op).await?.0;
219            super::fault!(pipe, AfterAuthorize, &op, a);
220        }
221        if mode == UploadMode::Fresh {
222            pipe.check_outbox_backpressure(&p, ahead.as_ref()).await?;
223            let credentials = super::admission::validate_credentials(&a.credential_capture)?;
224            let mut input = AdmissionInput::new(&op);
225            input.credential_headers = &credentials;
226            input.declared_bytes = declared;
227            input.pack_id = Some(key);
228            let allowance = pipe.admit_streaming(input).await?;
229            if let Some(rid) = &allowance.reservation {
230                pipe.abort_unsupported_stream(a, &p, rid).await?;
231                let refusal =
232                    ServerError::failed_precondition("admission reservations require BeginUpload");
233                return Err(if matches!(pipe.cfg.auth, AuthMode::TransportIdentity) {
234                    refusal.with_transport_admission_required()
235                } else {
236                    refusal
237                });
238            }
239            let charges = allowance.charges;
240            if op.auth.is_some() || !charges.is_empty() {
241                let req = WriteRequest {
242                    denial_ids: None,
243                    denial_packs: &[],
244                    authority_store: super::plan::AuthorityStore::from_capabilities(
245                        pipe.meta.capabilities(),
246                    ),
247                    authority_generation: op.authz.authority_generation,
248                    repo: &op.repo.name,
249                    kind: WriteKind::UploadReserve,
250                    refs: &[],
251                    ref_index: None,
252                    replay: replay_guard(&op),
253                    charges: &charges,
254                    namespace_charge: pipe.namespace_charge(
255                        &p,
256                        &charges,
257                        ahead.as_ref().and_then(|s| s.namespace_window),
258                    )?,
259                    grant: op.authz.grant.clone(),
260                    layout_version: pipe.meta.capabilities().implicit_layout_version.is_none(),
261                    mark_repo_known: false,
262                    lease: None,
263                    rejection: None,
264                    publication: None,
265                    pending: None,
266                    begin: None,
267                    advance: None,
268                    implicit: None,
269                };
270                pipe.apply_atomic(&op, a, &p, &req, ahead).await?;
271            }
272        }
273        super::fault!(pipe, AfterReserve, &op, a);
274        let target = if mode == UploadMode::Replay {
275            Target::Verify(Box::default())
276        } else {
277            let sink = pipe.blobs.begin(key.into(), declared).await;
278            Target::Sink(sink.map_err(|e| store_error(StorageOp::BlobPut, e))?)
279        };
280        Ok(Opened {
281            op,
282            p,
283            mode,
284            validator,
285            target,
286            ticket_id: None,
287        })
288    }
289
290    /// Open a replay-exempt ticketed upload with only authoritative fence-mode
291    /// reads, without reserving or writing business metadata.
292    pub(super) async fn begin_ticketed(
293        pipe: &'p Pipeline<B, N, H>,
294        a: &Authenticated,
295        pack_id: Option<&[u8]>,
296        total_bytes: Option<u64>,
297        token: &[u8],
298    ) -> Result<Self, ServerError> {
299        let mut outcome = pipe.outcome(a);
300        let opened = async {
301            let validator = UploadValidator::new(pack_id, total_bytes, pipe.cfg.upload_limits)?;
302            let AuthMode::AuthV2(cfg) = &pipe.cfg.auth else {
303                return Err(ServerError::new(
304                    Code::Unimplemented,
305                    "ticketed UploadPack requires auth v2",
306                ));
307            };
308            let keys = pipe.cfg.ticket_keys.as_ref().ok_or_else(|| {
309                ServerError::new(Code::Unimplemented, "upload tickets are not configured")
310            })?;
311            let key = validator.key();
312            let declared = validator.declared();
313            let mut op = pipe.identify(
314                a,
315                OpKind::UploadPack {
316                    key,
317                    declared_len: declared,
318                },
319            )?;
320            let auth = op
321                .auth
322                .as_ref()
323                .ok_or_else(|| internal("ticketed upload lacks auth"))?;
324            check_pack_commitment(auth, &key.0, declared)
325                .map_err(|e| ServerError::unauthenticated(e.to_string()))?;
326            let now_ms = ms(pipe.clock.now_ms().saturating_add(a.business_skew_ms));
327            let claims = verify_ticket(
328                keys,
329                token,
330                now_ms,
331                cfg.audience(),
332                &a.repo().identity,
333                &auth.signer,
334            )?;
335            pipe.check_ticket_generation(&op.repo.namespace, claims.authority_generation)
336                .await?;
337            op.authz.authority_generation = claims.authority_generation;
338            if claims.pack_id != key.0 || claims.bytes != declared {
339                return Err(ServerError::new(
340                    Code::PermissionDenied,
341                    "upload ticket binding mismatch",
342                ));
343            }
344            super::fault!(pipe, AfterAuthenticate, &op, a);
345            let sink = pipe
346                .blobs
347                .begin(key.into(), declared)
348                .await
349                .map_err(|e| store_error(StorageOp::BlobPut, e))?;
350            if let Err(err) = pipe
351                .check_ticket_generation(&op.repo.namespace, claims.authority_generation)
352                .await
353            {
354                sink.abort().await;
355                return Err(err);
356            }
357            Ok(Opened {
358                op,
359                p: pipe.shards.coordinator(&a.repo().repo.namespace),
360                mode: UploadMode::Ticketed,
361                validator,
362                target: Target::Sink(sink),
363                ticket_id: Some(claims.ticket_id),
364            })
365        }
366        .instrument(outcome.span.clone())
367        .await;
368        match opened {
369            Ok(o) => Ok(Self {
370                pipe,
371                a: a.clone(),
372                op: o.op,
373                p: o.p,
374                mode: o.mode,
375                validator: o.validator,
376                target: Some(o.target),
377                failed: None,
378                outcome,
379                ticket_id: o.ticket_id,
380                staging: super::staging::StagingBuffer::default(),
381            }),
382            Err(err) => {
383                outcome.record(Err(&err));
384                Err(err)
385            }
386        }
387    }
388
389    /// How this upload relates to its replay record.
390    #[must_use]
391    pub fn mode(&self) -> UploadMode {
392        self.mode
393    }
394
395    /// Validate one chunk's framing and write it to the blob sink (or, on a
396    /// replay, only hash it); returns whether the `last` chunk has arrived.
397    /// After an error the session is dead: every later call returns that
398    /// error.
399    ///
400    /// # Errors
401    /// The chunk's [`UploadError`]; `internal` for a failed blob write.
402    pub async fn push(
403        &mut self,
404        chunk_pack_id: Option<&[u8]>,
405        offset: Option<u64>,
406        mut data: Bytes,
407        last: bool,
408    ) -> Result<bool, ServerError> {
409        if let Some(err) = &self.failed {
410            return Err(err.clone());
411        }
412        let span = self.outcome.span.clone();
413        let result = async {
414            let progress = self
415                .validator
416                .push(chunk_pack_id, offset, data.len(), last)?;
417            if self.ticket_id.is_some() {
418                while !data.is_empty() {
419                    if let Some(bytes) = self.staging.push(&mut data) {
420                        self.stage(bytes).await?;
421                    }
422                }
423                if progress.complete
424                    && let Some(bytes) = self.staging.finish()
425                {
426                    self.stage(bytes).await?;
427                }
428            } else {
429                match self.target.as_mut() {
430                    Some(Target::Sink(sink)) => {
431                        let written = sink.write(data).await;
432                        written.map_err(|e| store_error(StorageOp::BlobPut, e))?;
433                    }
434                    Some(Target::Verify(hasher)) => {
435                        hasher.update(&data);
436                    }
437                    None => return Err(internal("upload sink gone")),
438                }
439            }
440            Ok(progress.complete)
441        }
442        .instrument(span)
443        .await;
444        if let Err(err) = &result {
445            self.fail(err);
446        }
447        result
448    }
449
450    /// End of stream: commit the blob (BLAKE3-verified by the sink), run
451    /// `pre_receive`, then commit the replay record. A replayed upload
452    /// stops once the stream verified.
453    ///
454    /// # Errors
455    /// The stream's [`UploadError`] (`DigestMismatch` when the bytes do not
456    /// hash to the pack id); a hook's error; the commit batch's error.
457    pub async fn finish(mut self) -> Result<(), ServerError> {
458        if let Some(err) = self.failed.take() {
459            return Err(err);
460        }
461        let span = self.outcome.span.clone();
462        let result = Box::pin(self.complete()).instrument(span).await;
463        self.outcome.record(result.as_ref().copied());
464        result
465    }
466
467    async fn stage(&mut self, bytes: Bytes) -> Result<(), ServerError> {
468        self.pipe
469            .check_ticket_generation(&self.op.repo.namespace, self.op.authz.authority_generation)
470            .await?;
471        let Some(Target::Sink(sink)) = self.target.as_mut() else {
472            return Err(internal("ticketed upload sink gone"));
473        };
474        sink.write(bytes)
475            .await
476            .map_err(|e| store_error(StorageOp::BlobPut, e))?;
477        self.pipe
478            .check_ticket_generation(&self.op.repo.namespace, self.op.authz.authority_generation)
479            .await
480    }
481
482    #[allow(clippy::too_many_lines)] // Upload commit and durable ticket completion share one lifecycle.
483    async fn complete(&mut self) -> Result<(), ServerError> {
484        let pipe = self.pipe;
485        let done = self.validator.clone().finish()?;
486        if self.ticket_id.is_some() {
487            pipe.check_ticket_generation(
488                &self.op.repo.namespace,
489                self.op.authz.authority_generation,
490            )
491            .await?;
492        }
493        match self.target.take() {
494            Some(Target::Sink(sink)) => match sink.commit().await {
495                Ok(_) => {}
496                Err(StoreError::Invalid(detail)) => {
497                    tracing::info!(%detail, "pack sink rejected the upload");
498                    return Err(UploadError::DigestMismatch.into());
499                }
500                Err(e) => return Err(store_error(StorageOp::BlobPut, e)),
501            },
502            Some(Target::Verify(hasher)) if hasher.finalize() == done.key.0 => {}
503            Some(Target::Verify(_)) => return Err(UploadError::DigestMismatch.into()),
504            None => return Err(internal("upload sink gone")),
505        }
506        super::fault!(pipe, AfterBlobCommit, &self.op, &self.a);
507        if let Some(ticket_id) = self.ticket_id {
508            pipe.check_ticket_generation(
509                &self.op.repo.namespace,
510                self.op.authz.authority_generation,
511            )
512            .await?;
513            write_upload_marker(&pipe.blobs, &ticket_id, &done.key.0)
514                .await
515                .map_err(|e| store_error(StorageOp::BlobPut, e))?;
516            pipe.check_ticket_generation(
517                &self.op.repo.namespace,
518                self.op.authz.authority_generation,
519            )
520            .await?;
521            pipe.metrics.incr(METRIC_UPLOAD_BYTES, &[], done.total);
522            return Ok(());
523        }
524        if self.mode == UploadMode::Replay {
525            return Ok(());
526        }
527        pipe.metrics.incr(METRIC_UPLOAD_BYTES, &[], done.total);
528        tracing::debug!(stage = "pre_receive");
529        let checked = pipe
530            .hooks
531            .pre_receive()
532            .check(&self.op, Some(&done.key.into()))
533            .await
534            .map_err(ServerError::strip_admission_shape);
535        let Some(replay) = replay_guard(&self.op) else {
536            return checked;
537        };
538        let rejection = match &checked {
539            Ok(()) => None,
540            Err(err) => match storable(err) {
541                Some(rejection) => Some(rejection),
542                None => return checked,
543            },
544        };
545        if self.now() > replay.expires_at_ms.saturating_add(MAX_CLOCK_LEAD_MS) {
546            return Self::lapsed(checked);
547        }
548        let req = WriteRequest {
549            denial_ids: None,
550            denial_packs: &[],
551            authority_store: super::plan::AuthorityStore::from_capabilities(
552                pipe.meta.capabilities(),
553            ),
554            authority_generation: self.op.authz.authority_generation,
555            repo: &self.op.repo.name,
556            kind: WriteKind::UploadCommit,
557            refs: &[],
558            ref_index: None,
559            replay: Some(replay),
560            charges: &[],
561            namespace_charge: None,
562            grant: self.op.authz.grant.clone(),
563            layout_version: pipe.meta.capabilities().implicit_layout_version.is_none(),
564            mark_repo_known: false,
565            lease: None,
566            rejection: rejection.as_ref(),
567            publication: None,
568            pending: None,
569            begin: None,
570            advance: None,
571            implicit: None,
572        };
573        match pipe
574            .apply_atomic(&self.op, &self.a, &self.p, &req, None)
575            .await
576        {
577            Ok(StoredResult::UploadPack) => Ok(()),
578            Ok(other) => Err(super::stored_mismatch(&other)),
579            // A commit that missed its deadline once the envelope expired:
580            // no retry of this nonce can authenticate either.
581            Err(e) if e.code() == Code::Unavailable && self.now() > replay.expires_at_ms => {
582                Self::lapsed(checked)
583            }
584            Err(e) => Err(e),
585        }
586    }
587
588    /// The real clock: the envelope cap bounds deadlines, which never use
589    /// the test clock skew.
590    fn now(&self) -> i64 {
591        self.pipe.clock.now_ms()
592    }
593
594    /// The upload outlived its envelope: answer without the record commit
595    /// (see the module docs).
596    fn lapsed(checked: Result<(), ServerError>) -> Result<(), ServerError> {
597        tracing::info!(
598            "envelope lapsed during the upload; the in-flight record is left to the pruner"
599        );
600        checked
601    }
602
603    /// Discard the upload: nothing becomes visible. An in-flight record
604    /// stays resumable. Recorded as `canceled`.
605    pub async fn abort(self) {
606        self.abort_with(&ServerError::new(Code::Canceled, "upload aborted"))
607            .await;
608    }
609
610    /// Discard the upload because of `err`, e.g. a malformed message the
611    /// binding decoded or a broken request stream: like [`Self::abort`],
612    /// but the request is recorded with `err`'s code, the one the client
613    /// receives. A session that already failed keeps its first error.
614    pub async fn abort_with(mut self, err: &ServerError) {
615        if let Some(Target::Sink(sink)) = self.target.take() {
616            sink.abort().await;
617        }
618        self.fail(err);
619    }
620
621    /// Record the session's one failure.
622    fn fail(&mut self, err: &ServerError) {
623        self.outcome.record(Err(err));
624        if self.failed.is_none() {
625            self.failed = Some(err.clone());
626        }
627    }
628}