1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
use bytes::Bytes;
use std::ops::RangeInclusive;
use std::time::Duration;
use derive_where::derive_where;
use ractor::{ActorRef, RpcReplyPort};
use malachitebft_core_consensus::{MisbehaviorEvidence, Role, VoteExtensionError};
use malachitebft_core_types::{CommitCertificate, Context, Round, ValueId, VoteExtensions};
use malachitebft_sync::{PeerId, RawDecidedValue};
use crate::util::streaming::StreamMessage;
pub use malachitebft_core_consensus::{LocallyProposedValue, ProposedValue};
pub use malachitebft_core_types::HeightParams;
/// A reference to the host actor.
pub type HostRef<Ctx> = ActorRef<HostMsg<Ctx>>;
/// What to do next after a decision.
#[derive_where(Debug)]
pub enum Next<Ctx: Context> {
/// Start at the given height with the provided parameters.
Start(Ctx::Height, HeightParams<Ctx>),
/// Restart at the given height with the provided parameters.
Restart(Ctx::Height, HeightParams<Ctx>),
}
/// Messages that need to be handled by the host actor.
#[derive_where(Debug)]
pub enum HostMsg<Ctx: Context> {
/// Notifies the application that consensus is ready.
///
/// The application MUST reply with a message to instruct
/// consensus to start at a given height.
ConsensusReady {
/// Use this reply port to instruct consensus to start the first height.
reply_to: RpcReplyPort<(Ctx::Height, HeightParams<Ctx>)>,
},
/// Consensus has started a new round.
StartedRound {
/// The height at which the round started.
height: Ctx::Height,
/// The round number that started.
round: Round,
/// The address of the proposer for this round.
proposer: Ctx::Address,
/// The role of the node in this round.
role: Role,
/// Use this reply port to send the undecided values that were already seen for this
/// round. This is needed when recovering from a crash.
///
/// The application MUST reply immediately with the values it has, or with an empty vector.
reply_to: RpcReplyPort<Vec<ProposedValue<Ctx>>>,
},
/// Request to build a local value to propose
///
/// The application MUST reply to this message with the requested value
/// within the specified timeout duration.
GetValue {
/// The height at which the value should be proposed.
height: Ctx::Height,
/// The round in which the value should be proposed.
round: Round,
/// The amount of time the application has to build the value.
timeout: Duration,
/// Use this reply port to send the value that was built.
reply_to: RpcReplyPort<LocallyProposedValue<Ctx>>,
},
/// ExtendVote allows the application to extend the pre-commit vote with arbitrary data.
///
/// When consensus is preparing to send a pre-commit vote, it first calls `ExtendVote`.
/// The application then returns a blob of data called a vote extension.
/// This data is opaque to the consensus algorithm but can contain application-specific information.
/// The proposer of the next block will receive all vote extensions along with the commit certificate.
ExtendVote {
/// The height at which the vote is being extended.
height: Ctx::Height,
/// The round in which the vote is being extended.
round: Round,
/// The ID of the value that is being voted on.
value_id: ValueId<Ctx>,
/// The vote extension to be added to the vote, if any.
reply_to: RpcReplyPort<Option<Ctx::Extension>>,
},
/// Verify a vote extension
///
/// If the vote extension is deemed invalid, the vote it was part of
/// will be discarded altogether.
VerifyVoteExtension {
/// The height for which the vote is.
height: Ctx::Height,
/// The round for which the vote is.
round: Round,
/// The ID of the value that the vote extension is for.
value_id: ValueId<Ctx>,
/// The vote extension to verify.
extension: Ctx::Extension,
/// Use this reply port to send the result of the verification.
reply_to: RpcReplyPort<Result<(), VoteExtensionError>>,
},
/// Requests the application to re-stream a proposal that it has already seen.
///
/// The application MUST re-publish again all the proposal parts pertaining
/// to that value by sending [`NetworkMsg::PublishProposalPart`] messages through
/// the [`Channels::network`] channel.
RestreamValue {
/// The height at which the value was proposed.
height: Ctx::Height,
/// The round in which the value was proposed.
round: Round,
/// The round in which the value was valid.
valid_round: Round,
/// The address of the proposer of the value.
address: Ctx::Address,
/// The ID of the value to restream.
value_id: ValueId<Ctx>,
},
/// Requests the earliest height available in the history maintained by the application.
///
/// The application MUST respond with its earliest available height.
GetHistoryMinHeight { reply_to: RpcReplyPort<Ctx::Height> },
/// Notifies the application that consensus has received a proposal part over the network.
///
/// If this part completes the full proposal, the application MUST respond
/// with the complete proposed value. Otherwise, it MUST respond with `None`.
ReceivedProposalPart {
from: PeerId,
part: StreamMessage<Ctx::ProposalPart>,
reply_to: RpcReplyPort<ProposedValue<Ctx>>,
},
/// Notifies the application that consensus has decided on a value.
///
/// This message includes a commit certificate containing the ID of
/// the value that was decided on, the height and round at which it was decided,
/// and the aggregated signatures of the validators that committed to it.
/// It also includes to the vote extensions received for that height.
Decided {
/// The commit certificate containing the ID of the value that was decided on,
/// the the height and round at which it was decided, and the aggregated signatures
/// of the validators that committed to it.
certificate: CommitCertificate<Ctx>,
/// Vote extensions that were received for this height.
extensions: VoteExtensions<Ctx>,
},
/// Notifies the application that consensus has finalized a height after collecting additional precommits.
///
/// This message is sent when the target time for the height has been reached,
/// which may include a delay between `Decided` and `Finalized` messages.
/// During this delay, additional precommits may have been collected. The certificate may contain more
/// signatures than the one sent in the initial Decided message.
///
/// In response to this message, the application MUST send a [`Next`]
/// message back to consensus, instructing it to either start the next height if
/// the application was able to commit the decided value, or to restart the current height
/// otherwise.
///
/// If the application does not reply, consensus will stall.
Finalized {
/// The commit certificate with extended signatures collected during finalization period.
certificate: CommitCertificate<Ctx>,
/// Vote extensions that were received for this height (including additional ones).
extensions: VoteExtensions<Ctx>,
/// Misbehavior evidence collected since last height was decided.
evidence: MisbehaviorEvidence<Ctx>,
/// Use this reply port to instruct consensus to start the next height.
reply_to: RpcReplyPort<Next<Ctx>>,
},
/// Requests a range of previously decided values from the application's storage.
///
/// The application MUST respond with those values if available, or `None` otherwise.
///
/// ## Important
/// The range is NOT checked for validity by consensus. It is the application's responsibility
/// to ensure that the the range is within valid bounds.
GetDecidedValues {
/// Range of decided values to retrieve
range: RangeInclusive<Ctx::Height>,
/// Channel for sending back the decided value
reply_to: RpcReplyPort<Vec<RawDecidedValue<Ctx>>>,
},
/// Notifies the application that a value has been synced from the network.
/// This may happen when the node is catching up with the network.
///
/// If a value can be decoded from the bytes provided, then the application MUST reply
/// to this message with the decoded value. Otherwise, it MUST reply with `None`.
ProcessSyncedValue {
/// Height of the synced value
height: Ctx::Height,
/// Round of the synced value
round: Round,
/// Address of the original proposer
proposer: Ctx::Address,
/// Raw encoded value data
value_bytes: Bytes,
/// Channel for sending back the proposed value, if successfully decoded
/// or `None` if the value could not be decoded
reply_to: RpcReplyPort<ProposedValue<Ctx>>,
},
}