1use 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 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 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 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 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}