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