Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_o3db_sync/src/dist/transport.rs

5.7 KiB, 158 runs

created by r1870400018:11394, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1//! Wire-format envelopes and message kinds for distributed Ozone.
2//!
3//! Transport itself lives *outside* this crate: the engine emits envelopes
4//! as a [`Commands`](crate::dist::Commands) return value, the caller
5//! dispatches them. The production adapter is Shield (UDP, signed-hash
6//! datagrams, AddressGuard rate-limiting); the test adapter is whatever the
7//! caller builds from a channel or a mock.
8//!
9//! Keeping transport out of the engine means distributed Ozone is a pure
10//! state machine: every decision it makes is a function of its inputs with
11//! no hidden I/O, which is the property that made the primitive crates
12//! (Kademlia, OAM, IBLT, HotStuff) easy to test and reason about.
13//!
14//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
15//! Anthropic Claude
16
17use super::record::{
18 Record,
19 RecordId,
20};
21use super::hotstuff::types::{
22 NewView,
23 Proposal,
24 Vote,
25};
26
27use oxedyne_fe2o3_core::prelude::*;
28use crate::kademlia::id::NodeId;
29
30
31/// Matches a [`MsgKind::GetRequest`] with its eventual
32/// [`MsgKind::GetResponse`]. Opaque 64-bit tokens, picked monotonically per
33/// process.
34pub type RequestId = u64;
35
36
37/// An envelope wraps a message body with its sender and intended recipient.
38///
39/// The engine consumes envelopes via [`DistOzone::handle_envelope`] and emits
40/// them in the `outbound` field of its outcome types. Signing, encryption and
41/// on-wire encoding are the transport adapter's responsibility -- the engine
42/// treats envelopes as opaque authenticated structures.
43///
44/// [`DistOzone::handle_envelope`]: crate::dist::DistOzone::handle_envelope
45#[derive(Clone, Debug, Eq, PartialEq)]
46pub struct Envelope {
47 pub from: NodeId,
48 pub to: NodeId,
49 pub body: MsgKind,
50}
51
52impl Envelope {
53 pub fn new(from: NodeId, to: NodeId, body: MsgKind) -> Self {
54 Self { from, to, body }
55 }
56}
57
58
59/// The distributed-Ozone message kinds.
60///
61/// This enum covers the replication-broadcast / read-routing cycle
62/// (`ReplicatePut`, `GetRequest`, `GetResponse`), the IBLT anti-entropy
63/// cycle (`AntiEntropyDigest`, `AntiEntropyReply`, `AntiEntropyPush`),
64/// and the HotStuff cohort cycle for strong-consistency tables
65/// (`CohortSubmit`, `CohortPropose`, `CohortVote`, `CohortNewView`).
66/// Brickyard backup messages are deferred until that layer lands.
67#[derive(Clone, Debug, Eq, PartialEq)]
68pub enum MsgKind {
69 // A write: "persist this record if you consider yourself a holder".
70 // Emitted for every remote holder. Each recipient re-checks its own
71 // placement decision, so a peer with a slightly different view of N may
72 // decline; the record is then dropped and the next anti-entropy round
73 // fills the gap.
74 ReplicatePut {
75 record: Record,
76 },
77 // A read request: "do you have this record?". Emitted when the local peer
78 // is not itself a holder of the record it wants to read. The recipient
79 // answers with GetResponse.
80 GetRequest {
81 request_id: RequestId,
82 table: String,
83 id: RecordId,
84 },
85 GetResponse {
86 request_id: RequestId,
87 record: Option<Record>, // None if the recipient did not have it
88 },
89
90 // Anti-entropy digest: "here is my IBLT for this table; reconcile against
91 // yours and reply with the symmetric difference". The sender builds the
92 // sketch from its own Storage::digests enumeration. The recipient builds
93 // its own sketch with matching parameters, subtracts, decodes, and replies
94 // with AntiEntropyReply.
95 AntiEntropyDigest {
96 table: String,
97 sketch: Vec<u8>, // Iblt::to_bytes output, opaque to transport
98 },
99 // Anti-entropy reply: "these records are what I have and you lack; please
100 // send me records with these identifiers". On decode failure -- sketch
101 // overload -- the recipient bulk-replies with every record it holds for
102 // the table and an empty requested-id list, and the originator absorbs
103 // what it lacks. Simple, at the cost of bandwidth on a fresh join.
104 AntiEntropyReply {
105 table: String,
106 records: Vec<Record>, // held here, and lacked by the originator
107 requested_ids: Vec<RecordId>, // wanted from the originator
108 bulk: bool, // set when the sketch could not decode
109 },
110 // Anti-entropy push: "here are the records you requested", sent by the
111 // originator in answer to an AntiEntropyReply's requested_ids list.
112 AntiEntropyPush {
113 table: String,
114 records: Vec<Record>,
115 },
116
117 // Forwarded write: "you are the HotStuff leader for this record; drive
118 // consensus on my behalf". Sent by a peer whose own DistOzone::put named a
119 // cohort-backed table for which it is not the initial round's leader. The
120 // recipient leader creates a CohortInstance if none exists and opens a
121 // CohortPropose round.
122 CohortSubmit {
123 record: Record,
124 },
125 // HotStuff leader-to-cohort proposal. The (table, id) pair selects the
126 // per-record HotStuff instance; the Proposal itself carries the view,
127 // phase, block hash, optional block payload and optional justify QC.
128 CohortPropose {
129 table: String,
130 id: RecordId,
131 proposal: Proposal,
132 },
133 // HotStuff replica-to-leader vote.
134 CohortVote {
135 table: String,
136 id: RecordId,
137 vote: Vote,
138 },
139 // HotStuff view change, sent by a replica to the incoming leader when its
140 // local timer fires without seeing progress in the current view.
141 CohortNewView {
142 table: String,
143 id: RecordId,
144 new_view: NewView,
145 },
146}
147
148impl MsgKind {
149 pub fn label(&self) -> &'static str {
150 match self {
151 Self::ReplicatePut { .. } => "ReplicatePut",
152 Self::GetRequest { .. } => "GetRequest",
153 Self::GetResponse { .. } => "GetResponse",
154 Self::AntiEntropyDigest { .. } => "AntiEntropyDigest",
155 Self::AntiEntropyReply { .. } => "AntiEntropyReply",
156 Self::AntiEntropyPush { .. } => "AntiEntropyPush",
157 Self::CohortSubmit { .. } => "CohortSubmit",
158 Self::CohortPropose { .. } => "CohortPropose",
159 Self::CohortVote { .. } => "CohortVote",
160 Self::CohortNewView { .. } => "CohortNewView",
161 }
162 }
163}