Skip to main content

heddle_thread_api/
replication.rs

1// SPDX-License-Identifier: Apache-2.0
2//! A bounded, transport-neutral causal exchange shared by device and hosted
3//! endpoints. The RPC adapter supplies the verified request scope.
4use 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/// Internal frames carry the same typed payloads in both directions.
41#[allow(clippy::large_enum_variant)] // heddle-api's inline import-authority bundle; boxing adds a heap hop per frame
42pub enum Frame {
43    Have(ReplicationHave),
44    Need(ReplicationNeed),
45    Operations(ReplicationOperations),
46    Receipt(ReplicationReceipt),
47}
48/// Validated originals retain shared evidence without re-encoding a wire batch
49/// for every committed-prefix unit.
50#[allow(clippy::large_enum_variant)] // heddle-api's inline import-authority bundle; boxing adds a heap hop per frame
51pub(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)] // heddle-api's inline import-authority bundle; boxing adds a heap hop per frame
107pub 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    /// `facets` is the intersection of authenticated admission scope and the
127    /// negotiated opening. Local export still checks current sharing policy.
128    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    /// At most one bounded frontier page. Call again until None. A generation
159    /// change during a paged announcement starts another round, closing gaps.
160    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        // A watcher may already have delivered a mutation that happened while
200        // this round was being paged. Keep the caller draining until a fresh
201        // round covers it; None must mean the announcement is caught up.
202        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    /// Validate the complete wire batch before committing any prefix. Each
211    /// returned unit can then receive its own durable acknowledgement.
212    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    /// Return the committed result before optional peer repair bookkeeping.
292    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    /// One bounded dependency window, refilled after each received operation.
447    /// Peer frontiers are durable; outstanding requests are session-local.
448    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    /// Resolve payload only when the writer is ready. Pending output queues
491    /// contain IDs, and policy is checked again immediately before disclosure.
492    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            // Native replication carries no import-authority proof bundle.
527            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
540/// Exact completion fields for one durably admitted original operation.
541/// Hosts can bind their acknowledgement exception to these canonical fields.
542pub 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;