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 | |
| 17 | use super::record::{ |
| 18 | Record, |
| 19 | RecordId, |
| 20 | }; |
| 21 | use super::hotstuff::types::{ |
| 22 | NewView, |
| 23 | Proposal, |
| 24 | Vote, |
| 25 | }; |
| 26 | |
| 27 | use oxedyne_fe2o3_core::prelude::*; |
| 28 | use 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. |
| 34 | pub 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)] |
| 46 | pub struct Envelope { |
| 47 | pub from: NodeId, |
| 48 | pub to: NodeId, |
| 49 | pub body: MsgKind, |
| 50 | } |
| 51 | |
| 52 | impl 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)] |
| 68 | pub 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 | |
| 148 | impl 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 | } |