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    protocol: Option<api::heddle::api::common::ProtocolCompatibility>,
125}
126impl<B: ReplicaStore> Session<B> {
127    /// `facets` is the intersection of authenticated admission scope and the
128    /// negotiated opening. Local export still checks current sharing policy.
129    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    /// Bind optional HYBRID capability to the authenticated Open/Ready pair.
153    /// Ordinary Sync sessions leave this unset until api#307's cutover.
154    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    /// At most one bounded frontier page. Call again until None. A generation
179    /// change during a paged announcement starts another round, closing gaps.
180    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        // A watcher may already have delivered a mutation that happened while
220        // this round was being paged. Keep the caller draining until a fresh
221        // round covers it; None must mean the announcement is caught up.
222        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    /// Validate the complete wire batch before committing any prefix. Each
231    /// returned unit can then receive its own durable acknowledgement.
232    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    /// Return the committed result before optional peer repair bookkeeping.
313    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    /// One bounded dependency window, refilled after each received operation.
469    /// Peer frontiers are durable; outstanding requests are session-local.
470    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    /// Resolve payload only when the writer is ready. Pending output queues
513    /// contain IDs, and policy is checked again immediately before disclosure.
514    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
566/// Exact completion fields for one durably admitted original operation.
567/// Hosts can bind their acknowledgement exception to these canonical fields.
568pub 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;