use crate::{
chain::quantus_subxt,
cli::exercise::{
report::Report,
runner::{submit_ok, ExerciseCtx},
},
error::{QuantusError, Result},
exercise_step,
};
use sp_runtime::traits::{BlakeTwo256, Hash};
use std::path::PathBuf;
use subxt::tx::Payload;
pub enum UpgradeMode {
SetCode(PathBuf),
SelfNoop,
}
pub async fn run(
ctx: &mut ExerciseCtx,
report: &mut Report,
phase: &str,
mode: UpgradeMode,
timeout_secs: u64,
) -> Result<()> {
match mode {
UpgradeMode::SetCode(wasm_path) => {
exercise_step!(
report,
phase,
"governance_set_code",
governance_set_code(ctx, &wasm_path, timeout_secs)
);
},
UpgradeMode::SelfNoop => {
exercise_step!(
report,
phase,
"governance_self_upgrade",
governance_self_upgrade(ctx, timeout_secs)
);
},
}
Ok(())
}
async fn submit_root_referendum(ctx: &mut ExerciseCtx, encoded_call: Vec<u8>) -> Result<u32> {
let preimage_hash: sp_core::H256 = BlakeTwo256::hash(&encoded_call);
let call_len = encoded_call.len() as u32;
crate::cli::common::submit_preimage(&ctx.client, &ctx.alice, encoded_call, ctx.wait_mode())
.await?;
let latest = ctx.client.get_latest_block().await?;
let count_addr = quantus_subxt::api::storage().tech_referenda().referendum_count();
let index = ctx.client.client().storage().at(latest).fetch(&count_addr).await?.unwrap_or(0);
let submit_call =
crate::cli::exercise::scenarios::governance::build_submit_call(preimage_hash, call_len);
let alice = ctx.alice.clone();
submit_ok(ctx, &alice, submit_call).await?;
crate::log_status!("⬆️ Referendum #{index} submitted");
let deposit_call = quantus_subxt::api::tx().tech_referenda().place_decision_deposit(index);
submit_ok(ctx, &alice, deposit_call).await?;
cast_collective_ayes(ctx, index).await?;
crate::log_status!("⬆️ Decision deposit placed and 3 aye votes cast; waiting for enactment…");
Ok(index)
}
async fn cast_collective_ayes(ctx: &ExerciseCtx, index: u32) -> Result<()> {
let client = &ctx.client;
let alice = ctx.alice.clone();
let bob = ctx.bob.clone();
let charlie = ctx.charlie.clone();
let mode = ctx.wait_mode();
let (r_a, r_b, r_c) = tokio::join!(
crate::cli::tech_collective::vote_on_referendum(client, &alice, index, true, mode),
crate::cli::tech_collective::vote_on_referendum(client, &bob, index, true, mode),
crate::cli::tech_collective::vote_on_referendum(client, &charlie, index, true, mode),
);
let mut failures = Vec::new();
for (who, result) in [("alice", r_a), ("bob", r_b), ("charlie", r_c)] {
if let Err(e) = result {
failures.push((who, e));
}
}
if failures.is_empty() {
return Ok(());
}
if referendum_is_approved(ctx, index).await? {
crate::log_status!(
"⬆️ Referendum #{index} already approved; ignoring {} vote(s) that raced the close \
({})",
failures.len(),
failures
.iter()
.map(|(who, e)| format!("{who}: {e}"))
.collect::<Vec<_>>()
.join("; ")
);
return Ok(());
}
let detail = failures
.into_iter()
.map(|(who, e)| format!("{who}: {e}"))
.collect::<Vec<_>>()
.join("; ");
Err(QuantusError::Generic(format!(
"collective aye votes failed on referendum #{index} (and it is not yet approved): {detail}"
)))
}
async fn referendum_is_approved(ctx: &ExerciseCtx, index: u32) -> Result<bool> {
use quantus_subxt::api::runtime_types::pallet_referenda::types::ReferendumInfo;
let info_addr = quantus_subxt::api::storage().tech_referenda().referendum_info_for(index);
let latest = ctx.client.get_latest_block().await?;
Ok(matches!(
ctx.client.client().storage().at(latest).fetch(&info_addr).await?,
Some(ReferendumInfo::Approved(..))
))
}
async fn refund_deposits(ctx: &mut ExerciseCtx, index: u32) -> &'static str {
let alice = ctx.alice.clone();
let refund_decision = quantus_subxt::api::tx().tech_referenda().refund_decision_deposit(index);
let refund_submission =
quantus_subxt::api::tx().tech_referenda().refund_submission_deposit(index);
let refund_result_a = submit_ok(ctx, &alice, refund_decision).await;
let refund_result_b = submit_ok(ctx, &alice, refund_submission).await;
match (refund_result_a, refund_result_b) {
(Ok(_), Ok(_)) => "deposits refunded",
_ => "deposit refund not yet available (referendum bookkeeping pending)",
}
}
async fn referendum_state(ctx: &ExerciseCtx, index: u32) -> Result<String> {
let info_addr = quantus_subxt::api::storage().tech_referenda().referendum_info_for(index);
let latest = ctx.client.get_latest_block().await?;
let info = ctx.client.client().storage().at(latest).fetch(&info_addr).await?;
Ok(format!("{info:?}"))
}
async fn governance_set_code(
ctx: &mut ExerciseCtx,
wasm_path: &std::path::Path,
timeout_secs: u64,
) -> Result<String> {
let (spec_before, _) = ctx.client.get_runtime_version().await?;
let wasm = std::fs::read(wasm_path).map_err(|e| {
QuantusError::Generic(format!("failed to read WASM {}: {e}", wasm_path.display()))
})?;
crate::log_status!(
"⬆️ Upgrade phase: proposing set_code with {} ({} bytes), current spec {}",
wasm_path.display(),
wasm.len(),
spec_before
);
let set_code = quantus_subxt::api::tx().system().set_code(wasm);
let encoded = set_code
.encode_call_data(&ctx.client.client().metadata())
.map_err(|e| QuantusError::Generic(format!("failed to encode set_code: {e:?}")))?;
let index = submit_root_referendum(ctx, encoded).await?;
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(timeout_secs);
loop {
tokio::time::sleep(std::time::Duration::from_secs(6)).await;
let (spec_now, _) = ctx.client.get_runtime_version().await?;
if spec_now > spec_before {
let refunds = refund_deposits(ctx, index).await;
return Ok(format!(
"runtime upgraded via referendum #{index}: spec {spec_before} → {spec_now}; {refunds}"
));
}
if std::time::Instant::now() > deadline {
let info = referendum_state(ctx, index).await?;
return Err(QuantusError::Generic(format!(
"spec version still {spec_before} after {timeout_secs}s; referendum #{index} \
state: {info}. Is the node built with the fast-governance feature and the \
WASM spec_version higher than {spec_before}?"
)));
}
}
}
async fn governance_self_upgrade(ctx: &mut ExerciseCtx, timeout_secs: u64) -> Result<String> {
let (spec, _) = ctx.client.get_runtime_version().await?;
let latest = ctx.client.get_latest_block().await?;
let code = ctx
.client
.client()
.storage()
.at(latest)
.fetch_raw(b":code".to_vec())
.await?
.ok_or_else(|| QuantusError::Generic(":code storage entry not found".to_string()))?;
let code_hash: sp_core::H256 = BlakeTwo256::hash(&code);
crate::log_status!(
"⬆️ Self-upgrade phase: re-installing the current runtime ({} bytes, spec {}, hash \
{code_hash:?})",
code.len(),
spec
);
let authorize = quantus_subxt::api::tx().system().authorize_upgrade_without_checks(code_hash);
let encoded = authorize.encode_call_data(&ctx.client.client().metadata()).map_err(|e| {
QuantusError::Generic(format!("failed to encode authorize_upgrade_without_checks: {e:?}"))
})?;
let index = submit_root_referendum(ctx, encoded).await?;
let authorized_addr = quantus_subxt::api::storage().system().authorized_upgrade();
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(timeout_secs);
loop {
tokio::time::sleep(std::time::Duration::from_secs(3)).await;
let latest = ctx.client.get_latest_block().await?;
if let Some(authorized) =
ctx.client.client().storage().at(latest).fetch(&authorized_addr).await?
{
if authorized.code_hash != code_hash {
return Err(QuantusError::Generic(format!(
"unexpected authorized upgrade hash: {:?} (expected {code_hash:?})",
authorized.code_hash
)));
}
break;
}
if std::time::Instant::now() > deadline {
let info = referendum_state(ctx, index).await?;
return Err(QuantusError::Generic(format!(
"upgrade not authorized within {timeout_secs}s; referendum #{index} state: \
{info}. Is the node built with the fast-governance feature?"
)));
}
}
crate::log_status!("⬆️ Upgrade authorized on-chain; applying the code blob…");
let mut last_seen = ctx.client.get_latest_block().await?;
let apply = quantus_subxt::api::tx().system().apply_authorized_upgrade(code);
let alice = ctx.alice.clone();
submit_ok(ctx, &alice, apply).await?;
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(timeout_secs);
loop {
if let Some(block_hash) = find_code_updated_since(ctx, &mut last_seen).await? {
let (spec_after, _) = ctx.client.get_runtime_version().await?;
if spec_after != spec {
return Err(QuantusError::Generic(format!(
"no-op self-upgrade unexpectedly changed the spec version: {spec} → {spec_after}"
)));
}
let refunds = refund_deposits(ctx, index).await;
return Ok(format!(
"current runtime (spec {spec}) re-installed via referendum #{index} \
(authorize_upgrade_without_checks + apply_authorized_upgrade); \
System::CodeUpdated observed in block {block_hash:?}; {refunds}"
));
}
if std::time::Instant::now() > deadline {
return Err(QuantusError::Generic(format!(
"apply_authorized_upgrade was included but no System::CodeUpdated event was \
observed within {timeout_secs}s"
)));
}
tokio::time::sleep(std::time::Duration::from_secs(3)).await;
}
}
async fn find_code_updated_since(
ctx: &ExerciseCtx,
last_seen: &mut subxt::utils::H256,
) -> Result<Option<subxt::utils::H256>> {
let tip = ctx.client.get_latest_block().await?;
if tip == *last_seen {
return Ok(None);
}
let mut chain = Vec::new();
let mut cursor = tip;
while cursor != *last_seen && chain.len() < 64 {
chain.push(cursor);
cursor = ctx.client.client().blocks().at(cursor).await?.header().parent_hash;
}
for block_hash in chain.into_iter().rev() {
let events = ctx.client.client().blocks().at(block_hash).await?.events().await?;
if events
.find_first::<quantus_subxt::api::system::events::CodeUpdated>()
.map_err(|e| QuantusError::Generic(format!("failed to decode events: {e:?}")))?
.is_some()
{
*last_seen = tip;
return Ok(Some(block_hash));
}
}
*last_seen = tip;
Ok(None)
}