Skip to main content

heddle_thread_api/fetch/
native.rs

1//! Install a verified hosted download through the existing local store and
2//! replica admission paths. This disk operation never advances a checkout.
3use 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
11/// Current endpoint possession under independently retained account authority.
12/// Retain the original credential privately; it never joins source proof packs.
13pub struct OwnedDeviceBinding<'a> {
14    pub attachment: &'a crate::contract::RootAttachment,
15    pub credential: &'a [u8],
16}
17
18impl StagedSource {
19    /// Verify a privately owned device through independent account authority.
20    /// Originals must already belong to this exact Spool.
21    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        // HRT1 is a disclosure proof, not a full tree object. Split it out before
395        // registering the visible immutable records in the shared object store.
396        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
527// Stage already negotiates Source alone. Keep that trust boundary explicit at
528// installation too: source verification is never original Metadata authority.
529fn 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}