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