1use super::{
2 actors::{batcher, resolver, voter},
3 config::{Config, SkipPolicy},
4 elector::{self, Elector as _},
5 types::{Activity, Context},
6};
7use crate::{
8 CertifiableAutomaton, Relay, Reporter,
9 simplex::{Lookahead, Plan, scheme::Scheme},
10};
11use commonware_cryptography::Digest;
12use commonware_macros::select;
13use commonware_p2p::{Blocker, Receiver, Sender};
14use commonware_parallel::Strategy;
15use commonware_runtime::{
16 BufferPooler, Clock, ContextCell, Handle, Metrics, Spawner, Storage, spawn_cell,
17};
18use rand_core::CryptoRng;
19use tracing::debug;
20
21pub struct Engine<
23 E: BufferPooler + Clock + CryptoRng + Spawner + Storage + Metrics,
24 S: Scheme<D>,
25 L: elector::Config<S>,
26 B: Blocker<PublicKey = S::PublicKey>,
27 D: Digest,
28 A: CertifiableAutomaton<Context = Context<D, S::PublicKey>, Digest = D>,
29 R: Relay<Digest = D, PublicKey = S::PublicKey, Plan = Plan<S::PublicKey>>,
30 F: Reporter<Activity = Activity<S, D>>,
31 T: Strategy,
32> {
33 context: ContextCell<E>,
34
35 voter: voter::Actor<E, S, L::Elector, B, D, A, R, F>,
36 voter_mailbox: voter::Mailbox<S, D>,
37
38 batcher: batcher::Actor<E, S, B, D, F, R, T>,
39 batcher_mailbox: batcher::Mailbox<S, D>,
40
41 resolver: resolver::Actor<E, S, B, D, T>,
42 resolver_mailbox: resolver::Mailbox<S, D>,
43}
44
45impl<
46 E: BufferPooler + Clock + CryptoRng + Spawner + Storage + Metrics,
47 S: Scheme<D>,
48 L: elector::Config<S>,
49 B: Blocker<PublicKey = S::PublicKey>,
50 D: Digest,
51 A: CertifiableAutomaton<Context = Context<D, S::PublicKey>, Digest = D>,
52 R: Relay<Digest = D, PublicKey = S::PublicKey, Plan = Plan<S::PublicKey>>,
53 F: Reporter<Activity = Activity<S, D>>,
54 T: Strategy,
55> Engine<E, S, L, B, D, A, R, F, T>
56{
57 pub fn new(mut context: E, cfg: Config<S, L, B, D, A, R, F, T>) -> Self {
59 cfg.assert(&mut context);
61 let skip_budget = match cfg.skip {
62 SkipPolicy::Disabled => 0,
63 SkipPolicy::Enabled { budget, .. } => budget.resolve(cfg.scheme.participants().len()),
64 };
65 let elector = cfg.elector.build(cfg.scheme.participants());
66 let terms = elector.terms();
67 let term_length = terms.length();
68 if let Some(stall_timeout) = terms.stall_timeout() {
69 assert!(
70 stall_timeout > cfg.certification_timeout,
71 "stall timeout must be greater than certification timeout"
72 );
73 }
74
75 let (batcher, batcher_mailbox) = batcher::Actor::new(
77 context.child("batcher"),
78 batcher::Config {
79 scheme: cfg.scheme.clone(),
80 blocker: cfg.blocker.clone(),
81 reporter: cfg.reporter.clone(),
82 track_historical_votes: cfg.track_historical_votes,
83 relay: cfg.relay.clone(),
84 strategy: cfg.strategy.clone(),
85 epoch: cfg.epoch,
86 mailbox_size: cfg.mailbox_size,
87 view_retention: cfg.view_retention,
88 skip: cfg.skip,
89 lookahead: Lookahead::new(&terms),
90 forward: cfg.forward,
91 floor: cfg.floor.view(),
92 },
93 );
94
95 let (voter, voter_mailbox) = voter::Actor::new(
97 context.child("voter"),
98 voter::Config {
99 scheme: cfg.scheme.clone(),
100 elector,
101 blocker: cfg.blocker.clone(),
102 automaton: cfg.automaton,
103 relay: cfg.relay,
104 reporter: cfg.reporter,
105 partition: cfg.partition,
106 mailbox_size: cfg.mailbox_size,
107 epoch: cfg.epoch,
108 floor: cfg.floor,
109 leader_timeout: cfg.leader_timeout,
110 certification_timeout: cfg.certification_timeout,
111 timeout_retry: cfg.timeout_retry,
112 skip_budget,
113 view_retention: cfg.view_retention,
114 replay_buffer: cfg.replay_buffer,
115 write_buffer: cfg.write_buffer,
116 page_cache: cfg.page_cache,
117 },
118 );
119
120 let (resolver, resolver_mailbox) = resolver::Actor::new(
122 context.child("resolver"),
123 resolver::Config {
124 blocker: cfg.blocker,
125 scheme: cfg.scheme,
126 strategy: cfg.strategy,
127 mailbox_size: cfg.mailbox_size,
128 epoch: cfg.epoch,
129 fetch_timeout: cfg.fetch_timeout,
130 term_length,
131 },
132 );
133
134 Self {
136 context: ContextCell::new(context),
137
138 voter,
139 voter_mailbox,
140
141 batcher,
142 batcher_mailbox,
143
144 resolver,
145 resolver_mailbox,
146 }
147 }
148
149 pub fn start(
187 mut self,
188 vote_network: (
189 impl Sender<PublicKey = S::PublicKey>,
190 impl Receiver<PublicKey = S::PublicKey>,
191 ),
192 certificate_network: (
193 impl Sender<PublicKey = S::PublicKey>,
194 impl Receiver<PublicKey = S::PublicKey>,
195 ),
196 resolver_network: (
197 impl Sender<PublicKey = S::PublicKey>,
198 impl Receiver<PublicKey = S::PublicKey>,
199 ),
200 ) -> Handle<()> {
201 spawn_cell!(
202 self.context,
203 self.run(vote_network, certificate_network, resolver_network)
204 )
205 }
206
207 async fn run(
208 self,
209 vote_network: (
210 impl Sender<PublicKey = S::PublicKey>,
211 impl Receiver<PublicKey = S::PublicKey>,
212 ),
213 certificate_network: (
214 impl Sender<PublicKey = S::PublicKey>,
215 impl Receiver<PublicKey = S::PublicKey>,
216 ),
217 resolver_network: (
218 impl Sender<PublicKey = S::PublicKey>,
219 impl Receiver<PublicKey = S::PublicKey>,
220 ),
221 ) {
222 let (vote_sender, vote_receiver) = vote_network;
225 let (certificate_sender, certificate_receiver) = certificate_network;
226 let mut batcher_task = self.batcher.start(
227 self.voter_mailbox.clone(),
228 vote_receiver,
229 certificate_receiver,
230 );
231
232 let (resolver_sender, resolver_receiver) = resolver_network;
234 let mut resolver_task =
235 self.resolver
236 .start(self.voter_mailbox, resolver_sender, resolver_receiver);
237
238 let mut voter_task = self.voter.start(
240 self.batcher_mailbox,
241 self.resolver_mailbox,
242 vote_sender,
243 certificate_sender,
244 );
245
246 let mut shutdown = self.context.stopped();
248 select! {
249 _ = &mut shutdown => {
250 debug!("context shutdown, stopping engine");
251 },
252 voter = &mut voter_task => {
253 debug!(?voter, "voter stopped, shutting down engine");
254 },
255 batcher = &mut batcher_task => {
256 debug!(?batcher, "batcher stopped, shutting down engine");
257 },
258 resolver = &mut resolver_task => {
259 debug!(?resolver, "resolver stopped, shutting down engine");
260 },
261 }
262 }
263}