1use crate::store::keys;
3
4use super::{
5 IndexedConfig,
6 classify::{self, UploadType},
7 entries::{FrameMeta, index_entries},
8 extract::{self, Extractor, Renew},
9 resolve,
10 state::{self, VerificationV1},
11};
12use crate::pipeline::ShardMap;
13use crate::repo::RepoId;
14use crate::store::{
15 codec::TicketV1,
16 index::{self, IndexEntry},
17 read,
18};
19use crate::telemetry::Metrics;
20use crate::{
21 Batch, BatchOutcome, BlobBody, BlobKey, BlobStore, BoxFuture, Clock, MultipartBlobStore,
22 NamespaceStore, Partition, Precondition, ServerError,
23};
24use futures::StreamExt as _;
25use mkit_core::hash::{Hash, hash};
26use mkit_core::object::Object;
27use mkit_core::ops::graph::{ClosureMode, children};
28use mkit_core::pack::{
29 DecodedEntry, DeltaBaseSource, PackDecodeCursor, PackError, decode_entries_with,
30 delta_base_hashes,
31};
32use mkit_core::sign::verify_object_signature;
33use mkit_core::transfer::decode_packlist;
34use mkit_core::verify::{ObjectSource, VerifyError, verify_push};
35use std::borrow::Cow;
36use std::collections::{BTreeMap, BTreeSet};
37use std::sync::Arc;
38
39fn storage_failed() -> ServerError {
40 ServerError::unavailable("object storage request failed")
41}
42fn bad_object() -> ServerError {
43 ServerError::invalid_argument("object hash mismatch")
44}
45fn decode_failure(error: ServerError, already_verified: bool, pack: &Hash) -> ServerError {
46 if already_verified {
47 tracing::error!(pack = %mkit_core::hash::to_hex(pack), "verified pack failed decode recheck");
48 ServerError::unavailable("verified pack content inconsistency")
49 } else {
50 error
51 }
52}
53pub(super) fn closure_error(now: u64, created: u64, bound: u64) -> ServerError {
54 if resolve::lagged(now, created, bound) {
55 ServerError::unavailable("repository membership not yet visible")
56 } else {
57 ServerError::invalid_argument("open closure")
58 }
59}
60pub(super) fn packlist_error(now: u64, created: u64, bound: u64) -> ServerError {
61 if resolve::lagged(now, created, bound) {
62 ServerError::unavailable("repository membership not yet visible")
63 } else {
64 ServerError::invalid_argument("packlist lists a pack that is not in this repository")
65 }
66}
67
68#[derive(Debug, Default, PartialEq, Eq)]
72pub struct StagedCommits {
73 pub parents: BTreeMap<Hash, Vec<Hash>>,
75 pub objects: usize,
77 pub bytes: u64,
79 pub external_bases: BTreeSet<Hash>,
81 pub denial_ids: BTreeSet<Hash>,
83 pub denial_packs: Vec<Hash>,
85 pub inspection: Option<super::inspection::InspectionSet>,
87}
88
89pub(crate) fn history_parents(object: &Object) -> Option<Vec<Hash>> {
91 match object {
92 Object::Commit(commit) => Some(commit.parents.clone()),
93 Object::Remix(remix) => Some(remix.parents.clone()),
94 Object::Tag(_) => Some(Vec::new()),
95 _ => None,
96 }
97}
98
99#[allow(clippy::too_many_arguments)]
102async fn check_member_head_type<B: BlobStore, S: NamespaceStore>(
103 blobs: &B,
104 store: &S,
105 shards: &dyn ShardMap,
106 repo: &RepoId,
107 head: Hash,
108 located: index::LocatedObject,
109 created: u64,
110 budget: u64,
111 (cfg, clock, metrics): (IndexedConfig, &dyn Clock, &dyn Metrics),
112) -> Result<(), ServerError> {
113 let mut memo = resolve::MemberCache::default();
114 let mut visiting = BTreeSet::new();
115 let (bytes, _) = resolve::member_object(
116 blobs,
117 store,
118 shards,
119 repo,
120 head,
121 located,
122 cfg.max_delta_chain_depth,
123 budget,
124 &mut memo,
125 &mut visiting,
126 metrics,
127 )
128 .await
129 .map_err(|failure| failure.public_error(now_ms(clock), created, cfg.relay_lag_bound_ms))?;
130 let object = mkit_core::serialize::deserialize(&bytes).map_err(|_| bad_object())?;
131 if history_parents(&object).is_none() {
132 return Err(ServerError::invalid_argument("open closure"));
133 }
134 Ok(())
135}
136
137#[allow(clippy::too_many_arguments)]
141pub async fn verify_member_head<B: BlobStore, S: NamespaceStore>(
142 blobs: &B,
143 store: &S,
144 shards: &dyn ShardMap,
145 repo: &RepoId,
146 head: Hash,
147 created: u64,
148 (cfg, clock, metrics): (IndexedConfig, &dyn Clock, &dyn Metrics),
149) -> Result<(), ServerError> {
150 let found = resolve::locate_split(store, shards, repo, &[head], metrics).await?;
151 match found.get(&head) {
152 Some(Ok(Some(located))) => {
153 check_member_head_type(
154 blobs,
155 store,
156 shards,
157 repo,
158 head,
159 *located,
160 created,
161 cfg.decode_budget,
162 (cfg, clock, metrics),
163 )
164 .await
165 }
166 Some(Err(_)) => Err(ServerError::invalid_argument("object index limit exceeded")),
167 _ => Err(closure_error(
168 now_ms(clock),
169 created,
170 cfg.relay_lag_bound_ms,
171 )),
172 }
173}
174
175async fn pack_bytes<B: BlobStore>(
176 blobs: &B,
177 ticket: &TicketV1,
178 cap: u64,
179) -> Result<Vec<u8>, ServerError> {
180 super::check_pack_cap(ticket.bytes, cap)?;
181 let body = blobs
182 .get(&BlobKey::pack(ticket.pack_id), None)
183 .await
184 .map_err(|_| storage_failed())?
185 .ok_or_else(storage_failed)?;
186 let capacity = usize::try_from(ticket.bytes).map_err(|_| bad_object())?;
187 let mut bytes = Vec::new();
188 match body {
189 BlobBody::Bytes(chunk) => {
190 if chunk.len() != capacity {
191 return Err(storage_failed());
192 }
193 bytes.extend_from_slice(&chunk);
194 }
195 BlobBody::Stream { len, mut stream } => {
196 if len != ticket.bytes {
197 return Err(storage_failed());
198 }
199 while let Some(chunk) = stream.next().await {
200 let chunk = chunk.map_err(|_| storage_failed())?;
201 if bytes.len().saturating_add(chunk.len()) > capacity {
202 return Err(storage_failed());
203 }
204 bytes.extend_from_slice(&chunk);
205 }
206 }
207 }
208 if bytes.len() != capacity || hash(&bytes) != ticket.pack_id {
209 return Err(bad_object());
210 }
211 Ok(bytes)
212}
213
214struct Bases(BTreeMap<Hash, Arc<[u8]>>);
215impl DeltaBaseSource for Bases {
216 const VERIFIED: bool = false;
217 fn base(&mut self, id: &Hash) -> Result<Option<Vec<u8>>, PackError> {
218 Ok(self.0.get(id).map(|bytes| bytes.to_vec()))
219 }
220}
221
222struct StagedSource<'a>(&'a BTreeMap<Hash, (Vec<u8>, Object, u64)>);
223impl ObjectSource for StagedSource<'_> {
224 fn fetch(&mut self, id: &Hash) -> Result<Option<Cow<'_, [u8]>>, VerifyError> {
225 Ok(self
226 .0
227 .get(id)
228 .map(|(bytes, _, _)| Cow::Borrowed(bytes.as_slice())))
229 }
230}
231
232struct PackWork {
233 ticket: TicketV1,
234 entries: Vec<IndexEntry>,
235 needs_index: bool,
236}
237
238struct HeldLease {
239 raw: crate::Value,
240 until_ms: u64,
241}
242
243fn now_ms(clock: &dyn Clock) -> u64 {
244 u64::try_from(clock.now_ms()).unwrap_or(0)
245}
246
247fn deadline(clock: &dyn Clock) -> u64 {
248 now_ms(clock).saturating_add(10_000)
249}
250
251async fn renew_pending<S: NamespaceStore>(
252 store: &S,
253 source: &Partition,
254 repo: &RepoId,
255 pack: &Hash,
256 acquired: &mut BTreeMap<Hash, HeldLease>,
257 clock: &dyn Clock,
258) -> Result<(), ServerError> {
259 let Some(held) = acquired.get_mut(pack) else {
260 return Ok(());
261 };
262 let now = now_ms(clock);
263 if held.until_ms.saturating_sub(now) >= state::VERIFICATION_LEASE_MS / 2 {
264 return Ok(());
265 }
266 let until_ms = now.saturating_add(state::VERIFICATION_LEASE_MS);
267 let next = VerificationV1::Pending {
268 lease_until_ms: until_ms,
269 };
270 if !state::write(
271 store,
272 source,
273 &repo.name,
274 pack,
275 Some(&held.raw),
276 &next,
277 deadline(clock),
278 )
279 .await
280 .map_err(|_| super::pending(1_000))?
281 {
282 return Err(super::pending(1_000));
283 }
284 held.raw = state::encode(&next);
285 held.until_ms = until_ms;
286 Ok(())
287}
288
289struct Lease<'a, S> {
291 store: &'a S,
292 source: &'a Partition,
293 repo: &'a RepoId,
294 acquired: &'a mut BTreeMap<Hash, HeldLease>,
295 clock: &'a dyn Clock,
296}
297
298impl<S: NamespaceStore> Renew for Lease<'_, S> {
299 fn renew(&mut self) -> BoxFuture<'_, Result<(), ServerError>> {
300 Box::pin(renew_all_pending(
301 self.store,
302 self.source,
303 self.repo,
304 &mut *self.acquired,
305 self.clock,
306 ))
307 }
308}
309
310async fn renew_all_pending<S: NamespaceStore>(
311 store: &S,
312 source: &Partition,
313 repo: &RepoId,
314 acquired: &mut BTreeMap<Hash, HeldLease>,
315 clock: &dyn Clock,
316) -> Result<(), ServerError> {
317 for pack in acquired.keys().copied().collect::<Vec<_>>() {
318 renew_pending(store, source, repo, &pack, acquired, clock).await?;
319 }
320 Ok(())
321}
322
323async fn reject_content<S: NamespaceStore>(
324 store: &S,
325 source: &Partition,
326 repo: &RepoId,
327 pack: &Hash,
328 acquired: &BTreeMap<Hash, HeldLease>,
329 message: &str,
330 clock: &dyn Clock,
331 metrics: &dyn Metrics,
332) -> ServerError {
333 let Some(held) = acquired.get(pack) else {
334 tracing::error!(pack = %mkit_core::hash::to_hex(pack), "verified pack failed content recheck");
335 return ServerError::unavailable("verified pack content inconsistency");
336 };
337 if !matches!(
338 state::write(
339 store,
340 source,
341 &repo.name,
342 pack,
343 Some(&held.raw),
344 &VerificationV1::Rejected {
345 code: "invalid_argument".into(),
346 message: message.into(),
347 },
348 deadline(clock),
349 )
350 .await,
351 Ok(true)
352 ) {
353 tracing::error!(pack = %mkit_core::hash::to_hex(pack), code = "invalid_argument", "failed to persist rejected verification state");
354 metrics.incr(crate::telemetry::METRIC_INDEX_REJECTED_WRITE_FAILED, &[], 1);
355 }
356 ServerError::invalid_argument(message.to_owned())
357}
358
359#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
365pub async fn verify_ticketed<B: MultipartBlobStore, S: NamespaceStore>(
366 blobs: &B,
367 store: &S,
368 shards: &dyn ShardMap,
369 repo: &RepoId,
370 source: &Partition,
371 tickets: &[TicketV1],
372 ticket_ids: &[Hash],
373 head: Hash,
374 cfg: IndexedConfig,
375 clock: &dyn Clock,
376 metrics: &dyn Metrics,
377) -> Result<StagedCommits, ServerError> {
378 verify_ticketed_optional(
379 blobs, store, shards, repo, source, tickets, ticket_ids, head, cfg, clock, metrics, None,
380 )
381 .await
382}
383
384#[allow(clippy::too_many_arguments)]
389pub async fn verify_ticketed_inspected<B: MultipartBlobStore, S: NamespaceStore>(
390 blobs: &B,
391 store: &S,
392 shards: &dyn ShardMap,
393 repo: &RepoId,
394 source: &Partition,
395 tickets: &[TicketV1],
396 ticket_ids: &[Hash],
397 head: Hash,
398 cfg: IndexedConfig,
399 clock: &dyn Clock,
400 metrics: &dyn Metrics,
401 inspection_limit: usize,
402) -> Result<StagedCommits, ServerError> {
403 verify_ticketed_optional(
404 blobs,
405 store,
406 shards,
407 repo,
408 source,
409 tickets,
410 ticket_ids,
411 head,
412 cfg,
413 clock,
414 metrics,
415 Some(inspection_limit),
416 )
417 .await
418}
419
420#[allow(clippy::too_many_arguments)]
421async fn verify_ticketed_optional<B: MultipartBlobStore, S: NamespaceStore>(
422 blobs: &B,
423 store: &S,
424 shards: &dyn ShardMap,
425 repo: &RepoId,
426 source: &Partition,
427 tickets: &[TicketV1],
428 ticket_ids: &[Hash],
429 head: Hash,
430 cfg: IndexedConfig,
431 clock: &dyn Clock,
432 metrics: &dyn Metrics,
433 inspection_limit: Option<usize>,
434) -> Result<StagedCommits, ServerError> {
435 for ticket in tickets {
438 super::check_pack_cap(ticket.bytes, cfg.max_pack_bytes)?;
439 }
440 let inspection_count = if let Some(limit) = inspection_limit {
441 Some(super::inspection::preflight_native(blobs, tickets, limit).await?)
442 } else {
443 None
444 };
445 let mut acquired = BTreeMap::new();
446 let result = verify_ticketed_inner(
447 blobs,
448 store,
449 shards,
450 repo,
451 source,
452 tickets,
453 ticket_ids,
454 head,
455 cfg,
456 clock,
457 metrics,
458 &mut acquired,
459 inspection_limit.zip(inspection_count),
460 )
461 .await;
462 if result.is_err() {
463 for (pack, held) in acquired {
464 if let Err(error) =
465 state::clear_pending(store, source, &repo.name, &pack, &held.raw, deadline(clock))
466 .await
467 {
468 tracing::error!(%error, "failed to release pending verification lease");
469 }
470 }
471 }
472 result
473}
474
475#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
476async fn verify_ticketed_inner<B: MultipartBlobStore, S: NamespaceStore>(
477 blobs: &B,
478 store: &S,
479 shards: &dyn ShardMap,
480 repo: &RepoId,
481 source: &Partition,
482 tickets: &[TicketV1],
483 ticket_ids: &[Hash],
484 head: Hash,
485 cfg: IndexedConfig,
486 clock: &dyn Clock,
487 metrics: &dyn Metrics,
488 acquired: &mut BTreeMap<Hash, HeldLease>,
489 inspection_limit: Option<(usize, u64)>,
490) -> Result<StagedCommits, ServerError> {
491 let now = now_ms(clock);
492 let consumed: BTreeSet<_> = tickets.iter().map(|ticket| ticket.pack_id).collect();
493 let mut staged: BTreeMap<Hash, (Vec<u8>, Object, u64)> = BTreeMap::new();
494 let mut staged_bytes = 0u64;
495 let mut external_bases = BTreeSet::new();
496 let mut staged_owner = BTreeMap::new();
497 let mut denial_ids = BTreeSet::new();
498 let mut work = Vec::with_capacity(tickets.len());
499 let mut raw_packs = BTreeSet::new();
500 let mut packlists = Vec::new();
501 for ticket in tickets {
502 renew_all_pending(store, source, repo, acquired, clock).await?;
503 let observed = state::read(store, source, &repo.name, &ticket.pack_id)
504 .await
505 .map_err(|_| storage_failed())?;
506 let already_verified = match &observed {
507 Some((VerificationV1::Verified { pack_len, .. }, _)) if *pack_len == ticket.bytes => {
508 true
509 }
510 Some((VerificationV1::Verified { .. }, _)) => {
511 tracing::error!(pack = %mkit_core::hash::to_hex(&ticket.pack_id), "verified pack length changed");
512 return Err(ServerError::unavailable(
513 "verified pack content inconsistency",
514 ));
515 }
516 Some((VerificationV1::Rejected { code, message }, _)) => {
517 return Err(if code == "invalid_argument" {
518 ServerError::invalid_argument(message.clone())
519 } else {
520 ServerError::failed_precondition(message.clone())
521 });
522 }
523 Some((state, _)) if state::concurrent_pending(state, now_ms(clock)).is_some() => {
524 return Err(super::pending(1_000));
525 }
526 _ => false,
527 };
528 let pending_until_ms = now_ms(clock).saturating_add(state::VERIFICATION_LEASE_MS);
529 let pending_raw = if already_verified {
530 None
531 } else {
532 let pending = VerificationV1::Pending {
533 lease_until_ms: pending_until_ms,
534 };
535 let prior = observed.as_ref().map(|(_, raw)| raw);
536 if !state::write(
537 store,
538 source,
539 &repo.name,
540 &ticket.pack_id,
541 prior,
542 &pending,
543 deadline(clock),
544 )
545 .await
546 .map_err(|_| storage_failed())?
547 {
548 return Err(super::pending(1_000));
549 }
550 Some(state::encode(&pending))
551 };
552 if let Some(raw) = &pending_raw {
553 acquired.insert(
554 ticket.pack_id,
555 HeldLease {
556 raw: raw.clone(),
557 until_ms: pending_until_ms,
558 },
559 );
560 }
561 let bytes = match pack_bytes(blobs, ticket, cfg.max_pack_bytes).await {
562 Ok(bytes) => bytes,
563 Err(_) if already_verified => {
564 tracing::error!(pack = %mkit_core::hash::to_hex(&ticket.pack_id), "verified pack bytes changed");
565 return Err(ServerError::unavailable(
566 "verified pack content inconsistency",
567 ));
568 }
569 Err(error) if error.public_message() == "object hash mismatch" => {
570 return Err(reject_content(
571 store,
572 source,
573 repo,
574 &ticket.pack_id,
575 acquired,
576 "object hash mismatch",
577 clock,
578 metrics,
579 )
580 .await);
581 }
582 Err(error) => return Err(error),
583 };
584 let Ok(kind) = classify::classify(&bytes) else {
585 return Err(reject_content(
586 store,
587 source,
588 repo,
589 &ticket.pack_id,
590 acquired,
591 "unknown upload type",
592 clock,
593 metrics,
594 )
595 .await);
596 };
597 match kind {
598 UploadType::Packlist => {
599 let Ok(list) = decode_packlist(&bytes) else {
600 return Err(reject_content(
601 store,
602 source,
603 repo,
604 &ticket.pack_id,
605 acquired,
606 "object hash mismatch",
607 clock,
608 metrics,
609 )
610 .await);
611 };
612 crate::takedown::inventory::stage_packlist(
613 store,
614 &ticket.pack_id,
615 ticket.bytes,
616 list.prev,
617 &list.packs,
618 now_ms(clock),
619 )
620 .await
621 .map_err(|_| storage_failed())?;
622 for child in &list.packs {
623 crate::takedown::inventory::dependency(
624 store,
625 &ticket.pack_id,
626 ticket.bytes,
627 child,
628 now_ms(clock),
629 )
630 .await
631 .map_err(|_| storage_failed())?;
632 }
633 packlists.push((ticket.created_at_ms, list.packs));
634 work.push(PackWork {
635 ticket: ticket.clone(),
636 entries: Vec::new(),
637 needs_index: pending_raw.is_some(),
638 });
639 }
640 UploadType::Pack => {
641 raw_packs.insert(ticket.pack_id);
642 let Ok(base_ids) = delta_base_hashes(&bytes) else {
643 return Err(reject_content(
644 store,
645 source,
646 repo,
647 &ticket.pack_id,
648 acquired,
649 "object hash mismatch",
650 clock,
651 metrics,
652 )
653 .await);
654 };
655 let prelocated = resolve::locate_split_quiet(store, shards, repo, &base_ids)
660 .await
661 .ok();
662 renew_all_pending(store, source, repo, acquired, clock).await?;
663 let mut bases = Bases(BTreeMap::new());
664 let mut depths = BTreeMap::new();
665 let mut memo = resolve::MemberCache::default();
666 let mut visiting = BTreeSet::new();
667 if let Ok(mut probe) = PackDecodeCursor::new(
671 &bytes,
672 super::geometry::entry_limits(cfg.decode_budget.saturating_sub(staged_bytes)),
673 ) {
674 let mut probe_depths = BTreeMap::new();
675 loop {
676 let result = probe.resume(&mut bases, |entry| {
677 let hops = entry.delta_base.map_or(0, |base| {
678 probe_depths
679 .get(&base)
680 .copied()
681 .unwrap_or(0_u32)
682 .saturating_add(1)
683 });
684 if hops > cfg.max_delta_chain_depth {
685 return Err(PackError::PackfileTooLarge);
686 }
687 probe_depths.entry(entry.id).or_insert(hops);
688 Ok(())
689 });
690 renew_all_pending(store, source, repo, acquired, clock).await?;
691 let Err(PackError::DeltaBaseMissing(hex)) = result else {
692 break;
693 };
694 let base = mkit_core::hash::from_hex(&hex).map_err(|_| {
695 decode_failure(bad_object(), already_verified, &ticket.pack_id)
696 })?;
697 if !base_ids.contains(&base) || bases.0.contains_key(&base) {
698 return Err(decode_failure(
699 bad_object(),
700 already_verified,
701 &ticket.pack_id,
702 ));
703 }
704 let cached = prelocated
705 .as_ref()
706 .and_then(|answers| answers.get(&base))
707 .copied();
708 let answer = if matches!(cached, Some(Ok(Some(_)))) {
709 cached
710 } else {
711 let found =
712 resolve::locate_split(store, shards, repo, &[base], metrics)
713 .await
714 .map_err(|error| {
715 decode_failure(error, already_verified, &ticket.pack_id)
716 })?;
717 found.get(&base).copied()
718 };
719 let located = match answer {
720 Some(Ok(Some(located))) => located,
721 Some(Err(_)) => {
722 return Err(decode_failure(
723 resolve::ResolveFailure::Capped.public_error(
724 now_ms(clock),
725 ticket.created_at_ms,
726 cfg.relay_lag_bound_ms,
727 ),
728 already_verified,
729 &ticket.pack_id,
730 ));
731 }
732 _ => {
733 return Err(decode_failure(
734 resolve::missing_base(
735 now_ms(clock),
736 ticket.created_at_ms,
737 cfg.relay_lag_bound_ms,
738 ),
739 already_verified,
740 &ticket.pack_id,
741 ));
742 }
743 };
744 super::geometry::check_entry(
745 located.value.decoded_size,
746 located.value.frame_length,
747 )
748 .map_err(|_| {
749 ServerError::invalid_argument(resolve::DECODE_BUDGET_MESSAGE)
750 })?;
751 let budget = cfg.decode_budget.saturating_sub(staged_bytes);
752 let (canonical, depth) = resolve::member_object(
753 blobs,
754 store,
755 shards,
756 repo,
757 base,
758 located,
759 cfg.max_delta_chain_depth,
760 budget,
761 &mut memo,
762 &mut visiting,
763 metrics,
764 )
765 .await
766 .map_err(|failure| {
767 decode_failure(
768 failure.public_error(
769 now_ms(clock),
770 ticket.created_at_ms,
771 cfg.relay_lag_bound_ms,
772 ),
773 already_verified,
774 &ticket.pack_id,
775 )
776 })?;
777 renew_all_pending(store, source, repo, acquired, clock).await?;
778 for ((id, _, _), _) in memo.rows() {
779 crate::takedown::inventory::dependency(
780 store,
781 &ticket.pack_id,
782 ticket.bytes,
783 id,
784 now_ms(clock),
785 )
786 .await
787 .map_err(|_| storage_failed())?;
788 }
789 bases.0.insert(base, canonical);
790 depths.insert(base, depth);
791 probe
792 .set_max_decoded_bytes(
793 cfg.decode_budget.saturating_sub(
794 staged_bytes.saturating_add(memo.retained_bytes()),
795 ),
796 )
797 .map_err(|_| {
798 decode_failure(
799 ServerError::invalid_argument(
800 "pack exceeds indexed decode budget",
801 ),
802 already_verified,
803 &ticket.pack_id,
804 )
805 })?;
806 }
807 }
808 let mut frames = Vec::new();
809 let mut local_depths: BTreeMap<Hash, (u32, Option<Hash>)> = BTreeMap::new();
810 let mut in_pack_depth_exceeded = false;
811 let mut external_depth_exceeded = false;
812 let decoded = decode_entries_with(
813 &bytes,
814 &mut bases,
815 super::geometry::entry_limits(
816 cfg.decode_budget
817 .saturating_sub(staged_bytes.saturating_add(memo.retained_bytes())),
818 ),
819 |entry: DecodedEntry<'_>| {
820 super::geometry::check_entry(entry.bytes.len() as u64, entry.frame_length)?;
821 let (hops, external) = match entry.delta_base {
822 None => (0, None),
823 Some(base) => match local_depths.get(&base) {
824 Some((depth, external)) => (depth.saturating_add(1), *external),
825 None => (1, Some(base)),
826 },
827 };
828 let total = hops.saturating_add(
829 external
830 .and_then(|base| depths.get(&base).copied())
831 .unwrap_or(0),
832 );
833 if hops > cfg.max_delta_chain_depth {
834 in_pack_depth_exceeded = true;
835 return Err(PackError::PackfileTooLarge);
836 }
837 if total > cfg.max_delta_chain_depth {
838 external_depth_exceeded = true;
839 return Err(PackError::PackfileTooLarge);
840 }
841 local_depths.entry(entry.id).or_insert((hops, external));
842 frames.push(FrameMeta {
843 id: entry.id,
844 frame_offset: entry.frame_offset,
845 frame_length: entry.frame_length,
846 wire_type: entry.wire_type,
847 delta_base: entry.delta_base,
848 decoded_size: entry.bytes.len() as u64,
849 });
850 if let std::collections::btree_map::Entry::Vacant(slot) =
851 staged.entry(entry.id)
852 {
853 staged_bytes = staged_bytes
854 .checked_add(entry.bytes.len() as u64)
855 .ok_or(PackError::PackfileTooLarge)?;
856 if staged_bytes.saturating_add(memo.retained_bytes())
857 > cfg.decode_budget
858 {
859 return Err(PackError::PackfileTooLarge);
860 }
861 slot.insert((entry.bytes.to_vec(), entry.object, ticket.created_at_ms));
862 }
863 staged_owner.entry(entry.id).or_insert(ticket.pack_id);
864 Ok(())
865 },
866 );
867 renew_all_pending(store, source, repo, acquired, clock).await?;
868 if let Err(error) = decoded {
869 if in_pack_depth_exceeded {
870 return Err(reject_content(
871 store,
872 source,
873 repo,
874 &ticket.pack_id,
875 acquired,
876 "delta chain too deep",
877 clock,
878 metrics,
879 )
880 .await);
881 }
882 if external_depth_exceeded {
883 return Err(decode_failure(
884 ServerError::invalid_argument("delta chain too deep"),
885 already_verified,
886 &ticket.pack_id,
887 ));
888 }
889 let mapped = if let PackError::DeltaBaseMissing(hex) = &error {
890 let _ = hex;
891 resolve::missing_base(
892 now_ms(clock),
893 ticket.created_at_ms,
894 cfg.relay_lag_bound_ms,
895 )
896 } else if matches!(error, PackError::PackfileTooLarge) {
897 ServerError::invalid_argument("pack exceeds indexed decode budget")
898 } else {
899 bad_object()
900 };
901 if mapped.public_message() == "object hash mismatch" {
902 return Err(reject_content(
903 store,
904 source,
905 repo,
906 &ticket.pack_id,
907 acquired,
908 "object hash mismatch",
909 clock,
910 metrics,
911 )
912 .await);
913 }
914 return Err(decode_failure(mapped, already_verified, &ticket.pack_id));
915 }
916 external_bases.extend(memo.rows().map(|((_, pack, _), _)| *pack));
917 let Ok(entries) = index_entries(&frames, cfg.max_delta_chain_depth) else {
918 return Err(reject_content(
919 store,
920 source,
921 repo,
922 &ticket.pack_id,
923 acquired,
924 "delta chain too deep",
925 clock,
926 metrics,
927 )
928 .await);
929 };
930 work.push(PackWork {
931 ticket: ticket.clone(),
932 entries,
933 needs_index: pending_raw.is_some(),
934 });
935 }
936 }
937 }
938 for id in staged.keys() {
939 crate::takedown::denial::require_clear(store, id).await?;
940 }
941 for (id, (_, object, _)) in &staged {
942 if verify_object_signature(object).is_err() {
943 let owner = staged_owner[id];
944 return Err(reject_content(
945 store,
946 source,
947 repo,
948 &owner,
949 acquired,
950 "bad signature",
951 clock,
952 metrics,
953 )
954 .await);
955 }
956 }
957 let mut needed: BTreeMap<Hash, u64> = BTreeMap::new();
958 for (_, object, created) in staged.values() {
959 for child in children(object, ClosureMode::History) {
960 if !staged.contains_key(&child) {
961 needed.entry(child).or_insert(*created);
962 }
963 }
964 }
965 if !staged.contains_key(&head) {
966 needed.entry(head).or_insert_with(|| {
967 tickets
968 .iter()
969 .map(|ticket| ticket.created_at_ms)
970 .min()
971 .unwrap_or(now)
972 });
973 }
974 let mut member_head = None;
975 if !needed.is_empty() {
976 let ids: Vec<_> = needed.keys().copied().collect();
977 let found = resolve::locate_split(store, shards, repo, &ids, metrics).await?;
978 for (id, created) in needed {
979 match found.get(&id) {
980 Some(Ok(Some(located))) => {
981 if id == head {
982 member_head = Some((*located, created));
983 }
984 }
985 Some(Err(_)) => {
986 return Err(ServerError::invalid_argument("object index limit exceeded"));
987 }
988 _ => {
989 return Err(closure_error(
990 u64::try_from(clock.now_ms()).unwrap_or(0),
991 created,
992 cfg.relay_lag_bound_ms,
993 ));
994 }
995 }
996 }
997 }
998 if let Some((located, created)) = member_head {
1001 check_member_head_type(
1002 blobs,
1003 store,
1004 shards,
1005 repo,
1006 head,
1007 located,
1008 created,
1009 cfg.decode_budget.saturating_sub(staged_bytes),
1010 (cfg, clock, metrics),
1011 )
1012 .await?;
1013 }
1014 let staged_ids: BTreeSet<_> = staged.keys().copied().collect();
1015 let report = verify_push(
1016 &[head],
1017 ClosureMode::History,
1018 &mut StagedSource(&staged),
1019 |id| !staged_ids.contains(id),
1020 )
1021 .map_err(|error| match error {
1022 VerifyError::TooManyClosureObjects => {
1023 ServerError::invalid_argument("object index limit exceeded")
1024 }
1025 VerifyError::ClosureRootWrongType(_) => ServerError::invalid_argument("open closure"),
1026 VerifyError::Store(_) => storage_failed(),
1027 _ => bad_object(),
1028 })?;
1029 if !report.bad_signatures.is_empty() {
1030 return Err(ServerError::invalid_argument("bad signature"));
1031 }
1032 if !report.corrupt.is_empty() {
1033 return Err(bad_object());
1034 }
1035 if !report.missing.is_empty() || !report.bad_tips.is_empty() {
1036 return Err(ServerError::invalid_argument("open closure"));
1037 }
1038 for (created, packs) in packlists {
1039 let missing: Vec<_> = packs
1040 .into_iter()
1041 .filter(|pack| !consumed.contains(pack))
1042 .collect();
1043 if missing.is_empty() {
1044 continue;
1045 }
1046 if missing.len() > index::MAX_LOOKUP_IDS {
1047 tracing::error!(reason = "ids", "packlist membership lookup capped");
1048 metrics.incr(
1049 crate::telemetry::METRIC_INDEX_LOOKUP_CAPPED,
1050 &[("reason", "ids")],
1051 1,
1052 );
1053 return Err(ServerError::invalid_argument("object index limit exceeded"));
1054 }
1055 let found = read::members_many(store, shards, repo, source, &missing)
1056 .await
1057 .map_err(|_| storage_failed())?;
1058 if found.iter().any(|member| !member) {
1059 return Err(packlist_error(
1060 u64::try_from(clock.now_ms()).unwrap_or(0),
1061 created,
1062 cfg.relay_lag_bound_ms,
1063 ));
1064 }
1065 }
1066 let selected = extract::select(&staged, cfg.extract_min_bytes);
1067 if extract::selected_bytes(&staged, &selected) > cfg.effective_max_extract_bytes() {
1068 return Err(ServerError::invalid_argument(
1069 "pack exceeds indexed decode budget",
1070 ));
1071 }
1072 let extractor = Extractor {
1075 blobs,
1076 store,
1077 shards,
1078 repo,
1079 cfg,
1080 clock,
1081 metrics,
1082 staged: &staged,
1083 staged_bytes,
1084 resolved: std::sync::atomic::AtomicU64::new(0),
1085 denial_pack: std::sync::Mutex::new(None),
1086 };
1087 for pack in &work {
1088 if !pack.needs_index {
1089 continue;
1090 }
1091 for entry in &pack.entries {
1092 let (_, object, _) = staged.get(&entry.object).ok_or_else(storage_failed)?;
1093 crate::takedown::inventory::stage(
1094 store,
1095 &pack.ticket.pack_id,
1096 pack.ticket.bytes,
1097 &entry.object,
1098 object,
1099 entry.value.delta_base,
1100 now_ms(clock),
1101 )
1102 .await
1103 .map_err(|_| storage_failed())?;
1104 }
1105 let plan = index::plan_index_rows_direct(
1106 shards,
1107 repo,
1108 source,
1109 &pack.ticket.pack_id,
1110 &pack.entries,
1111 now,
1112 )
1113 .map_err(|_| storage_failed())?;
1114 for direct in plan.direct {
1115 renew_all_pending(store, source, repo, acquired, clock).await?;
1116 let mut batch = Batch::new().require(Precondition::NotAfter(deadline(clock)));
1117 for (key, value) in direct.puts {
1118 if let Some(keys::ParsedKey::ObjectIndex { object, .. }) = keys::parse(&key) {
1119 crate::takedown::denial::require_clear(store, &object).await?;
1120 }
1121 batch = batch.put(key, value);
1122 }
1123 if !matches!(
1124 store.apply(&direct.target, batch).await,
1125 Ok(BatchOutcome::Committed)
1126 ) {
1127 return Err(super::pending(1_000));
1128 }
1129 }
1130 let ticket_id = tickets
1134 .iter()
1135 .zip(ticket_ids)
1136 .find_map(|(t, id)| (t.pack_id == pack.ticket.pack_id).then_some(id))
1137 .ok_or_else(|| {
1138 ServerError::internal(
1139 "object storage request failed",
1140 "missing consumed ticket id",
1141 )
1142 })?;
1143 let mut lease = Lease {
1144 store,
1145 source,
1146 repo,
1147 acquired,
1148 clock,
1149 };
1150 *extractor.denial_pack.lock().map_err(|_| storage_failed())? =
1151 Some((pack.ticket.pack_id, pack.ticket.bytes));
1152 for (id, kind) in &selected {
1153 if staged_owner.get(id) == Some(&pack.ticket.pack_id) {
1154 match Box::pin(extractor.extract(*id, *kind, ticket_id, &mut lease)).await {
1156 Ok(()) => {}
1157 Err(extract::ExtractError::Server(error)) => return Err(error),
1158 Err(extract::ExtractError::Content) => {
1161 return Err(reject_content(
1162 store,
1163 source,
1164 repo,
1165 &pack.ticket.pack_id,
1166 lease.acquired,
1167 extract::MALFORMED_MESSAGE,
1168 clock,
1169 metrics,
1170 )
1171 .await);
1172 }
1173 }
1174 }
1175 }
1176 crate::takedown::inventory::complete(
1177 store,
1178 &pack.ticket.pack_id,
1179 pack.ticket.bytes,
1180 now_ms(clock),
1181 )
1182 .await
1183 .map_err(|_| storage_failed())?;
1184 renew_all_pending(store, source, repo, acquired, clock).await?;
1185 let raw = &acquired[&pack.ticket.pack_id].raw;
1186 if !state::write(
1187 store,
1188 source,
1189 &repo.name,
1190 &pack.ticket.pack_id,
1191 Some(raw),
1192 &VerificationV1::Verified {
1193 pack_len: pack.ticket.bytes,
1194 verified_at_ms: now_ms(clock),
1195 publication: None,
1196 },
1197 deadline(clock),
1198 )
1199 .await
1200 .map_err(|_| super::pending(1_000))?
1201 {
1202 return Err(super::pending(1_000));
1203 }
1204 acquired.remove(&pack.ticket.pack_id);
1205 }
1206 let parents = staged
1207 .iter()
1208 .filter_map(|(id, (_, object, _))| Some((*id, history_parents(object)?)))
1209 .collect();
1210 denial_ids.extend(tickets.iter().map(|ticket| ticket.pack_id));
1211 let objects = staged.len();
1212 let mut inspection =
1213 inspection_limit.map(|(limit, _)| super::inspection::InspectionSet::new(limit));
1214 if let Some(set) = &mut inspection {
1215 if let Some((_, count)) = inspection_limit {
1216 set.reserve_added_count(count)?;
1217 }
1218 let entries = staged
1219 .into_iter()
1220 .map(|(id, (bytes, object, _))| {
1221 let object_type = object.object_type() as u8;
1222 super::inspection::NativeEntry {
1223 id,
1224 size: bytes.len() as u64,
1225 object_type,
1226 }
1227 })
1228 .collect();
1229 for pack in raw_packs {
1230 set.add_raw_pack(pack);
1231 }
1232 set.defer_native(entries);
1233 }
1234 Ok(StagedCommits {
1235 denial_packs: tickets.iter().map(|ticket| ticket.pack_id).collect(),
1236 denial_ids,
1237 parents,
1238 objects,
1239 bytes: staged_bytes,
1240 external_bases,
1241 inspection,
1242 })
1243}