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 self.is_complete() {
388 return repository
389 .store()
390 .install_pack_streaming(&pack, &index)
391 .map(|_| ())
392 .map_err(preparation);
393 }
394 let visible_pack = self.directory.path().join("visible.pack");
397 let visible_index = self.directory.path().join("visible.idx");
398 let output = std::fs::OpenOptions::new()
399 .read(true)
400 .write(true)
401 .create_new(true)
402 .open(&visible_pack)?;
403 let mut builder = heddle_pack::store::pack::StreamingPackBuilder::new(
404 output,
405 visible_index.clone(),
406 Default::default(),
407 self.directory.path().join("visible-buckets"),
408 )
409 .map_err(preparation)?;
410 let reader =
411 heddle_pack::store::pack::PackReader::open(&pack, &index).map_err(preparation)?;
412 reader
413 .visit_objects(|id, kind, bytes| {
414 if kind != heddle_pack::store::pack::ObjectType::Tree
415 || !objects::object::is_redacted_tree(bytes)
416 {
417 builder.add_id(id, kind, bytes)?;
418 }
419 Ok(())
420 })
421 .map_err(preparation)?;
422 let (output, _) = builder.finalize().map_err(preparation)?;
423 drop(output);
424 for partial in &self.partial_trees {
425 let bytes =
426 objects::object::encode_redacted_projection(partial).map_err(preparation)?;
427 repository
428 .store()
429 .put_partial_tree(&partial.declared_root(), &bytes)
430 .map_err(preparation)?;
431 }
432 repository
433 .store()
434 .install_pack_streaming(&visible_pack, &visible_index)
435 .map(|_| ())
436 .map_err(preparation)
437 }
438}
439struct PendingClaim {
440 original: crate::replication::ownership::OriginalClaim,
441 remaining: std::collections::BTreeSet<ContentHash>,
442}
443fn install_ready_resolution(
444 replica: &ThreadReplica,
445 resolution: &crate::replication::ownership::OriginalResolution,
446 authority: Option<&repo::device_authority::DeviceAuthority>,
447 spool_path: &str,
448 now: i64,
449) -> Result<(), Error> {
450 if let Some(existing) = replica.ownership_resolution().map_err(preparation)? {
451 if existing == resolution.original {
452 return Ok(());
453 }
454 return Err(Error::Invalid(
455 "incoming ownership resolution conflicts with retained history",
456 ));
457 }
458 let value = heddle_object_model::object::thread_replication::ownership_resolution::ThreadOwnershipResolution::decode(&resolution.original.canonical)
459 .map_err(preparation)?;
460 if replica.ownership_claims().map_err(preparation)?.len() != value.conflicting_claims.len() {
461 return Ok(());
462 }
463 for head in &value.frontier {
464 if replica
465 .operation(head)
466 .map_err(preparation)?
467 .is_none_or(|(_, status)| {
468 status != objects::object::thread_replication::Admission::Accepted
469 })
470 {
471 return Ok(());
472 }
473 }
474 if resolution.authority_admission.is_some() {
475 return Err(Error::HostedTrustRequired);
476 } else {
477 replica
478 .resolve_ownership(
479 &resolution.original,
480 authority.ok_or(Error::Invalid(
481 "new ownership resolution requires current recipient authority",
482 ))?,
483 spool_path,
484 now,
485 )
486 .map_err(preparation)?;
487 }
488 Ok(())
489}
490fn install_ready_claims(
491 replica: &ThreadReplica,
492 pending: &mut Vec<PendingClaim>,
493 authority: Option<&repo::device_authority::DeviceAuthority>,
494 spool_path: &str,
495 now: i64,
496) -> Result<(), Error> {
497 let mut index = 0;
498 while index < pending.len() {
499 if !pending[index].remaining.is_empty() {
500 index += 1;
501 continue;
502 }
503 let claim = pending.remove(index).original;
504 if claim.authority_admission.is_some() {
505 return Err(Error::HostedTrustRequired);
506 } else if !replica
507 .ownership_claims()
508 .map_err(preparation)?
509 .contains(&claim.original)
510 {
511 replica
512 .claim_ownership_for_import(
513 &claim.original,
514 authority.ok_or(Error::Invalid("new claim authority absent"))?,
515 spool_path,
516 now,
517 )
518 .map_err(preparation)?;
519 }
520 }
521 Ok(())
522}
523fn preparation(error: impl std::fmt::Display) -> Error {
524 Error::Preparation(error.to_string())
525}
526
527fn require_source_operation(
530 operation: &heddle_object_model::object::thread_replication::ThreadOperation,
531) -> repo::thread_replication::Result<()> {
532 if operation.facet() != heddle_object_model::object::thread_replication::ThreadFacet::Source {
533 return Err(repo::thread_replication::Error::Invalid(
534 "source installation cannot admit non-source authority".into(),
535 ));
536 }
537 Ok(())
538}
539
540#[cfg(test)]
541mod tests {
542 use crypto::{Ed25519Signer, Signer, thread_operation::SignedOperation};
543 use heddle_object_model::object::{
544 CollaborationActor,
545 thread_replication::{
546 ThreadGenesis, ThreadOperation, ThreadOperationBody,
547 metadata::{AUTHORITY_FORMAT, Control, ThreadControl},
548 },
549 };
550
551 use super::*;
552 #[test]
553 fn source_install_gate_cannot_create_metadata_original_authority() {
554 let directory = tempfile::tempdir().expect("repository");
555 let repository = Repository::init_default(directory.path()).expect("repo");
556 let signer = Ed25519Signer::from_seed(&[56; 32]).expect("signer");
557 let spool = uuid::Uuid::from_u128(11);
558 let genesis = ThreadGenesis {
559 owner: objects::object::thread_replication::GenesisOwner::LocalKey(
560 signer.public_key().try_into().expect("key"),
561 ),
562 version: 1,
563 spool: spool.to_string(),
564 parent: None,
565 base: repository.head().expect("head").expect("base"),
566 name: "source import".into(),
567 intent: "original authority".into(),
568 creator: signer.public_key().try_into().expect("key"),
569 nonce: vec![5],
570 };
571 let replica = ThreadReplica::create(
572 repository.heddle_dir(),
573 &SignedGenesis::sign(&genesis, &signer).expect("original genesis"),
574 )
575 .expect("replica");
576 let proof = b"unverified author evidence must not establish admission".to_vec();
577 let control = ThreadControl {
578 version: 1,
579 spool,
580 actor: CollaborationActor {
581 principal_id: uuid::Uuid::from_u128(22),
582 agent_id: None,
583 },
584 authority_digest: ContentHash::compute_typed(AUTHORITY_FORMAT, &proof),
585 authority_envelope: proof,
586 client_operation_id: uuid::Uuid::now_v7(),
587 occurred_at_ms: 0,
588 control: Control::Name("unproved author".into()),
589 };
590 let operation = ThreadOperation {
591 version: 1,
592 thread: replica.thread_id(),
593 parents: Default::default(),
594 publisher: signer.public_key().try_into().expect("key"),
595 body: ThreadOperationBody::Metadata(control.encode().expect("valid canonical control")),
596 };
597 let signed = SignedOperation::sign(&operation, &signer).expect("valid original signature");
598 signed.verify().expect("signature itself is valid");
599 let failure = replica
600 .receive(&signed, repository.store(), require_source_operation)
601 .expect_err("source-only gate denies unproved Metadata");
602 assert!(failure.to_string().contains("non-source authority"));
603 assert!(
604 replica
605 .operation(&operation.id().expect("ID"))
606 .expect("stored operation")
607 .is_none(),
608 "denial precedes immutable persistence"
609 );
610 assert!(
611 !replica
612 .original_authority_admitted(&signed)
613 .expect("admission marker"),
614 "source trust cannot manufacture author admission"
615 );
616 }
617}