Oregami
Repositories/oxedyne/fe2o3

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

30.1 KiB, 227 runs

created by r1870400018:11382, 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//! The distributed-Ozone engine.
2//!
3//! [`DistOzone`] composes the placement service, the peer set, and a
4//! [`Storage`] backend into one cohesive replication engine. It is a pure
5//! state machine: every public method either reads state or returns the
6//! outbound envelopes the caller should dispatch through its transport
7//! adapter. No method calls `send` itself.
8//!
9//! # Write path
10//!
11//! [`DistOzone::put`] runs the placement decision, persists locally if the
12//! local peer is a holder, and returns [`PutOutcome`] with a
13//! [`ReplicatePut`](crate::transport::MsgKind::ReplicatePut) envelope for
14//! every remote holder. The caller dispatches the envelopes; recipients
15//! re-check their own placement decision, so a put that a sender directed
16//! at a peer with a slightly different view of `N` can be silently dropped
17//! by the recipient without harm -- the next anti-entropy round fills it
18//! in.
19//!
20//! # Read path
21//!
22//! [`DistOzone::get`] reads from the local store if the local peer is a
23//! holder, returning [`GetOutcome::Local`] or [`GetOutcome::LocalMiss`]. If
24//! the local peer is *not* a holder, it returns [`GetOutcome::Remote`] with
25//! a request id and the
26//! [`GetRequest`](crate::transport::MsgKind::GetRequest) envelopes to
27//! dispatch. The caller polls [`DistOzone::poll_get`] to learn when a
28//! response has landed.
29//!
30//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
31//! Anthropic Claude
32
33use super::cohort;
34use super::config::{
35 Consistency,
36 DistOzoneConfig,
37 TableConfig,
38};
39use super::consensus::{
40 self,
41 CohortInstance,
42};
43use super::hotstuff::{
44 replica::Command as HsCommand,
45 types::{
46 NewView,
47 Proposal,
48 Vote,
49 },
50};
51use super::peer_set::PeerSet;
52use super::placement::Placement;
53use super::record::{
54 Record,
55 RecordId,
56};
57use super::storage::Storage;
58use super::transport::{
59 Envelope,
60 MsgKind,
61 RequestId,
62};
63
64use oxedyne_fe2o3_core::prelude::*;
65use oxedyne_fe2o3_data::iblt::{
66 DecodeOutcome,
67 Iblt,
68 IbltConfig,
69};
70use crate::kademlia::id::{
71 ID_LEN,
72 NodeId,
73};
74use crate::oam::config::OamConfig;
75
76use std::collections::HashMap;
77use std::sync::{
78 Mutex,
79 atomic::{
80 AtomicU64,
81 Ordering,
82 },
83};
84
85
86// Fixed key and value lengths for the anti-entropy IBLT sketches. The key is a
87// 32-byte RecordId and the value is the 32-byte content hash produced by
88// Storage::digests. Matching lengths across peers is mandatory for IBLT
89// subtraction; callers cannot override.
90const ANTI_ENTROPY_KEY_LEN: usize = ID_LEN;
91const ANTI_ENTROPY_VALUE_LEN: usize = 32;
92
93
94// Choosing more than one peer protects against a single straggler; choosing
95// many wastes bandwidth. Three is the operating-point default referenced in the
96// spec ("an OAM holder" -- plural in realistic deployments).
97pub const DEFAULT_READ_FANOUT: usize = 3;
98
99
100pub struct DistOzone<S: Storage> {
101 cfg: DistOzoneConfig,
102 placement: Placement,
103 peer_set: PeerSet,
104 storage: S,
105 next_rid: AtomicU64,
106 pending_gets: Mutex<HashMap<RequestId, PendingGet>>,
107 read_fanout: usize,
108 // An entry is created lazily on the first message, and kept after Decide so
109 // duplicate late messages are silently absorbed.
110 cohorts: Mutex<HashMap<(String, RecordId), CohortInstance>>,
111}
112
113/// The in-flight state of a remote read.
114#[derive(Clone, Debug)]
115struct PendingGet {
116 #[allow(dead_code)] // Preserved for diagnostics and future retry logic.
117 table: String,
118 #[allow(dead_code)]
119 id: RecordId,
120 outstanding: usize,
121 first_response: Option<Record>,
122 resolved_empty: bool,
123}
124
125impl<S: Storage> DistOzone<S> {
126 /// The bootstrap peer list is filtered to exclude the local peer, and the
127 /// placement service's threshold is precomputed.
128 pub fn new(cfg: DistOzoneConfig, storage: S) -> Outcome<Self> {
129 let placement = Placement::new(cfg.local_peer_id, cfg.oam);
130 let peer_set = PeerSet::from_bootstrap(
131 &cfg.local_peer_id,
132 cfg.bootstrap_peers.iter().cloned(),
133 );
134 Ok(Self {
135 cfg,
136 placement,
137 peer_set,
138 storage,
139 next_rid: AtomicU64::new(1),
140 pending_gets: Mutex::new(HashMap::new()),
141 read_fanout: DEFAULT_READ_FANOUT,
142 cohorts: Mutex::new(HashMap::new()),
143 })
144 }
145
146 /// One is the minimum; values above the peer-set size are clamped
147 /// automatically at request time.
148 pub fn set_read_fanout(&mut self, fanout: usize) {
149 self.read_fanout = fanout.max(1);
150 }
151
152 pub fn config(&self) -> &DistOzoneConfig {
153 &self.cfg
154 }
155
156 pub fn peer_set(&self) -> &PeerSet {
157 &self.peer_set
158 }
159
160 pub fn placement(&self) -> &Placement {
161 &self.placement
162 }
163
164 pub fn storage(&self) -> &S {
165 &self.storage
166 }
167
168 /// Ignores the local peer; true if the peer was new.
169 pub fn insert_peer(&mut self, peer: NodeId) -> bool {
170 if peer == self.cfg.local_peer_id {
171 return false;
172 }
173 self.peer_set.insert(peer)
174 }
175
176 pub fn remove_peer(&mut self, peer: &NodeId) -> bool {
177 self.peer_set.remove(peer)
178 }
179
180 /// Recomputes the cached placement threshold. The estimate typically comes
181 /// from a HyperLogLog merge.
182 pub fn update_network_size(&mut self, network_size: u64) -> Outcome<()> {
183 let oam = res!(OamConfig::new(self.cfg.oam.replication, network_size));
184 self.cfg.oam = oam;
185 self.placement.update_oam(oam);
186 Ok(())
187 }
188
189 fn table_or_err(&self, name: &str) -> Outcome<&TableConfig> {
190 match self.cfg.table(name) {
191 Some(t) => Ok(t),
192 None => Err(err!(
193 "Unknown table in DistOzone operation: {}.", name;
194 Invalid, Input, Missing)),
195 }
196 }
197
198 /// On an eventual table the record is persisted locally if the local
199 /// peer is a holder and a [`MsgKind::ReplicatePut`] envelope is emitted
200 /// for every remote holder.
201 ///
202 /// On a cohort-backed table the put enters a HotStuff round. If the
203 /// local peer is the cohort's initial leader the engine opens a
204 /// [`MsgKind::CohortPropose`] round and returns the proposal envelopes.
205 /// Otherwise it emits a [`MsgKind::CohortSubmit`] envelope to the
206 /// leader, who drives the round on the submitter's behalf. In both
207 /// cases [`PutOutcome::consensus_pending`] is set; the caller learns
208 /// when consensus completes through
209 /// [`InboundOutcome::completed_consensus_put`].
210 pub fn put(&self, record: Record) -> Outcome<PutOutcome> {
211 let tc = res!(self.table_or_err(&record.table));
212 match tc.consistency {
213 Consistency::Eventual => self.put_eventual(record),
214 Consistency::Cohort { lambda } =>
215 self.put_cohort(record, lambda),
216 }
217 }
218
219 fn put_eventual(&self, record: Record) -> Outcome<PutOutcome> {
220 let decision = self.placement.decide(&record.id, &self.peer_set);
221 let local_persisted = decision.local_is_holder;
222 if local_persisted {
223 res!(self.storage.put(&record));
224 }
225
226 let mut outbound = Vec::with_capacity(decision.remote_holders.len());
227 for peer in decision.remote_holders {
228 outbound.push(Envelope::new(
229 self.cfg.local_peer_id,
230 *peer,
231 MsgKind::ReplicatePut { record: record.clone() },
232 ));
233 }
234 Ok(PutOutcome {
235 local_persisted,
236 outbound,
237 consensus_pending: None,
238 })
239 }
240
241 fn put_cohort(&self, record: Record, lambda: u64) -> Outcome<PutOutcome> {
242 let table = record.table.clone();
243 let id = record.id;
244 let sel = res!(cohort::select(
245 &table,
246 &id,
247 &self.peer_set,
248 &self.cfg.local_peer_id,
249 lambda,
250 ));
251 if sel.local_is_leader {
252 // Drive the round locally.
253 let outbound = res!(self.leader_open_round(sel, record));
254 Ok(PutOutcome {
255 local_persisted: false,
256 outbound,
257 consensus_pending: Some((table, id)),
258 })
259 } else {
260 // Forward to the leader.
261 let env = Envelope::new(
262 self.cfg.local_peer_id,
263 sel.leader,
264 MsgKind::CohortSubmit { record },
265 );
266 Ok(PutOutcome {
267 local_persisted: false,
268 outbound: vec![env],
269 consensus_pending: Some((table, id)),
270 })
271 }
272 }
273
274 fn leader_open_round(
275 &self,
276 sel: cohort::Cohort,
277 record: Record,
278 )
279 -> Outcome<Vec<Envelope>>
280 {
281 let table = record.table.clone();
282 let id = record.id;
283 let block = consensus::encode_record(&record);
284 let block_hash = consensus::block_hash(&block);
285
286 let lambda = sel.members.len() as u64;
287 let mut cohorts = lock_mutex!(self.cohorts);
288 let instance = match cohorts.get_mut(&(table.clone(), id)) {
289 Some(i) => i,
290 None => {
291 let fresh = res!(CohortInstance::new(
292 sel,
293 &self.cfg.local_peer_id,
294 lambda,
295 ));
296 cohorts.insert((table.clone(), id), fresh);
297 cohorts.get_mut(&(table.clone(), id)).expect("just inserted")
298 },
299 };
300 if instance.has_decided() {
301 return Ok(Vec::new());
302 }
303 let cmds = res!(instance.replica.propose(block, block_hash));
304 let outbound = res!(self.translate_commands(&table, &id, instance, cmds));
305 Ok(outbound)
306 }
307
308 /// Applies the local side effects -- persistence on Decide -- as it goes.
309 fn translate_commands(
310 &self,
311 table: &str,
312 id: &RecordId,
313 instance: &mut CohortInstance,
314 cmds: Vec<HsCommand>,
315 )
316 -> Outcome<Vec<Envelope>>
317 {
318 let mut out = Vec::new();
319 for cmd in cmds {
320 match cmd {
321 HsCommand::BroadcastProposal(proposal) => {
322 // Broadcast to every cohort member *except* ourselves.
323 // HotStuff's on_proposal recipe expects the leader to
324 // record its own vote via on_proposal too, so we also
325 // feed the proposal back into our local replica.
326 let self_id = instance.replica.config().self_id;
327 for (i, member) in instance.members.iter().enumerate() {
328 if (i as u16) == self_id {
329 continue;
330 }
331 out.push(Envelope::new(
332 self.cfg.local_peer_id,
333 *member,
334 MsgKind::CohortPropose {
335 table: table.to_string(),
336 id: *id,
337 proposal: proposal.clone(),
338 },
339 ));
340 }
341 // Now fold the proposal into the local replica. We
342 // recurse on the commands produced, which lets a leader
343 // that is also a voter correctly emit its own SendVote.
344 let local_cmds = res!(instance.replica.on_proposal(proposal));
345 let more = res!(self.translate_commands(
346 table, id, instance, local_cmds,
347 ));
348 out.extend(more);
349 },
350 HsCommand::SendVote { to, vote } => {
351 let target = match instance.node_id(to) {
352 Some(n) => n,
353 None => return Err(err!(
354 "HotStuff SendVote targets replica id {} \
355 which is out of cohort range (size = {}).",
356 to, instance.members.len();
357 Bug, Invalid)),
358 };
359 // The leader of the current view is the usual target; a
360 // vote sent to ourselves needs to be folded in directly
361 // rather than sent on the wire.
362 if target == self.cfg.local_peer_id {
363 let local_cmds = res!(instance.replica.on_vote(vote));
364 let more = res!(self.translate_commands(
365 table, id, instance, local_cmds,
366 ));
367 out.extend(more);
368 } else {
369 out.push(Envelope::new(
370 self.cfg.local_peer_id,
371 target,
372 MsgKind::CohortVote {
373 table: table.to_string(),
374 id: *id,
375 vote,
376 },
377 ));
378 }
379 },
380 HsCommand::SendNewView { to, new_view } => {
381 let target = match instance.node_id(to) {
382 Some(n) => n,
383 None => return Err(err!(
384 "HotStuff SendNewView targets replica id {} \
385 which is out of cohort range (size = {}).",
386 to, instance.members.len();
387 Bug, Invalid)),
388 };
389 if target == self.cfg.local_peer_id {
390 let local_cmds = res!(instance.replica.on_new_view(new_view));
391 let more = res!(self.translate_commands(
392 table, id, instance, local_cmds,
393 ));
394 out.extend(more);
395 } else {
396 out.push(Envelope::new(
397 self.cfg.local_peer_id,
398 target,
399 MsgKind::CohortNewView {
400 table: table.to_string(),
401 id: *id,
402 new_view,
403 },
404 ));
405 }
406 },
407 HsCommand::Decide { view: _, block } => {
408 let record = res!(consensus::decode_record(&block));
409 if record.table != table || &record.id != id {
410 return Err(err!(
411 "HotStuff Decide block decoded to a record at \
412 ({}, {:?}) that does not match the consensus \
413 slot ({}, {:?}).",
414 record.table, record.id.as_bytes(),
415 table, id.as_bytes();
416 Bug, Invalid, Mismatch));
417 }
418 let hash = consensus::block_hash(&block);
419 res!(instance.mark_decided(hash));
420 res!(self.storage.put(&record));
421 },
422 }
423 }
424 Ok(out)
425 }
426
427 /// The caller owns the timer; on expiry it calls this and dispatches the
428 /// returned envelopes. Empty if the instance is absent -- already decided,
429 /// or never created.
430 pub fn cohort_timeout(
431 &self,
432 table: &str,
433 id: &RecordId,
434 )
435 -> Outcome<Vec<Envelope>>
436 {
437 let mut cohorts = lock_mutex!(self.cohorts);
438 let instance = match cohorts.get_mut(&(table.to_string(), *id)) {
439 Some(i) => i,
440 None => return Ok(Vec::new()),
441 };
442 if instance.has_decided() {
443 return Ok(Vec::new());
444 }
445 let cmds = res!(instance.replica.on_timeout());
446 self.translate_commands(table, id, instance, cmds)
447 }
448
449 /// Reads from local storage if the local peer is a holder; otherwise
450 /// dispatches a request to the nearest peers and returns the in-flight
451 /// request handle.
452 pub fn get(
453 &self,
454 table: &str,
455 id: &RecordId,
456 )
457 -> Outcome<GetOutcome>
458 {
459 res!(self.table_or_err(table));
460
461 if self.placement.i_am_holder(id) {
462 return Ok(match res!(self.storage.get(table, id)) {
463 Some(r) => GetOutcome::Local(r),
464 None => GetOutcome::LocalMiss,
465 });
466 }
467
468 let targets = self.placement.read_targets(
469 id,
470 &self.peer_set,
471 self.read_fanout,
472 );
473 if targets.is_empty() {
474 return Ok(GetOutcome::NoTargets);
475 }
476 let request_id = self.next_rid.fetch_add(1, Ordering::Relaxed);
477 let outbound = targets.iter()
478 .map(|peer| Envelope::new(
479 self.cfg.local_peer_id,
480 **peer,
481 MsgKind::GetRequest {
482 request_id,
483 table: table.to_string(),
484 id: *id,
485 },
486 ))
487 .collect::<Vec<_>>();
488 {
489 let mut pending = lock_mutex!(self.pending_gets);
490 pending.insert(request_id, PendingGet {
491 table: table.to_string(),
492 id: *id,
493 outstanding: outbound.len(),
494 first_response: None,
495 resolved_empty: false,
496 });
497 }
498 Ok(GetOutcome::Remote { request_id, outbound })
499 }
500
501 pub fn handle_envelope(&self, env: Envelope) -> Outcome<InboundOutcome> {
502 if env.to != self.cfg.local_peer_id {
503 // An envelope addressed to somebody else; ignore. This is mostly
504 // a belt-and-braces guard: the transport adapter should not
505 // deliver misaddressed envelopes in the first place.
506 return Ok(InboundOutcome::empty());
507 }
508 match env.body {
509 MsgKind::ReplicatePut { record } => {
510 // Re-check placement: sender's view of N may disagree with
511 // ours near the threshold. Drop silently if we do not
512 // consider ourselves a holder.
513 if !self.placement.i_am_holder(&record.id) {
514 return Ok(InboundOutcome::empty());
515 }
516 // Reject puts for tables we do not know about.
517 if self.cfg.table(&record.table).is_none() {
518 return Err(err!(
519 "ReplicatePut for unknown table '{}'.", record.table;
520 Invalid, Input, Missing));
521 }
522 res!(self.storage.put(&record));
523 Ok(InboundOutcome::empty())
524 }
525 MsgKind::GetRequest { request_id, table, id } => {
526 if self.cfg.table(&table).is_none() {
527 return Err(err!(
528 "GetRequest for unknown table '{}'.", table;
529 Invalid, Input, Missing));
530 }
531 let record = res!(self.storage.get(&table, &id));
532 let reply = Envelope::new(
533 self.cfg.local_peer_id,
534 env.from,
535 MsgKind::GetResponse { request_id, record },
536 );
537 Ok(InboundOutcome {
538 outbound: vec![reply],
539 completed_get: None,
540 completed_consensus_put: None,
541 })
542 }
543 MsgKind::GetResponse { request_id, record } => {
544 let completed = {
545 let mut pending = lock_mutex!(self.pending_gets);
546 let Some(slot) = pending.get_mut(&request_id) else {
547 // Unknown request id; stale or cancelled response.
548 return Ok(InboundOutcome::empty());
549 };
550 if slot.outstanding > 0 {
551 slot.outstanding -= 1;
552 }
553 match record {
554 Some(r) if slot.first_response.is_none() => {
555 slot.first_response = Some(r);
556 }
557 None => {
558 slot.resolved_empty = true;
559 }
560 _ => { /* later non-first response; ignore. */ }
561 }
562 // A pending get is "resolved" when either a response with
563 // a record has landed or every target has replied empty.
564 slot.first_response.is_some() || slot.outstanding == 0
565 };
566 Ok(InboundOutcome {
567 outbound: Vec::new(),
568 completed_get: if completed { Some(request_id) } else { None },
569 completed_consensus_put: None,
570 })
571 }
572 MsgKind::AntiEntropyDigest { table, sketch } => {
573 self.handle_anti_entropy_digest(env.from, table, sketch)
574 }
575 MsgKind::AntiEntropyReply { table, records, requested_ids, bulk } => {
576 self.handle_anti_entropy_reply(
577 env.from, table, records, requested_ids, bulk,
578 )
579 }
580 MsgKind::AntiEntropyPush { table, records } => {
581 self.handle_anti_entropy_push(table, records)
582 }
583 MsgKind::CohortSubmit { record } => {
584 self.handle_cohort_submit(record)
585 }
586 MsgKind::CohortPropose { table, id, proposal } => {
587 self.handle_cohort_propose(env.from, table, id, proposal)
588 }
589 MsgKind::CohortVote { table, id, vote } => {
590 self.handle_cohort_vote(table, id, vote)
591 }
592 MsgKind::CohortNewView { table, id, new_view } => {
593 self.handle_cohort_new_view(table, id, new_view)
594 }
595 }
596 }
597
598 pub fn poll_get(&self, request_id: RequestId) -> Outcome<PollOutcome> {
599 let pending = lock_mutex!(self.pending_gets);
600 let Some(slot) = pending.get(&request_id) else {
601 return Ok(PollOutcome::Unknown);
602 };
603 if let Some(r) = &slot.first_response {
604 return Ok(PollOutcome::Record(r.clone()));
605 }
606 if slot.outstanding == 0 && slot.resolved_empty {
607 return Ok(PollOutcome::NotFound);
608 }
609 Ok(PollOutcome::Pending)
610 }
611
612 /// Late responses that arrive after cancellation are ignored by the next
613 /// [`handle_envelope`](Self::handle_envelope) call.
614 pub fn cancel_get(&self, request_id: RequestId) -> Outcome<()> {
615 let mut pending = lock_mutex!(self.pending_gets);
616 pending.remove(&request_id);
617 Ok(())
618 }
619
620 /// The envelope carries a serialised IBLT built from the local storage's
621 /// [`Storage::digests`] enumeration for that table. The recipient
622 /// subtracts it against its own sketch, decodes the symmetric difference
623 /// and answers with [`MsgKind::AntiEntropyReply`]. Cohort-backed tables
624 /// reconcile through consensus rather than anti-entropy, and are rejected.
625 pub fn build_anti_entropy_request(
626 &self,
627 table: &str,
628 target: NodeId,
629 )
630 -> Outcome<Envelope>
631 {
632 let tc = res!(self.table_or_err(table));
633 if !matches!(tc.consistency, Consistency::Eventual) {
634 return Err(err!(
635 "anti-entropy is defined only for Eventual tables \
636 (table '{}').", table;
637 Invalid, Input, Unimplemented));
638 }
639 let iblt = res!(self.build_table_iblt(tc));
640 let sketch = iblt.to_bytes();
641 Ok(Envelope::new(
642 self.cfg.local_peer_id,
643 target,
644 MsgKind::AntiEntropyDigest {
645 table: table.to_string(),
646 sketch,
647 },
648 ))
649 }
650
651 /// Factored out so both the request builder and the inbound digest handler
652 /// use the same sketch shape.
653 fn build_table_iblt(&self, tc: &TableConfig) -> Outcome<Iblt> {
654 let cfg = IbltConfig {
655 num_cells: tc.iblt_cells,
656 num_hashes: TableConfig::IBLT_NUM_HASHES,
657 key_len: ANTI_ENTROPY_KEY_LEN,
658 value_len: ANTI_ENTROPY_VALUE_LEN,
659 seed: tc.iblt_seed(),
660 };
661 let mut iblt = res!(Iblt::new(cfg));
662 let digests = res!(self.storage.digests(&tc.name));
663 for d in digests {
664 res!(iblt.insert(d.id.as_bytes(), &d.content));
665 }
666 Ok(iblt)
667 }
668
669 /// Decodes the symmetric difference against the local sketch and returns an
670 /// [`AntiEntropyReply`][ar] envelope carrying records the sender lacks and
671 /// a list of record identifiers the recipient lacks. On sketch overload it
672 /// falls back to a bulk reply of every record the recipient holds for the
673 /// table.
674 ///
675 /// [ar]: MsgKind::AntiEntropyReply
676 fn handle_anti_entropy_digest(
677 &self,
678 from: NodeId,
679 table: String,
680 sketch: Vec<u8>,
681 )
682 -> Outcome<InboundOutcome>
683 {
684 let tc = res!(self.table_or_err(&table));
685 if !matches!(tc.consistency, Consistency::Eventual) {
686 return Err(err!(
687 "anti-entropy digest received for non-Eventual table '{}'.",
688 table;
689 Invalid, Input, Unimplemented));
690 }
691 let expected_cfg = IbltConfig {
692 num_cells: tc.iblt_cells,
693 num_hashes: TableConfig::IBLT_NUM_HASHES,
694 key_len: ANTI_ENTROPY_KEY_LEN,
695 value_len: ANTI_ENTROPY_VALUE_LEN,
696 seed: tc.iblt_seed(),
697 };
698 let their_iblt = res!(Iblt::from_bytes(&sketch));
699 if their_iblt.config() != expected_cfg {
700 return Err(err!(
701 "anti-entropy sketch config mismatch for table '{}'.",
702 table;
703 Invalid, Input, Mismatch));
704 }
705 let mut mine = res!(self.build_table_iblt(tc));
706 res!(mine.subtract(&their_iblt));
707 let decode = res!(mine.decode());
708 let (records_for_sender, requested_ids, bulk) = match decode {
709 DecodeOutcome::Complete { inserted, deleted } => {
710 // `inserted` = keys in mine not in theirs -> records I
711 // should send. `deleted` = keys in theirs not in mine ->
712 // ids I should request.
713 let mut records_for_sender = Vec::with_capacity(inserted.len());
714 for (key_bytes, _value_hash) in inserted {
715 let rid = res!(RecordId::from_slice(&key_bytes));
716 if let Some(r) = res!(self.storage.get(&table, &rid)) {
717 records_for_sender.push(r);
718 }
719 }
720 let mut requested_ids = Vec::with_capacity(deleted.len());
721 for (key_bytes, _value_hash) in deleted {
722 requested_ids.push(res!(RecordId::from_slice(&key_bytes)));
723 }
724 (records_for_sender, requested_ids, false)
725 }
726 DecodeOutcome::Incomplete { .. } => {
727 // Sketch overloaded. Fall back to bulk: send everything I
728 // have for this table; the sender absorbs what it lacks.
729 // This is simple and correct; a later optimisation can
730 // teach the sender to retry with a larger sketch.
731 let digests = res!(self.storage.digests(&table));
732 let mut records = Vec::with_capacity(digests.len());
733 for d in digests {
734 if let Some(r) = res!(self.storage.get(&table, &d.id)) {
735 records.push(r);
736 }
737 }
738 (records, Vec::new(), true)
739 }
740 };
741 let reply = Envelope::new(
742 self.cfg.local_peer_id,
743 from,
744 MsgKind::AntiEntropyReply {
745 table,
746 records: records_for_sender,
747 requested_ids,
748 bulk,
749 },
750 );
751 Ok(InboundOutcome {
752 outbound: vec![reply],
753 completed_get: None,
754 completed_consensus_put: None,
755 })
756 }
757
758 /// Applies the records the recipient was missing, and builds an
759 /// [`AntiEntropyPush`][ap] envelope for any records requested in return.
760 ///
761 /// [ap]: MsgKind::AntiEntropyPush
762 fn handle_anti_entropy_reply(
763 &self,
764 from: NodeId,
765 table: String,
766 records: Vec<Record>,
767 requested_ids: Vec<RecordId>,
768 _bulk: bool,
769 )
770 -> Outcome<InboundOutcome>
771 {
772 res!(self.table_or_err(&table));
773
774 // Apply every record the peer sent us, re-checking placement so a
775 // stale-N sender cannot push a record to a peer that has since
776 // stopped considering itself a holder.
777 for record in records {
778 if record.table != table {
779 continue;
780 }
781 if self.placement.i_am_holder(&record.id) {
782 res!(self.storage.put(&record));
783 }
784 }
785
786 // Build a push for every requested id we actually have.
787 let mut to_push = Vec::with_capacity(requested_ids.len());
788 for rid in requested_ids {
789 if let Some(r) = res!(self.storage.get(&table, &rid)) {
790 to_push.push(r);
791 }
792 }
793 if to_push.is_empty() {
794 return Ok(InboundOutcome::empty());
795 }
796 let push = Envelope::new(
797 self.cfg.local_peer_id,
798 from,
799 MsgKind::AntiEntropyPush {
800 table,
801 records: to_push,
802 },
803 );
804 Ok(InboundOutcome {
805 outbound: vec![push],
806 completed_get: None,
807 completed_consensus_put: None,
808 })
809 }
810
811 /// Opens a HotStuff round if the local peer is the initial leader for the
812 /// `(table, record_id)` pair. If it is not, the submission is dropped
813 /// silently -- almost always a stale-cohort race, where the submitter's
814 /// view of the peer set differed from the leader's.
815 fn handle_cohort_submit(
816 &self,
817 record: Record,
818 )
819 -> Outcome<InboundOutcome>
820 {
821 let tc = res!(self.table_or_err(&record.table));
822 let lambda = match tc.consistency {
823 Consistency::Cohort { lambda } => lambda,
824 Consistency::Eventual => return Err(err!(
825 "CohortSubmit received for eventual-consistency table '{}'.",
826 record.table;
827 Invalid, Input, Mismatch)),
828 };
829 let sel = res!(cohort::select(
830 &record.table,
831 &record.id,
832 &self.peer_set,
833 &self.cfg.local_peer_id,
834 lambda,
835 ));
836 if !sel.local_is_leader {
837 // Not our job -- drop silently.
838 return Ok(InboundOutcome::empty());
839 }
840 let outbound = res!(self.leader_open_round(sel, record));
841 Ok(InboundOutcome {
842 outbound,
843 completed_get: None,
844 completed_consensus_put: None,
845 })
846 }
847
848 /// Creates a fresh per-record replica if one does not yet exist. A proposal
849 /// addressed to a peer that is not a cohort member is dropped silently.
850 fn handle_cohort_propose(
851 &self,
852 from: NodeId,
853 table: String,
854 id: RecordId,
855 proposal: Proposal,
856 )
857 -> Outcome<InboundOutcome>
858 {
859 let tc = res!(self.table_or_err(&table));
860 let lambda = match tc.consistency {
861 Consistency::Cohort { lambda } => lambda,
862 Consistency::Eventual => return Err(err!(
863 "CohortPropose received for eventual-consistency table '{}'.",
864 table;
865 Invalid, Input, Mismatch)),
866 };
867 let sel = res!(cohort::select(
868 &table,
869 &id,
870 &self.peer_set,
871 &self.cfg.local_peer_id,
872 lambda,
873 ));
874 if !sel.local_is_member {
875 // Proposal arrived but we are not in the cohort -- drop. This
876 // can happen if peer-set views disagree; the sender will retry
877 // after its own set updates.
878 return Ok(InboundOutcome::empty());
879 }
880 // Verify the proposal came from someone plausibly in the cohort.
881 // Stronger origin checks (view-aware leader matching) live in the
882 // HotStuff replica itself.
883 if !sel.members.iter().any(|m| m == &from) {
884 return Ok(InboundOutcome::empty());
885 }
886 let mut cohorts = lock_mutex!(self.cohorts);
887 let key = (table.clone(), id);
888 let instance = match cohorts.get_mut(&key) {
889 Some(i) => i,
890 None => {
891 let fresh = res!(CohortInstance::new(
892 sel,
893 &self.cfg.local_peer_id,
894 lambda,
895 ));
896 cohorts.insert(key.clone(), fresh);
897 cohorts.get_mut(&key).expect("just inserted")
898 },
899 };
900 if instance.has_decided() {
901 return Ok(InboundOutcome::empty());
902 }
903 let before_decided = instance.has_decided();
904 let cmds = res!(instance.replica.on_proposal(proposal));
905 let outbound = res!(self.translate_commands(&table, &id, instance, cmds));
906 let completed_consensus_put = if !before_decided && instance.has_decided() {
907 Some((table, id))
908 } else {
909 None
910 };
911 Ok(InboundOutcome {
912 outbound,
913 completed_get: None,
914 completed_consensus_put,
915 })
916 }
917
918 /// Leader-only: a non-leader replica silently ignores a vote, per the
919 /// HotStuff spec.
920 fn handle_cohort_vote(
921 &self,
922 table: String,
923 id: RecordId,
924 vote: Vote,
925 )
926 -> Outcome<InboundOutcome>
927 {
928 let mut cohorts = lock_mutex!(self.cohorts);
929 let key = (table.clone(), id);
930 let instance = match cohorts.get_mut(&key) {
931 Some(i) => i,
932 None => {
933 // Vote arrived before we created an instance. Drop -- a
934 // well-behaved cohort member only votes after seeing the
935 // leader's Propose, so by the time a vote reaches us we
936 // should already have an instance. This case is most
937 // likely an adversarial or replayed envelope.
938 return Ok(InboundOutcome::empty());
939 },
940 };
941 if instance.has_decided() {
942 return Ok(InboundOutcome::empty());
943 }
944 let before_decided = instance.has_decided();
945 let cmds = res!(instance.replica.on_vote(vote));
946 let outbound = res!(self.translate_commands(&table, &id, instance, cmds));
947 let completed_consensus_put = if !before_decided && instance.has_decided() {
948 Some((table, id))
949 } else {
950 None
951 };
952 Ok(InboundOutcome {
953 outbound,
954 completed_get: None,
955 completed_consensus_put,
956 })
957 }
958
959 fn handle_cohort_new_view(
960 &self,
961 table: String,
962 id: RecordId,
963 new_view: NewView,
964 )
965 -> Outcome<InboundOutcome>
966 {
967 let mut cohorts = lock_mutex!(self.cohorts);
968 let key = (table.clone(), id);
969 let instance = match cohorts.get_mut(&key) {
970 Some(i) => i,
971 None => return Ok(InboundOutcome::empty()),
972 };
973 if instance.has_decided() {
974 return Ok(InboundOutcome::empty());
975 }
976 let cmds = res!(instance.replica.on_new_view(new_view));
977 let outbound = res!(self.translate_commands(&table, &id, instance, cmds));
978 Ok(InboundOutcome {
979 outbound,
980 completed_get: None,
981 completed_consensus_put: None,
982 })
983 }
984
985 /// Each record is placement-checked before persistence.
986 fn handle_anti_entropy_push(
987 &self,
988 table: String,
989 records: Vec<Record>,
990 )
991 -> Outcome<InboundOutcome>
992 {
993 res!(self.table_or_err(&table));
994 for record in records {
995 if record.table != table {
996 continue;
997 }
998 if self.placement.i_am_holder(&record.id) {
999 res!(self.storage.put(&record));
1000 }
1001 }
1002 Ok(InboundOutcome::empty())
1003 }
1004}
1005
1006
1007/// The result of a [`DistOzone::put`] call.
1008#[derive(Clone, Debug)]
1009pub struct PutOutcome {
1010 // Always false for a cohort-backed write, which persists on Decide and is
1011 // signalled through InboundOutcome::completed_consensus_put.
1012 pub local_persisted: bool,
1013 pub outbound: Vec<Envelope>,
1014 // Set when the put entered a HotStuff consensus round.
1015 pub consensus_pending: Option<(String, RecordId)>,
1016}
1017
1018
1019/// The result of a [`DistOzone::get`] call.
1020#[derive(Clone, Debug)]
1021pub enum GetOutcome {
1022 Local(Record),
1023 LocalMiss, // a holder, but no record at that id
1024 // A remote read has been initiated; completion is reported through
1025 // DistOzone::poll_get.
1026 Remote {
1027 request_id: RequestId,
1028 outbound: Vec<Envelope>,
1029 },
1030 NoTargets, // not a holder, and no remote targets are known
1031}
1032
1033
1034/// The result of a [`DistOzone::handle_envelope`] call.
1035#[derive(Clone, Debug)]
1036pub struct InboundOutcome {
1037 pub outbound: Vec<Envelope>,
1038 // Set when this envelope completed a pending remote read; the caller polls
1039 // DistOzone::poll_get to collect the record itself.
1040 pub completed_get: Option<RequestId>,
1041 // Set when this envelope drove a cohort round to Decide and the record was
1042 // persisted locally.
1043 pub completed_consensus_put: Option<(String, RecordId)>,
1044}
1045
1046impl InboundOutcome {
1047 fn empty() -> Self {
1048 Self {
1049 outbound: Vec::new(),
1050 completed_get: None,
1051 completed_consensus_put: None,
1052 }
1053 }
1054}
1055
1056
1057/// The result of a [`DistOzone::poll_get`] query.
1058#[derive(Clone, Debug)]
1059pub enum PollOutcome {
1060 Pending,
1061 Record(Record),
1062 NotFound, // every outstanding holder replied that the record is absent
1063 Unknown, // unknown, or cancelled
1064}