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        ready: ReplicationReady {
107            thread: Some(thread.clone()),
108            endpoint: Some(local.clone()),
109            facets: facets.into_iter().map(super::wire_facet).collect(),
110            sharing_policy_version,
111            budget: Some(ReadBudget {
112                max_items: if requested == 0 {
113                    MAX_ITEMS
114                } else {
115                    requested.min(MAX_ITEMS)
116                },
117                max_frame_bytes: FRAME_LIMIT as u32,
118                max_snapshot_bytes: 0,
119            }),
120            record_formats: vec![OPERATION_FORMAT.into()],
121            // No HYBRID import-authority support is claimed; no proof bundle.
122            protocol: None,
123            import_authority: None,
124        },
125    })
126}
127
128/// Parsed proposal; authorization and durable installation belong to the host.
129pub struct AcceptedOpening {
130    pub ready: ReplicationReady,
131    pub genesis: Option<ThreadGenesis>,
132    pub genesis_record: Option<ThreadGenesisRecord>,
133}
134
135/// Structural and original-signature validation only. Account authority and
136/// hosted admission receipts still require independently retained trust.
137pub fn verify_genesis_record(
138    record: &ThreadGenesisRecord,
139    thread: &ThreadRef,
140) -> Result<ThreadGenesis, Error> {
141    use prost::Message;
142    if record.boundary_acceptances.len() > crate::boundary_acceptance::MAX_ACCEPTANCES
143        || record.encoded_len() > 256 * 1024
144    {
145        return Err(Error::Protocol("genesis wrapper evidence exceeds bounds"));
146    }
147    let signed = record
148        .genesis
149        .as_ref()
150        .ok_or(Error::Protocol("original signed genesis missing"))?;
151    let genesis = verify_genesis(signed, thread)?;
152    if record.creator_authority.len() > 64 * 1024 {
153        return Err(Error::Protocol("creator authority exceeds bound"));
154    }
155    use heddle_object_model::object::thread_replication::GenesisOwner;
156    match genesis.owner {
157        GenesisOwner::LocalKey(_)
158            if !record.creator_authority.is_empty() || record.admission.is_some() =>
159        {
160            return Err(Error::Protocol(
161                "local-key ownership requires an explicit claim, not an account envelope",
162            ));
163        }
164        GenesisOwner::Account(_) if record.creator_authority.is_empty() => {
165            return Err(Error::Protocol(
166                "account-owned genesis requires original creator authority",
167            ));
168        }
169        _ => {}
170    }
171    super::ownership::verify_claims(record, &genesis)?;
172    Ok(genesis)
173}
174
175pub fn sign_genesis(genesis: &ThreadGenesis, signer: &impl Signer) -> Result<SignedRecord, Error> {
176    let signed =
177        SignedGenesis::sign(genesis, signer).map_err(|error| Error::Io(error.to_string()))?;
178    Ok(SignedRecord {
179        format: GENESIS_FORMAT.into(),
180        canonical_record: signed.canonical,
181        signatures: vec![RecordSignature {
182            public_key: genesis.creator.to_vec(),
183            signature: signed.signature,
184        }],
185    })
186}
187
188pub fn verify_genesis(record: &SignedRecord, thread: &ThreadRef) -> Result<ThreadGenesis, Error> {
189    if record.format != GENESIS_FORMAT || record.signatures.len() != 1 {
190        return Err(Error::Protocol("unsupported Thread genesis record"));
191    }
192    let genesis = SignedGenesis {
193        canonical: record.canonical_record.clone(),
194        signature: record.signatures[0].signature.clone(),
195    }
196    .verify()
197    .map_err(|_| Error::Protocol("invalid Thread genesis signature"))?;
198    let id = genesis
199        .id()
200        .map_err(|_| Error::Protocol("invalid Thread genesis"))?;
201    if record.signatures[0].public_key != genesis.creator
202        || thread
203            .spool
204            .as_ref()
205            .is_none_or(|spool| spool.id != genesis.spool)
206        || thread
207            .id
208            .as_ref()
209            .is_none_or(|thread| thread.value != id.as_bytes())
210    {
211        return Err(Error::Protocol(
212            "Thread genesis identity does not match opening",
213        ));
214    }
215    Ok(genesis)
216}
217
218pub fn validate_ready(
219    ready: &ReplicationReady,
220    thread: &ThreadRef,
221    destination: &EndpointRef,
222    requested_facets: &BTreeSet<ThreadFacet>,
223    requested_max_items: u32,
224) -> Result<(BTreeSet<ThreadFacet>, usize), Error> {
225    crate::hybrid::replication_ready(ready).map_err(Error::Protocol)?;
226    validate_endpoint(destination)?;
227    if ready.endpoint.as_ref() != Some(destination)
228        || ready.thread.as_ref() != Some(thread)
229        || ready.record_formats != [OPERATION_FORMAT]
230    {
231        return Err(Error::Protocol(
232            "replication Ready binding differs from opening",
233        ));
234    }
235    let facets = parse_facets(&ready.facets)?;
236    if !facets.is_subset(requested_facets) {
237        return Err(Error::Protocol("replication Ready widened admission scope"));
238    }
239    let budget = ready
240        .budget
241        .as_ref()
242        .ok_or(Error::Protocol("replication Ready requires budget"))?;
243    let ceiling = if requested_max_items == 0 {
244        MAX_ITEMS
245    } else {
246        requested_max_items.min(MAX_ITEMS)
247    };
248    if budget.max_items == 0
249        || budget.max_items > ceiling
250        || budget.max_frame_bytes != FRAME_LIMIT as u32
251    {
252        return Err(Error::Protocol("unsupported replication Ready budget"));
253    }
254    Ok((facets, budget.max_items as usize))
255}
256
257#[cfg(test)]
258mod tests {
259    use super::*;
260
261    #[test]
262    fn first_publication_retains_signed_genesis_and_rejects_changed_identity() {
263        use crypto::{Ed25519Signer, Signer};
264        use heddle_object_model::object::{
265            StateId,
266            thread_replication::{GENESIS_FORMAT, ThreadGenesis},
267        };
268        let signer = Ed25519Signer::from_seed(&[23; 32]).expect("origin device");
269        let genesis = ThreadGenesis {
270            version: 1,
271            spool: "01980000-0000-7000-8000-000000000001".into(),
272            parent: None,
273            base: StateId::from_bytes([1; 32]),
274            name: "original".into(),
275            intent: "publish once".into(),
276            creator: signer.public_key().try_into().expect("public key"),
277            owner: heddle_object_model::object::thread_replication::GenesisOwner::LocalKey(
278                signer.public_key().try_into().expect("public key"),
279            ),
280            nonce: vec![2; 16],
281        };
282        let canonical = genesis.encode().expect("canonical genesis");
283        let mut signing = GENESIS_FORMAT.as_bytes().to_vec();
284        signing.push(0);
285        signing.extend(&canonical);
286        let signed = SignedRecord {
287            format: GENESIS_FORMAT.into(),
288            canonical_record: canonical,
289            signatures: vec![RecordSignature {
290                public_key: signer.public_key().to_vec(),
291                signature: signer.sign(&signing).expect("origin signature"),
292            }],
293        };
294        let thread = ThreadRef {
295            spool: Some(SpoolRef {
296                id: genesis.spool.clone(),
297            }),
298            id: Some(ThreadId {
299                value: genesis.id().expect("Thread ID").as_bytes().to_vec(),
300            }),
301        };
302        let local = EndpointRef {
303            public_key: vec![7; 32],
304            kind: EndpointKind::Weft as i32,
305        };
306        let open = ReplicationOpen {
307            thread: Some(thread.clone()),
308            thread_genesis: Some(ThreadGenesisRecord {
309                boundary_acceptances: Vec::new(),
310                ownership_claims: vec![],
311                ownership_claim_admissions: vec![],
312                ownership_resolutions: vec![],
313                ownership_resolution_admissions: vec![],
314                genesis: Some(signed),
315                creator_authority: vec![],
316                admission: None,
317            }),
318            source: Some(EndpointRef {
319                public_key: vec![8; 32],
320                kind: EndpointKind::Device as i32,
321            }),
322            destination: Some(local.clone()),
323            facets: vec![SharedFacet::Source as i32],
324            session_nonce: vec![9; 16],
325            record_formats: vec![OPERATION_FORMAT.into()],
326            ..Default::default()
327        };
328        let facets = BTreeSet::from([ThreadFacet::Source]);
329        assert!(
330            accept(&open, &thread, &local, [8; 32], &facets, vec![]).is_ok(),
331            "first publication must accept the original signed genesis"
332        );
333        let mut changed = open.clone();
334        changed
335            .thread_genesis
336            .as_mut()
337            .expect("genesis")
338            .genesis
339            .as_mut()
340            .expect("signed genesis")
341            .signatures[0]
342            .signature[0] ^= 1;
343        assert!(accept(&changed, &thread, &local, [8; 32], &facets, vec![]).is_err());
344        changed = open.clone();
345        changed
346            .thread
347            .as_mut()
348            .expect("Thread")
349            .id
350            .as_mut()
351            .expect("ID")
352            .value[0] ^= 1;
353        let claimed = changed.thread.clone().expect("claimed Thread");
354        assert!(accept(&changed, &claimed, &local, [8; 32], &facets, vec![]).is_err());
355        changed = open;
356        changed
357            .thread
358            .as_mut()
359            .expect("Thread")
360            .spool
361            .as_mut()
362            .expect("spool")
363            .id = "another".into();
364        let claimed = changed.thread.clone().expect("claimed Thread");
365        assert!(accept(&changed, &claimed, &local, [8; 32], &facets, vec![]).is_err());
366    }
367
368    #[test]
369    fn opening_binds_both_endpoints_thread_formats_and_negotiated_limits() {
370        let thread = ThreadRef {
371            spool: Some(SpoolRef { id: "spool".into() }),
372            id: Some(ThreadId { value: vec![3; 32] }),
373        };
374        let local = EndpointRef {
375            kind: EndpointKind::Weft as i32,
376            public_key: vec![1; 32],
377        };
378        let source = EndpointRef {
379            kind: EndpointKind::Device as i32,
380            public_key: vec![2; 32],
381        };
382        let facets = BTreeSet::from([ThreadFacet::Source, ThreadFacet::Discussion]);
383        let open = ReplicationOpen {
384            thread: Some(thread.clone()),
385            source: Some(source),
386            destination: Some(local.clone()),
387            facets: facets
388                .iter()
389                .copied()
390                .map(super::super::wire_facet)
391                .collect(),
392            session_nonce: vec![4; 16],
393            record_formats: vec![OPERATION_FORMAT.into()],
394            budget: Some(ReadBudget {
395                max_items: 1,
396                max_frame_bytes: FRAME_LIMIT as u32,
397                max_snapshot_bytes: 0,
398            }),
399            ..Default::default()
400        };
401        let allowed = BTreeSet::from([ThreadFacet::Source]);
402        let ready = accept(&open, &thread, &local, [2; 32], &allowed, vec![5; 32])
403            .expect("authorized opening")
404            .ready;
405        assert_eq!(
406            validate_ready(&ready, &thread, &local, &facets, 1).expect("bound ready"),
407            (allowed.clone(), 1)
408        );
409        assert_eq!(ready.sharing_policy_version, vec![5; 32]);
410        assert!(accept(&open, &thread, &local, [7; 32], &allowed, vec![]).is_err());
411        let mut changed = open.clone();
412        changed
413            .destination
414            .as_mut()
415            .expect("destination")
416            .public_key = vec![8; 32];
417        assert!(accept(&changed, &thread, &local, [2; 32], &allowed, vec![]).is_err());
418        changed = open.clone();
419        changed
420            .thread
421            .as_mut()
422            .expect("thread")
423            .id
424            .as_mut()
425            .expect("id")
426            .value = vec![8; 32];
427        assert!(accept(&changed, &thread, &local, [2; 32], &allowed, vec![]).is_err());
428        changed = open.clone();
429        changed.facets = vec![SharedFacet::Source as i32; 2];
430        assert!(accept(&changed, &thread, &local, [2; 32], &allowed, vec![]).is_err());
431        changed = open.clone();
432        changed.record_formats.clear();
433        assert!(accept(&changed, &thread, &local, [2; 32], &allowed, vec![]).is_err());
434        changed = open.clone();
435        changed.budget.as_mut().expect("budget").max_frame_bytes = 128;
436        assert!(accept(&changed, &thread, &local, [2; 32], &allowed, vec![]).is_err());
437        let mut widened = ready.clone();
438        widened.budget.as_mut().expect("budget").max_items = 2;
439        assert!(validate_ready(&widened, &thread, &local, &facets, 1).is_err());
440        widened = ready;
441        widened.facets.push(SharedFacet::Collaboration as i32);
442        assert!(validate_ready(&widened, &thread, &local, &allowed, 1).is_err());
443    }
444
445    #[test]
446    fn hybrid_opening_and_ready_fields_are_rejected_before_admission() {
447        use api::heddle::api::common::ProtocolCompatibility;
448        let thread = ThreadRef {
449            spool: Some(SpoolRef { id: "spool".into() }),
450            id: Some(ThreadId { value: vec![3; 32] }),
451        };
452        let local = EndpointRef {
453            kind: EndpointKind::Weft as i32,
454            public_key: vec![1; 32],
455        };
456        let facets = BTreeSet::from([ThreadFacet::Source]);
457        let open = ReplicationOpen {
458            thread: Some(thread.clone()),
459            source: Some(EndpointRef {
460                kind: EndpointKind::Device as i32,
461                public_key: vec![2; 32],
462            }),
463            destination: Some(local.clone()),
464            facets: vec![SharedFacet::Source as i32],
465            session_nonce: vec![4; 16],
466            record_formats: vec![OPERATION_FORMAT.into()],
467            budget: Some(ReadBudget {
468                max_items: 1,
469                max_frame_bytes: FRAME_LIMIT as u32,
470                max_snapshot_bytes: 0,
471            }),
472            ..Default::default()
473        };
474        let ready = accept(&open, &thread, &local, [2; 32], &facets, vec![5; 32])
475            .expect("a non-HYBRID opening is accepted")
476            .ready;
477        validate_ready(&ready, &thread, &local, &facets, 1).expect("non-HYBRID ready");
478        let hybrid = |result: Result<(), Error>| {
479            assert!(
480                matches!(result, Err(Error::Protocol(message)) if message.contains("api#307")),
481                "HYBRID fields must be refused, never ignored"
482            );
483        };
484        let mut changed = open.clone();
485        changed.import_authority = Some(ImportPublicProofBundleV1::default());
486        hybrid(accept(&changed, &thread, &local, [2; 32], &facets, vec![]).map(|_| ()));
487        changed = open;
488        changed.protocol = Some(ProtocolCompatibility::default());
489        hybrid(accept(&changed, &thread, &local, [2; 32], &facets, vec![]).map(|_| ()));
490        let mut widened = ready.clone();
491        widened.import_authority = Some(ImportPublicProofBundleV1::default());
492        hybrid(validate_ready(&widened, &thread, &local, &facets, 1).map(|_| ()));
493        widened = ready;
494        widened.protocol = Some(ProtocolCompatibility::default());
495        hybrid(validate_ready(&widened, &thread, &local, &facets, 1).map(|_| ()));
496    }
497}