/home/runner/work/lyquor/lyquor/toolchain/shaker/src/availability.rs
Line | Count | Source |
1 | | use std::future::Future; |
2 | | use std::sync::Arc; |
3 | | use std::time::Duration; |
4 | | |
5 | | use alloy_sol_types::{SolCall, sol}; |
6 | | use anyhow::Context; |
7 | | use lyquor_eth::{EthSubmitter, Signer}; |
8 | | use lyquor_jsonrpc::types::{ |
9 | | BlockNumber, EthCall, EthCallResp, EthCallTx, EthGetTransactionReceipt, EthGetTransactionReceiptResp, |
10 | | }; |
11 | | use lyquor_primitives::alloy_primitives::B256; |
12 | | use lyquor_primitives::oracle::OracleConfig; |
13 | | use lyquor_primitives::{Address, NodeID, U256, decode_object}; |
14 | | |
15 | | use crate::Client; |
16 | | |
17 | | const AVAILABILITY_TOPIC: &str = "availability"; |
18 | | const DEFAULT_POLL_INTERVAL: Duration = Duration::from_millis(200); |
19 | | const DEFAULT_POLL_TIMEOUT: Duration = Duration::from_secs(60); |
20 | | const COMMITTEE_TIMEOUT_HINT: &str = "a committee node may be offline, the threshold may be unreachable, or a committee node may be missing its Ed25519 key"; |
21 | | |
22 | | sol! { |
23 | | interface BartenderAvailability { |
24 | | function __lyquor_oracle_initialize( |
25 | | string topic, |
26 | | address targetAddr, |
27 | | bool isEvm, |
28 | | bytes32[] committee, |
29 | | uint16 threshold |
30 | | ) external; |
31 | | function __lyquor_oracle_advance_epoch( |
32 | | string topic, |
33 | | address targetAddr, |
34 | | bool isEvm |
35 | | ) external returns (bool); |
36 | | function __lyquor_oracle_finalize_epoch( |
37 | | string topic, |
38 | | address targetAddr, |
39 | | bool isEvm |
40 | | ) external returns (bool); |
41 | | function __lyquor_oracle_dest_epoch_info( |
42 | | string topic, |
43 | | bool fullConfig |
44 | | ) external returns (uint64 epoch, bytes32 configHash, uint32 changeCount, bytes config); |
45 | | function get_availability_epoch() external returns (uint32); |
46 | | function get_availability_counts() external returns (uint64 admittedImages, uint64 pendingDeployments); |
47 | | function get_availability_deadline_blocks() external returns (uint64); |
48 | | } |
49 | 0 | } |
50 | | |
51 | | /// Observable availability-gate state. |
52 | | #[derive(Debug, Clone, PartialEq, Eq, serde::Serialize)] |
53 | | pub struct AvailabilityStatus { |
54 | | pub source_epoch: u32, |
55 | | pub dest_epoch: u32, |
56 | | pub committee: Vec<NodeID>, |
57 | | pub threshold: u16, |
58 | | pub deadline_blocks: u64, |
59 | | pub admitted_image_count: u64, |
60 | | pub pending_deployment_count: u64, |
61 | | } |
62 | | |
63 | | #[derive(Debug)] |
64 | | struct EpochState { |
65 | | source_epoch: u32, |
66 | | dest_epoch: u32, |
67 | | config: OracleConfig, |
68 | | } |
69 | | |
70 | 0 | async fn eth_call<C: SolCall>(client: &Client, bartender: Address, call: C) -> anyhow::Result<C::Return> { |
71 | 0 | let response: EthCallResp = client |
72 | 0 | .request(EthCall { |
73 | 0 | tx: EthCallTx { |
74 | 0 | from: None, |
75 | 0 | to: Some(bartender), |
76 | 0 | gas: None, |
77 | 0 | gas_price: None, |
78 | 0 | value: None, |
79 | 0 | data: Some(call.abi_encode().into()), |
80 | 0 | }, |
81 | 0 | block_number: BlockNumber::Latest, |
82 | 0 | }) |
83 | 0 | .await |
84 | 0 | .context("Availability call failed")?; |
85 | 0 | C::abi_decode_returns(&response.0).context("Failed to decode availability call response") |
86 | 0 | } |
87 | | |
88 | 0 | async fn epoch_state(client: &Client, bartender: Address) -> anyhow::Result<EpochState> { |
89 | 0 | let (source_epoch, dest) = tokio::try_join!( |
90 | 0 | eth_call(client, bartender, BartenderAvailability::get_availability_epochCall {}), |
91 | 0 | eth_call( |
92 | 0 | client, |
93 | 0 | bartender, |
94 | 0 | BartenderAvailability::__lyquor_oracle_dest_epoch_infoCall { |
95 | 0 | topic: AVAILABILITY_TOPIC.to_owned(), |
96 | 0 | fullConfig: true, |
97 | 0 | }, |
98 | | ) |
99 | 0 | )?; |
100 | 0 | let dest_epoch = u32::try_from(dest.epoch).context("Availability destination epoch does not fit in u32")?; |
101 | 0 | let config = if dest.config.is_empty() { |
102 | 0 | OracleConfig { |
103 | 0 | committee: Vec::new(), |
104 | 0 | threshold: 0, |
105 | 0 | } |
106 | | } else { |
107 | 0 | decode_object(&dest.config).context("Bartender returned an invalid availability committee config")? |
108 | | }; |
109 | 0 | Ok(EpochState { |
110 | 0 | source_epoch, |
111 | 0 | dest_epoch, |
112 | 0 | config, |
113 | 0 | }) |
114 | 0 | } |
115 | | |
116 | 0 | async fn poll<T, F, Fut>(timeout_message: impl Into<String>, mut operation: F) -> anyhow::Result<T> |
117 | 0 | where |
118 | 0 | F: FnMut() -> Fut, |
119 | 0 | Fut: Future<Output = anyhow::Result<Option<T>>>, |
120 | 0 | { |
121 | 0 | let timeout_message = timeout_message.into(); |
122 | 0 | tokio::time::timeout(DEFAULT_POLL_TIMEOUT, async { |
123 | | loop { |
124 | 0 | if let Some(value) = operation().await? { |
125 | 0 | return Ok(value); |
126 | 0 | } |
127 | 0 | tokio::time::sleep(DEFAULT_POLL_INTERVAL).await; |
128 | | } |
129 | 0 | }) |
130 | 0 | .await |
131 | 0 | .map_err(|_| anyhow::Error::msg(timeout_message))? |
132 | 0 | } |
133 | | |
134 | 0 | async fn wait_for_transaction(client: &Client, tx_hash: B256) -> anyhow::Result<()> { |
135 | 0 | let receipt = poll( |
136 | 0 | format!("Timed out waiting for availability initialization transaction {tx_hash} to be sequenced"), |
137 | 0 | || async { |
138 | 0 | let receipt: EthGetTransactionReceiptResp = client.request(EthGetTransactionReceipt(tx_hash)).await?; |
139 | 0 | Ok(receipt) |
140 | 0 | }, |
141 | | ) |
142 | 0 | .await?; |
143 | 0 | if receipt.status != U256::from(1_u8) { |
144 | 0 | anyhow::bail!( |
145 | | "Availability initialization transaction {tx_hash} failed with status {}", |
146 | | receipt.status |
147 | | ); |
148 | 0 | } |
149 | 0 | Ok(()) |
150 | 0 | } |
151 | | |
152 | 5 | fn validate_committee(committee: &[NodeID], threshold: u16) -> anyhow::Result<()> { |
153 | 5 | if committee.is_empty() { |
154 | 1 | anyhow::bail!("Availability committee cannot be empty"); |
155 | 4 | } |
156 | 4 | if threshold == 0 || usize::from3 (threshold3 ) > committee.len() { |
157 | 2 | anyhow::bail!( |
158 | | "Availability threshold {} must be between 1 and the committee size {}", |
159 | | threshold, |
160 | 2 | committee.len() |
161 | | ); |
162 | 2 | } |
163 | 2 | let mut unique = committee.to_vec(); |
164 | 2 | unique.sort_unstable(); |
165 | 2 | unique.dedup(); |
166 | 2 | if unique.len() != committee.len() { |
167 | 1 | anyhow::bail!("Availability committee contains duplicate node IDs"); |
168 | 1 | } |
169 | 1 | Ok(()) |
170 | 5 | } |
171 | | |
172 | | /// Complete bartender's first availability epoch. |
173 | 0 | pub async fn activate<S: Signer + Clone + Send + Sync + 'static>( |
174 | 0 | client: &Client, signer: &S, bartender: Address, target: Address, committee: &[NodeID], threshold: u16, |
175 | 0 | ) -> anyhow::Result<AvailabilityStatus> { |
176 | 0 | validate_committee(committee, threshold)?; |
177 | | |
178 | 0 | let call = BartenderAvailability::__lyquor_oracle_initializeCall { |
179 | 0 | topic: AVAILABILITY_TOPIC.to_owned(), |
180 | 0 | targetAddr: target, |
181 | | isEvm: false, |
182 | 0 | committee: committee.iter().map(|id| <[u8; 32]>::from(*id).into()).collect(), |
183 | 0 | threshold, |
184 | | }; |
185 | 0 | let submitter = EthSubmitter::new(client.clone(), Arc::new(signer.clone())); |
186 | 0 | let tx_hash = submitter |
187 | 0 | .submit_contract_call(bartender, call.abi_encode().into()) |
188 | 0 | .await |
189 | 0 | .context("Failed to submit availability committee initialization")?; |
190 | 0 | wait_for_transaction(client, tx_hash).await?; |
191 | | |
192 | 0 | poll( |
193 | 0 | format!("Timed out submitting the availability epoch advance; {COMMITTEE_TIMEOUT_HINT}"), |
194 | 0 | || async { |
195 | 0 | let advanced = eth_call( |
196 | 0 | client, |
197 | 0 | bartender, |
198 | 0 | BartenderAvailability::__lyquor_oracle_advance_epochCall { |
199 | 0 | topic: AVAILABILITY_TOPIC.to_owned(), |
200 | 0 | targetAddr: target, |
201 | 0 | isEvm: false, |
202 | 0 | }, |
203 | 0 | ) |
204 | 0 | .await?; |
205 | 0 | Ok(advanced.then_some(())) |
206 | 0 | }, |
207 | | ) |
208 | 0 | .await?; |
209 | 0 | let mut state = poll( |
210 | 0 | format!("Timed out waiting for availability destination epoch 1; {COMMITTEE_TIMEOUT_HINT}"), |
211 | 0 | || async { |
212 | 0 | let state = epoch_state(client, bartender).await?; |
213 | 0 | Ok((state.dest_epoch >= 1).then_some(state)) |
214 | 0 | }, |
215 | | ) |
216 | 0 | .await?; |
217 | | |
218 | 0 | let finalized = eth_call( |
219 | 0 | client, |
220 | 0 | bartender, |
221 | 0 | BartenderAvailability::__lyquor_oracle_finalize_epochCall { |
222 | 0 | topic: AVAILABILITY_TOPIC.to_owned(), |
223 | 0 | targetAddr: target, |
224 | 0 | isEvm: false, |
225 | 0 | }, |
226 | 0 | ) |
227 | 0 | .await?; |
228 | 0 | if !finalized { |
229 | 0 | anyhow::bail!("Bartender did not accept availability epoch finalization"); |
230 | 0 | } |
231 | 0 | let expected = state.dest_epoch; |
232 | 0 | state = poll( |
233 | 0 | format!("Timed out waiting for availability source epoch {expected}; {COMMITTEE_TIMEOUT_HINT}"), |
234 | 0 | || async { |
235 | 0 | let state = epoch_state(client, bartender).await?; |
236 | 0 | Ok((state.source_epoch >= expected).then_some(state)) |
237 | 0 | }, |
238 | | ) |
239 | 0 | .await?; |
240 | | |
241 | 0 | if state.source_epoch == 0 || state.source_epoch != state.dest_epoch { |
242 | 0 | anyhow::bail!( |
243 | | "Availability ceremony did not settle: source epoch {}, destination epoch {}", |
244 | | state.source_epoch, |
245 | | state.dest_epoch |
246 | | ); |
247 | 0 | } |
248 | | |
249 | 0 | let status = status_from_state(client, bartender, state).await?; |
250 | 0 | if status.committee != committee || status.threshold != threshold { |
251 | 0 | anyhow::bail!( |
252 | | "Availability committee is already active with a different configuration; active committee {:?}, threshold {}", |
253 | | status.committee, |
254 | | status.threshold |
255 | | ); |
256 | 0 | } |
257 | 0 | Ok(status) |
258 | 0 | } |
259 | | |
260 | | /// Read bartender's availability committee and deployment-admission status. |
261 | 0 | pub async fn status(client: &Client, bartender: Address) -> anyhow::Result<AvailabilityStatus> { |
262 | 0 | let state = epoch_state(client, bartender).await?; |
263 | 0 | status_from_state(client, bartender, state).await |
264 | 0 | } |
265 | | |
266 | 0 | async fn status_from_state( |
267 | 0 | client: &Client, bartender: Address, state: EpochState, |
268 | 0 | ) -> anyhow::Result<AvailabilityStatus> { |
269 | 0 | let (counts, deadline_blocks) = tokio::try_join!( |
270 | 0 | eth_call(client, bartender, BartenderAvailability::get_availability_countsCall {}), |
271 | 0 | eth_call( |
272 | 0 | client, |
273 | 0 | bartender, |
274 | 0 | BartenderAvailability::get_availability_deadline_blocksCall {}, |
275 | | ) |
276 | 0 | )?; |
277 | 0 | let committee = state |
278 | 0 | .config |
279 | 0 | .committee |
280 | 0 | .into_iter() |
281 | 0 | .map(|signer| { |
282 | 0 | NodeID::try_from(signer.key.as_ref()) |
283 | 0 | .map_err(|err| anyhow::anyhow!("Invalid node ID in availability committee: {err:?}")) |
284 | 0 | }) |
285 | 0 | .collect::<anyhow::Result<Vec<_>>>()?; |
286 | 0 | Ok(AvailabilityStatus { |
287 | 0 | source_epoch: state.source_epoch, |
288 | 0 | dest_epoch: state.dest_epoch, |
289 | 0 | committee, |
290 | 0 | threshold: state.config.threshold, |
291 | 0 | deadline_blocks, |
292 | 0 | admitted_image_count: counts.admittedImages, |
293 | 0 | pending_deployment_count: counts.pendingDeployments, |
294 | 0 | }) |
295 | 0 | } |
296 | | |
297 | | #[cfg(test)] |
298 | | mod tests { |
299 | | use super::*; |
300 | | use lyquor_test::test; |
301 | | |
302 | | #[test] |
303 | | fn activation_rejects_invalid_committee_thresholds_and_duplicates() { |
304 | | let node = NodeID::from(1); |
305 | | assert!(validate_committee(&[], 1).is_err()); |
306 | | assert!(validate_committee(&[node], 0).is_err()); |
307 | | assert!(validate_committee(&[node], 2).is_err()); |
308 | | assert!(validate_committee(&[node, node], 1).is_err()); |
309 | | assert!(validate_committee(&[node], 1).is_ok()); |
310 | | } |
311 | | } |