1use crypto::thread_operation::SignedGenesis;
4use heddle_object_model::object::{
5 ContentHash, StateId,
6 thread_replication::integration::{SPOOL_GENESIS_TRUST_FORMAT, TrustedHostedExecutor},
7};
8use objects::store::ObjectStore;
9use prost::Message;
10use repo::{Repository, thread_replication::ThreadReplica};
11
12use super::{Error, StagedSource};
13use crate::contract::EndpointKind;
14
15pub struct OwnedDeviceBinding<'a> {
18 pub attachment: &'a crate::contract::RootAttachment,
19 pub credential: &'a [u8],
20}
21
22impl StagedSource {
23 pub fn install(self, repository: &Repository, now_unix_seconds: i64) -> Result<StateId, Error> {
27 let endpoint = self
28 .ready
29 .endpoint
30 .as_ref()
31 .ok_or(Error::Invalid("endpoint absent"))?;
32 if endpoint.kind != EndpointKind::Weft as i32 {
33 return Err(Error::Invalid(
34 "hosted installation requires the selected Weft endpoint",
35 ));
36 }
37 let executor = endpoint
38 .public_key
39 .as_slice()
40 .try_into()
41 .map_err(|_| Error::Invalid("invalid hosted endpoint key"))?;
42 let spool = uuid::Uuid::parse_str(
43 &self
44 .ready
45 .thread
46 .as_ref()
47 .and_then(|t| t.spool.as_ref())
48 .ok_or(Error::Invalid("Spool absent"))?
49 .id,
50 )
51 .map_err(preparation)?;
52 let genesis = self
53 .ready
54 .owner_genesis
55 .as_ref()
56 .ok_or(Error::Invalid("owner genesis absent"))?;
57 let owner = self
58 .ready
59 .ownership
60 .as_ref()
61 .ok_or(Error::Invalid("owner history absent"))?;
62 let verified =
63 repo::verify_spool_owner_observation(genesis, owner, spool, now_unix_seconds)
64 .map_err(preparation)?;
65 let body = genesis
66 .genesis
67 .as_ref()
68 .ok_or(Error::Invalid("owner genesis body absent"))?;
69 let trust = TrustedHostedExecutor {
70 spool,
71 spool_genesis: ContentHash::compute_typed(
72 SPOOL_GENESIS_TRUST_FORMAT,
73 &body.encode_to_vec(),
74 ),
75 executor,
76 };
77 for signed in &self.operations {
78 let operation = signed.verify().map_err(preparation)?;
79 require_source_operation(&operation).map_err(preparation)?;
80 if operation
81 .hosted_execution_binding()
82 .map_err(preparation)?
83 .is_some()
84 {
85 trust.authorize(&operation).map_err(preparation)?;
86 }
87 }
88 repository
91 .install_native_spool_id(spool)
92 .map_err(preparation)?;
93 repository
94 .verify_and_pin_owner_observation(
95 genesis,
96 owner,
97 spool,
98 &verified.wire().canonical_spool_path_segments,
99 now_unix_seconds,
100 )
101 .map_err(preparation)?;
102 self.install_replicas(repository, Some(&trust), None, "", now_unix_seconds)?;
103 Ok(self.state.id())
104 }
105 pub fn install_owned_device(
110 self,
111 repository: &Repository,
112 authority: &repo::device_authority::DeviceAuthority,
113 binding: OwnedDeviceBinding<'_>,
114 spool_path: &str,
115 now_unix_seconds: i64,
116 ) -> Result<StateId, Error> {
117 let endpoint = self
118 .ready
119 .endpoint
120 .as_ref()
121 .ok_or(Error::Invalid("endpoint absent"))?;
122 if endpoint.kind != EndpointKind::Device as i32 {
123 return Err(Error::Invalid(
124 "owned-device installation requires device endpoint",
125 ));
126 }
127 let owner = repo::verify_account_owner_observation(&authority.owner, now_unix_seconds)
128 .map_err(preparation)?;
129 let account = owner
130 .signed_root()
131 .root
132 .as_ref()
133 .ok_or(Error::Invalid("account owner root absent"))?;
134 let account_id = uuid::Uuid::from_slice(&account.account_uuid)
135 .map_err(preparation)?
136 .to_string();
137 authority
138 .verify_mint_root(&binding.attachment.root_public_key, now_unix_seconds)
139 .map_err(preparation)?;
140 authority
141 .verify_publisher(&binding.attachment.subject_public_key)
142 .map_err(preparation)?;
143 authority
144 .verify_publisher(&endpoint.public_key)
145 .map_err(preparation)?;
146 let roots = biscuit_verifier::parse_ed25519_public_keys_hex(
147 &hex::encode(&binding.attachment.root_public_key),
148 1,
149 )
150 .map_err(preparation)?;
151 let verified = crate::root_attachment::verify(
152 binding.attachment,
153 binding.credential,
154 &roots,
155 &account_id,
156 endpoint,
157 chrono::DateTime::from_timestamp(now_unix_seconds, 0)
158 .ok_or(Error::Invalid("invalid endpoint verification time"))?,
159 )?;
160 if verified
161 .credential_revocation_ids()
162 .iter()
163 .any(|id| authority.revoked_ids.contains(id))
164 {
165 return Err(Error::Invalid(
166 "endpoint binding credential is explicitly revoked",
167 ));
168 }
169 let spool = self
170 .ready
171 .thread
172 .as_ref()
173 .and_then(|thread| thread.spool.as_ref())
174 .ok_or(Error::Invalid("Spool absent"))?;
175 repository
176 .install_native_spool_id(spool.id.parse().map_err(preparation)?)
177 .map_err(preparation)?;
178 self.install_replicas(
179 repository,
180 None,
181 Some(authority),
182 spool_path,
183 now_unix_seconds,
184 )?;
185 Ok(self.state.id())
186 }
187 fn install_replicas(
188 &self,
189 repository: &Repository,
190 trust: Option<&TrustedHostedExecutor>,
191 authority: Option<&repo::device_authority::DeviceAuthority>,
192 spool_path: &str,
193 now: i64,
194 ) -> Result<(), Error> {
195 let main = self
196 .ready
197 .thread_genesis
198 .as_ref()
199 .ok_or(Error::Invalid("Thread genesis absent"))?;
200 let mut replicas = std::collections::BTreeMap::new();
201 let mut claims = std::collections::BTreeMap::<ContentHash, Vec<PendingClaim>>::new();
202 let mut resolutions = std::collections::BTreeMap::<
203 ContentHash,
204 crate::replication::ownership::OriginalResolution,
205 >::new();
206 for wrapper in std::iter::once(main).chain(&self.dependencies) {
207 let original = wrapper
208 .genesis
209 .as_ref()
210 .ok_or(Error::Invalid("original signed genesis absent"))?;
211 let [signature] = original.signatures.as_slice() else {
212 return Err(Error::Invalid("one original creator signature required"));
213 };
214 let signed = SignedGenesis {
215 canonical: original.canonical_record.clone(),
216 signature: signature.signature.clone(),
217 };
218 let genesis = signed.verify().map_err(preparation)?;
219 let replica = match &genesis.owner {
220 heddle_object_model::object::thread_replication::GenesisOwner::LocalKey(_) => {
221 if !wrapper.creator_authority.is_empty() || wrapper.admission.is_some() {
222 return Err(Error::Invalid(
223 "local ownership cannot carry implicit account admission",
224 ));
225 }
226 ThreadReplica::create(repository.heddle_dir(), &signed).map_err(preparation)?
227 }
228 heddle_object_model::object::thread_replication::GenesisOwner::Account(_) => {
229 if let Some(receipt) = &wrapper.admission {
230 if receipt.format
231 != heddle_object_model::object::thread_genesis_admission::FORMAT
232 {
233 return Err(Error::Invalid("unknown genesis admission format"));
234 }
235 let [_signature] = receipt.signatures.as_slice() else {
236 return Err(Error::Invalid(
237 "one hosted genesis admission signature required",
238 ));
239 };
240 let admission = crate::boundary_acceptance::genesis_admission(wrapper)?
241 .ok_or(Error::Invalid("genesis admission absent"))?;
242 match trust {
243 Some(trust) => ThreadReplica::create_from_genesis_admission(
244 repository.heddle_dir(),
245 &signed,
246 &wrapper.creator_authority,
247 &admission,
248 trust,
249 ),
250 None => ThreadReplica::create_from_pinned_genesis_admission(
251 repository.heddle_dir(),
252 &signed,
253 &wrapper.creator_authority,
254 &admission,
255 ),
256 }
257 .map_err(preparation)?
258 } else {
259 let authority = authority.ok_or(Error::Invalid(
260 "original account genesis admission required",
261 ))?;
262 ThreadReplica::create_authorized(
263 repository.heddle_dir(),
264 &signed,
265 &wrapper.creator_authority,
266 authority,
267 spool_path,
268 "/heddle.api.v1alpha2.ThreadService/StartThread",
269 now,
270 )
271 .map_err(preparation)?
272 }
273 }
274 };
275 if let Some(trust) = trust {
276 let executor = trust.executor;
277 repository
278 .pin_thread_hosted_executor(&replica, executor)
279 .map_err(preparation)?;
280 }
281 let mut pending = Vec::new();
282 let retained = replica.ownership_claims().map_err(preparation)?;
283 for claim in crate::replication::ownership::verify_claims(wrapper, &genesis)? {
284 let value = claim.original.verify().map_err(preparation)?;
285 if let Some(receipt) = &claim.authority_admission {
286 let trust = replica
287 .authority_admission_trust(receipt)
288 .map_err(preparation)?;
289 receipt
290 .verify_claim(&claim.original, &genesis, &trust)
291 .map_err(preparation)?;
292 } else if !retained.contains(&claim.original) {
293 let authority = authority.ok_or(Error::Invalid("new claim requires current original acceptance or independently pinned admission"))?;
294 repo::thread_replication::ownership_claim::verify_claim_authority(
295 &claim.original,
296 &genesis,
297 authority,
298 spool_path,
299 now,
300 )
301 .map_err(preparation)?;
302 }
303 pending.push(PendingClaim {
304 original: claim,
305 remaining: value.source_frontier,
306 });
307 }
308 for resolution in crate::replication::ownership::verify_resolutions(wrapper, &genesis)?
309 {
310 if resolutions
311 .insert(replica.thread_id(), resolution)
312 .is_some()
313 {
314 return Err(Error::Invalid("duplicate ownership resolution"));
315 }
316 }
317 claims.insert(replica.thread_id(), pending);
318 replicas.insert(replica.thread_id(), replica);
319 }
320 for signed in &self.operations {
324 let operation = signed.verify().map_err(preparation)?;
325 let replica = replicas
326 .get(&operation.thread)
327 .ok_or(Error::Invalid("source dependency replica absent"))?;
328 if let Some(receipt) = self
329 .authority_admissions
330 .get(&operation.id().map_err(preparation)?)
331 {
332 replica
333 .require_authority_admission(signed, receipt)
334 .map_err(preparation)?;
335 }
336 }
337 self.install_source_objects(repository)?;
338 for (thread, pending) in &mut claims {
339 install_ready_claims(
340 replicas
341 .get(thread)
342 .ok_or(Error::Invalid("claim replica absent"))?,
343 pending,
344 authority,
345 spool_path,
346 now,
347 )?;
348 }
349 for (thread, resolution) in &resolutions {
350 install_ready_resolution(
351 replicas
352 .get(thread)
353 .ok_or(Error::Invalid("resolution replica absent"))?,
354 resolution,
355 authority,
356 spool_path,
357 now,
358 )?;
359 }
360 for signed in &self.operations {
361 let operation = signed.verify().map_err(preparation)?;
362 let replica = replicas
363 .get(&operation.thread)
364 .ok_or(Error::Invalid("source dependency replica absent"))?;
365 let id = operation.id().map_err(preparation)?;
366 if !self.authority_admissions.contains_key(&id) {
367 let prior = replica
368 .operation_with_authority_admission(&id)
369 .map_err(preparation)?;
370 if !prior.is_some_and(|prior| {
371 prior.original == *signed
372 && prior.status == objects::object::thread_replication::Admission::Accepted
373 }) && let Some(author) = operation.source_author().map_err(preparation)?
374 {
375 match author {
376 objects::object::thread_replication::SourceAuthor::LocalKey => replica
377 .verify_local_source_owner(&operation)
378 .map_err(preparation)?,
379 objects::object::thread_replication::SourceAuthor::Account { .. } => {
380 replica.verify_source_authority(&operation, authority.ok_or(Error::Invalid("fresh source requires original authority or retained admission"))?, spool_path, now).map_err(preparation)?;
381 }
382 }
383 }
384 }
385 let admission = if !self.is_complete() {
386 replica.receive_source_metadata(
387 signed,
388 repository.store(),
389 self.authority_admissions.get(&id),
390 require_source_operation,
391 )
392 } else if let Some(receipt) = self
393 .authority_admissions
394 .get(&operation.id().map_err(preparation)?)
395 {
396 replica.receive_with_authority_admission(
397 signed,
398 receipt,
399 repository.store(),
400 require_source_operation,
401 )
402 } else {
403 replica.receive(signed, repository.store(), require_source_operation)
404 }
405 .map_err(preparation)?;
406 if admission != objects::object::thread_replication::Admission::Accepted {
407 return Err(Error::Invalid(
408 "source proof did not settle in dependency order",
409 ));
410 }
411 if let Some(pending) = claims.get_mut(&operation.thread) {
412 for claim in pending.iter_mut() {
413 claim.remaining.remove(&id);
414 }
415 install_ready_claims(replica, pending, authority, spool_path, now)?;
416 }
417 if let Some(resolution) = resolutions.get(&operation.thread) {
418 install_ready_resolution(replica, resolution, authority, spool_path, now)?;
419 }
420 }
421 if claims.values().any(|claims| !claims.is_empty()) {
422 return Err(Error::Invalid(
423 "ownership cutoff did not settle before source completion",
424 ));
425 }
426 for (thread, resolution) in &resolutions {
427 let replica = replicas
428 .get(thread)
429 .ok_or(Error::Invalid("resolution replica absent"))?;
430 install_ready_resolution(replica, resolution, authority, spool_path, now)?;
431 if replica
432 .ownership_resolution()
433 .map_err(preparation)?
434 .is_none()
435 {
436 return Err(Error::Invalid(
437 "ownership resolution frontier did not settle",
438 ));
439 }
440 }
441 for thread in claims
442 .keys()
443 .filter(|thread| !resolutions.contains_key(*thread))
444 {
445 replicas
446 .get(thread)
447 .ok_or(Error::Invalid("claim replica absent"))?
448 .effective_owner()
449 .map_err(preparation)?;
450 }
451 let main_id = crate::replication::opening::verify_genesis(
452 main.genesis
453 .as_ref()
454 .ok_or(Error::Invalid("signed genesis absent"))?,
455 self.ready
456 .thread
457 .as_ref()
458 .ok_or(Error::Invalid("Thread absent"))?,
459 )?
460 .id()
461 .map_err(preparation)?;
462 let selected = replicas
463 .get(&main_id)
464 .ok_or(Error::Invalid("selected replica absent"))?;
465 if self.operations.is_empty() {
466 let genesis = selected.genesis().map_err(preparation)?;
467 let canonical = self.state.encode_current_msgpack().map_err(preparation)?;
468 objects::object::thread_replication::hosted_import::initial_base_state(
469 &genesis, &canonical,
470 )
471 .map_err(preparation)?;
472 } else {
473 let mut source = selected;
474 let mut proved = false;
475 let mut possession = Vec::new();
476 for _ in 0..128 {
477 if source
478 .accepted_source_revision(self.state.id())
479 .map_err(preparation)?
480 .is_some()
481 {
482 possession.push(source);
483 proved = true;
484 break;
485 }
486 let genesis = source.genesis().map_err(preparation)?;
487 if genesis.base != self.state.id() {
488 break;
489 }
490 possession.push(source);
491 let Some(parent) = genesis.parent else { break };
492 let Some(next) = replicas.get(&parent) else {
493 break;
494 };
495 if next.genesis().map_err(preparation)?.spool != genesis.spool {
496 break;
497 }
498 source = next;
499 }
500 if !proved {
501 return Err(Error::Invalid("selected source proof did not settle"));
502 }
503 if self.is_complete() {
504 for replica in possession {
505 replica
506 .record_source_possession(self.state.id())
507 .map_err(preparation)?;
508 }
509 }
510 }
511 if self.is_complete() {
512 selected
513 .record_source_possession(self.state.id())
514 .map_err(preparation)?;
515 }
516 Ok(())
517 }
518
519 pub(super) fn install_source_objects(&self, repository: &Repository) -> Result<(), Error> {
520 let pack = self.directory.path().join("source.pack");
521 let index = self.directory.path().join("source.idx");
522 if self.is_complete() {
523 return repository
524 .store()
525 .install_pack_streaming(&pack, &index)
526 .map(|_| ())
527 .map_err(preparation);
528 }
529 let visible_pack = self.directory.path().join("visible.pack");
532 let visible_index = self.directory.path().join("visible.idx");
533 let output = std::fs::OpenOptions::new()
534 .read(true)
535 .write(true)
536 .create_new(true)
537 .open(&visible_pack)?;
538 let mut builder = heddle_pack::store::pack::StreamingPackBuilder::new(
539 output,
540 visible_index.clone(),
541 Default::default(),
542 self.directory.path().join("visible-buckets"),
543 )
544 .map_err(preparation)?;
545 let reader =
546 heddle_pack::store::pack::PackReader::open(&pack, &index).map_err(preparation)?;
547 reader
548 .visit_objects(|id, kind, bytes| {
549 if kind != heddle_pack::store::pack::ObjectType::Tree
550 || !objects::object::is_redacted_tree(bytes)
551 {
552 builder.add_id(id, kind, bytes)?;
553 }
554 Ok(())
555 })
556 .map_err(preparation)?;
557 let (output, _) = builder.finalize().map_err(preparation)?;
558 drop(output);
559 for partial in &self.partial_trees {
560 let bytes =
561 objects::object::encode_redacted_projection(partial).map_err(preparation)?;
562 repository
563 .store()
564 .put_partial_tree(&partial.declared_root(), &bytes)
565 .map_err(preparation)?;
566 }
567 repository
568 .store()
569 .install_pack_streaming(&visible_pack, &visible_index)
570 .map(|_| ())
571 .map_err(preparation)
572 }
573}
574struct PendingClaim {
575 original: crate::replication::ownership::OriginalClaim,
576 remaining: std::collections::BTreeSet<ContentHash>,
577}
578fn install_ready_resolution(
579 replica: &ThreadReplica,
580 resolution: &crate::replication::ownership::OriginalResolution,
581 authority: Option<&repo::device_authority::DeviceAuthority>,
582 spool_path: &str,
583 now: i64,
584) -> Result<(), Error> {
585 if let Some(existing) = replica.ownership_resolution().map_err(preparation)? {
586 if existing == resolution.original {
587 return Ok(());
588 }
589 return Err(Error::Invalid(
590 "incoming ownership resolution conflicts with retained history",
591 ));
592 }
593 let value = heddle_object_model::object::thread_replication::ownership_resolution::ThreadOwnershipResolution::decode(&resolution.original.canonical)
594 .map_err(preparation)?;
595 if replica.ownership_claims().map_err(preparation)?.len() != value.conflicting_claims.len() {
596 return Ok(());
597 }
598 for head in &value.frontier {
599 if replica
600 .operation(head)
601 .map_err(preparation)?
602 .is_none_or(|(_, status)| {
603 status != objects::object::thread_replication::Admission::Accepted
604 })
605 {
606 return Ok(());
607 }
608 }
609 if let Some(receipt) = &resolution.authority_admission {
610 replica
611 .resolve_ownership_with_admission(&resolution.original, receipt)
612 .map_err(preparation)?;
613 } else {
614 replica
615 .resolve_ownership(
616 &resolution.original,
617 authority.ok_or(Error::Invalid(
618 "new ownership resolution requires current recipient authority",
619 ))?,
620 spool_path,
621 now,
622 )
623 .map_err(preparation)?;
624 }
625 Ok(())
626}
627fn install_ready_claims(
628 replica: &ThreadReplica,
629 pending: &mut Vec<PendingClaim>,
630 authority: Option<&repo::device_authority::DeviceAuthority>,
631 spool_path: &str,
632 now: i64,
633) -> Result<(), Error> {
634 let mut index = 0;
635 while index < pending.len() {
636 if !pending[index].remaining.is_empty() {
637 index += 1;
638 continue;
639 }
640 let claim = pending.remove(index).original;
641 if let Some(receipt) = &claim.authority_admission {
642 replica
643 .claim_ownership_with_admission(&claim.original, receipt)
644 .map_err(preparation)?;
645 } else if !replica
646 .ownership_claims()
647 .map_err(preparation)?
648 .contains(&claim.original)
649 {
650 replica
651 .claim_ownership_for_import(
652 &claim.original,
653 authority.ok_or(Error::Invalid("new claim authority absent"))?,
654 spool_path,
655 now,
656 )
657 .map_err(preparation)?;
658 }
659 }
660 Ok(())
661}
662fn preparation(error: impl std::fmt::Display) -> Error {
663 Error::Preparation(error.to_string())
664}
665
666fn require_source_operation(
669 operation: &heddle_object_model::object::thread_replication::ThreadOperation,
670) -> repo::thread_replication::Result<()> {
671 if operation.facet() != heddle_object_model::object::thread_replication::ThreadFacet::Source {
672 return Err(repo::thread_replication::Error::Invalid(
673 "source installation cannot admit non-source authority".into(),
674 ));
675 }
676 Ok(())
677}
678
679#[cfg(test)]
680mod tests {
681 use crypto::{Ed25519Signer, Signer, thread_operation::SignedOperation};
682 use heddle_object_model::object::{
683 CollaborationActor,
684 thread_replication::{
685 ThreadGenesis, ThreadOperation, ThreadOperationBody,
686 metadata::{AUTHORITY_FORMAT, Control, ThreadControl},
687 },
688 };
689
690 use super::*;
691 #[test]
692 fn source_install_gate_cannot_create_metadata_original_authority() {
693 let directory = tempfile::tempdir().expect("repository");
694 let repository = Repository::init_default(directory.path()).expect("repo");
695 let signer = Ed25519Signer::from_seed(&[56; 32]).expect("signer");
696 let spool = uuid::Uuid::from_u128(11);
697 let genesis = ThreadGenesis {
698 owner: objects::object::thread_replication::GenesisOwner::LocalKey(
699 signer.public_key().try_into().expect("key"),
700 ),
701 version: 1,
702 spool: spool.to_string(),
703 parent: None,
704 base: repository.head().expect("head").expect("base"),
705 name: "source import".into(),
706 intent: "original authority".into(),
707 creator: signer.public_key().try_into().expect("key"),
708 nonce: vec![5],
709 };
710 let replica = ThreadReplica::create(
711 repository.heddle_dir(),
712 &SignedGenesis::sign(&genesis, &signer).expect("original genesis"),
713 )
714 .expect("replica");
715 let proof = b"unverified author evidence must not establish admission".to_vec();
716 let control = ThreadControl {
717 version: 1,
718 spool,
719 actor: CollaborationActor {
720 principal_id: uuid::Uuid::from_u128(22),
721 agent_id: None,
722 },
723 authority_digest: ContentHash::compute_typed(AUTHORITY_FORMAT, &proof),
724 authority_envelope: proof,
725 client_operation_id: uuid::Uuid::now_v7(),
726 occurred_at_ms: 0,
727 control: Control::Name("unproved author".into()),
728 };
729 let operation = ThreadOperation {
730 version: 1,
731 thread: replica.thread_id(),
732 parents: Default::default(),
733 publisher: signer.public_key().try_into().expect("key"),
734 body: ThreadOperationBody::Metadata(control.encode().expect("valid canonical control")),
735 };
736 let signed = SignedOperation::sign(&operation, &signer).expect("valid original signature");
737 signed.verify().expect("signature itself is valid");
738 let failure = replica
739 .receive(&signed, repository.store(), require_source_operation)
740 .expect_err("source-only gate denies unproved Metadata");
741 assert!(failure.to_string().contains("non-source authority"));
742 assert!(
743 replica
744 .operation(&operation.id().expect("ID"))
745 .expect("stored operation")
746 .is_none(),
747 "denial precedes immutable persistence"
748 );
749 assert!(
750 !replica
751 .original_authority_admitted(&signed)
752 .expect("admission marker"),
753 "source trust cannot manufacture author admission"
754 );
755 }
756}