-
Notifications
You must be signed in to change notification settings - Fork 6
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
Receiving of the block proposed event from the driver (#61)
* block proposed JSON parsing * block proposed event receiver * timeout for the error while waiting for the block proposed event * test fixed
- Loading branch information
1 parent
6e3800f
commit 9239df4
Showing
10 changed files
with
189 additions
and
18 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,46 @@ | ||
use crate::{taiko::Taiko, utils::node_message::NodeMessage}; | ||
use std::sync::Arc; | ||
use tokio::sync::mpsc::Sender; | ||
use tracing::{error, info}; | ||
|
||
pub struct BlockProposedEventReceiver { | ||
taiko: Arc<Taiko>, | ||
node_tx: Sender<NodeMessage>, | ||
} | ||
|
||
impl BlockProposedEventReceiver { | ||
pub fn new(taiko: Arc<Taiko>, node_tx: Sender<NodeMessage>) -> Self { | ||
Self { taiko, node_tx } | ||
} | ||
|
||
pub async fn start(receiver: Self) { | ||
tokio::spawn(async move { | ||
receiver.check_for_events().await; | ||
}); | ||
} | ||
|
||
pub async fn check_for_events(&self) { | ||
loop { | ||
let block_proposed_event = self.taiko.wait_for_block_proposed_event().await; | ||
match block_proposed_event { | ||
Ok(block_proposed) => { | ||
info!( | ||
"Received block proposed event for block: {}", | ||
block_proposed.block_id | ||
); | ||
if let Err(e) = self | ||
.node_tx | ||
.send(NodeMessage::BlockProposed(block_proposed)) | ||
.await | ||
{ | ||
error!("Error sending block proposed event by channel: {:?}", e); | ||
} | ||
} | ||
Err(e) => { | ||
error!("Error receiving block proposed event: {:?}", e); | ||
tokio::time::sleep(std::time::Duration::from_secs(10)).await; | ||
} | ||
} | ||
} | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,63 @@ | ||
use anyhow::Error; | ||
use serde::{Deserialize, Deserializer, Serialize}; | ||
use serde_json::Value; | ||
|
||
#[derive(Serialize, Deserialize, Debug)] | ||
#[serde(rename_all = "PascalCase")] | ||
pub struct BlockProposed { | ||
#[serde(rename = "BlockID")] | ||
pub block_id: u64, | ||
pub tx_list_hash: [u8; 32], | ||
#[serde(deserialize_with = "deserialize_proposer")] | ||
pub proposer: [u8; 20], | ||
} | ||
|
||
fn deserialize_proposer<'de, D>(deserializer: D) -> Result<[u8; 20], D::Error> | ||
where | ||
D: Deserializer<'de>, | ||
{ | ||
let s: String = Deserialize::deserialize(deserializer)?; | ||
let s = s.trim_start_matches("0x"); | ||
let bytes = hex::decode(s).map_err(serde::de::Error::custom)?; | ||
if bytes.len() != 20 { | ||
return Err(serde::de::Error::custom( | ||
"Invalid length for proposer address", | ||
)); | ||
} | ||
let mut array = [0u8; 20]; | ||
array.copy_from_slice(&bytes); | ||
Ok(array) | ||
} | ||
|
||
pub fn decompose_block_proposed_json(json_data: Value) -> Result<BlockProposed, Error> { | ||
let block_proposed: BlockProposed = serde_json::from_value(json_data)?; | ||
Ok(block_proposed) | ||
} | ||
|
||
#[cfg(test)] | ||
mod tests { | ||
use super::*; | ||
use serde_json::json; | ||
|
||
#[test] | ||
fn test_decompose_block_proposed_json() { | ||
let json_data = json!({ | ||
"BlockID":4321,"TxListHash":[12,34,56,78,90,12,34,56,78,90,12,34,56,78,90,12,34,56,78,90,12,34,56,78,90,12,34,56,78,90,12,34],"Proposer":"0x0000000000000000000000000000000000000008" | ||
}); | ||
|
||
let result = decompose_block_proposed_json(json_data).unwrap(); | ||
|
||
assert_eq!(result.block_id, 4321); | ||
assert_eq!( | ||
result.tx_list_hash, | ||
[ | ||
12, 34, 56, 78, 90, 12, 34, 56, 78, 90, 12, 34, 56, 78, 90, 12, 34, 56, 78, 90, 12, | ||
34, 56, 78, 90, 12, 34, 56, 78, 90, 12, 34, | ||
] | ||
); | ||
assert_eq!( | ||
result.proposer, | ||
[0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 0, 8,] | ||
); | ||
} | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,4 +1,6 @@ | ||
pub mod block_proposed; | ||
pub mod commit; | ||
pub mod config; | ||
pub mod node_message; | ||
pub mod rpc_client; | ||
pub mod rpc_server; |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,7 @@ | ||
use super::block_proposed::BlockProposed; | ||
|
||
#[derive(Debug)] | ||
pub enum NodeMessage { | ||
BlockProposed(BlockProposed), | ||
P2P(String), | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters