diff --git a/src/connection.rs b/src/connection.rs index 0197602..191365e 100644 --- a/src/connection.rs +++ b/src/connection.rs @@ -221,7 +221,7 @@ impl Connection { // Should do the loop here - let (tx_tcp, rx_tcp) = tokio::sync::mpsc::channel::>(32); + let (tx_tcp, rx_tcp) = tokio::sync::mpsc::channel::>(256); let (tx_write, mut rx_write) = tokio::sync::mpsc::channel::>(32); let handle = tokio::spawn(async move { diff --git a/src/database.rs b/src/database.rs index 6aecd03..9c50816 100644 --- a/src/database.rs +++ b/src/database.rs @@ -77,13 +77,13 @@ pub fn create_tables(schema_name: &String, postgres_client: &mut Client) { ); CREATE TABLE IF NOT EXISTS {schema_name}.withdrawals ( block_hash BYTEA NOT NULL, - index BIGINT NOT NULL, - validator_index BIGINT NOT NULL, + index NUMERIC(78) NOT NULL, + validator_index NUMERIC(78) NOT NULL, address BYTEA NOT NULL, - amount BIGINT NOT NULL + amount NUMERIC(78) NOT NULL ); CREATE TABLE IF NOT EXISTS {schema_name}.receipts ( - txid BYTEA NOT NULL, + txid BYTEA, tx_type SMALLINT NOT NULL, post_state_or_status BYTEA NOT NULL, cumulative_gas BIGINT NOT NULL, @@ -240,17 +240,31 @@ pub fn save_blocks( ); withdrawals_string.push_str(&tmp); }); - receipts.iter().enumerate().for_each(|(i, r)| { + let (system_receipts, user_receipts): (Vec<_>, Vec<_>) = + receipts.iter().partition(|r| r.tx_type >= 128); + + for r in system_receipts { + let tmp = format!( + "\\N;{};\\\\x{};{};{}\n", // Important! We don't end with a ';' + r.tx_type, + hex::encode(&r.post_state_or_status), + r.cumulative_gas, + serde_json::to_value(&r.logs).unwrap(), + ); + receipts_string.push_str(&tmp); + } + + for (r, t) in user_receipts.iter().zip(txs.iter()) { let tmp = format!( "\\\\x{};{};\\\\x{};{};{}\n", // Important! We don't end with a ';' - hex::encode(&txs[i].txid), + hex::encode(&t.txid), r.tx_type, hex::encode(&r.post_state_or_status), r.cumulative_gas, serde_json::to_value(&r.logs).unwrap(), ); receipts_string.push_str(&tmp); - }); + } }); transactions_strings.push(transactions_string.clone()); diff --git a/src/main.rs b/src/main.rs index eb05a8f..4cbe336 100644 --- a/src/main.rs +++ b/src/main.rs @@ -318,23 +318,24 @@ fn main() { continue; } - let headers_body = loop { - match conn.rx_tcp.recv().await { - Some(body) if body[0].saturating_sub(16) == 4 => break Some(body), - Some(_) => continue, - None => break None, + let headers_body = tokio::time::timeout(Duration::from_secs(30), async { + loop { + match conn.rx_tcp.recv().await { + Some(body) if body[0].saturating_sub(16) == 4 => break Some(body), + Some(_) => continue, + None => break None, + } } - }; + }).await.unwrap_or(None); let headers_body = match headers_body { Some(b) => b, - None => continue, + None => { warn!("Timed out or connection closed waiting for block headers, retrying"); continue; } }; let headers = eth::parse_block_headers(&headers_body[1..]); if headers.is_empty() { - warn!("No block headers received"); - pool.push_back(conn); - break; + warn!("No block headers received, dropping peer and retrying"); + continue; } current_hash = headers.last().unwrap().parent_hash.to_vec(); @@ -364,21 +365,21 @@ fn main() { warn!("Connection lost during GetBlockBodies, will retry batch"); return Err(block_headers); } - loop { - match conn.rx_tcp.recv().await { - Some(body) if body[0].saturating_sub(16) == 6 => { - transactions - .extend(eth::parse_block_bodies(&body[1..])); - break; - } - Some(_) => continue, - None => { - warn!( - "Connection closed during GetBlockBodies, will retry batch" - ); - return Err(block_headers); + let ok = tokio::time::timeout(Duration::from_secs(30), async { + loop { + match conn.rx_tcp.recv().await { + Some(body) if body[0].saturating_sub(16) == 6 => { + transactions.extend(eth::parse_block_bodies(&body[1..])); + break true; + } + Some(_) => continue, + None => break false, } } + }).await.unwrap_or(false); + if !ok { + warn!("Timed out or connection closed during GetBlockBodies, will retry batch"); + return Err(block_headers); } } @@ -397,18 +398,21 @@ fn main() { warn!("Connection lost during GetReceipts, will retry batch"); return Err(block_headers); } - loop { - match conn.rx_tcp.recv().await { - Some(body) if body[0].saturating_sub(16) == 16 => { - receipts.extend(eth::parse_receipts(&body[1..])); - break; - } - Some(_) => continue, - None => { - warn!("Connection closed during GetReceipts, will retry batch"); - return Err(block_headers); + let ok = tokio::time::timeout(Duration::from_secs(30), async { + loop { + match conn.rx_tcp.recv().await { + Some(body) if body[0].saturating_sub(16) == 16 => { + receipts.extend(eth::parse_receipts(&body[1..])); + break true; + } + Some(_) => continue, + None => break false, } } + }).await.unwrap_or(false); + if !ok { + warn!("Timed out or connection closed during GetReceipts, will retry batch"); + return Err(block_headers); } } } @@ -481,16 +485,18 @@ fn main() { continue; } - let headers_body = loop { - match conn.rx_tcp.recv().await { - Some(body) if body[0].saturating_sub(16) == 4 => break Some(body), - Some(_) => continue, - None => break None, + let headers_body = tokio::time::timeout(Duration::from_secs(30), async { + loop { + match conn.rx_tcp.recv().await { + Some(body) if body[0].saturating_sub(16) == 4 => break Some(body), + Some(_) => continue, + None => break None, + } } - }; + }).await.unwrap_or(None); let headers_body = match headers_body { Some(b) => b, - None => continue, + None => { warn!("Timed out or connection closed waiting for block headers, retrying"); continue; } }; let headers = eth::parse_block_headers(&headers_body[1..]); @@ -621,19 +627,22 @@ fn main() { retry_headers = Some(block_headers); continue 'outer; } - loop { - match conn.rx_tcp.recv().await { - Some(body) if body[0].saturating_sub(16) == 6 => { - transactions.extend(eth::parse_block_bodies(&body[1..])); - break; - } - Some(_) => continue, - None => { - warn!("Connection closed during GetBlockBodies, will retry batch"); - retry_headers = Some(block_headers); - continue 'outer; + let ok = tokio::time::timeout(Duration::from_secs(30), async { + loop { + match conn.rx_tcp.recv().await { + Some(body) if body[0].saturating_sub(16) == 6 => { + transactions.extend(eth::parse_block_bodies(&body[1..])); + break true; + } + Some(_) => continue, + None => break false, } } + }).await.unwrap_or(false); + if !ok { + warn!("Timed out or connection closed during GetBlockBodies, will retry batch"); + retry_headers = Some(block_headers); + continue 'outer; } } @@ -649,19 +658,22 @@ fn main() { retry_headers = Some(block_headers); continue 'outer; } - loop { - match conn.rx_tcp.recv().await { - Some(body) if body[0].saturating_sub(16) == 16 => { - receipts.extend(eth::parse_receipts(&body[1..])); - break; - } - Some(_) => continue, - None => { - warn!("Connection closed during GetReceipts, will retry batch"); - retry_headers = Some(block_headers); - continue 'outer; + let ok = tokio::time::timeout(Duration::from_secs(30), async { + loop { + match conn.rx_tcp.recv().await { + Some(body) if body[0].saturating_sub(16) == 16 => { + receipts.extend(eth::parse_receipts(&body[1..])); + break true; + } + Some(_) => continue, + None => break false, } } + }).await.unwrap_or(false); + if !ok { + warn!("Timed out or connection closed during GetReceipts, will retry batch"); + retry_headers = Some(block_headers); + continue 'outer; } } } diff --git a/src/networks.rs b/src/networks.rs index 6b20942..0fcd631 100644 --- a/src/networks.rs +++ b/src/networks.rs @@ -120,6 +120,28 @@ impl Network { network_id: 0x2105, }; + // Berachain Mainnet + pub const BERACHAIN_MAINNET: Network = Network { + genesis_hash: [ + 213, 120, 25, 66, 33, 40, 218, 28, 68, 51, 159, 199, 149, 102, 98, 55, 140, 23, 226, + 33, 62, 102, 155, 66, 122, 201, 28, 209, 29, 252, 251, 56, + ], + head_td: 0, + fork_id: [0x701a097f, 0], + network_id: 0x138de, + }; + + // Berachain Bepolia + pub const BERACHAIN_BEPOLIA: Network = Network { + genesis_hash: [ + 2, 7, 102, 29, 227, 143, 14, 84, 186, 145, 200, 40, 96, 150, 231, 36, 134, 120, 76, + 121, 220, 106, 150, 129, 252, 72, 107, 56, 51, 92, 4, 47, + ], + head_td: 0, + fork_id: [0x2edd8d57, 0], + network_id: 0x138c5, + }; + pub fn find(network: &str) -> Result> { match network { "ethereum_ropsten" => Ok(Self::ETHEREUM_ROPSTEN), @@ -132,6 +154,8 @@ impl Network { "binance_mainnet" => Ok(Self::BINANCE_MAINNET), "polygon_mainnet" => Ok(Self::POLYGON_MAINNET), "base_mainnet" => Ok(Self::BASE_MAINNET), + "berachain_mainnet" => Ok(Self::BERACHAIN_MAINNET), + "berachain_bepolia" => Ok(Self::BERACHAIN_BEPOLIA), _ => Err("not matching available networks.".into()), } } @@ -148,6 +172,8 @@ impl Network { &Self::BINANCE_MAINNET => "binance_mainnet", &Self::POLYGON_MAINNET => "polygon_mainnet", &Self::BASE_MAINNET => "base_mainnet", + &Self::BERACHAIN_MAINNET => "berachain_mainnet", + &Self::BERACHAIN_BEPOLIA => "berachain_bepolia", _ => panic!("Unknown network"), } .to_string()