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).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
535// Stage already negotiates Source alone. Keep that trust boundary explicit at
536// installation too: source verification is never original Metadata authority.
537fn 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}