1use std::collections::{HashSet, VecDeque};
3#[cfg(feature = "async-source")]
4use std::sync::{
5 Arc,
6 atomic::{AtomicBool, Ordering},
7};
8#[cfg(feature = "async-source")]
9use std::time::Instant;
10
11use serde::{Deserialize, Serialize};
12
13#[cfg(feature = "async-source")]
14use crate::store::AsyncObjectSource;
15use crate::{
16 error::{HeddleError, Result},
17 object::{
18 AnnotatedTag, BindingDelta, ContentHash, RedactionsBlob, ReverseDependencyIndex,
19 SemanticEntryKind, SemanticIndexRoot, SemanticTreeNode, State, StateAttachment,
20 StateAttachmentBody, StateAttachmentId, StateAttachmentKind, StateId, TreeEntryTarget,
21 decode_tree_delta_header, is_delta_tree,
22 },
23 store::{ObjectSource, ObjectStore, pack::ObjectType as PackObjectType},
24};
25
26#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
27pub enum ObjectId {
28 Hash(ContentHash),
29 StateId(StateId),
30 StateAttachment {
31 state: StateId,
32 id: StateAttachmentId,
33 kind: StateAttachmentKind,
39 },
40}
41
42#[derive(Debug, Clone, Serialize, Deserialize)]
43pub struct ObjectInfo {
44 pub id: ObjectId,
45 pub obj_type: ObjectType,
46 pub size: u64,
47 pub delta_base: Option<ContentHash>,
48}
49
50#[derive(Debug, Clone, PartialEq, Eq, Hash, Serialize, Deserialize)]
51pub struct PlannedObject {
52 pub id: ObjectId,
53 pub obj_type: ObjectType,
54}
55
56#[derive(Debug, Clone)]
57pub struct StateClosureTransferObjects {
58 pub planned_objects: Vec<PlannedObject>,
59 pub full_objects: Option<Vec<ObjectInfo>>,
60}
61
62#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Serialize, Deserialize)]
63pub enum ObjectType {
64 Blob,
65 Tree,
66 State,
67 Action,
68 AnnotatedTag,
69 Redaction,
75 Purge,
79 StateVisibility,
86 StateAttachment,
87 KeyBinding,
92}
93
94#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
95pub enum ObjectTypeBucket {
96 Blob,
97 Tree,
98 State,
99 Action,
100 AnnotatedTag,
101 Redaction,
102 Purge,
103 StateVisibility,
104 StateAttachment,
105 KeyBinding,
106}
107
108impl ObjectType {
109 pub fn wire_name(self) -> &'static str {
110 match self {
111 ObjectType::Blob => "blob",
112 ObjectType::Tree => "tree",
113 ObjectType::State => "state",
114 ObjectType::Action => "action",
115 ObjectType::AnnotatedTag => "annotated_tag",
116 ObjectType::Redaction => "redaction",
117 ObjectType::Purge => "purge",
118 ObjectType::StateVisibility => "state_visibility",
119 ObjectType::StateAttachment => "state_attachment",
120 ObjectType::KeyBinding => "key_binding",
121 }
122 }
123
124 pub fn from_wire(value: &str) -> Result<Self> {
125 match value {
126 "blob" => Ok(ObjectType::Blob),
127 "tree" => Ok(ObjectType::Tree),
128 "state" => Ok(ObjectType::State),
129 "action" => Ok(ObjectType::Action),
130 "annotated_tag" => Ok(ObjectType::AnnotatedTag),
131 "redaction" => Ok(ObjectType::Redaction),
132 "purge" => Ok(ObjectType::Purge),
133 "state_visibility" => Ok(ObjectType::StateVisibility),
134 "state_attachment" => Ok(ObjectType::StateAttachment),
135 "key_binding" => Ok(ObjectType::KeyBinding),
136 _ => Err(HeddleError::InvalidObject(format!(
137 "unknown object type: {value}"
138 ))),
139 }
140 }
141
142 pub fn packable(self) -> bool {
153 !matches!(
154 self,
155 ObjectType::Redaction
156 | ObjectType::Purge
157 | ObjectType::StateVisibility
158 | ObjectType::KeyBinding
159 )
160 }
161
162 pub fn packable_for_push(self) -> bool {
172 self.packable() && !matches!(self, ObjectType::StateAttachment)
173 }
174
175 pub fn packable_for_pull(self) -> bool {
183 self.packable()
184 }
185
186 pub fn pack_object_type(self) -> Result<PackObjectType> {
187 match self {
188 ObjectType::Blob => Ok(PackObjectType::Blob),
189 ObjectType::Tree => Ok(PackObjectType::Tree),
190 ObjectType::State => Ok(PackObjectType::State),
191 ObjectType::Action => Ok(PackObjectType::Action),
192 ObjectType::AnnotatedTag => Ok(PackObjectType::AnnotatedTag),
193 ObjectType::StateAttachment => Ok(PackObjectType::StateAttachment),
194 ObjectType::Redaction => Err(HeddleError::InvalidObject(
195 "Redaction sidecar records cannot be packed into the content-addressed object pack"
196 .to_string(),
197 )),
198 ObjectType::Purge => Err(HeddleError::InvalidObject(
199 "Purge sidecar records cannot be packed into the content-addressed object pack"
200 .to_string(),
201 )),
202 ObjectType::StateVisibility => Err(HeddleError::InvalidObject(
203 "StateVisibility sidecar records cannot be packed into the content-addressed object pack"
204 .to_string(),
205 )),
206 ObjectType::KeyBinding => Err(HeddleError::InvalidObject(
207 "KeyBinding registry objects cannot be packed into the content-addressed object pack"
208 .to_string(),
209 )),
210 }
211 }
212
213 pub fn bucket(self) -> ObjectTypeBucket {
214 match self {
215 ObjectType::Blob => ObjectTypeBucket::Blob,
216 ObjectType::Tree => ObjectTypeBucket::Tree,
217 ObjectType::State => ObjectTypeBucket::State,
218 ObjectType::Action => ObjectTypeBucket::Action,
219 ObjectType::AnnotatedTag => ObjectTypeBucket::AnnotatedTag,
220 ObjectType::Redaction => ObjectTypeBucket::Redaction,
221 ObjectType::Purge => ObjectTypeBucket::Purge,
222 ObjectType::StateVisibility => ObjectTypeBucket::StateVisibility,
223 ObjectType::StateAttachment => ObjectTypeBucket::StateAttachment,
224 ObjectType::KeyBinding => ObjectTypeBucket::KeyBinding,
225 }
226 }
227}
228
229#[derive(Debug, Clone, Default)]
230pub struct StateClosureOptions {
231 pub depth: Option<u32>,
232 pub exclude_states: Vec<StateId>,
233}
234
235pub fn enumerate_state_closure(
236 store: &impl ObjectStore,
237 state_id: StateId,
238) -> Result<Vec<ObjectInfo>> {
239 enumerate_state_closure_with_options(store, state_id, StateClosureOptions::default())
240}
241
242pub fn enumerate_state_closure_with_options(
243 store: &impl ObjectStore,
244 state_id: StateId,
245 options: StateClosureOptions,
246) -> Result<Vec<ObjectInfo>> {
247 let mut out = Vec::new();
248 walk_state_closure(store, state_id, options, |event| {
249 if let Some(info) = object_info_from_event(store, event)? {
250 out.push(info);
251 }
252 Ok(())
253 })?;
254 for (hash, tag) in annotated_tags_for_state(store, state_id)? {
255 out.push(annotated_tag_info(hash, &tag));
256 }
257
258 Ok(out)
259}
260
261pub fn enumerate_state_closure_plan(
262 store: &impl ObjectStore,
263 state_id: StateId,
264) -> Result<Vec<PlannedObject>> {
265 enumerate_state_closure_plan_with_options(store, state_id, StateClosureOptions::default())
266}
267
268pub fn enumerate_state_closure_plan_with_options(
269 store: &impl ObjectStore,
270 state_id: StateId,
271 options: StateClosureOptions,
272) -> Result<Vec<PlannedObject>> {
273 let mut out = Vec::new();
274 walk_state_closure(store, state_id, options, |event| {
275 if let Some(object) = planned_object_from_event(store, event)? {
276 out.push(object);
277 }
278 Ok(())
279 })?;
280 out.extend(
281 annotated_tags_for_state(store, state_id)?
282 .into_iter()
283 .map(|(hash, _)| PlannedObject {
284 id: ObjectId::Hash(hash),
285 obj_type: ObjectType::AnnotatedTag,
286 }),
287 );
288
289 Ok(out)
290}
291
292pub fn enumerate_state_closure_transfer_with_options(
293 store: &impl ObjectStore,
294 state_id: StateId,
295 options: StateClosureOptions,
296 full_descriptor_object_threshold: usize,
297) -> Result<StateClosureTransferObjects> {
298 let mut planned_objects = Vec::new();
299 let mut full_objects = Some(Vec::new());
300
301 walk_state_closure(store, state_id, options, |event| {
302 if let Some(object) = planned_object_from_event(store, event)? {
303 planned_objects.push(object);
304 }
305
306 if full_objects.is_some() && planned_objects.len() > full_descriptor_object_threshold {
307 full_objects = None;
308 }
309 if let Some(objects) = full_objects.as_mut()
310 && let Some(info) = object_info_from_event(store, event)?
311 {
312 objects.push(info);
313 }
314
315 Ok(())
316 })?;
317 let tags = annotated_tags_for_state(store, state_id)?;
318 planned_objects.extend(tags.iter().map(|(hash, _)| PlannedObject {
319 id: ObjectId::Hash(*hash),
320 obj_type: ObjectType::AnnotatedTag,
321 }));
322 if full_objects.is_some() && planned_objects.len() > full_descriptor_object_threshold {
323 full_objects = None;
324 }
325 if let Some(objects) = full_objects.as_mut() {
326 objects.extend(
327 tags.iter()
328 .map(|(hash, tag)| annotated_tag_info(*hash, tag)),
329 );
330 }
331
332 Ok(StateClosureTransferObjects {
333 planned_objects,
334 full_objects,
335 })
336}
337
338fn annotated_tags_for_state(
339 store: &impl ObjectStore,
340 state_id: StateId,
341) -> Result<Vec<(ContentHash, AnnotatedTag)>> {
342 let mut roots = Vec::new();
343 for hash in store.list_annotated_tags()? {
344 let Some(tag) = store.get_annotated_tag(&hash)? else {
345 continue;
346 };
347 if tag
348 .marker()
349 .is_some_and(|marker| marker.peeled_state == state_id)
350 {
351 roots.push((hash, tag));
352 }
353 }
354
355 let mut tags = Vec::new();
356 let mut seen = HashSet::new();
357 let mut stack = roots;
358 while let Some((hash, tag)) = stack.pop() {
359 if !seen.insert(hash) {
360 continue;
361 }
362 if let Some(inner_hash) = tag.target_tag() {
363 let inner = store.get_annotated_tag(&inner_hash)?.ok_or_else(|| {
364 HeddleError::NotFound(format!(
365 "annotated tag {hash} references missing inner tag {inner_hash}"
366 ))
367 })?;
368 stack.push((inner_hash, inner));
369 }
370 tags.push((hash, tag));
371 }
372 Ok(tags)
373}
374
375fn annotated_tag_info(hash: ContentHash, tag: &AnnotatedTag) -> ObjectInfo {
376 ObjectInfo {
377 id: ObjectId::Hash(hash),
378 obj_type: ObjectType::AnnotatedTag,
379 size: tag.encode_current_msgpack().len() as u64,
380 delta_base: None,
381 }
382}
383
384pub fn enumerate_state_closure_transfer_from_boundaries(
391 store: &impl ObjectStore,
392 state_id: StateId,
393 boundary_states: &[StateId],
394 full_descriptor_object_threshold: usize,
395) -> Result<StateClosureTransferObjects> {
396 let mut planned_objects = Vec::new();
397 let mut full_objects = Some(Vec::new());
398 let excluded_states = boundary_states.iter().copied().collect();
399
400 walk_state_closure_with_exclusions(
401 store,
402 state_id,
403 None,
404 excluded_states,
405 HashSet::new(),
406 |event| {
407 if let Some(object) = planned_object_from_event(store, event)? {
408 planned_objects.push(object);
409 }
410
411 if full_objects.is_some() && planned_objects.len() > full_descriptor_object_threshold {
412 full_objects = None;
413 }
414 if let Some(objects) = full_objects.as_mut()
415 && let Some(info) = object_info_from_event(store, event)?
416 {
417 objects.push(info);
418 }
419
420 Ok(())
421 },
422 )?;
423
424 Ok(StateClosureTransferObjects {
425 planned_objects,
426 full_objects,
427 })
428}
429
430#[derive(Debug, Clone, Copy)]
431enum StateClosureEvent<'a> {
432 State {
433 id: StateId,
434 state: &'a State,
435 },
436 Tree {
437 hash: ContentHash,
438 storage_size: u64,
439 delta_base: Option<ContentHash>,
440 },
441 Blob {
442 hash: ContentHash,
443 },
444 Redaction {
445 blob: ContentHash,
446 },
447 StateVisibility {
448 state: StateId,
449 },
450 StateAttachment {
451 state: StateId,
452 attachment: &'a StateAttachment,
453 },
454 ExcludedState {
455 id: StateId,
456 },
457 ExcludedHash {
458 hash: ContentHash,
459 },
460}
461
462fn walk_state_closure(
463 store: &impl ObjectStore,
464 state_id: StateId,
465 options: StateClosureOptions,
466 visit: impl for<'event> FnMut(StateClosureEvent<'event>) -> Result<()>,
467) -> Result<()> {
468 let (excluded_states, excluded_hashes) = collect_excluded(store, &options.exclude_states)?;
469
470 walk_state_closure_with_exclusions(
471 store,
472 state_id,
473 options.depth,
474 excluded_states,
475 excluded_hashes,
476 visit,
477 )
478}
479
480fn walk_state_closure_with_exclusions(
481 store: &impl ObjectStore,
482 state_id: StateId,
483 max_depth: Option<u32>,
484 excluded_states: HashSet<StateId>,
485 excluded_hashes: HashSet<ContentHash>,
486 mut visit: impl for<'event> FnMut(StateClosureEvent<'event>) -> Result<()>,
487) -> Result<()> {
488 let mut seen_states: HashSet<StateId> = HashSet::new();
489 let mut seen_hashes: HashSet<ContentHash> = HashSet::new();
490 let mut queue: VecDeque<(StateId, u32)> = VecDeque::new();
491 queue.push_back((state_id, 0));
492
493 while let Some((id, depth)) = queue.pop_front() {
494 if excluded_states.contains(&id) {
495 visit(StateClosureEvent::ExcludedState { id })?;
496 continue;
497 }
498 if !seen_states.insert(id) {
499 continue;
500 }
501
502 let state = store
503 .get_state(&id)?
504 .ok_or_else(|| HeddleError::MissingObject {
505 object_type: "state".to_string(),
506 id: id.to_string(),
507 })?;
508
509 visit(StateClosureEvent::State { id, state: &state })?;
510 if store.has_state_visibility_for_state(&id)? {
511 visit(StateClosureEvent::StateVisibility { state: id })?;
512 }
513 for attachment in store.list_state_attachments(&id)? {
514 visit(StateClosureEvent::StateAttachment {
515 state: id,
516 attachment: &attachment,
517 })?;
518 match attachment.body {
519 StateAttachmentBody::Context(root) => walk_tree_closure_filtered(
520 store,
521 root,
522 &excluded_hashes,
523 &mut seen_hashes,
524 &mut visit,
525 )?,
526 StateAttachmentBody::RiskSignals(hash)
527 | StateAttachmentBody::ReviewSignatures(hash)
528 | StateAttachmentBody::Discussions(hash)
529 | StateAttachmentBody::StructuredConflicts(hash)
530 | StateAttachmentBody::ContextConsumption(hash) => {
531 walk_blob_filtered(store, hash, &excluded_hashes, &mut seen_hashes, &mut visit)?
532 }
533 StateAttachmentBody::SemanticIndex(root) => walk_semantic_index_closure(
534 store,
535 root,
536 &excluded_hashes,
537 &mut seen_hashes,
538 &mut visit,
539 )?,
540 StateAttachmentBody::Signature(_) => {}
541 }
542 }
543
544 if max_depth.map(|max| depth < max).unwrap_or(true) {
545 for parent in &state.parents {
546 queue.push_back((*parent, depth + 1));
547 }
548 }
549
550 walk_tree_closure_filtered(
551 store,
552 state.tree,
553 &excluded_hashes,
554 &mut seen_hashes,
555 &mut visit,
556 )?;
557 if let Some(provenance_root) = state.provenance {
558 walk_tree_closure_filtered(
559 store,
560 provenance_root,
561 &excluded_hashes,
562 &mut seen_hashes,
563 &mut visit,
564 )?;
565 }
566 }
567
568 Ok(())
569}
570
571fn walk_tree_closure_filtered(
572 store: &impl ObjectStore,
573 tree_hash: ContentHash,
574 excluded: &HashSet<ContentHash>,
575 seen: &mut HashSet<ContentHash>,
576 visit: &mut impl for<'event> FnMut(StateClosureEvent<'event>) -> Result<()>,
577) -> Result<()> {
578 if excluded.contains(&tree_hash) {
579 visit(StateClosureEvent::ExcludedHash { hash: tree_hash })?;
580 return Ok(());
581 }
582 if !seen.insert(tree_hash) {
583 return Ok(());
584 }
585
586 let tree = store
587 .get_tree(&tree_hash)?
588 .ok_or_else(|| HeddleError::MissingObject {
589 object_type: "tree".to_string(),
590 id: tree_hash.to_hex(),
591 })?;
592
593 let (storage_size, delta_base) = tree_storage_metadata(store, tree_hash)?;
594 visit(StateClosureEvent::Tree {
595 hash: tree_hash,
596 storage_size,
597 delta_base,
598 })?;
599
600 if let Some(anchor) = delta_base {
601 walk_tree_storage_anchor(store, anchor, seen, visit)?;
602 }
603
604 for entry in tree.entries() {
605 match entry.target() {
606 TreeEntryTarget::Blob { hash, .. } | TreeEntryTarget::Symlink { hash } => {
607 walk_blob_filtered(store, *hash, excluded, seen, visit)?;
608 }
609 TreeEntryTarget::Tree { hash } => {
610 walk_tree_closure_filtered(store, *hash, excluded, seen, visit)?;
611 }
612 TreeEntryTarget::Gitlink { .. } => {}
613 TreeEntryTarget::Spoollink { .. } => {}
616 }
617 }
618
619 Ok(())
620}
621
622fn tree_storage_metadata(
623 store: &impl ObjectStore,
624 tree_hash: ContentHash,
625) -> Result<(u64, Option<ContentHash>)> {
626 let body =
627 store
628 .get_tree_serialized(&tree_hash)?
629 .ok_or_else(|| HeddleError::MissingObject {
630 object_type: "tree".to_string(),
631 id: tree_hash.to_hex(),
632 })?;
633 let delta_base = if is_delta_tree(&body) {
634 let header = decode_tree_delta_header(&body)?;
635 if header.anchor == tree_hash {
636 return Err(HeddleError::InvalidObject(
637 "HDC1 result id must differ from its anchor id".to_string(),
638 ));
639 }
640 Some(header.anchor)
641 } else {
642 None
643 };
644 Ok((body.len() as u64, delta_base))
645}
646
647fn walk_tree_storage_anchor(
648 store: &impl ObjectStore,
649 anchor_hash: ContentHash,
650 seen: &mut HashSet<ContentHash>,
651 visit: &mut impl for<'event> FnMut(StateClosureEvent<'event>) -> Result<()>,
652) -> Result<()> {
653 if !seen.insert(anchor_hash) {
654 return Ok(());
655 }
656 if store.get_tree(&anchor_hash)?.is_none() {
657 return Err(HeddleError::MissingObject {
658 object_type: "tree delta anchor".to_string(),
659 id: anchor_hash.to_hex(),
660 });
661 }
662 let (storage_size, delta_base) = tree_storage_metadata(store, anchor_hash)?;
663 if delta_base.is_some() {
664 return Err(HeddleError::InvalidObject(
665 "HDC1 anchor must be materialized; delta chains are forbidden".to_string(),
666 ));
667 }
668 visit(StateClosureEvent::Tree {
669 hash: anchor_hash,
670 storage_size,
671 delta_base: None,
672 })
673}
674
675fn walk_blob_filtered(
676 store: &impl ObjectStore,
677 blob_hash: ContentHash,
678 excluded: &HashSet<ContentHash>,
679 seen: &mut HashSet<ContentHash>,
680 visit: &mut impl for<'event> FnMut(StateClosureEvent<'event>) -> Result<()>,
681) -> Result<()> {
682 if excluded.contains(&blob_hash) {
683 visit(StateClosureEvent::ExcludedHash { hash: blob_hash })?;
684 return Ok(());
685 }
686 if !seen.insert(blob_hash) {
687 return Ok(());
688 }
689 visit(StateClosureEvent::Blob { hash: blob_hash })?;
690 if store.has_redactions_for_blob(&blob_hash)? {
691 visit(StateClosureEvent::Redaction { blob: blob_hash })?;
692 }
693 Ok(())
694}
695
696fn walk_semantic_index_closure(
709 store: &impl ObjectStore,
710 root_hash: ContentHash,
711 excluded: &HashSet<ContentHash>,
712 seen: &mut HashSet<ContentHash>,
713 visit: &mut impl for<'event> FnMut(StateClosureEvent<'event>) -> Result<()>,
714) -> Result<()> {
715 let mut stack: Vec<ContentHash> = vec![root_hash];
718 while let Some(node_hash) = stack.pop() {
719 if !emit_semantic_blob(store, node_hash, excluded, seen, visit)? {
720 continue; }
722 let blob = store
723 .get_blob(&node_hash)?
724 .ok_or_else(|| missing_blob(node_hash))?;
725 let node = decode_semantic_container(&blob, node_hash)?;
729 for child in node {
730 match child {
731 SemanticChild::Interior(hash) => stack.push(hash),
732 SemanticChild::Leaf(hash) => {
733 emit_semantic_blob(store, hash, excluded, seen, visit)?;
735 }
736 SemanticChild::BindingDelta(hash) => {
737 emit_binding_delta(store, hash, excluded, seen, visit)?;
743 }
744 SemanticChild::ImporterIndex(hash) => {
745 emit_importer_index(store, hash, excluded, seen, visit)?;
746 }
747 }
748 }
749 }
750 Ok(())
751}
752
753enum SemanticChild {
755 Interior(ContentHash),
756 Leaf(ContentHash),
757 BindingDelta(ContentHash),
758 ImporterIndex(ContentHash),
759}
760
761fn decode_semantic_container(
764 blob: &crate::object::Blob,
765 node_hash: ContentHash,
766) -> Result<Vec<SemanticChild>> {
767 if let Ok(root) = SemanticIndexRoot::decode(blob.content()) {
771 let mut children = vec![SemanticChild::Interior(root.tree)];
772 if let Some(binding_delta) = root.binding_delta {
773 children.push(SemanticChild::BindingDelta(binding_delta));
774 }
775 if let Some(importer_index) = root.importer_index {
776 children.push(SemanticChild::ImporterIndex(importer_index));
777 }
778 return Ok(children);
779 }
780 let node = SemanticTreeNode::decode(blob.content())
781 .map_err(|err| HeddleError::Serialization(format!("semantic node {node_hash}: {err}")))?;
782 Ok(node
783 .entries
784 .iter()
785 .filter_map(|entry| match entry.kind {
786 SemanticEntryKind::Dir => Some(SemanticChild::Interior(entry.node)),
787 SemanticEntryKind::File => Some(SemanticChild::Leaf(entry.node)),
788 SemanticEntryKind::Opaque => None,
790 })
791 .collect())
792}
793
794fn emit_binding_delta(
795 store: &impl ObjectStore,
796 hash: ContentHash,
797 excluded: &HashSet<ContentHash>,
798 seen: &mut HashSet<ContentHash>,
799 visit: &mut impl for<'event> FnMut(StateClosureEvent<'event>) -> Result<()>,
800) -> Result<()> {
801 if !emit_semantic_blob(store, hash, excluded, seen, visit)? {
802 return Ok(());
803 }
804 let blob = store.get_blob(&hash)?.ok_or_else(|| missing_blob(hash))?;
805 BindingDelta::decode(blob.content()).map_err(|err| {
806 HeddleError::Serialization(format!("semantic binding delta {hash}: {err}"))
807 })?;
808 Ok(())
809}
810
811fn emit_importer_index(
812 store: &impl ObjectStore,
813 hash: ContentHash,
814 excluded: &HashSet<ContentHash>,
815 seen: &mut HashSet<ContentHash>,
816 visit: &mut impl for<'event> FnMut(StateClosureEvent<'event>) -> Result<()>,
817) -> Result<()> {
818 if !emit_semantic_blob(store, hash, excluded, seen, visit)? {
819 return Ok(());
820 }
821 let blob = store.get_blob(&hash)?.ok_or_else(|| missing_blob(hash))?;
822 ReverseDependencyIndex::decode(blob.content()).map_err(|err| {
823 HeddleError::Serialization(format!("semantic reverse-dependency index {hash}: {err}"))
824 })?;
825 Ok(())
826}
827
828fn emit_semantic_blob(
833 store: &impl ObjectStore,
834 hash: ContentHash,
835 excluded: &HashSet<ContentHash>,
836 seen: &mut HashSet<ContentHash>,
837 visit: &mut impl for<'event> FnMut(StateClosureEvent<'event>) -> Result<()>,
838) -> Result<bool> {
839 if excluded.contains(&hash) {
840 visit(StateClosureEvent::ExcludedHash { hash })?;
841 return Ok(false);
842 }
843 if !seen.insert(hash) {
844 return Ok(false);
845 }
846 if store.get_blob(&hash)?.is_none() {
847 return Err(missing_blob(hash));
848 }
849 visit(StateClosureEvent::Blob { hash })?;
850 Ok(true)
851}
852
853fn collect_semantic_hashes(
858 store: &impl ObjectStore,
859 root_hash: ContentHash,
860 excluded: &mut HashSet<ContentHash>,
861) -> Result<()> {
862 let mut stack: Vec<ContentHash> = vec![root_hash];
863 while let Some(node_hash) = stack.pop() {
864 if !excluded.insert(node_hash) {
865 continue;
866 }
867 let Some(blob) = store.get_blob(&node_hash)? else {
868 continue;
869 };
870 let children = match decode_semantic_container(&blob, node_hash) {
871 Ok(children) => children,
872 Err(_) => continue,
873 };
874 for child in children {
875 match child {
876 SemanticChild::Interior(hash) => stack.push(hash),
877 SemanticChild::Leaf(hash) => {
878 excluded.insert(hash);
879 }
880 SemanticChild::BindingDelta(hash) | SemanticChild::ImporterIndex(hash) => {
881 excluded.insert(hash);
882 }
883 }
884 }
885 }
886 Ok(())
887}
888
889fn object_info_from_event(
890 store: &impl ObjectStore,
891 event: StateClosureEvent<'_>,
892) -> Result<Option<ObjectInfo>> {
893 match event {
894 StateClosureEvent::State { id, state } => {
895 let state_bytes = state.encode_current_msgpack()?;
896 Ok(Some(ObjectInfo {
897 id: ObjectId::StateId(id),
898 obj_type: ObjectType::State,
899 size: state_bytes.len() as u64,
900 delta_base: None,
901 }))
902 }
903 StateClosureEvent::Tree {
904 hash,
905 storage_size,
906 delta_base,
907 } => Ok(Some(ObjectInfo {
908 id: ObjectId::Hash(hash),
909 obj_type: ObjectType::Tree,
910 size: storage_size,
911 delta_base,
912 })),
913 StateClosureEvent::Blob { hash, .. } => {
914 let Some(blob) = store.get_blob(&hash)? else {
915 if blob_has_purge_evidence(store, &hash)? {
916 return Ok(None);
917 }
918 return Err(missing_blob(hash));
919 };
920 Ok(Some(ObjectInfo {
921 id: ObjectId::Hash(hash),
922 obj_type: ObjectType::Blob,
923 size: blob.size() as u64,
924 delta_base: None,
925 }))
926 }
927 StateClosureEvent::Redaction { blob } => Ok(store
928 .get_redactions_bytes_for_blob(&blob)?
929 .map(|bytes| ObjectInfo {
930 id: ObjectId::Hash(blob),
931 obj_type: ObjectType::Redaction,
932 size: bytes.len() as u64,
933 delta_base: None,
934 })),
935 StateClosureEvent::StateVisibility { state } => Ok(store
936 .get_state_visibility_bytes_for_state(&state)?
937 .map(|bytes| ObjectInfo {
938 id: ObjectId::StateId(state),
939 obj_type: ObjectType::StateVisibility,
940 size: bytes.len() as u64,
941 delta_base: None,
942 })),
943 StateClosureEvent::StateAttachment { state, attachment } => {
944 let bytes = attachment.encode_current_msgpack()?;
945 Ok(Some(ObjectInfo {
946 id: ObjectId::StateAttachment {
947 state,
948 id: attachment.id(),
949 kind: attachment.body.kind(),
950 },
951 obj_type: ObjectType::StateAttachment,
952 size: bytes.len() as u64,
953 delta_base: None,
954 }))
955 }
956 StateClosureEvent::ExcludedState { id } => {
957 let _ = id;
958 Ok(None)
959 }
960 StateClosureEvent::ExcludedHash { hash } => {
961 let _ = hash;
962 Ok(None)
963 }
964 }
965}
966
967fn planned_object_from_event(
968 store: &impl ObjectStore,
969 event: StateClosureEvent<'_>,
970) -> Result<Option<PlannedObject>> {
971 match event {
972 StateClosureEvent::State { id, .. } => Ok(Some(PlannedObject {
973 id: ObjectId::StateId(id),
974 obj_type: ObjectType::State,
975 })),
976 StateClosureEvent::Tree { hash, .. } => Ok(Some(PlannedObject {
977 id: ObjectId::Hash(hash),
978 obj_type: ObjectType::Tree,
979 })),
980 StateClosureEvent::Blob { hash, .. } => {
981 if store.get_blob(&hash)?.is_none() {
982 if blob_has_purge_evidence(store, &hash)? {
983 return Ok(None);
984 }
985 return Err(missing_blob(hash));
986 }
987 Ok(Some(PlannedObject {
988 id: ObjectId::Hash(hash),
989 obj_type: ObjectType::Blob,
990 }))
991 }
992 StateClosureEvent::Redaction { blob } => Ok(Some(PlannedObject {
993 id: ObjectId::Hash(blob),
994 obj_type: ObjectType::Redaction,
995 })),
996 StateClosureEvent::StateVisibility { state } => Ok(Some(PlannedObject {
997 id: ObjectId::StateId(state),
998 obj_type: ObjectType::StateVisibility,
999 })),
1000 StateClosureEvent::StateAttachment { state, attachment } => Ok(Some(PlannedObject {
1001 id: ObjectId::StateAttachment {
1002 state,
1003 id: attachment.id(),
1004 kind: attachment.body.kind(),
1005 },
1006 obj_type: ObjectType::StateAttachment,
1007 })),
1008 StateClosureEvent::ExcludedState { id } => {
1009 let _ = id;
1010 Ok(None)
1011 }
1012 StateClosureEvent::ExcludedHash { hash } => {
1013 let _ = hash;
1014 Ok(None)
1015 }
1016 }
1017}
1018
1019fn missing_blob(hash: ContentHash) -> HeddleError {
1023 HeddleError::MissingObject {
1024 object_type: "blob".to_string(),
1025 id: hash.to_hex(),
1026 }
1027}
1028
1029fn blob_has_purge_evidence(store: &impl ObjectStore, hash: &ContentHash) -> Result<bool> {
1030 let Some(bytes) = store.get_redactions_bytes_for_blob(hash)? else {
1031 return Ok(false);
1032 };
1033 let redactions = RedactionsBlob::decode(&bytes).map_err(|error| {
1034 HeddleError::InvalidObject(format!(
1035 "invalid redaction sidecar for missing blob {}: {error}",
1036 hash.to_hex()
1037 ))
1038 })?;
1039 Ok(redactions
1040 .redactions
1041 .iter()
1042 .any(|redaction| redaction.redacted_blob == *hash && redaction.is_purged()))
1043}
1044
1045pub fn missing_blobs_in_tree(
1046 store: &impl ObjectStore,
1047 tree_hash: ContentHash,
1048) -> Result<Vec<ContentHash>> {
1049 let mut missing = Vec::new();
1050 collect_missing_blobs_recursive(store, &tree_hash, &mut missing)?;
1051 Ok(missing)
1052}
1053
1054fn collect_missing_blobs_recursive(
1055 store: &impl ObjectStore,
1056 tree_hash: &ContentHash,
1057 missing: &mut Vec<ContentHash>,
1058) -> Result<()> {
1059 let Some(tree) = store.get_tree(tree_hash).map_err(|err| {
1060 HeddleError::InvalidObject(format!(
1061 "load tree {} while collecting lazy hydration missing blobs: {err}",
1062 tree_hash.to_hex()
1063 ))
1064 })?
1065 else {
1066 return Ok(());
1067 };
1068
1069 for entry in tree.entries() {
1070 match entry.target() {
1071 TreeEntryTarget::Blob { hash, .. } | TreeEntryTarget::Symlink { hash } => {
1072 if !store.has_blob(hash).map_err(|err| {
1073 HeddleError::InvalidObject(format!(
1074 "check blob {} while collecting lazy hydration missing blobs: {err}",
1075 hash.to_hex()
1076 ))
1077 })? {
1078 missing.push(*hash);
1079 }
1080 }
1081 TreeEntryTarget::Tree { hash } => {
1082 collect_missing_blobs_recursive(store, hash, missing)?;
1083 }
1084 TreeEntryTarget::Gitlink { .. } => {}
1085 TreeEntryTarget::Spoollink { .. } => {}
1088 }
1089 }
1090 Ok(())
1091}
1092
1093fn collect_excluded(
1094 store: &impl ObjectStore,
1095 roots: &[StateId],
1096) -> Result<(HashSet<StateId>, HashSet<ContentHash>)> {
1097 if roots.is_empty() {
1098 return Ok((HashSet::new(), HashSet::new()));
1099 }
1100
1101 let mut excluded_states: HashSet<StateId> = HashSet::new();
1102 let mut excluded_hashes: HashSet<ContentHash> = HashSet::new();
1103 let mut queue: VecDeque<StateId> = VecDeque::new();
1104
1105 for id in roots {
1106 queue.push_back(*id);
1107 }
1108
1109 while let Some(id) = queue.pop_front() {
1110 if !excluded_states.insert(id) {
1111 continue;
1112 }
1113
1114 let state = match store.get_state(&id)? {
1115 Some(state) => state,
1116 None => continue,
1117 };
1118
1119 for parent in &state.parents {
1120 queue.push_back(*parent);
1121 }
1122
1123 collect_tree_hashes(store, state.tree, &mut excluded_hashes)?;
1124 if let Some(provenance_root) = state.provenance {
1125 collect_tree_hashes(store, provenance_root, &mut excluded_hashes)?;
1126 }
1127 for attachment in store.list_state_attachments(&id)? {
1128 match attachment.body {
1129 StateAttachmentBody::Context(root) => {
1130 collect_tree_hashes(store, root, &mut excluded_hashes)?
1131 }
1132 StateAttachmentBody::RiskSignals(hash)
1133 | StateAttachmentBody::ReviewSignatures(hash)
1134 | StateAttachmentBody::Discussions(hash)
1135 | StateAttachmentBody::StructuredConflicts(hash)
1136 | StateAttachmentBody::ContextConsumption(hash) => {
1137 excluded_hashes.insert(hash);
1138 }
1139 StateAttachmentBody::SemanticIndex(root) => {
1140 collect_semantic_hashes(store, root, &mut excluded_hashes)?;
1141 }
1142 StateAttachmentBody::Signature(_) => {}
1143 }
1144 }
1145 }
1146
1147 Ok((excluded_states, excluded_hashes))
1148}
1149
1150fn collect_tree_hashes(
1151 store: &impl ObjectStore,
1152 tree_hash: ContentHash,
1153 excluded: &mut HashSet<ContentHash>,
1154) -> Result<()> {
1155 if !excluded.insert(tree_hash) {
1156 return Ok(());
1157 }
1158
1159 let tree = match store.get_tree(&tree_hash)? {
1160 Some(tree) => tree,
1161 None => return Ok(()),
1162 };
1163
1164 for entry in tree.entries() {
1165 match entry.target() {
1166 TreeEntryTarget::Blob { hash, .. } | TreeEntryTarget::Symlink { hash } => {
1167 excluded.insert(*hash);
1168 }
1169 TreeEntryTarget::Tree { hash } => {
1170 collect_tree_hashes(store, *hash, excluded)?;
1171 }
1172 TreeEntryTarget::Gitlink { .. } => {}
1173 TreeEntryTarget::Spoollink { .. } => {}
1176 }
1177 }
1178
1179 Ok(())
1180}
1181
1182pub fn is_ancestor(
1183 store: &impl ObjectStore,
1184 ancestor: StateId,
1185 descendant: StateId,
1186) -> Result<bool> {
1187 walk_ancestor(ancestor, descendant, |id| ObjectStore::get_state(store, id))
1188}
1189
1190pub fn is_ancestor_from_source(
1192 source: &(impl ObjectSource + ?Sized),
1193 ancestor: StateId,
1194 descendant: StateId,
1195) -> Result<bool> {
1196 walk_ancestor(ancestor, descendant, |id| source.get_state(id))
1197}
1198
1199#[cfg(feature = "async-source")]
1202#[derive(Clone, Debug)]
1203pub struct AncestryBudget {
1204 pub max_states: usize,
1205 pub cancelled: Arc<AtomicBool>,
1206 pub deadline: Option<Instant>,
1207}
1208
1209#[cfg(feature = "async-source")]
1210#[derive(Debug, thiserror::Error)]
1211pub enum AncestryError {
1212 #[error("ancestry traversal cancelled")]
1213 Cancelled,
1214 #[error("ancestry traversal deadline exceeded")]
1215 Deadline,
1216 #[error("ancestry work limit exhausted")]
1217 WorkLimit,
1218 #[error(transparent)]
1219 Source(#[from] HeddleError),
1220}
1221
1222#[cfg(feature = "async-source")]
1226pub async fn is_ancestor_async_bounded<S>(
1227 source: &S,
1228 ancestor: &StateId,
1229 descendant: &StateId,
1230 budget: &AncestryBudget,
1231) -> std::result::Result<bool, AncestryError>
1232where
1233 S: AsyncObjectSource + ?Sized,
1234{
1235 if ancestor == descendant {
1236 return Ok(true);
1237 }
1238 let check_interruption = || {
1239 if budget.cancelled.load(Ordering::Acquire) {
1240 return Err(AncestryError::Cancelled);
1241 }
1242 if budget
1243 .deadline
1244 .is_some_and(|deadline| Instant::now() >= deadline)
1245 {
1246 return Err(AncestryError::Deadline);
1247 }
1248 Ok(())
1249 };
1250 let mut seen = HashSet::new();
1251 let mut stack = vec![*descendant];
1252 while let Some(id) = stack.pop() {
1253 check_interruption()?;
1254 if !seen.insert(id) {
1255 continue;
1256 }
1257 if seen.len() > budget.max_states {
1258 return Err(AncestryError::WorkLimit);
1259 }
1260 let state = source.get_state(&id).await?;
1261 check_interruption()?;
1262 let Some(state) = state else {
1263 continue;
1264 };
1265 heddle_perf_contract::record_history_object_decode();
1266 if id == *ancestor {
1267 return Ok(true);
1268 }
1269 stack.extend(state.parents);
1270 }
1271 Ok(false)
1272}
1273
1274#[cfg(feature = "async-source")]
1277pub async fn is_ancestor_async<S>(
1278 source: &S,
1279 ancestor: &StateId,
1280 descendant: &StateId,
1281) -> Result<bool>
1282where
1283 S: AsyncObjectSource + ?Sized,
1284{
1285 let budget = AncestryBudget {
1286 max_states: usize::MAX,
1287 cancelled: Arc::new(AtomicBool::new(false)),
1288 deadline: None,
1289 };
1290 is_ancestor_async_bounded(source, ancestor, descendant, &budget)
1291 .await
1292 .map_err(|error| match error {
1293 AncestryError::Source(error) => error,
1294 other => HeddleError::InvalidObject(other.to_string()),
1295 })
1296}
1297
1298fn walk_ancestor(
1299 ancestor: StateId,
1300 descendant: StateId,
1301 mut get_state: impl FnMut(&StateId) -> Result<Option<State>>,
1302) -> Result<bool> {
1303 if ancestor == descendant {
1304 return Ok(true);
1305 }
1306
1307 let mut seen: HashSet<StateId> = HashSet::new();
1308 let mut queue: VecDeque<StateId> = VecDeque::new();
1309 queue.push_back(descendant);
1310
1311 while let Some(id) = queue.pop_front() {
1312 if !seen.insert(id) {
1313 continue;
1314 }
1315 let state = match get_state(&id)? {
1316 Some(s) => s,
1317 None => return Ok(false),
1318 };
1319 for parent in state.parents {
1320 if parent == ancestor {
1321 return Ok(true);
1322 }
1323 queue.push_back(parent);
1324 }
1325 }
1326
1327 Ok(false)
1328}
1329
1330#[cfg(all(test, feature = "async-source"))]
1331mod async_ancestry_tests {
1332 use std::{
1333 collections::HashMap,
1334 future::Future,
1335 sync::atomic::AtomicUsize,
1336 task::{Context, Poll, Waker},
1337 };
1338
1339 use super::*;
1340 use crate::object::{Attribution, Blob, Principal, Tree};
1341
1342 struct Source {
1343 states: HashMap<StateId, State>,
1344 reads: AtomicUsize,
1345 fail: Option<StateId>,
1346 cancel: Option<Arc<AtomicBool>>,
1347 }
1348
1349 impl AsyncObjectSource for Source {
1350 async fn get_tree(&self, _: &ContentHash) -> Result<Option<Tree>> {
1351 Ok(None)
1352 }
1353 async fn get_blob(&self, _: &ContentHash) -> Result<Option<Blob>> {
1354 Ok(None)
1355 }
1356 async fn get_state(&self, id: &StateId) -> Result<Option<State>> {
1357 self.reads.fetch_add(1, Ordering::Relaxed);
1358 if self.fail == Some(*id) {
1359 return Err(HeddleError::InvalidObject("source failed".into()));
1360 }
1361 if let Some(cancel) = &self.cancel {
1362 cancel.store(true, Ordering::Release);
1363 }
1364 Ok(self.states.get(id).cloned())
1365 }
1366 }
1367
1368 fn run<T>(future: impl Future<Output = T>) -> T {
1369 let mut future = Box::pin(future);
1370 let mut context = Context::from_waker(Waker::noop());
1371 loop {
1372 match future.as_mut().poll(&mut context) {
1373 Poll::Ready(value) => return value,
1374 Poll::Pending => std::thread::yield_now(),
1375 }
1376 }
1377 }
1378
1379 fn state(name: &str, parents: Vec<StateId>) -> State {
1380 State::new(
1381 Tree::new().hash(),
1382 parents,
1383 Attribution::human(Principal::new(name, "test@example.com")),
1384 )
1385 }
1386
1387 fn budget(limit: usize) -> AncestryBudget {
1388 AncestryBudget {
1389 max_states: limit,
1390 cancelled: Arc::new(AtomicBool::new(false)),
1391 deadline: None,
1392 }
1393 }
1394
1395 #[test]
1396 fn bounded_walk_distinguishes_complete_missing_error_and_interruption() {
1397 let root = state("root", vec![]);
1398 let left = state("left", vec![root.id()]);
1399 let right = state("right", vec![root.id()]);
1400 let merge = state("merge", vec![left.id(), right.id()]);
1401 let missing = state("missing", vec![]).id();
1402 let states = [root.clone(), left.clone(), right.clone(), merge.clone()]
1403 .into_iter()
1404 .map(|state| (state.id(), state))
1405 .collect();
1406 let mut source = Source {
1407 states,
1408 reads: AtomicUsize::new(0),
1409 fail: None,
1410 cancel: None,
1411 };
1412 assert!(
1413 run(is_ancestor_async_bounded(
1414 &source,
1415 &merge.id(),
1416 &merge.id(),
1417 &budget(0)
1418 ))
1419 .unwrap()
1420 );
1421 assert_eq!(source.reads.load(Ordering::Relaxed), 0);
1422
1423 assert!(
1424 run(is_ancestor_async_bounded(
1425 &source,
1426 &root.id(),
1427 &merge.id(),
1428 &budget(4)
1429 ))
1430 .unwrap()
1431 );
1432 source.states.remove(&root.id());
1433 assert!(
1434 !run(is_ancestor_async_bounded(
1435 &source,
1436 &root.id(),
1437 &merge.id(),
1438 &budget(4)
1439 ))
1440 .unwrap()
1441 );
1442 source.states.insert(root.id(), root.clone());
1443 assert!(
1444 !run(is_ancestor_async_bounded(
1445 &source,
1446 &merge.id(),
1447 &left.id(),
1448 &budget(4)
1449 ))
1450 .unwrap()
1451 );
1452 assert!(
1453 !run(is_ancestor_async_bounded(
1454 &source,
1455 &missing,
1456 &merge.id(),
1457 &budget(5)
1458 ))
1459 .unwrap()
1460 );
1461 assert!(matches!(
1462 run(is_ancestor_async_bounded(
1463 &source,
1464 &root.id(),
1465 &merge.id(),
1466 &budget(1)
1467 )),
1468 Err(AncestryError::WorkLimit)
1469 ));
1470
1471 source.fail = Some(merge.id());
1472 assert!(matches!(
1473 run(is_ancestor_async_bounded(
1474 &source,
1475 &root.id(),
1476 &merge.id(),
1477 &budget(4)
1478 )),
1479 Err(AncestryError::Source(_))
1480 ));
1481 source.fail = None;
1482
1483 let cancelled = budget(4);
1484 cancelled.cancelled.store(true, Ordering::Release);
1485 let before = source.reads.load(Ordering::Relaxed);
1486 assert!(matches!(
1487 run(is_ancestor_async_bounded(
1488 &source,
1489 &root.id(),
1490 &merge.id(),
1491 &cancelled
1492 )),
1493 Err(AncestryError::Cancelled)
1494 ));
1495 assert_eq!(source.reads.load(Ordering::Relaxed), before);
1496
1497 let mut expired = budget(4);
1498 expired.deadline = Some(Instant::now());
1499 assert!(matches!(
1500 run(is_ancestor_async_bounded(
1501 &source,
1502 &root.id(),
1503 &merge.id(),
1504 &expired
1505 )),
1506 Err(AncestryError::Deadline)
1507 ));
1508
1509 let mid_cancel = budget(4);
1510 source.cancel = Some(mid_cancel.cancelled.clone());
1511 assert!(matches!(
1512 run(is_ancestor_async_bounded(
1513 &source,
1514 &root.id(),
1515 &merge.id(),
1516 &mid_cancel
1517 )),
1518 Err(AncestryError::Cancelled)
1519 ));
1520 }
1521}