Coverage Report

Created: 2026-10-02 21:01

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/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
}