Skip to main content

heddle_thread_api/replication/
opening.rs

1//! Shared negotiation for already-authorized, already-resolved Threads.
2//! This validates transport bindings and limits; it does not grant authority.
3use std::collections::BTreeSet;
4
5use crypto::{Signer, thread_operation::SignedGenesis};
6use heddle_object_model::object::thread_replication::{
7    GENESIS_FORMAT, OPERATION_FORMAT, ThreadFacet, ThreadGenesis,
8};
9
10use crate::{contract::*, transport::Error};
11
12pub const FRAME_LIMIT: usize = 512 * 1024;
13pub const MAX_ITEMS: u32 = 64;
14
15pub fn validate_endpoint(endpoint: &EndpointRef) -> Result<(), Error> {
16    if endpoint.public_key.len() != 32
17        || !matches!(
18            EndpointKind::try_from(endpoint.kind),
19            Ok(EndpointKind::Device | EndpointKind::Weft)
20        )
21    {
22        return Err(Error::Protocol("invalid replication endpoint"));
23    }
24    Ok(())
25}
26
27pub fn parse_facets(values: &[i32]) -> Result<BTreeSet<ThreadFacet>, Error> {
28    if values.is_empty() || values.len() > ThreadFacet::ALL.len() {
29        return Err(Error::Protocol(
30            "replication requires a bounded set of distinct facets",
31        ));
32    }
33    let facets = values
34        .iter()
35        .map(|value| {
36            super::native_facet(*value)
37                .map_err(|_| Error::Protocol("unsupported replication facet"))
38        })
39        .collect::<Result<BTreeSet<_>, _>>()?;
40    if facets.len() != values.len() {
41        return Err(Error::Protocol("duplicate replication facet"));
42    }
43    Ok(facets)
44}
45
46/// The remote key must come from the authenticated transport, never the frame.
47/// An attached genesis retains the original creator's signature and identity.
48/// The host must authorize publication and commit that genesis before Ready;
49/// parsing it grants no authority. PoP covers the entire exact opening.
50pub fn accept(
51    open: &ReplicationOpen,
52    thread: &ThreadRef,
53    local: &EndpointRef,
54    remote_key: [u8; 32],
55    admission: &BTreeSet<ThreadFacet>,
56    sharing_policy_version: Vec<u8>,
57) -> Result<AcceptedOpening, Error> {
58    crate::hybrid::replication_open(open).map_err(Error::Protocol)?;
59    validate_endpoint(local)?;
60    if open.thread.as_ref() != Some(thread) {
61        return Err(Error::Protocol(
62            "replication requires an already resolved Thread",
63        ));
64    }
65    let source = open
66        .source
67        .as_ref()
68        .ok_or(Error::Protocol("opening requires source endpoint"))?;
69    validate_endpoint(source)?;
70    if source.public_key != remote_key || open.destination.as_ref() != Some(local) {
71        return Err(Error::Protocol(
72            "opening endpoints differ from Iroh connection",
73        ));
74    }
75    if open.session_nonce.len() != 16
76        || !open
77            .record_formats
78            .iter()
79            .any(|format| format == OPERATION_FORMAT)
80    {
81        return Err(Error::Protocol("unsupported replication session or format"));
82    }
83    // Native records are indivisible. Negotiate the supported frame size
84    // explicitly instead of promising a smaller budget and exceeding it.
85    if open.budget.as_ref().is_some_and(|budget| {
86        budget.max_frame_bytes != 0 && budget.max_frame_bytes != FRAME_LIMIT as u32
87    }) {
88        return Err(Error::Protocol("unsupported replication frame budget"));
89    }
90    let facets: BTreeSet<_> = parse_facets(&open.facets)?
91        .intersection(admission)
92        .copied()
93        .collect();
94    if facets.is_empty() {
95        return Err(Error::Protocol("no authorized replication facets"));
96    }
97    let requested = open.budget.as_ref().map_or(0, |budget| budget.max_items);
98    let genesis = open
99        .thread_genesis
100        .as_ref()
101        .map(|record| verify_genesis_record(record, thread))
102        .transpose()?;
103    Ok(AcceptedOpening {
104        genesis,
105        genesis_record: open.thread_genesis.clone(),
106        import_authority: open.import_authority.clone(),
107        native_authority: open.native_authority.clone(),
108        ready: ReplicationReady {
109            native_authority: None,
110            thread: Some(thread.clone()),
111            endpoint: Some(local.clone()),
112            facets: facets.into_iter().map(super::wire_facet).collect(),
113            sharing_policy_version,
114            budget: Some(ReadBudget {
115                max_items: if requested == 0 {
116                    MAX_ITEMS
117                } else {
118                    requested.min(MAX_ITEMS)
119                },
120                max_frame_bytes: FRAME_LIMIT as u32,
121                max_snapshot_bytes: 0,
122            }),
123            record_formats: vec![OPERATION_FORMAT.into()],
124            // Echo only this exact opening's optional capable path. The
125            // mandatory Sync switch stays off until api#307.
126            protocol: open.protocol.clone(),
127            import_authority: None,
128        },
129    })
130}
131
132/// Parsed proposal; authorization and durable installation belong to the host.
133pub struct AcceptedOpening {
134    pub ready: ReplicationReady,
135    pub genesis: Option<ThreadGenesis>,
136    pub genesis_record: Option<ThreadGenesisRecord>,
137    /// Retained structural evidence; the receiver must verify before mutation.
138    pub import_authority: Option<crate::contract::ImportPublicProofBundleV1>,
139    pub native_authority: Option<crate::contract::NativePublicProofBundleV1>,
140}
141
142/// Structural and original-signature validation only. Account authority and
143/// hosted admission receipts still require independently retained trust.
144pub fn verify_genesis_record(
145    record: &ThreadGenesisRecord,
146    thread: &ThreadRef,
147) -> Result<ThreadGenesis, Error> {
148    use prost::Message;
149    if record.boundary_acceptances.len() > crate::boundary_acceptance::MAX_ACCEPTANCES
150        || record.encoded_len() > 256 * 1024
151    {
152        return Err(Error::Protocol("genesis wrapper evidence exceeds bounds"));
153    }
154    let signed = record
155        .genesis
156        .as_ref()
157        .ok_or(Error::Protocol("original signed genesis missing"))?;
158    if let Some(binding) = &record.native_genesis_authority {
159        api::native_witness::verify_genesis_authority(
160            binding,
161            record
162                .genesis
163                .as_ref()
164                .ok_or(Error::Protocol("original genesis absent"))?,
165            &record.creator_authority,
166        )
167        .map_err(|_| Error::Protocol("native creator binding differs from original"))?;
168    }
169    let genesis = verify_genesis(signed, thread)?;
170    if record.creator_authority.len() > 64 * 1024 {
171        return Err(Error::Protocol("creator authority exceeds bound"));
172    }
173    use heddle_object_model::object::thread_replication::GenesisOwner;
174    match genesis.owner {
175        GenesisOwner::LocalKey(_)
176            if !record.creator_authority.is_empty() || record.admission.is_some() =>
177        {
178            return Err(Error::Protocol(
179                "local-key ownership requires an explicit claim, not an account envelope",
180            ));
181        }
182        GenesisOwner::Account(_) if record.creator_authority.is_empty() => {
183            return Err(Error::Protocol(
184                "account-owned genesis requires original creator authority",
185            ));
186        }
187        _ => {}
188    }
189    super::ownership::verify_claims(record, &genesis)?;
190    Ok(genesis)
191}
192
193pub fn sign_genesis(genesis: &ThreadGenesis, signer: &impl Signer) -> Result<SignedRecord, Error> {
194    let signed =
195        SignedGenesis::sign(genesis, signer).map_err(|error| Error::Io(error.to_string()))?;
196    Ok(SignedRecord {
197        format: GENESIS_FORMAT.into(),
198        canonical_record: signed.canonical,
199        signatures: vec![RecordSignature {
200            public_key: genesis.creator.to_vec(),
201            signature: signed.signature,
202        }],
203    })
204}
205
206pub fn verify_genesis(record: &SignedRecord, thread: &ThreadRef) -> Result<ThreadGenesis, Error> {
207    if record.format != GENESIS_FORMAT || record.signatures.len() != 1 {
208        return Err(Error::Protocol("unsupported Thread genesis record"));
209    }
210    let genesis = SignedGenesis {
211        canonical: record.canonical_record.clone(),
212        signature: record.signatures[0].signature.clone(),
213    }
214    .verify()
215    .map_err(|_| Error::Protocol("invalid Thread genesis signature"))?;
216    let id = genesis
217        .id()
218        .map_err(|_| Error::Protocol("invalid Thread genesis"))?;
219    if record.signatures[0].public_key != genesis.creator
220        || thread
221            .spool
222            .as_ref()
223            .is_none_or(|spool| spool.id != genesis.spool)
224        || thread
225            .id
226            .as_ref()
227            .is_none_or(|thread| thread.value != id.as_bytes())
228    {
229        return Err(Error::Protocol(
230            "Thread genesis identity does not match opening",
231        ));
232    }
233    Ok(genesis)
234}
235
236pub fn validate_ready(
237    ready: &ReplicationReady,
238    thread: &ThreadRef,
239    destination: &EndpointRef,
240    requested_facets: &BTreeSet<ThreadFacet>,
241    requested_max_items: u32,
242) -> Result<(BTreeSet<ThreadFacet>, usize), Error> {
243    crate::hybrid::replication_ready(ready).map_err(Error::Protocol)?;
244    validate_endpoint(destination)?;
245    if ready.endpoint.as_ref() != Some(destination)
246        || ready.thread.as_ref() != Some(thread)
247        || ready.record_formats != [OPERATION_FORMAT]
248    {
249        return Err(Error::Protocol(
250            "replication Ready binding differs from opening",
251        ));
252    }
253    let facets = parse_facets(&ready.facets)?;
254    if !facets.is_subset(requested_facets) {
255        return Err(Error::Protocol("replication Ready widened admission scope"));
256    }
257    let budget = ready
258        .budget
259        .as_ref()
260        .ok_or(Error::Protocol("replication Ready requires budget"))?;
261    let ceiling = if requested_max_items == 0 {
262        MAX_ITEMS
263    } else {
264        requested_max_items.min(MAX_ITEMS)
265    };
266    if budget.max_items == 0
267        || budget.max_items > ceiling
268        || budget.max_frame_bytes != FRAME_LIMIT as u32
269    {
270        return Err(Error::Protocol("unsupported replication Ready budget"));
271    }
272    Ok((facets, budget.max_items as usize))
273}
274
275#[cfg(test)]
276mod tests {
277    use super::*;
278
279    #[test]
280    fn first_publication_retains_signed_genesis_and_rejects_changed_identity() {
281        use crypto::{Ed25519Signer, Signer};
282        use heddle_object_model::object::{
283            StateId,
284            thread_replication::{GENESIS_FORMAT, ThreadGenesis},
285        };
286        let signer = Ed25519Signer::from_seed(&[23; 32]).expect("origin device");
287        let genesis = ThreadGenesis {
288            version: 1,
289            spool: "01980000-0000-7000-8000-000000000001".into(),
290            parent: None,
291            base: StateId::from_bytes([1; 32]),
292            name: "original".into(),
293            intent: "publish once".into(),
294            creator: signer.public_key().try_into().expect("public key"),
295            owner: heddle_object_model::object::thread_replication::GenesisOwner::LocalKey(
296                signer.public_key().try_into().expect("public key"),
297            ),
298            nonce: vec![2; 16],
299        };
300        let canonical = genesis.encode().expect("canonical genesis");
301        let mut signing = GENESIS_FORMAT.as_bytes().to_vec();
302        signing.push(0);
303        signing.extend(&canonical);
304        let signed = SignedRecord {
305            format: GENESIS_FORMAT.into(),
306            canonical_record: canonical,
307            signatures: vec![RecordSignature {
308                public_key: signer.public_key().to_vec(),
309                signature: signer.sign(&signing).expect("origin signature"),
310            }],
311        };
312        let thread = ThreadRef {
313            spool: Some(SpoolRef {
314                id: genesis.spool.clone(),
315            }),
316            id: Some(ThreadId {
317                value: genesis.id().expect("Thread ID").as_bytes().to_vec(),
318            }),
319        };
320        let local = EndpointRef {
321            public_key: vec![7; 32],
322            kind: EndpointKind::Weft as i32,
323        };
324        let open = ReplicationOpen {
325            thread: Some(thread.clone()),
326            thread_genesis: Some(ThreadGenesisRecord {
327                native_genesis_authority: None,
328                boundary_acceptances: Vec::new(),
329                ownership_claims: vec![],
330                ownership_claim_admissions: vec![],
331                ownership_resolutions: vec![],
332                ownership_resolution_admissions: vec![],
333                genesis: Some(signed),
334                creator_authority: vec![],
335                admission: None,
336            }),
337            source: Some(EndpointRef {
338                public_key: vec![8; 32],
339                kind: EndpointKind::Device as i32,
340            }),
341            destination: Some(local.clone()),
342            facets: vec![SharedFacet::Source as i32],
343            session_nonce: vec![9; 16],
344            record_formats: vec![OPERATION_FORMAT.into()],
345            ..Default::default()
346        };
347        let facets = BTreeSet::from([ThreadFacet::Source]);
348        assert!(
349            accept(&open, &thread, &local, [8; 32], &facets, vec![]).is_ok(),
350            "first publication must accept the original signed genesis"
351        );
352        let mut changed = open.clone();
353        changed
354            .thread_genesis
355            .as_mut()
356            .expect("genesis")
357            .genesis
358            .as_mut()
359            .expect("signed genesis")
360            .signatures[0]
361            .signature[0] ^= 1;
362        assert!(accept(&changed, &thread, &local, [8; 32], &facets, vec![]).is_err());
363        changed = open.clone();
364        changed
365            .thread
366            .as_mut()
367            .expect("Thread")
368            .id
369            .as_mut()
370            .expect("ID")
371            .value[0] ^= 1;
372        let claimed = changed.thread.clone().expect("claimed Thread");
373        assert!(accept(&changed, &claimed, &local, [8; 32], &facets, vec![]).is_err());
374        changed = open;
375        changed
376            .thread
377            .as_mut()
378            .expect("Thread")
379            .spool
380            .as_mut()
381            .expect("spool")
382            .id = "another".into();
383        let claimed = changed.thread.clone().expect("claimed Thread");
384        assert!(accept(&changed, &claimed, &local, [8; 32], &facets, vec![]).is_err());
385    }
386
387    #[test]
388    fn opening_binds_both_endpoints_thread_formats_and_negotiated_limits() {
389        let thread = ThreadRef {
390            spool: Some(SpoolRef { id: "spool".into() }),
391            id: Some(ThreadId { value: vec![3; 32] }),
392        };
393        let local = EndpointRef {
394            kind: EndpointKind::Weft as i32,
395            public_key: vec![1; 32],
396        };
397        let source = EndpointRef {
398            kind: EndpointKind::Device as i32,
399            public_key: vec![2; 32],
400        };
401        let facets = BTreeSet::from([ThreadFacet::Source, ThreadFacet::Discussion]);
402        let open = ReplicationOpen {
403            thread: Some(thread.clone()),
404            source: Some(source),
405            destination: Some(local.clone()),
406            facets: facets
407                .iter()
408                .copied()
409                .map(super::super::wire_facet)
410                .collect(),
411            session_nonce: vec![4; 16],
412            record_formats: vec![OPERATION_FORMAT.into()],
413            budget: Some(ReadBudget {
414                max_items: 1,
415                max_frame_bytes: FRAME_LIMIT as u32,
416                max_snapshot_bytes: 0,
417            }),
418            ..Default::default()
419        };
420        let allowed = BTreeSet::from([ThreadFacet::Source]);
421        let ready = accept(&open, &thread, &local, [2; 32], &allowed, vec![5; 32])
422            .expect("authorized opening")
423            .ready;
424        assert_eq!(
425            validate_ready(&ready, &thread, &local, &facets, 1).expect("bound ready"),
426            (allowed.clone(), 1)
427        );
428        assert_eq!(ready.sharing_policy_version, vec![5; 32]);
429        assert!(accept(&open, &thread, &local, [7; 32], &allowed, vec![]).is_err());
430        let mut changed = open.clone();
431        changed
432            .destination
433            .as_mut()
434            .expect("destination")
435            .public_key = vec![8; 32];
436        assert!(accept(&changed, &thread, &local, [2; 32], &allowed, vec![]).is_err());
437        changed = open.clone();
438        changed
439            .thread
440            .as_mut()
441            .expect("thread")
442            .id
443            .as_mut()
444            .expect("id")
445            .value = vec![8; 32];
446        assert!(accept(&changed, &thread, &local, [2; 32], &allowed, vec![]).is_err());
447        changed = open.clone();
448        changed.facets = vec![SharedFacet::Source as i32; 2];
449        assert!(accept(&changed, &thread, &local, [2; 32], &allowed, vec![]).is_err());
450        changed = open.clone();
451        changed.record_formats.clear();
452        assert!(accept(&changed, &thread, &local, [2; 32], &allowed, vec![]).is_err());
453        changed = open.clone();
454        changed.budget.as_mut().expect("budget").max_frame_bytes = 128;
455        assert!(accept(&changed, &thread, &local, [2; 32], &allowed, vec![]).is_err());
456        let mut widened = ready.clone();
457        widened.budget.as_mut().expect("budget").max_items = 2;
458        assert!(validate_ready(&widened, &thread, &local, &facets, 1).is_err());
459        widened = ready;
460        widened.facets.push(SharedFacet::Collaboration as i32);
461        assert!(validate_ready(&widened, &thread, &local, &allowed, 1).is_err());
462    }
463
464    #[test]
465    fn hybrid_opening_and_ready_fields_are_rejected_before_admission() {
466        use api::heddle::api::common::ProtocolCompatibility;
467        let thread = ThreadRef {
468            spool: Some(SpoolRef { id: "spool".into() }),
469            id: Some(ThreadId { value: vec![3; 32] }),
470        };
471        let local = EndpointRef {
472            kind: EndpointKind::Weft as i32,
473            public_key: vec![1; 32],
474        };
475        let facets = BTreeSet::from([ThreadFacet::Source]);
476        let open = ReplicationOpen {
477            thread: Some(thread.clone()),
478            source: Some(EndpointRef {
479                kind: EndpointKind::Device as i32,
480                public_key: vec![2; 32],
481            }),
482            destination: Some(local.clone()),
483            facets: vec![SharedFacet::Source as i32],
484            session_nonce: vec![4; 16],
485            record_formats: vec![OPERATION_FORMAT.into()],
486            budget: Some(ReadBudget {
487                max_items: 1,
488                max_frame_bytes: FRAME_LIMIT as u32,
489                max_snapshot_bytes: 0,
490            }),
491            ..Default::default()
492        };
493        let ready = accept(&open, &thread, &local, [2; 32], &facets, vec![5; 32])
494            .expect("a non-HYBRID opening is accepted")
495            .ready;
496        validate_ready(&ready, &thread, &local, &facets, 1).expect("non-HYBRID ready");
497        let hybrid = |result: Result<(), Error>| {
498            assert!(
499                matches!(result, Err(Error::Protocol(message)) if message.contains("api#307")),
500                "HYBRID fields must be refused, never ignored"
501            );
502        };
503        let mut changed = open.clone();
504        changed.import_authority = Some(ImportPublicProofBundleV1::default());
505        hybrid(accept(&changed, &thread, &local, [2; 32], &facets, vec![]).map(|_| ()));
506        changed = open;
507        changed.protocol = Some(ProtocolCompatibility::default());
508        hybrid(accept(&changed, &thread, &local, [2; 32], &facets, vec![]).map(|_| ()));
509        let mut widened = ready.clone();
510        widened.import_authority = Some(ImportPublicProofBundleV1::default());
511        hybrid(validate_ready(&widened, &thread, &local, &facets, 1).map(|_| ()));
512        widened = ready;
513        widened.protocol = Some(ProtocolCompatibility::default());
514        hybrid(validate_ready(&widened, &thread, &local, &facets, 1).map(|_| ()));
515    }
516}