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}
125impl<B: ReplicaStore> Session<B> {
126 pub fn new(
129 replica: B,
130 destination: [u8; 32],
131 facets: BTreeSet<ThreadFacet>,
132 max_items: usize,
133 ) -> Result<Self> {
134 if max_items == 0 || max_items > 64 {
135 return Err(Error::Protocol("replication batch must be 1..64"));
136 }
137 Ok(Self {
138 replica,
139 destination,
140 facets,
141 max_items,
142 in_flight: BTreeSet::new(),
143 generation: -1,
144 announce_generation: -1,
145 announce_facet: 0,
146 after: None,
147 pending_input_bookkeeping: None,
148 })
149 }
150 pub async fn export_facets(&self) -> StoreResult<BTreeSet<ThreadFacet>, B::Error> {
151 let sharing = self
152 .replica
153 .sharing(self.destination)
154 .await
155 .map_err(StoreError::Store)?;
156 Ok(sharing.intersection(&self.facets).copied().collect())
157 }
158 pub async fn announcement(&mut self) -> StoreResult<Option<Frame>, B::Error> {
161 let current = self.replica.generation().await.map_err(StoreError::Store)?;
162 if self.generation == current && self.announce_generation < 0 {
163 return Ok(None);
164 }
165 if self.announce_generation < 0 {
166 self.announce_generation = current;
167 self.announce_facet = 0;
168 self.after = None;
169 }
170 let sharing = self.export_facets().await?;
171 let facets = ThreadFacet::ALL;
172 while self.announce_facet < facets.len() {
173 let facet = facets[self.announce_facet];
174 if !sharing.contains(&facet) {
175 self.announce_facet += 1;
176 self.after = None;
177 continue;
178 }
179 let page = self
180 .replica
181 .frontier_page(facet, self.after, self.max_items)
182 .await
183 .map_err(StoreError::Store)?;
184 if page.is_empty() {
185 self.announce_facet += 1;
186 self.after = None;
187 continue;
188 }
189 self.after = page.last().copied();
190 return Ok(Some(Frame::Have(ReplicationHave {
191 frontiers: vec![CausalFrontier {
192 facet: wire_facet(facet),
193 heads: page.into_iter().map(|id| id.as_bytes().to_vec()).collect(),
194 }],
195 })));
196 }
197 self.generation = self.announce_generation;
198 self.announce_generation = -1;
199 Ok(
203 (self.replica.generation().await.map_err(StoreError::Store)? != self.generation)
204 .then(|| Frame::Have(ReplicationHave::default())),
205 )
206 }
207 pub async fn handle(&mut self, frame: Frame) -> StoreResult<Vec<Outbound>, B::Error> {
208 self.handle_frame(frame, true).await
209 }
210 pub(crate) fn input_units(&self, frame: Frame) -> StoreResult<Vec<InputUnit>, B::Error> {
213 let Frame::Operations(batch) = frame else {
214 return Ok(vec![InputUnit::Frame(frame)]);
215 };
216 self.check_count(batch.operations.len())?;
217 self.check_count(batch.authority_admissions.len())?;
218 self.check_count(batch.boundary_acceptances.len())?;
219 let originals = crate::authority_admission::match_batch(&batch)
220 .map_err(|_| Error::Protocol("invalid original authority batch"))?;
221 let mut units = Vec::with_capacity(originals.len());
222 for received in originals {
223 let operation = received.original.verify().map_err(Error::from)?;
224 if !self.facets.contains(&operation.facet()) {
225 return Err(Error::Protocol("operation is outside admission scope").into());
226 }
227 units.push(InputUnit::Operation(received));
228 }
229 Ok(units)
230 }
231 pub(crate) async fn handle_unit(
232 &mut self,
233 unit: InputUnit,
234 ) -> StoreResult<Vec<Outbound>, B::Error> {
235 match unit {
236 InputUnit::Frame(frame) => self.handle_input(frame).await,
237 InputUnit::Operation(received) => Ok(vec![Outbound::Frame(Frame::Receipt(
238 self.commit_original(received, false).await?,
239 ))]),
240 }
241 }
242 async fn commit_original(
243 &mut self,
244 received: store::ReceivedOperation,
245 maintenance: bool,
246 ) -> StoreResult<ReplicationReceipt, B::Error> {
247 if self.pending_input_bookkeeping.is_some() {
248 return Err(Error::Protocol("prior input bookkeeping not completed").into());
249 }
250 let operation = received.original.verify().map_err(Error::from)?;
251 if !self.facets.contains(&operation.facet()) {
252 return Err(Error::Protocol("operation is outside admission scope").into());
253 }
254 let id = operation.id().map_err(Error::from)?;
255 self.in_flight.remove(&id);
256 let mut receipt = ReplicationReceipt::default();
257 match self
258 .replica
259 .receive(received)
260 .await
261 .map_err(StoreError::Store)?
262 {
263 Admission::Accepted => receipt.accepted_operation_ids.push(id.as_bytes().to_vec()),
264 Admission::Pending => {
265 receipt.pending_operation_ids.push(id.as_bytes().to_vec());
266 if maintenance {
267 self.replica
268 .remember_peer_heads(self.destination, vec![(operation.facet(), id)])
269 .await
270 .map_err(StoreError::Store)?;
271 } else {
272 self.pending_input_bookkeeping = Some((operation.facet(), id));
273 }
274 }
275 Admission::Rejected(message) => receipt.rejected.push(rejection(id, message)),
276 }
277 Ok(receipt)
278 }
279 pub(crate) fn has_input_bookkeeping(&self) -> bool {
280 self.pending_input_bookkeeping.is_some()
281 }
282 pub(crate) async fn finish_input_bookkeeping(&mut self) -> StoreResult<(), B::Error> {
283 if let Some(head) = self.pending_input_bookkeeping.take() {
284 self.replica
285 .remember_peer_heads(self.destination, vec![head])
286 .await
287 .map_err(StoreError::Store)?;
288 }
289 Ok(())
290 }
291 pub(crate) async fn handle_input(
293 &mut self,
294 frame: Frame,
295 ) -> StoreResult<Vec<Outbound>, B::Error> {
296 self.handle_frame(frame, false).await
297 }
298 async fn handle_frame(
299 &mut self,
300 frame: Frame,
301 maintenance: bool,
302 ) -> StoreResult<Vec<Outbound>, B::Error> {
303 let mut responses = Vec::new();
304 match frame {
305 Frame::Have(have) => {
306 let count: usize = have.frontiers.iter().map(|f| f.heads.len()).sum();
307 self.check_count(count)?;
308 self.check_count(have.frontiers.len())?;
309 let mut heads = Vec::new();
310 let mut receipt = ReplicationReceipt::default();
311 for frontier in have.frontiers {
312 let facet = native_facet(frontier.facet)?;
313 if !self.facets.contains(&facet) {
314 return Err(Error::Protocol("unnegotiated facet").into());
315 }
316 for bytes in frontier.heads {
317 let id = hash(&bytes)?;
318 if let Some((record, status)) = self
319 .replica
320 .operation(id)
321 .await
322 .map_err(StoreError::Store)?
323 {
324 if record.original.verify().map_err(Error::from)?.facet() != facet {
325 return Err(Error::Protocol(
326 "advertised frontier has the wrong facet",
327 )
328 .into());
329 }
330 match status {
331 Admission::Accepted => receipt.accepted_operation_ids.push(bytes),
332 Admission::Rejected(message) => {
333 receipt.rejected.push(rejection(id, message))
334 }
335 Admission::Pending => heads.push((facet, id)),
336 }
337 } else {
338 heads.push((facet, id));
339 }
340 }
341 }
342 self.replica
343 .remember_peer_heads(self.destination, heads)
344 .await
345 .map_err(StoreError::Store)?;
346 if !receipt.accepted_operation_ids.is_empty() || !receipt.rejected.is_empty() {
347 responses.push(Outbound::Frame(Frame::Receipt(receipt)));
348 }
349 }
350 Frame::Need(need) => {
351 self.check_count(need.operation_ids.len())?;
352 let sharing = self.export_facets().await?;
353 for bytes in need.operation_ids {
354 let id = hash(&bytes)?;
355 let Some((record, Admission::Accepted)) = self
356 .replica
357 .operation(id)
358 .await
359 .map_err(StoreError::Store)?
360 else {
361 return Err(Error::Protocol("requested operation unavailable").into());
362 };
363 let operation = record.original.verify().map_err(Error::from)?;
364 if !sharing.contains(&operation.facet()) {
365 return Err(
366 Error::Protocol("operation is outside current sharing policy").into(),
367 );
368 }
369 responses.push(Outbound::Operation(id));
370 }
371 }
372 Frame::Operations(batch) => {
373 if self.pending_input_bookkeeping.is_some() {
374 return Err(Error::Protocol("prior input bookkeeping not completed").into());
375 }
376 self.check_count(batch.operations.len())?;
377 self.check_count(batch.authority_admissions.len())?;
378 self.check_count(batch.boundary_acceptances.len())?;
379 let decoded =
380 crate::authority_admission::match_batch(&batch).map_err(
381 |error| match error {
382 crate::transport::Error::Protocol(message) => Error::Protocol(message),
383 _ => Error::Protocol("invalid original authority batch"),
384 },
385 )?;
386 let mut operations = Vec::with_capacity(decoded.len());
387 for received in decoded {
388 let operation = received.original.verify().map_err(Error::from)?;
389 let id = operation.id().map_err(Error::from)?;
390 if !self.facets.contains(&operation.facet()) {
391 return Err(Error::Protocol("operation is outside admission scope").into());
392 }
393 operations.push((received, operation, id));
394 }
395 let mut receipt = ReplicationReceipt::default();
396 for (received, _, _) in operations {
397 let next = self.commit_original(received, maintenance).await?;
398 receipt
399 .accepted_operation_ids
400 .extend(next.accepted_operation_ids);
401 receipt
402 .pending_operation_ids
403 .extend(next.pending_operation_ids);
404 receipt.rejected.extend(next.rejected);
405 }
406 responses.push(Outbound::Frame(Frame::Receipt(receipt)));
407 }
408 Frame::Receipt(receipt) => {
409 self.check_count(
410 receipt.accepted_operation_ids.len()
411 + receipt.pending_operation_ids.len()
412 + receipt.rejected.len(),
413 )?;
414 for bytes in receipt.accepted_operation_ids {
415 self.replica
416 .record_peer_receipt(self.destination, hash(&bytes)?, Admission::Accepted)
417 .await
418 .map_err(StoreError::Store)?;
419 }
420 for bytes in receipt.pending_operation_ids {
421 self.replica
422 .record_peer_receipt(self.destination, hash(&bytes)?, Admission::Pending)
423 .await
424 .map_err(StoreError::Store)?;
425 }
426 for rejected in receipt.rejected {
427 self.replica
428 .record_peer_receipt(
429 self.destination,
430 hash(&rejected.operation_id)?,
431 Admission::Rejected(
432 rejected.failure.map(|f| f.message).unwrap_or_default(),
433 ),
434 )
435 .await
436 .map_err(StoreError::Store)?;
437 }
438 }
439 }
440 if maintenance && let Some(repair) = self.control().await? {
441 responses.push(Outbound::Frame(repair));
442 }
443 Ok(responses)
444 }
445
446 pub async fn control(&mut self) -> StoreResult<Option<Frame>, B::Error> {
449 let settled = self
450 .replica
451 .settled_peer_heads(self.destination, self.facets.clone(), self.max_items)
452 .await
453 .map_err(StoreError::Store)?;
454 if !settled.is_empty() {
455 let mut receipt = ReplicationReceipt::default();
456 for (id, admission) in settled {
457 self.in_flight.remove(&id);
458 match admission {
459 Admission::Accepted => {
460 receipt.accepted_operation_ids.push(id.as_bytes().to_vec())
461 }
462 Admission::Rejected(message) => receipt.rejected.push(rejection(id, message)),
463 Admission::Pending => {}
464 }
465 }
466 return Ok(Some(Frame::Receipt(receipt)));
467 }
468 let available = self.max_items.saturating_sub(self.in_flight.len());
469 if available == 0 {
470 return Ok(None);
471 }
472 let candidates = self
473 .replica
474 .needed_from_peer(self.destination, self.facets.clone(), self.max_items)
475 .await
476 .map_err(StoreError::Store)?;
477 let ids: Vec<_> = candidates
478 .into_iter()
479 .filter(|id| !self.in_flight.contains(id))
480 .take(available)
481 .collect();
482 self.in_flight.extend(ids.iter().copied());
483 Ok((!ids.is_empty()).then(|| {
484 Frame::Need(ReplicationNeed {
485 operation_ids: ids.into_iter().map(|id| id.as_bytes().to_vec()).collect(),
486 })
487 }))
488 }
489
490 pub async fn export_operation(&self, id: ContentHash) -> StoreResult<Frame, B::Error> {
493 let Some((record, Admission::Accepted)) = self
494 .replica
495 .operation(id)
496 .await
497 .map_err(StoreError::Store)?
498 else {
499 return Err(Error::Protocol("requested operation unavailable").into());
500 };
501 let operation = record.original.verify().map_err(Error::from)?;
502 if !self.export_facets().await?.contains(&operation.facet()) {
503 return Err(Error::Protocol("operation is outside current sharing policy").into());
504 }
505 Ok(Frame::Operations(ReplicationOperations {
506 boundary_acceptances: crate::boundary_acceptance::authority_evidence(
507 record.authority_admission.as_ref(),
508 )
509 .map_err(|_| Error::Protocol("invalid boundary evidence"))?,
510 authority_admissions: record
511 .authority_admission
512 .as_ref()
513 .map(crate::authority_admission::encode)
514 .transpose()
515 .map_err(|_| Error::Protocol("invalid retained authority admission"))?
516 .into_iter()
517 .collect(),
518 operations: vec![SignedRecord {
519 format: OPERATION_FORMAT.into(),
520 canonical_record: record.original.canonical,
521 signatures: vec![RecordSignature {
522 public_key: operation.publisher.to_vec(),
523 signature: record.original.signature,
524 }],
525 }],
526 import_authority: None,
528 }))
529 }
530
531 fn check_count(&self, count: usize) -> Result<()> {
532 if count > self.max_items {
533 Err(Error::Protocol("replication item budget exceeded"))
534 } else {
535 Ok(())
536 }
537 }
538}
539
540pub fn completion_receipt(id: ContentHash, admission: Admission) -> ReplicationReceipt {
543 let mut receipt = ReplicationReceipt::default();
544 match admission {
545 Admission::Accepted => receipt.accepted_operation_ids.push(id.as_bytes().to_vec()),
546 Admission::Pending => receipt.pending_operation_ids.push(id.as_bytes().to_vec()),
547 Admission::Rejected(message) => receipt.rejected.push(rejection(id, message)),
548 }
549 receipt
550}
551
552fn rejection(id: ContentHash, message: String) -> ReplicationRejection {
553 ReplicationRejection {
554 operation_id: id.as_bytes().to_vec(),
555 failure: Some(api::heddle::api::common::CallFailure {
556 code: 9,
557 message,
558 ..Default::default()
559 }),
560 }
561}
562
563pub fn decode_record(record: SignedRecord) -> Result<SignedOperation> {
564 if record.format != OPERATION_FORMAT || record.signatures.len() != 1 {
565 return Err(Error::Protocol("unsupported signed operation format"));
566 }
567 let signature = &record.signatures[0];
568 let signed = SignedOperation {
569 canonical: record.canonical_record,
570 signature: signature.signature.clone(),
571 };
572 if signed.verify()?.publisher.as_slice() != signature.public_key {
573 return Err(Error::Protocol(
574 "record signature key differs from publisher",
575 ));
576 }
577 Ok(signed)
578}
579pub fn native_facet(facet: i32) -> Result<ThreadFacet> {
580 match SharedFacet::try_from(facet) {
581 Ok(SharedFacet::Source) => Ok(ThreadFacet::Source),
582 Ok(SharedFacet::Collaboration) => Ok(ThreadFacet::Discussion),
583 Ok(SharedFacet::Metadata) => Ok(ThreadFacet::Metadata),
584 _ => Err(Error::Protocol("unsupported replication facet")),
585 }
586}
587pub fn wire_facet(facet: ThreadFacet) -> i32 {
588 match facet {
589 ThreadFacet::Source => SharedFacet::Source as i32,
590 ThreadFacet::Discussion => SharedFacet::Collaboration as i32,
591 ThreadFacet::Metadata => SharedFacet::Metadata as i32,
592 }
593}
594fn hash(bytes: &[u8]) -> Result<ContentHash> {
595 Ok(ContentHash::from_bytes(bytes.try_into().map_err(|_| {
596 Error::Protocol("operation ID must be 32 bytes")
597 })?))
598}
599
600#[cfg(all(test, feature = "native"))]
601#[path = "replication_tests.rs"]
602mod tests;
603
604#[cfg(all(test, feature = "native"))]
605#[path = "replication_admission_tests.rs"]
606mod admission_tests;