Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion src/connection.rs
Original file line number Diff line number Diff line change
Expand Up @@ -221,7 +221,7 @@ impl Connection {

// Should do the loop here

let (tx_tcp, rx_tcp) = tokio::sync::mpsc::channel::<Vec<u8>>(32);
let (tx_tcp, rx_tcp) = tokio::sync::mpsc::channel::<Vec<u8>>(256);
let (tx_write, mut rx_write) = tokio::sync::mpsc::channel::<Vec<u8>>(32);

let handle = tokio::spawn(async move {
Expand Down
28 changes: 21 additions & 7 deletions src/database.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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());
Expand Down
136 changes: 74 additions & 62 deletions src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down Expand Up @@ -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);
}
}

Expand All @@ -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);
}
}
}
Expand Down Expand Up @@ -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..]);
Expand Down Expand Up @@ -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;
}
}

Expand All @@ -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;
}
}
}
Expand Down
26 changes: 26 additions & 0 deletions src/networks.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<Self, Box<dyn Error>> {
match network {
"ethereum_ropsten" => Ok(Self::ETHEREUM_ROPSTEN),
Expand All @@ -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()),
}
}
Expand All @@ -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()
Expand Down
Loading