1use std::collections::BTreeSet;
5
6use crypto::thread_operation::SignedOperation;
7use heddle_object_model::object::{
8 ContentHash,
9 thread_replication::{Admission, OPERATION_FORMAT, ThreadFacet},
10};
11#[cfg(feature = "native")]
12pub mod native;
13pub mod opening;
14pub mod ownership;
15pub mod store;
16use store::ReplicaStore;
17
18use crate::contract::*;
19
20#[derive(Debug, thiserror::Error)]
21pub enum Error {
22 #[error(transparent)]
23 Operation(#[from] heddle_object_model::error::HeddleError),
24 #[error(transparent)]
25 Signature(#[from] crypto::thread_operation::Error),
26 #[error("replication protocol: {0}")]
27 Protocol(&'static str),
28}
29pub type Result<T> = std::result::Result<T, Error>;
30
31#[derive(Debug, thiserror::Error)]
32pub enum StoreError<E: std::error::Error + 'static> {
33 #[error("replica store: {0}")]
34 Store(#[source] E),
35 #[error(transparent)]
36 Protocol(#[from] Error),
37}
38pub type StoreResult<T, E> = std::result::Result<T, StoreError<E>>;
39
40#[allow(clippy::large_enum_variant)] pub enum Frame {
43 Have(ReplicationHave),
44 Need(ReplicationNeed),
45 Operations(ReplicationOperations),
46 Receipt(ReplicationReceipt),
47}
48#[allow(clippy::large_enum_variant)] pub(crate) enum InputUnit {
52 Frame(Frame),
53 Operation(store::ReceivedOperation),
54}
55impl Frame {
56 pub fn request(self) -> ReplicateThreadRequest {
57 use replicate_thread_request::Body;
58 ReplicateThreadRequest {
59 body: Some(match self {
60 Self::Have(v) => Body::Have(v),
61 Self::Need(v) => Body::Need(v),
62 Self::Operations(v) => Body::Operations(v),
63 Self::Receipt(v) => Body::Receipt(v),
64 }),
65 }
66 }
67 pub fn response(self) -> ReplicateThreadResponse {
68 use replicate_thread_response::Body;
69 ReplicateThreadResponse {
70 body: Some(match self {
71 Self::Have(v) => Body::Have(v),
72 Self::Need(v) => Body::Need(v),
73 Self::Operations(v) => Body::Operations(v),
74 Self::Receipt(v) => Body::Receipt(v),
75 }),
76 }
77 }
78 pub fn from_request(request: ReplicateThreadRequest) -> Result<Self> {
79 use replicate_thread_request::Body;
80 Ok(match request.body {
81 Some(Body::Have(v)) => Self::Have(v),
82 Some(Body::Need(v)) => Self::Need(v),
83 Some(Body::Operations(v)) => {
84 crate::hybrid::operations(&v).map_err(Error::Protocol)?;
85 Self::Operations(v)
86 }
87 Some(Body::Receipt(v)) => Self::Receipt(v),
88 _ => return Err(Error::Protocol("unexpected replication opening")),
89 })
90 }
91 pub fn from_response(response: ReplicateThreadResponse) -> Result<Self> {
92 use replicate_thread_response::Body;
93 Ok(match response.body {
94 Some(Body::Have(v)) => Self::Have(v),
95 Some(Body::Need(v)) => Self::Need(v),
96 Some(Body::Operations(v)) => {
97 crate::hybrid::operations(&v).map_err(Error::Protocol)?;
98 Self::Operations(v)
99 }
100 Some(Body::Receipt(v)) => Self::Receipt(v),
101 _ => return Err(Error::Protocol("unexpected replication ready")),
102 })
103 }
104}
105
106#[allow(clippy::large_enum_variant)] pub enum Outbound {
108 Frame(Frame),
109 Operation(ContentHash),
110}
111
112#[derive(Clone)]
113pub struct Session<B: ReplicaStore> {
114 pub replica: B,
115 destination: [u8; 32],
116 facets: BTreeSet<ThreadFacet>,
117 max_items: usize,
118 in_flight: BTreeSet<ContentHash>,
119 generation: i64,
120 announce_generation: i64,
121 announce_facet: usize,
122 after: Option<ContentHash>,
123 pending_input_bookkeeping: Option<(ThreadFacet, ContentHash)>,
124 protocol: Option<api::heddle::api::common::ProtocolCompatibility>,
125}
126impl<B: ReplicaStore> Session<B> {
127 pub fn new(
130 replica: B,
131 destination: [u8; 32],
132 facets: BTreeSet<ThreadFacet>,
133 max_items: usize,
134 ) -> Result<Self> {
135 if max_items == 0 || max_items > 64 {
136 return Err(Error::Protocol("replication batch must be 1..64"));
137 }
138 Ok(Self {
139 replica,
140 destination,
141 facets,
142 max_items,
143 in_flight: BTreeSet::new(),
144 generation: -1,
145 announce_generation: -1,
146 announce_facet: 0,
147 after: None,
148 pending_input_bookkeeping: None,
149 protocol: None,
150 })
151 }
152 pub fn with_protocol(
155 mut self,
156 open: Option<&api::heddle::api::common::ProtocolCompatibility>,
157 ready: Option<&api::heddle::api::common::ProtocolCompatibility>,
158 ) -> Result<Self> {
159 crate::hybrid::negotiated(open, ready).map_err(Error::Protocol)?;
160 self.protocol = open.cloned();
161 Ok(self)
162 }
163 fn require_bundle_protocol(&self, batch: &ReplicationOperations) -> Result<()> {
164 if batch.import_authority.is_some() || batch.native_authority.is_some() {
165 api::import_authority::require_hybrid_peer(self.protocol.as_ref())
166 .map_err(|_| Error::Protocol("HYBRID operations require negotiated protocol"))?;
167 }
168 Ok(())
169 }
170 pub async fn export_facets(&self) -> StoreResult<BTreeSet<ThreadFacet>, B::Error> {
171 let sharing = self
172 .replica
173 .sharing(self.destination)
174 .await
175 .map_err(StoreError::Store)?;
176 Ok(sharing.intersection(&self.facets).copied().collect())
177 }
178 pub async fn announcement(&mut self) -> StoreResult<Option<Frame>, B::Error> {
181 let current = self.replica.generation().await.map_err(StoreError::Store)?;
182 if self.generation == current && self.announce_generation < 0 {
183 return Ok(None);
184 }
185 if self.announce_generation < 0 {
186 self.announce_generation = current;
187 self.announce_facet = 0;
188 self.after = None;
189 }
190 let sharing = self.export_facets().await?;
191 let facets = ThreadFacet::ALL;
192 while self.announce_facet < facets.len() {
193 let facet = facets[self.announce_facet];
194 if !sharing.contains(&facet) {
195 self.announce_facet += 1;
196 self.after = None;
197 continue;
198 }
199 let page = self
200 .replica
201 .frontier_page(facet, self.after, self.max_items)
202 .await
203 .map_err(StoreError::Store)?;
204 if page.is_empty() {
205 self.announce_facet += 1;
206 self.after = None;
207 continue;
208 }
209 self.after = page.last().copied();
210 return Ok(Some(Frame::Have(ReplicationHave {
211 frontiers: vec![CausalFrontier {
212 facet: wire_facet(facet),
213 heads: page.into_iter().map(|id| id.as_bytes().to_vec()).collect(),
214 }],
215 })));
216 }
217 self.generation = self.announce_generation;
218 self.announce_generation = -1;
219 Ok(
223 (self.replica.generation().await.map_err(StoreError::Store)? != self.generation)
224 .then(|| Frame::Have(ReplicationHave::default())),
225 )
226 }
227 pub async fn handle(&mut self, frame: Frame) -> StoreResult<Vec<Outbound>, B::Error> {
228 self.handle_frame(frame, true).await
229 }
230 pub(crate) fn input_units(&self, frame: Frame) -> StoreResult<Vec<InputUnit>, B::Error> {
233 let Frame::Operations(batch) = frame else {
234 return Ok(vec![InputUnit::Frame(frame)]);
235 };
236 self.require_bundle_protocol(&batch)?;
237 self.check_count(batch.operations.len())?;
238 self.check_count(batch.authority_admissions.len())?;
239 self.check_count(batch.boundary_acceptances.len())?;
240 let originals = crate::authority_admission::match_batch(&batch)
241 .map_err(|_| Error::Protocol("invalid original authority batch"))?;
242 let mut units = Vec::with_capacity(originals.len());
243 for received in originals {
244 let operation = received.original.verify().map_err(Error::from)?;
245 if !self.facets.contains(&operation.facet()) {
246 return Err(Error::Protocol("operation is outside admission scope").into());
247 }
248 units.push(InputUnit::Operation(received));
249 }
250 Ok(units)
251 }
252 pub(crate) async fn handle_unit(
253 &mut self,
254 unit: InputUnit,
255 ) -> StoreResult<Vec<Outbound>, B::Error> {
256 match unit {
257 InputUnit::Frame(frame) => self.handle_input(frame).await,
258 InputUnit::Operation(received) => Ok(vec![Outbound::Frame(Frame::Receipt(
259 self.commit_original(received, false).await?,
260 ))]),
261 }
262 }
263 async fn commit_original(
264 &mut self,
265 received: store::ReceivedOperation,
266 maintenance: bool,
267 ) -> StoreResult<ReplicationReceipt, B::Error> {
268 if self.pending_input_bookkeeping.is_some() {
269 return Err(Error::Protocol("prior input bookkeeping not completed").into());
270 }
271 let operation = received.original.verify().map_err(Error::from)?;
272 if !self.facets.contains(&operation.facet()) {
273 return Err(Error::Protocol("operation is outside admission scope").into());
274 }
275 let id = operation.id().map_err(Error::from)?;
276 self.in_flight.remove(&id);
277 let mut receipt = ReplicationReceipt::default();
278 match self
279 .replica
280 .receive(received)
281 .await
282 .map_err(StoreError::Store)?
283 {
284 Admission::Accepted => receipt.accepted_operation_ids.push(id.as_bytes().to_vec()),
285 Admission::Pending => {
286 receipt.pending_operation_ids.push(id.as_bytes().to_vec());
287 if maintenance {
288 self.replica
289 .remember_peer_heads(self.destination, vec![(operation.facet(), id)])
290 .await
291 .map_err(StoreError::Store)?;
292 } else {
293 self.pending_input_bookkeeping = Some((operation.facet(), id));
294 }
295 }
296 Admission::Rejected(message) => receipt.rejected.push(rejection(id, message)),
297 }
298 Ok(receipt)
299 }
300 pub(crate) fn has_input_bookkeeping(&self) -> bool {
301 self.pending_input_bookkeeping.is_some()
302 }
303 pub(crate) async fn finish_input_bookkeeping(&mut self) -> StoreResult<(), B::Error> {
304 if let Some(head) = self.pending_input_bookkeeping.take() {
305 self.replica
306 .remember_peer_heads(self.destination, vec![head])
307 .await
308 .map_err(StoreError::Store)?;
309 }
310 Ok(())
311 }
312 pub(crate) async fn handle_input(
314 &mut self,
315 frame: Frame,
316 ) -> StoreResult<Vec<Outbound>, B::Error> {
317 self.handle_frame(frame, false).await
318 }
319 async fn handle_frame(
320 &mut self,
321 frame: Frame,
322 maintenance: bool,
323 ) -> StoreResult<Vec<Outbound>, B::Error> {
324 let mut responses = Vec::new();
325 match frame {
326 Frame::Have(have) => {
327 let count: usize = have.frontiers.iter().map(|f| f.heads.len()).sum();
328 self.check_count(count)?;
329 self.check_count(have.frontiers.len())?;
330 let mut heads = Vec::new();
331 let mut receipt = ReplicationReceipt::default();
332 for frontier in have.frontiers {
333 let facet = native_facet(frontier.facet)?;
334 if !self.facets.contains(&facet) {
335 return Err(Error::Protocol("unnegotiated facet").into());
336 }
337 for bytes in frontier.heads {
338 let id = hash(&bytes)?;
339 if let Some((record, status)) = self
340 .replica
341 .operation(id)
342 .await
343 .map_err(StoreError::Store)?
344 {
345 if record.original.verify().map_err(Error::from)?.facet() != facet {
346 return Err(Error::Protocol(
347 "advertised frontier has the wrong facet",
348 )
349 .into());
350 }
351 match status {
352 Admission::Accepted => receipt.accepted_operation_ids.push(bytes),
353 Admission::Rejected(message) => {
354 receipt.rejected.push(rejection(id, message))
355 }
356 Admission::Pending => heads.push((facet, id)),
357 }
358 } else {
359 heads.push((facet, id));
360 }
361 }
362 }
363 self.replica
364 .remember_peer_heads(self.destination, heads)
365 .await
366 .map_err(StoreError::Store)?;
367 if !receipt.accepted_operation_ids.is_empty() || !receipt.rejected.is_empty() {
368 responses.push(Outbound::Frame(Frame::Receipt(receipt)));
369 }
370 }
371 Frame::Need(need) => {
372 self.check_count(need.operation_ids.len())?;
373 let sharing = self.export_facets().await?;
374 for bytes in need.operation_ids {
375 let id = hash(&bytes)?;
376 let Some((record, Admission::Accepted)) = self
377 .replica
378 .operation(id)
379 .await
380 .map_err(StoreError::Store)?
381 else {
382 return Err(Error::Protocol("requested operation unavailable").into());
383 };
384 let operation = record.original.verify().map_err(Error::from)?;
385 if !sharing.contains(&operation.facet()) {
386 return Err(
387 Error::Protocol("operation is outside current sharing policy").into(),
388 );
389 }
390 responses.push(Outbound::Operation(id));
391 }
392 }
393 Frame::Operations(batch) => {
394 self.require_bundle_protocol(&batch)?;
395 if self.pending_input_bookkeeping.is_some() {
396 return Err(Error::Protocol("prior input bookkeeping not completed").into());
397 }
398 self.check_count(batch.operations.len())?;
399 self.check_count(batch.authority_admissions.len())?;
400 self.check_count(batch.boundary_acceptances.len())?;
401 let decoded =
402 crate::authority_admission::match_batch(&batch).map_err(
403 |error| match error {
404 crate::transport::Error::Protocol(message) => Error::Protocol(message),
405 _ => Error::Protocol("invalid original authority batch"),
406 },
407 )?;
408 let mut operations = Vec::with_capacity(decoded.len());
409 for received in decoded {
410 let operation = received.original.verify().map_err(Error::from)?;
411 let id = operation.id().map_err(Error::from)?;
412 if !self.facets.contains(&operation.facet()) {
413 return Err(Error::Protocol("operation is outside admission scope").into());
414 }
415 operations.push((received, operation, id));
416 }
417 let mut receipt = ReplicationReceipt::default();
418 for (received, _, _) in operations {
419 let next = self.commit_original(received, maintenance).await?;
420 receipt
421 .accepted_operation_ids
422 .extend(next.accepted_operation_ids);
423 receipt
424 .pending_operation_ids
425 .extend(next.pending_operation_ids);
426 receipt.rejected.extend(next.rejected);
427 }
428 responses.push(Outbound::Frame(Frame::Receipt(receipt)));
429 }
430 Frame::Receipt(receipt) => {
431 self.check_count(
432 receipt.accepted_operation_ids.len()
433 + receipt.pending_operation_ids.len()
434 + receipt.rejected.len(),
435 )?;
436 for bytes in receipt.accepted_operation_ids {
437 self.replica
438 .record_peer_receipt(self.destination, hash(&bytes)?, Admission::Accepted)
439 .await
440 .map_err(StoreError::Store)?;
441 }
442 for bytes in receipt.pending_operation_ids {
443 self.replica
444 .record_peer_receipt(self.destination, hash(&bytes)?, Admission::Pending)
445 .await
446 .map_err(StoreError::Store)?;
447 }
448 for rejected in receipt.rejected {
449 self.replica
450 .record_peer_receipt(
451 self.destination,
452 hash(&rejected.operation_id)?,
453 Admission::Rejected(
454 rejected.failure.map(|f| f.message).unwrap_or_default(),
455 ),
456 )
457 .await
458 .map_err(StoreError::Store)?;
459 }
460 }
461 }
462 if maintenance && let Some(repair) = self.control().await? {
463 responses.push(Outbound::Frame(repair));
464 }
465 Ok(responses)
466 }
467
468 pub async fn control(&mut self) -> StoreResult<Option<Frame>, B::Error> {
471 let settled = self
472 .replica
473 .settled_peer_heads(self.destination, self.facets.clone(), self.max_items)
474 .await
475 .map_err(StoreError::Store)?;
476 if !settled.is_empty() {
477 let mut receipt = ReplicationReceipt::default();
478 for (id, admission) in settled {
479 self.in_flight.remove(&id);
480 match admission {
481 Admission::Accepted => {
482 receipt.accepted_operation_ids.push(id.as_bytes().to_vec())
483 }
484 Admission::Rejected(message) => receipt.rejected.push(rejection(id, message)),
485 Admission::Pending => {}
486 }
487 }
488 return Ok(Some(Frame::Receipt(receipt)));
489 }
490 let available = self.max_items.saturating_sub(self.in_flight.len());
491 if available == 0 {
492 return Ok(None);
493 }
494 let candidates = self
495 .replica
496 .needed_from_peer(self.destination, self.facets.clone(), self.max_items)
497 .await
498 .map_err(StoreError::Store)?;
499 let ids: Vec<_> = candidates
500 .into_iter()
501 .filter(|id| !self.in_flight.contains(id))
502 .take(available)
503 .collect();
504 self.in_flight.extend(ids.iter().copied());
505 Ok((!ids.is_empty()).then(|| {
506 Frame::Need(ReplicationNeed {
507 operation_ids: ids.into_iter().map(|id| id.as_bytes().to_vec()).collect(),
508 })
509 }))
510 }
511
512 pub async fn export_operation(&self, id: ContentHash) -> StoreResult<Frame, B::Error> {
515 let Some((record, Admission::Accepted)) = self
516 .replica
517 .operation(id)
518 .await
519 .map_err(StoreError::Store)?
520 else {
521 return Err(Error::Protocol("requested operation unavailable").into());
522 };
523 if record.import_authority.is_some() || record.native_authority.is_some() {
524 api::import_authority::require_hybrid_peer(self.protocol.as_ref())
525 .map_err(|_| Error::Protocol("HYBRID relay requires negotiated protocol"))?;
526 }
527 let operation = record.original.verify().map_err(Error::from)?;
528 if !self.export_facets().await?.contains(&operation.facet()) {
529 return Err(Error::Protocol("operation is outside current sharing policy").into());
530 }
531 Ok(Frame::Operations(ReplicationOperations {
532 native_authority: record.native_authority.as_deref().cloned(),
533 boundary_acceptances: crate::boundary_acceptance::authority_evidence(
534 record.authority_admission.as_ref(),
535 )
536 .map_err(|_| Error::Protocol("invalid boundary evidence"))?,
537 authority_admissions: record
538 .authority_admission
539 .as_ref()
540 .map(crate::authority_admission::encode)
541 .transpose()
542 .map_err(|_| Error::Protocol("invalid retained authority admission"))?
543 .into_iter()
544 .collect(),
545 operations: vec![SignedRecord {
546 format: OPERATION_FORMAT.into(),
547 canonical_record: record.original.canonical,
548 signatures: vec![RecordSignature {
549 public_key: operation.publisher.to_vec(),
550 signature: record.original.signature,
551 }],
552 }],
553 import_authority: record.import_authority.as_deref().cloned(),
554 }))
555 }
556
557 fn check_count(&self, count: usize) -> Result<()> {
558 if count > self.max_items {
559 Err(Error::Protocol("replication item budget exceeded"))
560 } else {
561 Ok(())
562 }
563 }
564}
565
566pub fn completion_receipt(id: ContentHash, admission: Admission) -> ReplicationReceipt {
569 let mut receipt = ReplicationReceipt::default();
570 match admission {
571 Admission::Accepted => receipt.accepted_operation_ids.push(id.as_bytes().to_vec()),
572 Admission::Pending => receipt.pending_operation_ids.push(id.as_bytes().to_vec()),
573 Admission::Rejected(message) => receipt.rejected.push(rejection(id, message)),
574 }
575 receipt
576}
577
578fn rejection(id: ContentHash, message: String) -> ReplicationRejection {
579 ReplicationRejection {
580 operation_id: id.as_bytes().to_vec(),
581 failure: Some(api::heddle::api::common::CallFailure {
582 code: 9,
583 message,
584 ..Default::default()
585 }),
586 }
587}
588
589pub fn decode_record(record: SignedRecord) -> Result<SignedOperation> {
590 if record.format != OPERATION_FORMAT || record.signatures.len() != 1 {
591 return Err(Error::Protocol("unsupported signed operation format"));
592 }
593 let signature = &record.signatures[0];
594 let signed = SignedOperation {
595 canonical: record.canonical_record,
596 signature: signature.signature.clone(),
597 };
598 if signed.verify()?.publisher.as_slice() != signature.public_key {
599 return Err(Error::Protocol(
600 "record signature key differs from publisher",
601 ));
602 }
603 Ok(signed)
604}
605pub fn native_facet(facet: i32) -> Result<ThreadFacet> {
606 match SharedFacet::try_from(facet) {
607 Ok(SharedFacet::Source) => Ok(ThreadFacet::Source),
608 Ok(SharedFacet::Collaboration) => Ok(ThreadFacet::Discussion),
609 Ok(SharedFacet::Metadata) => Ok(ThreadFacet::Metadata),
610 _ => Err(Error::Protocol("unsupported replication facet")),
611 }
612}
613pub fn wire_facet(facet: ThreadFacet) -> i32 {
614 match facet {
615 ThreadFacet::Source => SharedFacet::Source as i32,
616 ThreadFacet::Discussion => SharedFacet::Collaboration as i32,
617 ThreadFacet::Metadata => SharedFacet::Metadata as i32,
618 }
619}
620fn hash(bytes: &[u8]) -> Result<ContentHash> {
621 Ok(ContentHash::from_bytes(bytes.try_into().map_err(|_| {
622 Error::Protocol("operation ID must be 32 bytes")
623 })?))
624}
625
626#[cfg(all(test, feature = "native"))]
627#[path = "replication_tests.rs"]
628mod tests;
629
630#[cfg(all(test, feature = "native"))]
631#[path = "replication_admission_tests.rs"]
632mod admission_tests;