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 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 protocol: None,
123 import_authority: None,
124 },
125 })
126}
127
128pub struct AcceptedOpening {
130 pub ready: ReplicationReady,
131 pub genesis: Option<ThreadGenesis>,
132 pub genesis_record: Option<ThreadGenesisRecord>,
133}
134
135pub 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}