1use 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
46pub 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 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 protocol: open.protocol.clone(),
127 import_authority: None,
128 },
129 })
130}
131
132pub struct AcceptedOpening {
134 pub ready: ReplicationReady,
135 pub genesis: Option<ThreadGenesis>,
136 pub genesis_record: Option<ThreadGenesisRecord>,
137 pub import_authority: Option<crate::contract::ImportPublicProofBundleV1>,
139 pub native_authority: Option<crate::contract::NativePublicProofBundleV1>,
140}
141
142pub 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}