1use 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
48pub enum UploadMode {
49 Fresh,
52 Resume,
55 Replay,
58 Ticketed,
61}
62
63enum Target<S> {
66 Sink(S),
67 Verify(Box<Hasher>),
68}
69
70pub 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
98struct 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
108pub(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
117fn 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
139fn 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 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)] 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 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 #[must_use]
391 pub fn mode(&self) -> UploadMode {
392 self.mode
393 }
394
395 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 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)] 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 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 fn now(&self) -> i64 {
591 self.pipe.clock.now_ms()
592 }
593
594 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 pub async fn abort(self) {
606 self.abort_with(&ServerError::new(Code::Canceled, "upload aborted"))
607 .await;
608 }
609
610 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 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}