Skip to main content

mkit_server/indexed/
verify.rs

1//! Inline verification before a ticketed advance commits.
2use 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/// What a verified advance staged: the history edges of its commits, remixes
69/// and tags (a tag has none), and the decoded bytes it holds, which the
70/// fast-forward walk charges against the decode budget (WP-4.17).
71#[derive(Debug, Default, PartialEq, Eq)]
72pub struct StagedCommits {
73    /// Object id to its `parents`; only remix sources are never followed.
74    pub parents: BTreeMap<Hash, Vec<Hash>>,
75    /// Every staged object, of any type.
76    pub objects: usize,
77    /// Total decoded bytes of the staged objects.
78    pub bytes: u64,
79    /// Every external source pack used by any consumed entry, including surplus objects.
80    pub external_bases: BTreeSet<Hash>,
81    /// Current decoded IDs and all reused chain sources for fresh denial at apply.
82    pub denial_ids: BTreeSet<Hash>,
83    /// Immutable verified inventories, including canonical manifest pages.
84    pub denial_packs: Vec<Hash>,
85    /// Complete added-pack inspection metadata, only for configured inspection.
86    pub inspection: Option<super::inspection::InspectionSet>,
87}
88
89/// The `parents` of a history object; `None` for a blob, tree or manifest.
90pub(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/// Reconstruct a located member head and require a commit, remix or tag,
100/// which `verify_push` does not check for a known frontier.
101#[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/// A ticketless non-delete head must be a commit, remix or tag member of
138/// this repository (SPEC-SERVER §9.7). The miss answer is repository-scoped
139/// and `created` opens the §9.4 lag window.
140#[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
289/// The verification leases this call holds, renewed by the extractor.
290struct 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/// Verify every staged object, then extract the large ones into the object
360/// store before each pack's `Verified` state is written (WP-4.10), so
361/// `Verified` implies extracted, held and holder recorded. `ticket_ids` are
362/// the consuming tickets' ids, parallel to `tickets`. No ticket is consumed
363/// and no advance row is written.
364#[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/// Verify ticketed packs while gathering their bounded inspection metadata.
385///
386/// # Errors
387/// The same verification errors, plus the whole-advance inspection size refusal.
388#[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    // Check before inspection preflight or verification-state handling, so a
436    // cached Verified pack returns the same cap error as a fresh pack.
437    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                // Locate all syntactic candidates in bounded batches. A
656                // candidate's answer is used only if the decoder actually
657                // asks for it; in-pack links never fetch member bytes or
658                // emit capped-lookup metrics.
659                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                // The resumable core decoder identifies actual external
668                // bases in one pass. An earlier in-pack frame wins over any
669                // matching member, so no such member frame is fetched.
670                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    // `verify_push` skips known frontiers, including a known root's type.
999    // Reconstruct a member head once so a blob/tree cannot become a tip.
1000    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    // One extractor per advance: its resolution counter is shared by every
1073    // manifest and chunk (R-163).
1074    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        // Extract this pack's objects (those it introduced) after its index
1131        // rows and before `Verified`. A retry re-takes each hold and skips
1132        // what is stored.
1133        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                // Boxed: the extraction future is large and rarely awaited.
1155                match Box::pin(extractor.extract(*id, *kind, ticket_id, &mut lease)).await {
1156                    Ok(()) => {}
1157                    Err(extract::ExtractError::Server(error)) => return Err(error),
1158                    // A manifest this push carries does not match its chunks:
1159                    // content-intrinsic, so the verdict is persisted (§9.8).
1160                    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}