diff --git a/.gitignore b/.gitignore index 6d9263a0..13c8f4ab 100644 --- a/.gitignore +++ b/.gitignore @@ -34,6 +34,4 @@ circuits/*/*/*.bin.in **/*.out **/*/output.data* -proof-builder-rpc/*.ckpt - -runner \ No newline at end of file +proof-builder-rpc/*.ckpt \ No newline at end of file diff --git a/node/README.md b/node/README.md index 9e65ac18..8bf705ea 100644 --- a/node/README.md +++ b/node/README.md @@ -872,6 +872,7 @@ real graph raw data in the database. | `GOAT_PRIVATE_KEY` | Conditional | GOAT chain private key (required for Committee) | - | | `GOAT_ADDRESS` | Conditional | GOAT address (required for Operator/Verifier) | - | | `ENABLE_RELAYER` | No | Enable relayer mode for Committee nodes | `false` | +| `ENABLE_BABE_SETUP_STATE_CLEANUP` | No | Enable scheduled BABE setup state cleanup for Operator/Verifier nodes | `false` | | `BTC_CHAIN_URL` | No | Bitcoin Esplora API endpoint | Public Esplora | | `MARA_SLIPSTREAM_API_URL` | No | MARA slipstream API base URL (used for non-standard tx broadcast) | mainnet: `https://slipstream.mara.com/api`; testnet4: `https://teststream.mara.com/api` | | `GOAT_PROOF_BUILD_URL` | No | Proof Builder RPC endpoint | - | diff --git a/node/src/action.rs b/node/src/action.rs index b350f87a..96095162 100644 --- a/node/src/action.rs +++ b/node/src/action.rs @@ -112,7 +112,7 @@ pub struct CutCircuits { pub verifier_index: usize, pub selected_circuit_indexes: Vec, } -#[derive(Serialize, Deserialize, Clone)] +#[derive(Debug, Serialize, Deserialize, Clone, PartialEq, Eq)] pub struct SolderingProofReady { pub instance_id: Uuid, pub graph_id: Uuid, diff --git a/node/src/env.rs b/node/src/env.rs index 355a3c46..072c7161 100644 --- a/node/src/env.rs +++ b/node/src/env.rs @@ -34,6 +34,7 @@ pub const ENV_GOAT_SEQUENCER_SET_MULTI_SIG_VERIFIER_ADDRESS: &str = "ENV_GOAT_SEQUENCER_SET_MULTI_SIG_VERIFIER_ADDRESS"; pub const ENV_ENABLE_RELAYER: &str = "ENABLE_RELAYER"; pub const ENV_ENABLE_UPDATE_SPV_CONTRACT: &str = "ENABLE_UPDATE_SPV_CONTRACT"; +pub const ENV_ENABLE_BABE_SETUP_STATE_CLEANUP: &str = "ENABLE_BABE_SETUP_STATE_CLEANUP"; pub const ENV_BTC_BLOCK_CONFIRMS: &str = "BTC_BLOCK_CONFIRMS"; pub const ENV_MARA_SLIPSTREAM_API_URL: &str = "MARA_SLIPSTREAM_API_URL"; pub const DEFAULT_MARA_SLIPSTREAM_MAINNET_API_URL: &str = "https://slipstream.mara.com/api"; @@ -241,6 +242,13 @@ pub fn is_enable_update_spv_contract() -> bool { } } +pub fn is_enable_babe_setup_state_cleanup() -> bool { + match std::env::var(ENV_ENABLE_BABE_SETUP_STATE_CLEANUP) { + Ok(value) => value.eq_ignore_ascii_case("true"), + Err(_) => false, + } +} + pub fn get_btc_block_confirms(network: Network) -> u32 { std::env::var(ENV_BTC_BLOCK_CONFIRMS) .ok() diff --git a/node/src/handle.rs b/node/src/handle.rs index 0bd345f5..e4de0287 100644 --- a/node/src/handle.rs +++ b/node/src/handle.rs @@ -17,12 +17,12 @@ use bitcoin::{PublicKey, XOnlyPublicKey}; use bitvm_lib::actors::Actor; use bitvm_lib::babe_adapter::{ BABE_M_CC, BABE_N_CC, BabeBundleBuilder, BabeChallengeAssertWitness, BabeProverState, - CACSetupPackage, ChallengeAssertWitnessRaw, CompactSolderingProofPayload, TxAssertWitness, - WOTS_SIG_COUNT, assert_wots_message, build_assert_witness, build_real_challenge_assert_witness, - build_real_setup_package, derive_finalized_indices, expand_compact_soldering_proof_payload, - extract_gc_circuit_data, open_real_setup_and_solder, - recover_operator_proof_from_assert_witness, recover_real_wrongly_challenged_witness, - verify_real_setup, + CACSetupPackage, ChallengeAssertWitnessRaw, CompactSolderingProofPayload, + FinalizedInstanceData, SolderingData, TxAssertWitness, WOTS_SIG_COUNT, assert_wots_message, + build_assert_witness, build_real_challenge_assert_witness, build_real_setup_package, + derive_finalized_indices, expand_compact_soldering_proof_payload, extract_gc_circuit_data, + open_real_setup_and_solder, recover_operator_proof_from_assert_witness, + recover_real_wrongly_challenged_witness, verify_real_setup, }; use bitvm_lib::committee::*; use bitvm_lib::keys::*; @@ -641,7 +641,8 @@ fn record_candidate_gc_data( verifier_index: usize, setup_package: &CACSetupPackage, gc_data: BitvmGcCircuitData, - prover_state: BabeProverState, + prover_state: &BabeProverState, + soldering_proof_ready: SolderingProofReady, ) -> Result>> { let frozen = state .frozen_verifier_pubkeys @@ -673,6 +674,9 @@ fn record_candidate_gc_data( if candidate.selected_circuit_indexes != prover_state.soldering.finalized_indices { bail!("BABE prover state finalized indices do not match selected verifier cut"); } + if soldering_proof_ready.verifier_index != verifier_index { + bail!("soldering proof reference verifier index does not match selected verifier slot"); + } if prover_state.finalized.len() != BABE_M_CC || prover_state.h_msgs.len() != BABE_M_CC { bail!("BABE prover state must contain exactly {BABE_M_CC} finalized instances and hashes"); } @@ -684,16 +688,16 @@ fn record_candidate_gc_data( { bail!("conflicting GC slot received for selected verifier"); } - if let Some(existing) = &candidate.prover_state - && existing != &prover_state + if let Some(existing) = &candidate.soldering_proof_ready + && existing != &soldering_proof_ready { - bail!("conflicting BABE prover state received for selected verifier"); + bail!("conflicting soldering proof reference received for selected verifier"); } if candidate.gc_data.is_none() { candidate.gc_data = Some(gc_data); } - if candidate.prover_state.is_none() { - candidate.prover_state = Some(prover_state); + if candidate.soldering_proof_ready.is_none() { + candidate.soldering_proof_ready = Some(soldering_proof_ready); } if state.candidates.iter().any(|candidate| candidate.gc_data.is_none()) { return Ok(None); @@ -701,6 +705,24 @@ fn record_candidate_gc_data( Ok(Some(state.candidates.iter().map(|candidate| candidate.gc_data.clone().unwrap()).collect())) } +fn build_babe_prover_state( + setup_package: &CACSetupPackage, + finalized: Vec, + soldering: SolderingData, +) -> Result { + let h_msgs = finalized + .iter() + .map(|finalized| { + setup_package + .commits + .get(finalized.index) + .map(|commit| commit.h_msg) + .ok_or_else(|| anyhow!("finalized BABE instance index is out of range")) + }) + .collect::>>()?; + Ok(BabeProverState { package: setup_package.clone(), finalized, soldering, h_msgs }) +} + fn validate_verifier_slot_lengths( verifier_index: usize, gc_data_len: usize, @@ -1213,11 +1235,8 @@ async fn handle_init_graph_verifier( verifier_pubkey, setup_package, private_state, - verifier_index: None, finalized_indices: vec![], - opened: vec![], - finalized: vec![], - soldering: None, + soldering_proof_ready: None, } }; @@ -1262,8 +1281,11 @@ async fn handle_gen_circuits_operator( return Ok(()); } - if setup_package.commits.is_empty() { - bail!("GenCircuits setup package has no commitments"); + if setup_package.commits.len() != BABE_N_CC { + bail!( + "invalid GenCircuits setup package commitment count: expected {BABE_N_CC}, got {}", + setup_package.commits.len() + ); } let mut state = load_babe_setup_state(ctx.local_db, instance_id, graph_id)?.unwrap_or_default(); @@ -1295,7 +1317,7 @@ async fn handle_gen_circuits_operator( verifier_index: None, selected_circuit_indexes: vec![], gc_data: None, - prover_state: None, + soldering_proof_ready: None, }); } @@ -1371,30 +1393,26 @@ async fn handle_cut_circuits_verifier( return Ok(()); } - if let Some(saved_index) = verifier_state.verifier_index { - if saved_index != verifier_index { + if let Some(soldering_proof_ready) = verifier_state.soldering_proof_ready.clone() { + if soldering_proof_ready.verifier_index != verifier_index { bail!( - "CutCircuits verifier index {verifier_index} conflicts with persisted slot {saved_index}" + "CutCircuits verifier index {verifier_index} conflicts with persisted slot {}", + soldering_proof_ready.verifier_index ); } if verifier_state.finalized_indices != *selected_circuit_indexes { bail!("CutCircuits finalized indices conflict with persisted selection"); } - if let Some(soldering) = verifier_state.soldering.clone() { - let setup_state = verifier_state; - send_soldering_proof_to_operator( - ctx.swarm, - instance_id, - graph_id, - verifier_index, - &setup_state.opened, - &setup_state.finalized, - &soldering, - ) - .await?; + send_to_peer( + ctx.swarm, + GOATMessage::new( + Actor::Operator, + GOATMessageContent::SolderingProofReady(soldering_proof_ready), + ), + ) + .await?; - return Ok(()); - } + return Ok(()); } let setup_package = verifier_state.setup_package.clone(); @@ -1423,24 +1441,28 @@ async fn handle_cut_circuits_verifier( .await .context("real BABE opening task failed")??; - verifier_state.verifier_index = Some(verifier_index); + let soldering_proof_ready = save_soldering_proof_payload( + instance_id, + graph_id, + verifier_index, + &opened, + &finalized, + &soldering, + ) + .await?; verifier_state.finalized_indices = selected_circuit_indexes.clone(); - verifier_state.opened = opened.clone(); - verifier_state.finalized = finalized.clone(); - verifier_state.soldering = Some(soldering.clone()); + verifier_state.soldering_proof_ready = Some(soldering_proof_ready.clone()); update_babe_setup_state(ctx.local_db, instance_id, graph_id, |state| { state.verifier = Some(verifier_state); })?; - send_soldering_proof_to_operator( + send_to_peer( ctx.swarm, - instance_id, - graph_id, - verifier_index, - &opened, - &finalized, - &soldering, + GOATMessage::new( + Actor::Operator, + GOATMessageContent::SolderingProofReady(soldering_proof_ready), + ), ) .await?; @@ -1459,6 +1481,8 @@ async fn handle_soldering_proof_ready_operator( if total_len == 0 { bail!("SolderingProofReady total_len must be greater than zero"); } + let soldering_proof_ready = + SolderingProofReady { instance_id, graph_id, verifier_index, payload_hash, total_len }; let operator_master_key = OperatorMasterKey::new(get_bitvm_key()?); let local_operator_pubkey = operator_master_key.master_keypair().public_key().into(); if !pending_graph_belongs_to_operator( @@ -1548,29 +1572,16 @@ async fn handle_soldering_proof_ready_operator( payload_hash = %soldering_payload_hash_hex(&payload_hash), "start processing soldering proof payload" ); - handle_soldering_proof_payload_operator( - ctx, - instance_id, - graph_id, - verifier_index, - payload_hash, - total_len, - &payload, - ) - .await + handle_soldering_proof_payload_operator(ctx, &soldering_proof_ready, &payload).await } -#[allow(clippy::too_many_arguments)] -pub(crate) async fn handle_soldering_proof_payload_operator( - ctx: &mut HandlerContext<'_>, - instance_id: Uuid, - graph_id: Uuid, - verifier_index: usize, - payload_hash: [u8; 32], - total_len: usize, +fn decode_soldering_proof_payload( + soldering_proof_ready: &SolderingProofReady, payload: &[u8], -) -> Result<()> { - if payload.len() != total_len { +) -> Result { + let SolderingProofReady { instance_id, graph_id, verifier_index, payload_hash, total_len } = + soldering_proof_ready; + if payload.len() != *total_len { tracing::warn!( instance_id = %instance_id, graph_id = %graph_id, @@ -1586,7 +1597,7 @@ pub(crate) async fn handle_soldering_proof_payload_operator( ); } let actual_hash = soldering_payload_hash(payload); - if actual_hash != payload_hash { + if actual_hash != *payload_hash { tracing::warn!( instance_id = %instance_id, graph_id = %graph_id, @@ -1597,21 +1608,28 @@ pub(crate) async fn handle_soldering_proof_payload_operator( ); bail!("SolderingProof payload hash mismatch"); } - let payload: CompactSolderingProofPayload = - bincode::deserialize(payload).context("deserialize compact soldering proof payload")?; - handle_compact_soldering_proof_operator(ctx, instance_id, graph_id, verifier_index, payload) - .await + bincode::deserialize(payload).context("deserialize compact soldering proof payload") +} + +pub(crate) async fn handle_soldering_proof_payload_operator( + ctx: &mut HandlerContext<'_>, + soldering_proof_ready: &SolderingProofReady, + payload: &[u8], +) -> Result<()> { + let payload = decode_soldering_proof_payload(soldering_proof_ready, payload)?; + handle_compact_soldering_proof_operator(ctx, soldering_proof_ready, payload).await } // verify Verifier SolderingProof, build Graph and broadcast CreateGraph. -#[tracing::instrument(level = "info", skip_all, fields(instance_id = %instance_id, graph_id = %graph_id))] +#[tracing::instrument(level = "info", skip_all, fields(instance_id = %soldering_proof_ready.instance_id, graph_id = %soldering_proof_ready.graph_id))] async fn handle_compact_soldering_proof_operator( ctx: &mut HandlerContext<'_>, - instance_id: Uuid, - graph_id: Uuid, - verifier_index: usize, + soldering_proof_ready: &SolderingProofReady, payload: CompactSolderingProofPayload, ) -> Result<()> { + let instance_id = soldering_proof_ready.instance_id; + let graph_id = soldering_proof_ready.graph_id; + let verifier_index = soldering_proof_ready.verifier_index; let operator_master_key = OperatorMasterKey::new(get_bitvm_key()?); let local_operator_pubkey = operator_master_key.master_keypair().public_key().into(); if !pending_graph_belongs_to_operator( @@ -1685,22 +1703,16 @@ async fn handle_compact_soldering_proof_operator( bail!("each verifier must contribute exactly {BABE_M_CC} finalized BABE instances"); } let epk = &setup_package.commits[finalized[0].index].epk; - let h_msgs: Vec<[u8; 20]> = - finalized.iter().map(|f| setup_package.commits[f.index].h_msg).collect(); - let gc_data = extract_gc_circuit_data(verifier_pubkey, epk, &h_msgs)?; - let prover_state = BabeProverState { - package: setup_package.clone(), - finalized, - soldering, - h_msgs: gc_data.final_msg_hashlocks.clone(), - }; + let prover_state = build_babe_prover_state(&setup_package, finalized, soldering)?; + let gc_data = extract_gc_circuit_data(verifier_pubkey, epk, &prover_state.h_msgs)?; let Some(bitvm_gc_circuit_datas) = record_candidate_gc_data( operator_state, verifier_pubkey, verifier_index, &setup_package, gc_data, - prover_state, + &prover_state, + soldering_proof_ready.clone(), )? else { save_babe_setup_state(ctx.local_db, instance_id, graph_id, &state)?; @@ -4267,7 +4279,13 @@ async fn handle_assert_sent_verifier( ); return Ok(()); }; - if saved_verifier_state.verifier_index != Some(verifier_index) { + let Some(soldering_proof_ready) = saved_verifier_state.soldering_proof_ready.as_ref() else { + tracing::warn!( + "Ignore AssertSent for {instance_id}:{graph_id}: missing soldering proof reference" + ); + return Ok(()); + }; + if soldering_proof_ready.verifier_index != verifier_index { tracing::warn!( "Ignore AssertSent for {instance_id}:{graph_id}: local setup slot does not match graph owner slot" ); @@ -4398,10 +4416,24 @@ async fn handle_challenge_assert_sent_operator( .iter() .find(|candidate| candidate.verifier_index == Some(verifier_index)) .ok_or_else(|| anyhow!("missing BABE prover state for verifier slot {verifier_index}"))?; - let prover_state = candidate - .prover_state - .as_ref() - .ok_or_else(|| anyhow!("missing BABE prover state for verifier slot {verifier_index}"))?; + let soldering_proof_ready = candidate.soldering_proof_ready.as_ref().ok_or_else(|| { + anyhow!("missing soldering proof reference for verifier slot {verifier_index}") + })?; + let store_base_path = get_soldering_proof_payload_store_path()?; + let payload_path = soldering_proof_payload_store_path( + &store_base_path, + instance_id, + graph_id, + verifier_index, + &soldering_proof_ready.payload_hash, + )?; + let payload = read_soldering_proof_store_payload(&payload_path) + .await + .with_context(|| format!("read soldering proof payload from {payload_path}"))?; + let payload = decode_soldering_proof_payload(soldering_proof_ready, &payload)?; + let (_, finalized, soldering) = expand_compact_soldering_proof_payload(payload) + .context("expand compact soldering proof payload for challenge")?; + let prover_state = build_babe_prover_state(&candidate.setup_package, finalized, soldering)?; let assert_witness = TxAssertWitness { wots_sig: challenge_witness.witness.wots_sig.clone(), pi2: operator_assert_witness.pi2, @@ -4416,7 +4448,7 @@ async fn handle_challenge_assert_sent_operator( anyhow!("cannot recover dynamic input from challenge witness WOTS signature") })?; let wrongly_challenged_witness = recover_real_wrongly_challenged_witness( - prover_state, + &prover_state, &challenge_witness, &proof, vk, @@ -5109,6 +5141,16 @@ mod tests { (package, gc_data, prover_state) } + fn soldering_proof_ready(hash_byte: u8) -> SolderingProofReady { + SolderingProofReady { + instance_id: Uuid::nil(), + graph_id: Uuid::nil(), + verifier_index: 0, + payload_hash: [hash_byte; 32], + total_len: 1, + } + } + fn operator_state(package: CACSetupPackage) -> OperatorBabeSetupState { OperatorBabeSetupState { frozen_verifier_pubkeys: Some(vec![verifier_pubkey()]), @@ -5118,7 +5160,7 @@ mod tests { verifier_index: Some(0), selected_circuit_indexes: (0..BABE_M_CC).collect(), gc_data: None, - prover_state: None, + soldering_proof_ready: None, }], asserted_operator_proof: None, } @@ -5135,7 +5177,7 @@ mod tests { verifier_index: None, selected_circuit_indexes: vec![], gc_data: None, - prover_state: None, + soldering_proof_ready: None, }], asserted_operator_proof: None, }; @@ -5146,9 +5188,10 @@ mod tests { } #[test] - fn candidate_records_one_slot_and_full_prover_state_idempotently() { + fn candidate_records_one_slot_and_payload_reference_idempotently() { let (package, gc_data, prover_state) = gc_submission(); let mut state = operator_state(package.clone()); + let soldering_proof_ready = soldering_proof_ready(1); let graph_data = record_candidate_gc_data( &mut state, @@ -5156,14 +5199,18 @@ mod tests { 0, &package, gc_data.clone(), - prover_state.clone(), + &prover_state, + soldering_proof_ready.clone(), ) .unwrap() .unwrap(); assert_eq!(graph_data.len(), 1); assert_eq!(graph_data[0].final_msg_hashlocks.len(), BABE_M_CC); - assert_eq!(state.candidates[0].prover_state.as_ref().unwrap().finalized.len(), BABE_M_CC); + assert_eq!( + state.candidates[0].soldering_proof_ready.as_ref(), + Some(&soldering_proof_ready) + ); let duplicate = record_candidate_gc_data( &mut state, @@ -5171,7 +5218,8 @@ mod tests { 0, &package, gc_data.clone(), - prover_state.clone(), + &prover_state, + soldering_proof_ready.clone(), ) .unwrap() .unwrap(); @@ -5186,21 +5234,26 @@ mod tests { 0, &package, gc_data, - conflict, + &conflict, + soldering_proof_ready, ) .is_err() ); } #[test] - fn legacy_operator_candidate_without_prover_state_deserializes() { - let (package, _, _) = gc_submission(); - let mut value = serde_json::to_value(operator_state(package)).unwrap(); - value["candidates"][0].as_object_mut().unwrap().remove("prover_state"); + fn rebuilds_prover_state_from_persisted_payload_data() { + let package = build_setup_package(BABE_M_CC + 1).unwrap(); + let selected = (0..BABE_M_CC).collect::>(); + let (_, finalized, soldering) = open_and_solder(&package, &selected).unwrap(); - let restored: OperatorBabeSetupState = serde_json::from_value(value).unwrap(); + let restored = + build_babe_prover_state(&package, finalized.clone(), soldering.clone()).unwrap(); - assert!(restored.candidates[0].prover_state.is_none()); + assert_eq!(restored.package, package); + assert_eq!(restored.finalized, finalized); + assert_eq!(restored.soldering, soldering); + assert_eq!(restored.h_msgs.len(), BABE_M_CC); } // ── collect_ack_txins ───────────────────────────────────────────────────── diff --git a/node/src/middleware/behaviour.rs b/node/src/middleware/behaviour.rs index f7773fe7..378d5371 100644 --- a/node/src/middleware/behaviour.rs +++ b/node/src/middleware/behaviour.rs @@ -21,7 +21,7 @@ impl AllBehaviours { // .unwrap(); let gossipsub_config = gossipsub::ConfigBuilder::default() - .max_transmit_size(4194304) // 4 MB + .max_transmit_size(16 * 1024 * 1024) // 16 MB .build() .map_err(io::Error::other) .unwrap(); diff --git a/node/src/scheduled_tasks/babe_setup_state_cleanup_task.rs b/node/src/scheduled_tasks/babe_setup_state_cleanup_task.rs new file mode 100644 index 00000000..086fca53 --- /dev/null +++ b/node/src/scheduled_tasks/babe_setup_state_cleanup_task.rs @@ -0,0 +1,13 @@ +use crate::env::get_soldering_proof_payload_store_path; +use crate::utils::cleanup_babe_setup_states; +use store::localdb::LocalDB; +use tracing::info; + +pub async fn babe_setup_state_cleanup_monitor(local_db: &LocalDB) -> anyhow::Result<()> { + let payload_store = get_soldering_proof_payload_store_path()?; + let deleted = cleanup_babe_setup_states(local_db, &payload_store).await?; + if deleted > 0 { + info!(deleted, "BABE setup state cleanup monitor finished"); + } + Ok(()) +} diff --git a/node/src/scheduled_tasks/mod.rs b/node/src/scheduled_tasks/mod.rs index 7544cb99..e590752e 100644 --- a/node/src/scheduled_tasks/mod.rs +++ b/node/src/scheduled_tasks/mod.rs @@ -1,3 +1,4 @@ +mod babe_setup_state_cleanup_task; mod event_watch_task; pub mod graph_maintenance_tasks; pub mod instance_maintenance_tasks; @@ -6,7 +7,11 @@ mod sequencer_set_hash_monitor_task; mod spv_maintenance_tasks; use crate::action::GOATMessageContent; -use crate::env::{get_maintenance_run_timeout_secs, is_enable_update_spv_contract, is_relayer}; +use crate::env::{ + get_maintenance_run_timeout_secs, is_enable_babe_setup_state_cleanup, + is_enable_update_spv_contract, is_relayer, +}; +use crate::scheduled_tasks::babe_setup_state_cleanup_task::babe_setup_state_cleanup_monitor; use crate::scheduled_tasks::graph_maintenance_tasks::{ detect_init_withdraw_call, detect_kickoff, detect_take1_or_challenge, process_graph_challenge, }; @@ -55,6 +60,13 @@ async fn run( let btc_client = btc_client.as_ref(); let goat_client = goat_client.as_ref(); + if is_enable_babe_setup_state_cleanup() + && matches!(&actor, Actor::Verifier | Actor::Operator | Actor::All) + && let Err(err) = babe_setup_state_cleanup_monitor(local_db).await + { + warn!("babe_setup_state_cleanup_monitor, err {:?}", err) + } + if (actor == Actor::Operator || is_relayer()) && let Err(err) = node_available_pbtc_update_monitor(local_db, goat_client).await { diff --git a/node/src/soldering_payload_store.rs b/node/src/soldering_payload_store.rs index fae3edc6..aa683546 100644 --- a/node/src/soldering_payload_store.rs +++ b/node/src/soldering_payload_store.rs @@ -109,6 +109,19 @@ pub(crate) async fn read_soldering_proof_store_payload(path: &str) -> Result Result<()> { + if is_soldering_proof_s3_path(path) { + bail!("S3 soldering proof cleanup is not supported"); + } + match tokio::fs::remove_file(path.trim()).await { + Ok(()) => Ok(()), + Err(err) if err.kind() == std::io::ErrorKind::NotFound => Ok(()), + Err(err) => { + Err(err).with_context(|| format!("delete soldering proof payload {}", path.trim())) + } + } +} + #[cfg(test)] mod tests { use super::*; diff --git a/node/src/utils.rs b/node/src/utils.rs index 2b1c30f4..f2e61058 100644 --- a/node/src/utils.rs +++ b/node/src/utils.rs @@ -8,8 +8,8 @@ use crate::error::SpecialError; use crate::middleware::AllBehaviours; use crate::rpc_service::current_time_secs; use crate::soldering_payload_store::{ - is_soldering_proof_s3_path, soldering_proof_payload_store_path, - write_soldering_proof_store_payload, + delete_soldering_proof_store_payload, is_soldering_proof_s3_path, + soldering_proof_payload_store_path, write_soldering_proof_store_payload, }; use alloy::primitives::{Address as EvmAddress, Signature as EvmSignature}; use alloy::signers::Signer; @@ -85,8 +85,8 @@ use crate::scheduled_tasks::graph_maintenance_tasks::{ }; use bitcoin_light_client_circuit::hash_operator_constant; use bitvm_lib::babe_adapter::{ - BabeProverState, BabeVerifierPrivateState, CACSetupPackage, FinalizedInstanceData, - SolderingData, compact_soldering_proof_payload, + BabeVerifierPrivateState, CACSetupPackage, FinalizedInstanceData, SolderingData, + compact_soldering_proof_payload, }; use bitvm_lib::transactions::base::BaseTransaction; use client::goat_chain::{DisproveTxType, GraphData, PeginStatus, WithdrawStatus}; @@ -4959,7 +4959,7 @@ pub(crate) fn babe_setup_state_root(local_db: &LocalDB) -> PathBuf { } fn babe_setup_state_path(local_db: &LocalDB, instance_id: Uuid, graph_id: Uuid) -> PathBuf { - babe_setup_state_root(local_db).join(instance_id.to_string()).join(format!("{graph_id}.json")) + babe_setup_state_root(local_db).join(instance_id.to_string()).join(format!("{graph_id}.bin")) } pub(crate) fn soldering_payload_hash(payload: &[u8]) -> [u8; 32] { @@ -4983,15 +4983,14 @@ pub(crate) async fn pending_graph_belongs_to_operator( } #[allow(clippy::too_many_arguments)] -pub(crate) async fn send_soldering_proof_to_operator( - swarm: &mut Swarm, +pub(crate) async fn save_soldering_proof_payload( instance_id: Uuid, graph_id: Uuid, verifier_index: usize, opened: &[(usize, u64)], finalized: &[FinalizedInstanceData], soldering: &SolderingData, -) -> Result<()> { +) -> Result { let compact_payload = compact_soldering_proof_payload(opened, finalized, soldering)?; let payload = bincode::serialize(&compact_payload) .context("serialize compact soldering proof payload")?; @@ -5026,23 +5025,17 @@ pub(crate) async fn send_soldering_proof_to_operator( total_len, payload_hash = %soldering_payload_hash_hex(&payload_hash), payload_path = %payload_path, - "send compact soldering proof ready from payload store" + "saved compact soldering proof to payload store" ); - let message_content = GOATMessageContent::SolderingProofReady(SolderingProofReady { - instance_id, - graph_id, - verifier_index, - payload_hash, - total_len, - }); - send_to_peer(swarm, GOATMessage::new(Actor::Operator, message_content)).await?; - - Ok(()) + Ok(SolderingProofReady { instance_id, graph_id, verifier_index, payload_hash, total_len }) } fn load_babe_setup_state_from_path(path: &Path) -> Result> { match std::fs::read(path) { - Ok(bytes) => Ok(Some(serde_json::from_slice(&bytes)?)), + Ok(bytes) => Ok(Some( + bincode::deserialize(&bytes) + .with_context(|| format!("decode BABE setup state {}", path.display()))?, + )), Err(err) if err.kind() == std::io::ErrorKind::NotFound => Ok(None), Err(err) => Err(err).with_context(|| format!("read BABE setup state {}", path.display())), } @@ -5053,9 +5046,15 @@ fn save_babe_setup_state_to_path(path: &Path, state: &BabeSetupState) -> Result< std::fs::create_dir_all(parent) .with_context(|| format!("create BABE setup state dir {}", parent.display()))?; } - let bytes = serde_json::to_vec_pretty(state)?; - std::fs::write(path, bytes) - .with_context(|| format!("write BABE setup state {}", path.display())) + let bytes = bincode::serialize(state)?; + let mut temp_path = path.as_os_str().to_os_string(); + temp_path.push(".tmp"); + let temp_path = PathBuf::from(temp_path); + std::fs::write(&temp_path, bytes) + .with_context(|| format!("write BABE setup state {}", temp_path.display()))?; + std::fs::rename(&temp_path, path).with_context(|| { + format!("move BABE setup state {} to {}", temp_path.display(), path.display()) + }) } pub(crate) fn load_babe_setup_state( @@ -5087,6 +5086,92 @@ pub(crate) fn update_babe_setup_state( Ok(state) } +pub(crate) async fn cleanup_babe_setup_states( + local_db: &LocalDB, + payload_store: &str, +) -> Result { + let root = babe_setup_state_root(local_db); + let instance_dirs = match std::fs::read_dir(&root) { + Ok(entries) => entries, + Err(err) if err.kind() == std::io::ErrorKind::NotFound => return Ok(0), + Err(err) => { + return Err(err) + .with_context(|| format!("read BABE setup state root {}", root.display())); + } + }; + let mut storage = local_db.acquire().await?; + let mut deleted = 0; + for instance_dir in instance_dirs { + let instance_dir = instance_dir?; + if !instance_dir.file_type()?.is_dir() { + continue; + } + let Ok(instance_id) = Uuid::parse_str(&instance_dir.file_name().to_string_lossy()) else { + continue; + }; + for state_file in std::fs::read_dir(instance_dir.path())? { + let state_file = state_file?; + let path = state_file.path(); + if path.extension().and_then(|extension| extension.to_str()) != Some("bin") { + continue; + } + let Some(graph_id) = path + .file_stem() + .and_then(|stem| stem.to_str()) + .and_then(|stem| Uuid::parse_str(stem).ok()) + else { + continue; + }; + let Some(graph) = storage.find_graph(&graph_id).await? else { + continue; + }; + let Ok(status) = GraphStatus::from_str(&graph.status) else { + warn!( + "Skip BABE setup cleanup for {graph_id}: invalid graph status {}", + graph.status + ); + continue; + }; + if graph.instance_id != instance_id || !status.is_closed() { + continue; + } + let Some(state) = load_babe_setup_state_from_path(&path)? else { + continue; + }; + let mut payloads = Vec::new(); + if let Some(ready) = state.verifier.and_then(|verifier| verifier.soldering_proof_ready) + { + payloads.push(ready); + } + if let Some(operator) = state.operator { + payloads.extend( + operator + .candidates + .into_iter() + .filter_map(|candidate| candidate.soldering_proof_ready), + ); + } + for ready in payloads { + if ready.instance_id != instance_id || ready.graph_id != graph_id { + bail!("BABE setup payload reference does not match state path"); + } + let payload_path = soldering_proof_payload_store_path( + payload_store, + instance_id, + graph_id, + ready.verifier_index, + &ready.payload_hash, + )?; + delete_soldering_proof_store_payload(&payload_path).await?; + } + std::fs::remove_file(&path) + .with_context(|| format!("delete BABE setup state {}", path.display()))?; + deleted += 1; + } + } + Ok(deleted) +} + #[derive(Clone, Default, Serialize, Deserialize)] pub struct BabeSetupState { pub verifier: Option, @@ -5098,11 +5183,8 @@ pub struct VerifierBabeSetupState { pub verifier_pubkey: PublicKey, pub setup_package: CACSetupPackage, pub private_state: BabeVerifierPrivateState, - pub verifier_index: Option, pub finalized_indices: Vec, - pub opened: Vec<(usize, u64)>, - pub finalized: Vec, - pub soldering: Option, + pub soldering_proof_ready: Option, } #[derive(Clone, Serialize, Deserialize)] @@ -5112,8 +5194,7 @@ pub struct OperatorVerifierCandidate { pub verifier_index: Option, pub selected_circuit_indexes: Vec, pub gc_data: Option, - #[serde(default)] - pub prover_state: Option, + pub soldering_proof_ready: Option, } #[derive(Clone, Serialize, Deserialize)]