Skip to main content

mkit_server/pipeline/
begin.rs

1//! `BeginUpload`'s pre-admission decisions and pure ref-shard fragment.
2use super::{
3    Addressing, AuthMode, Authenticated, BeginUploadResult, HookSet, Key, MultipartBlobStore,
4    NamespaceStore, OpKind, Operation, PackKey, Partition, Pipeline, PlanClock, ServerError,
5    Sharding, Snapshot, StorageOp, StoredResult, TicketCaps, check_ref_name, codec, internal, keys,
6    meta_error, ms, store_error, stored_mismatch,
7};
8use crate::replay::StoredRejection;
9use crate::store::tickets;
10use crate::store::tickets::{TicketPlanError, TicketSpec};
11use crate::upload::token::{TicketClaims, TicketKeys};
12use mkit_core::hash::Hash;
13
14pub(super) const CAP_MESSAGE: &str = "too many open upload tickets";
15
16#[derive(Debug, Clone)]
17pub(super) enum BeginWrite {
18    Return(BeginUploadResult),
19    Open(Box<TicketOpen>),
20}
21
22#[derive(Debug, Clone)]
23pub(super) struct TicketOpen {
24    pub spec: TicketSpec,
25    keys: TicketKeys,
26    caps: TicketCaps,
27    audience: String,
28    repository: String,
29    reserved: bool,
30}
31
32impl TicketOpen {
33    pub(super) fn reserved(&self) -> bool {
34        self.reserved
35    }
36}
37
38pub(super) fn decision_keys(
39    repo: &crate::repo::RepoName,
40    name: &str,
41    pack: &Hash,
42    signer: &Hash,
43) -> Result<Vec<Key>, ServerError> {
44    Ok(vec![
45        keys::ticket_index(repo, name, pack, signer).map_err(meta_error)?,
46        keys::tickets_per_ref(repo, name).map_err(meta_error)?,
47        keys::tickets_per_signer(repo, name, signer).map_err(meta_error)?,
48        keys::membership(repo, pack),
49    ])
50}
51
52pub(super) fn open_keys(spec: &TicketSpec) -> Vec<Key> {
53    let k = tickets::keys(spec);
54    vec![k.ticket, k.index, k.per_ref, k.per_signer, k.reservation]
55}
56
57pub(super) async fn read_indexed<N: NamespaceStore>(
58    meta: &N,
59    p: &Partition,
60    spec: &TicketSpec,
61    snap: &mut Snapshot,
62) -> Result<(), ServerError> {
63    let k = tickets::keys(spec);
64    if let Some(index) = snap.get(&k.index) {
65        let id = codec::decode_ref_id(index).map_err(meta_error)?;
66        let key = keys::ticket(&id);
67        if !snap.contains(&key) {
68            let value = meta.get(p, &key).await.map_err(meta_error)?;
69            snap.insert(key, value);
70        }
71    }
72    Ok(())
73}
74
75fn result(
76    keys: &TicketKeys,
77    audience: &str,
78    repository: &str,
79    ticket: &codec::TicketV1,
80) -> BeginUploadResult {
81    let id = tickets::ticket_id(&ticket.reservation_id);
82    let claims = TicketClaims {
83        authority_generation: ticket.authority_generation,
84        ticket_id: id,
85        audience: audience.to_owned(),
86        repository: repository.to_owned(),
87        signer: ticket.signer,
88        pack_id: ticket.pack_id,
89        bytes: ticket.bytes,
90        part_size: ticket.part_size,
91        expires_at_ms: ticket.expires_at_ms,
92        upload_session: ticket.upload_session.clone().unwrap_or_default(),
93    };
94    BeginUploadResult::Ticket {
95        id,
96        part_size: ticket.part_size,
97        expires_at_ms: ticket.expires_at_ms,
98        token: keys.mint(&claims),
99    }
100}
101
102impl<B: MultipartBlobStore, N: NamespaceStore, H: HookSet> Pipeline<B, N, H> {
103    /// Open a stateless authenticated upload ticket in the target ref shard.
104    ///
105    /// # Errors
106    /// Invalid geometry/ref, unsupported auth or missing keys, policy or cap
107    /// refusal, and the usual replay, quota, lease and storage errors.
108    pub async fn begin_upload(
109        &self,
110        a: &Authenticated,
111        ref_name: &str,
112        pack_id: &[u8],
113        bytes: u64,
114    ) -> Result<BeginUploadResult, ServerError> {
115        self.begin_upload_with_meta(a, ref_name, pack_id, bytes)
116            .await
117            .map(|(result, _)| result)
118    }
119
120    /// Begin a resumable upload and return success-only admission headers.
121    pub async fn begin_upload_with_meta(
122        &self,
123        a: &Authenticated,
124        ref_name: &str,
125        pack_id: &[u8],
126        bytes: u64,
127    ) -> Result<(BeginUploadResult, super::ResponseMeta), ServerError> {
128        self.observe(a, async {
129            check_ref_name(ref_name)?;
130            if ref_name.starts_with(mkit_core::refs::PACKMAP_REF_PREFIX) {
131                return Err(ServerError::invalid_argument(
132                    "BeginUpload names a branch or tag, not its packmap",
133                ));
134            }
135            if self.cfg.sharding == Sharding::D34 && !ref_name.starts_with("refs/heads/") {
136                return Err(ServerError::invalid_argument(
137                    "BeginUpload requires refs/heads/ on this server",
138                ));
139            }
140            let id: Hash = pack_id
141                .try_into()
142                .map_err(|_| ServerError::invalid_argument("pack_id must be 32 bytes"))?;
143            if bytes == 0 || bytes > self.cfg.upload_limits.max_total_bytes {
144                return Err(ServerError::invalid_argument(
145                    "upload bytes outside server limit",
146                ));
147            }
148            if !matches!(self.cfg.auth, AuthMode::AuthV2(_)) {
149                return Err(ServerError::new(
150                    crate::Code::Unimplemented,
151                    "BeginUpload requires auth v2",
152                ));
153            }
154            if self.cfg.ticket_keys.is_none() {
155                return Err(ServerError::new(
156                    crate::Code::Unimplemented,
157                    "upload tickets are not configured",
158                ));
159            }
160            if bytes > self.cfg.part_size && !self.blobs.supports_multipart() {
161                return Err(ServerError::unimplemented(
162                    "multipart uploads are not supported by this storage backend",
163                ));
164            }
165            match self
166                .write(
167                    a,
168                    OpKind::BeginUpload {
169                        ref_name: ref_name.into(),
170                        key: PackKey(id),
171                        bytes,
172                    },
173                )
174                .await?
175            {
176                (StoredResult::BeginUpload(answer), meta) => Ok((answer, meta)),
177                (other, _) => Err(stored_mismatch(&other)),
178            }
179        })
180        .await
181    }
182
183    pub(super) async fn begin_decision(
184        &self,
185        op: &Operation,
186        a: &Authenticated,
187        ahead: Option<&mut Snapshot>,
188    ) -> Result<Option<BeginUploadResult>, ServerError> {
189        let OpKind::BeginUpload { ref_name, key, .. } = &op.kind else {
190            return Ok(None);
191        };
192        let snap = ahead.ok_or_else(|| internal("ticket write lacks snapshot"))?;
193        let auth = op
194            .auth
195            .as_ref()
196            .ok_or_else(|| internal("ticket write lacks signer"))?;
197        let ks = decision_keys(&op.repo.name, ref_name, &key.0, &auth.signer)?;
198        let now = ms(self.clock.now_ms().saturating_add(a.business_skew_ms));
199        if let Some(index) = snap.get(&ks[0]) {
200            let id = codec::decode_ref_id(index).map_err(meta_error)?;
201            let k = keys::ticket(&id);
202            let value = self
203                .meta
204                .get(&self.shards.ref_shard(&op.repo, ref_name), &k)
205                .await
206                .map_err(meta_error)?;
207            snap.insert(k.clone(), value);
208            if let Some(raw) = snap.get(&k) {
209                let ticket = codec::decode_ticket(raw).map_err(meta_error)?;
210                if ticket.repo != op.repo.name
211                    || ticket.ref_name != *ref_name
212                    || ticket.pack_id != key.0
213                    || ticket.signer != auth.signer
214                    || tickets::ticket_id(&ticket.reservation_id) != id
215                {
216                    return Err(internal("ticket index binding mismatch"));
217                }
218                if ticket.expires_at_ms > now {
219                    if self.cfg.authority_fence.is_some()
220                        && ticket.authority_generation != op.authz.authority_generation
221                    {
222                        return Err(crate::authority::moved());
223                    }
224                    let (keys, audience) = self.ticket_config()?;
225                    return Ok(Some(result(keys, audience, &a.repo().identity, &ticket)));
226                }
227            }
228        }
229        let present = if matches!(self.cfg.addressing, Addressing::Multi(_)) {
230            snap.get(&ks[3]).is_some()
231        } else {
232            self.blobs
233                .head(&(*key).into())
234                .await
235                .map_err(|e| store_error(StorageOp::BlobHead, e))?
236                .is_some()
237        };
238        let present = if present && self.cfg.indexed.is_some() {
239            let clear = if self.cfg.takedown_denial {
240                crate::takedown::denial::require_pack_clear(
241                    &self.meta,
242                    self.shards.as_ref(),
243                    &op.repo,
244                    &key.0,
245                )
246                .await
247            } else {
248                crate::takedown::denial::require_clear(&self.meta, &key.0).await
249            };
250            match clear {
251                Ok(()) => true,
252                Err(e) if e.public_message() == "object blocked" => false,
253                Err(e) => return Err(e),
254            }
255        } else {
256            present
257        };
258        if present {
259            return Ok(Some(BeginUploadResult::AlreadyPresent));
260        }
261        for (k, cap) in [
262            (&ks[1], self.cfg.ticket_caps.per_ref),
263            (&ks[2], self.cfg.ticket_caps.per_signer),
264        ] {
265            let count = snap
266                .get(k)
267                .map(codec::decode_u64)
268                .transpose()
269                .map_err(meta_error)?
270                .unwrap_or(0);
271            if count >= cap {
272                return Err(ServerError::failed_precondition(CAP_MESSAGE));
273            }
274        }
275        Ok(None)
276    }
277
278    fn ticket_config(&self) -> Result<(&TicketKeys, &str), ServerError> {
279        let keys = self
280            .cfg
281            .ticket_keys
282            .as_ref()
283            .ok_or_else(|| internal("missing ticket keys"))?;
284        let AuthMode::AuthV2(auth) = &self.cfg.auth else {
285            return Err(internal("missing ticket audience"));
286        };
287        // The token encodes the audience with a u16 length; never let minting panic.
288        if u16::try_from(auth.audience().len()).is_err() {
289            return Err(internal("ticket audience too long"));
290        }
291        Ok((keys, auth.audience()))
292    }
293
294    pub(super) fn begin_write(
295        &self,
296        op: &Operation,
297        a: &Authenticated,
298        existing: Option<BeginUploadResult>,
299        reservation: Option<String>,
300    ) -> Result<Option<BeginWrite>, ServerError> {
301        let OpKind::BeginUpload {
302            ref_name,
303            key,
304            bytes,
305        } = &op.kind
306        else {
307            return Ok(None);
308        };
309        if let Some(answer) = existing {
310            return Ok(Some(BeginWrite::Return(answer)));
311        }
312        let auth = op
313            .auth
314            .as_ref()
315            .ok_or_else(|| internal("missing ticket signer"))?;
316        let (keys, audience) = self.ticket_config()?;
317        let now = ms(self.clock.now_ms().saturating_add(a.business_skew_ms));
318        let expires = now
319            .checked_add(self.cfg.ticket_ttl_ms)
320            .ok_or_else(|| internal("ticket expiry overflow"))?;
321        let reserved = reservation.is_some();
322        let rid = reservation
323            .unwrap_or_else(|| crate::store::outbox::synthetic_reservation_id(&auth.replay_scope));
324        // Validate admission-supplied ids before any infallible key constructor.
325        keys::reservation(&rid).map_err(meta_error)?;
326        Ok(Some(BeginWrite::Open(Box::new(TicketOpen {
327            spec: TicketSpec {
328                authority_generation: op.authz.authority_generation,
329                repo: op.repo.name.clone(),
330                ref_name: ref_name.clone(),
331                signer: auth.signer,
332                pack_id: key.0,
333                bytes: *bytes,
334                part_size: self.cfg.part_size,
335                expires_at_ms: expires,
336                created_at_ms: now,
337                now_ms: now,
338                reservation_id: rid,
339                upload_session: None,
340            },
341            keys: keys.clone(),
342            caps: self.cfg.ticket_caps,
343            audience: audience.into(),
344            repository: a.repo().identity.clone(),
345            reserved,
346        }))))
347    }
348}
349
350pub(super) fn plan(
351    begin: &BeginWrite,
352    snap: &Snapshot,
353    clock: &PlanClock,
354    pre: &mut Vec<crate::store::Precondition>,
355    writes: &mut Vec<crate::store::Write>,
356) -> Result<StoredResult, ServerError> {
357    let BeginWrite::Open(open) = begin else {
358        let BeginWrite::Return(answer) = begin else {
359            unreachable!()
360        };
361        if let BeginUploadResult::Ticket { id, .. } = answer {
362            let key = keys::ticket(id);
363            let raw = snap
364                .get(&key)
365                .ok_or_else(|| ServerError::aborted_retryable("upload ticket race"))?;
366            pre.push(crate::store::Precondition::Equals(key, raw.clone()));
367        }
368        return Ok(StoredResult::BeginUpload(answer.clone()));
369    };
370    let mut spec = open.spec.clone();
371    spec.now_ms = ms(clock.business_now_ms);
372    let k = tickets::keys(&spec);
373    let index = snap.get(&k.index).cloned();
374    let indexed_ticket = index
375        .as_ref()
376        .map(codec::decode_ref_id)
377        .transpose()
378        .map_err(meta_error)?
379        .and_then(|id| snap.get(&keys::ticket(&id)).cloned());
380    let reads = tickets::TicketReads {
381        ticket: snap.get(&k.ticket).cloned(),
382        indexed_ticket,
383        index,
384        per_ref: snap.get(&k.per_ref).cloned(),
385        per_signer: snap.get(&k.per_signer).cloned(),
386        reservation: snap.get(&k.reservation).cloned(),
387    };
388    match tickets::plan_ticket_open(&spec, &reads, open.caps, pre, writes) {
389        Ok(id) => Ok(StoredResult::BeginUpload(BeginUploadResult::Ticket {
390            id,
391            part_size: spec.part_size,
392            expires_at_ms: spec.expires_at_ms,
393            token: open.keys.mint(&TicketClaims {
394                authority_generation: spec.authority_generation,
395                ticket_id: id,
396                audience: open.audience.clone(),
397                repository: open.repository.clone(),
398                signer: spec.signer,
399                pack_id: spec.pack_id,
400                bytes: spec.bytes,
401                part_size: spec.part_size,
402                expires_at_ms: spec.expires_at_ms,
403                upload_session: spec.upload_session.clone().unwrap_or_default(),
404            }),
405        })),
406        Err(TicketPlanError::Existing(ticket)) if !open.reserved => {
407            if spec.authority_generation.is_some()
408                && ticket.authority_generation != spec.authority_generation
409            {
410                return Err(crate::authority::moved());
411            }
412            let key = keys::ticket(&tickets::ticket_id(&ticket.reservation_id));
413            let raw = snap
414                .get(&key)
415                .ok_or_else(|| ServerError::aborted_retryable("upload ticket race"))?;
416            pre.push(crate::store::Precondition::Equals(key, raw.clone()));
417            Ok(StoredResult::BeginUpload(result(
418                &open.keys,
419                &open.audience,
420                &open.repository,
421                &ticket,
422            )))
423        }
424        Err(TicketPlanError::CapExceeded { .. }) if !open.reserved => Ok(StoredResult::Rejected(
425            StoredRejection::new(crate::Code::FailedPrecondition, CAP_MESSAGE)
426                .expect("final cap error"),
427        )),
428        Err(TicketPlanError::Existing(_)) => {
429            Err(ServerError::aborted_retryable("upload ticket race"))
430        }
431        Err(TicketPlanError::CapExceeded { .. }) => {
432            Err(ServerError::failed_precondition(CAP_MESSAGE))
433        }
434        Err(TicketPlanError::Corrupt(err)) => Err(meta_error(err)),
435        Err(TicketPlanError::Invalid(detail)) => Err(internal(detail)),
436    }
437}