Skip to content
This repository was archived by the owner on Jul 5, 2026. It is now read-only.
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
55 changes: 28 additions & 27 deletions config/keys/private.pem
Original file line number Diff line number Diff line change
@@ -1,27 +1,28 @@
-----BEGIN RSA PRIVATE KEY-----
MIIEpAIBAAKCAQEA62uXcN/aTl6ZJF745GPRXvISnlJs/Uw58jb5eTnFswmNUzhY
Achyzfgcu8fLS4VKh9eVwHnbySh1szprYCWo+hIYAtDZmYYIBi1cTUkxtMTBpxDk
Be2Ghy8usIYyA/VJvytgkyhTlXxEF/uFquzteQD38UxNyd4zUKub4e4E+deVo9yt
/HL4UVCt8v+CCraQ0s86V5APJhotmmV7BKxXs31lqwtnNOLdBtASr+d+k7Ln3wuo
0CdK6doTjBRl7cMYRYoPOo4Izh8hYSLW6YMQ5KQiEqe2cGcCKgj5R1IMwnh3QVO+
l3jvacMeHCO3U1NVH09CPO5mm5KWBANmv/iC3QIDAQABAoIBACRmRUsRgXp+i+Ug
vhDqEhRBD3nlOq7LW2ZE87u3oAa3ol9MpebYrE+GXkL2eEtb95MbVS8maEIo/FHS
5Yk/KWpI4+eDjTF8lL8Hwm68s2/EwEBpjygPeq5qMCjhBtiv01A4j70RDiNdzFV8
8UTlTy5XZP6tEpX0wjBl6Ds9hw1t6OAnEaBbuK1sZqU8L8d0GgKryWIAE2HwPIav
L2dW4FXE/Sbm08hXzhlNX9Fq3vMrSZ+I3kiUuYuiBRtbUjCDBIFMlIaaJP1y1Bec
ZAwXpl/NtW4MsCRPKgILVEuYrMYr1A0PVWS66JK4ZavoaQq8JY7pNjmWMNUd06Ql
LFzbAM0CgYEA/t5sdUJ/Tu3ric+HiE6PY1zAazi2U0SUvNlMIMftXb28wmEC3IWd
C6LV2srQBG8CrIP1kpVGTIM+eZB/Q3cfUq97Qe91YhI3CLlfifobuqyAbZ1j0hAb
N5L/fxrxx2Tg2LpTMJycPkLtO078cHigNxPPiqNr5kOI58YHZY7Q8tcCgYEA7HcS
JOIjVXufC2OXyFPkPdaFWKm9h5WJa6CxUCtBmXXiHVegPBWGySbaFFqNRhk2agXb
wcLZdj5/5Lq4FhAjdKHhFBqMF8AHfFx2yNfDo742nyqoXLVnF/5bs/Q3B+B+uoXO
F0nA/ri15E88UqyWj8gtGbrl9H/9IWveVPd3tWsCgYBmys5zfJ5b9xlIO6suDoFG
UeJJXFYsvzw97mYF0pypchzvSLEev8TXLJWT6Lh9EUjCy3X/6LSxpz1LSjwJucGo
V54eubVeGHqZyin+PCFy6J/jldbsohJYF7F0Uimxgb4tqvhiYsehVNzZTsIBmqUD
kbni8IZUGGjfEb9p9m/PgQKBgQDsFDzDIhqQv6kb38SrtkXLDx92U5Drinn2QCqG
lYkawzyKeu94zS0SKn3TkEw3TfirhUnPes9NZDyfiWM8c8RSL0PdpFt1YryWhmH5
RqEGG2PBKP+J/3n71HCNiyZd8N3VLr2BNps+M/80/36EM9blmb6dT6FBp357HYyN
W7viHQKBgQCTKGqbwtmUzEK6lblevIFmj+XE42Jwpfnvvh79h34U/GL9JawrwgfR
Chfj+4AwdL+Uhsr6FmiygGFz/P6D+GdGkNlIIaEzdcLsuw1Z1Ia1DL3o7XNQ5C3F
fd+zSxdczYY4nH0GwNF6Na6o5s6Pn3wFvhaI3K4Twe5OrOs2AWsLjg==
-----END RSA PRIVATE KEY-----
-----BEGIN PRIVATE KEY-----
MIIEvgIBADANBgkqhkiG9w0BAQEFAASCBKgwggSkAgEAAoIBAQDM14zJAhClXJp6
7KqNDHtv13NUlnz+mjJZII+MPMXqnSDXzKgPS20OS03C3bVdRSrA16Lo7lU4EJwM
I/cR+Bc2//QF/O6Sg/Kw4Egyr1P2zceyYykfotsK9fRs+2vmwJgtvcZIvJW2i6Vp
L9I6r27hcn6onEKPMrFesOoV57xitlXARPFYyBrIRrvlpa5/+7/PdtTbsAH91Z7y
FLgtEbBvS/9S2R4x5HZSznTtjgwlJquJiQo3+wgxsq66RswyiDGm7eVYGM9J6QmV
et5Qdh0w9MBy4xxII4UaUbZdJqSz+mWEqgc6rS1LnwqpLDqa69y0z4qye48mTS4y
s7ycSaEVAgMBAAECggEAAUR2KQo7uyIzDH6pYX0JyHvfSU8zD8o5dIa4jKgVm2mE
egFYqtuPHa8GmKWRiTWz2YScC+/plBK6PHL+hNxxnFQCGQVjHoH1fvWsTK/8B4Nn
cGmfqAP0cgFqlUAK/18CsgnCD9Im5P3BNMDofpd2SqvQL8/js4ofQdQ7Zo5MAppW
Xf+21+DUkMPx5iBuTAep4BegN5n2PIh65qtno9qnX4WWKVmz2saXjwtbRLLYf85b
xuAuoXVETN0OfnUJnqo2pi0RLcBY7XjD8UwxuMxgBAGeG+JBz4iclnUBmp91XImN
0smULh/cAx9FHXXPwspQLbasKL/otxoLRmaRahMOIQKBgQD67p//yxB6A6mFKkCI
FUqGI96OeGO6NuZkZaU+1qrFHbBUgPwGa5J+KUCy+7TQizb8EYIF+Z8RSM+qokh2
X0BRKJael6D/MvMIrmJ2d5UQDtNjI3Xb8p3DII0gik9w4Zj3MSqVTXNVhh5nngyJ
XyRYSNPgMpW4yXVT/NR5Bw5apQKBgQDQ+qDwjFhVlKCVEa17/Dx9sPXUm19Sgz8h
1Sji64YGq5i16Fx9bkjrcyfq8hebyBzFu5q9w1ERqJUjVdfGtCK/GnS5c2465Zif
Uvbu9EfCFMnaoFH7VjPs3Kr6b//usevHIGw10q9/Zf7AQSXdrV686HJvNGl9Vv9j
PdzbAngRsQKBgEoDDAokWM3EOsHePn5k2UBLYB9hfvizrKy8Fks8gc31/cZO7Qbv
v5uai0y/VQuVpDgg6drdT3+HnEjV6M2RNqU5dYN9ca0T1/8dgEk06DB+TvcUxHSF
UOb2uOl6IghHYhi21bqHx5bYIiupwETcXRn1ERk1kleYhBSro/e2jxNJAoGBAJLw
dzNMa1wZemP2nxY7wEjcoa3RZc/9yuk+GVadJosQIvtdG5NydUFgoiO3/9OgfGKo
S+C8MgeJkvvagzMLPBdFQeeX+1zcTVlRm6FfEAmuVlQsQBjKfw5ABtS65ajvX4qP
CKc7sfyROfPymu5o1eFcTAJXRwlDn6UnPWCdNtGxAoGBAKTLZY61SQuv+/k7GS8w
TH1J1iqFbbKS6Yx8mXrFjakpp3Bgqn4utNjEJCbGvW50JSlZmJjZ5q6d9uz0LUnw
MokzRnxrXwc2lqa67L2NNZvLolOzExJve3s+UVkrnx9jd8dP0ocodh7tISJCsVfq
6TPNelWvZ3Ym8i5YF/Otc86U
-----END PRIVATE KEY-----
14 changes: 7 additions & 7 deletions config/keys/public.pem
Original file line number Diff line number Diff line change
@@ -1,9 +1,9 @@
-----BEGIN PUBLIC KEY-----
MIIBIjANBgkqhkiG9w0BAQEFAAOCAQ8AMIIBCgKCAQEA62uXcN/aTl6ZJF745GPR
XvISnlJs/Uw58jb5eTnFswmNUzhYAchyzfgcu8fLS4VKh9eVwHnbySh1szprYCWo
+hIYAtDZmYYIBi1cTUkxtMTBpxDkBe2Ghy8usIYyA/VJvytgkyhTlXxEF/uFquzt
eQD38UxNyd4zUKub4e4E+deVo9yt/HL4UVCt8v+CCraQ0s86V5APJhotmmV7BKxX
s31lqwtnNOLdBtASr+d+k7Ln3wuo0CdK6doTjBRl7cMYRYoPOo4Izh8hYSLW6YMQ
5KQiEqe2cGcCKgj5R1IMwnh3QVO+l3jvacMeHCO3U1NVH09CPO5mm5KWBANmv/iC
3QIDAQAB
MIIBIjANBgkqhkiG9w0BAQEFAAOCAQ8AMIIBCgKCAQEAzNeMyQIQpVyaeuyqjQx7
b9dzVJZ8/poyWSCPjDzF6p0g18yoD0ttDktNwt21XUUqwNei6O5VOBCcDCP3EfgX
Nv/0BfzukoPysOBIMq9T9s3HsmMpH6LbCvX0bPtr5sCYLb3GSLyVtoulaS/SOq9u
4XJ+qJxCjzKxXrDqFee8YrZVwETxWMgayEa75aWuf/u/z3bU27AB/dWe8hS4LRGw
b0v/UtkeMeR2Us507Y4MJSariYkKN/sIMbKuukbMMogxpu3lWBjPSekJlXreUHYd
MPTAcuMcSCOFGlG2XSaks/plhKoHOq0tS58KqSw6muvctM+KsnuPJk0uMrO8nEmh
FQIDAQAB
-----END PUBLIC KEY-----
1 change: 1 addition & 0 deletions config/solana/relayer-keypair.json
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
[11,8,74,129,166,152,151,110,156,221,26,193,11,98,182,253,60,53,238,201,96,195,190,209,227,76,42,143,3,51,83,10,141,37,202,66,40,97,168,42,38,147,80,186,130,238,205,201,72,168,188,191,101,248,238,157,249,155,55,137,84,8,69,81]
21 changes: 17 additions & 4 deletions relayer/src/relayer.rs
Original file line number Diff line number Diff line change
@@ -1,12 +1,12 @@
use std::{
collections::{hash_map::Entry, HashMap, HashSet},
net::IpAddr,
str::FromStr,
sync::{
atomic::{AtomicBool, AtomicU64, Ordering},
Arc, RwLock,
},
thread,
thread::JoinHandle,
thread::{self, JoinHandle},
time::{Duration, Instant, SystemTime},
};

Expand Down Expand Up @@ -426,6 +426,7 @@ impl RelayerImpl {
ofac_addresses: HashSet<Pubkey>,
address_lookup_table_cache: Arc<DashMap<Pubkey, AddressLookupTableAccount>>,
validator_packet_batch_size: usize,
forward_transaction_signal: Sender<()>,
forward_all: bool,
) -> Self {
const LEADER_LOOKAHEAD: u64 = 2;
Expand Down Expand Up @@ -454,6 +455,7 @@ impl RelayerImpl {
ofac_addresses,
address_lookup_table_cache,
validator_packet_batch_size,
forward_transaction_signal,
forward_all,
);
warn!("RelayerImpl thread exited with result {res:?}")
Expand Down Expand Up @@ -490,6 +492,7 @@ impl RelayerImpl {
ofac_addresses: HashSet<Pubkey>,
address_lookup_table_cache: Arc<DashMap<Pubkey, AddressLookupTableAccount>>,
validator_packet_batch_size: usize,
forward_transaction_signal: Sender<()>,
forward_all: bool,
) -> RelayerResult<()> {
let mut highest_slot = Slot::default();
Expand All @@ -504,6 +507,8 @@ impl RelayerImpl {
);

let mut slot_leaders = HashSet::new();
let mut is_my_validator_leader = false;
let my_validator_pubkey = Pubkey::from_str("validator_pubkey").unwrap(); // get from env ?

while !exit.load(Ordering::Relaxed) {
crossbeam_channel::select! {
Expand All @@ -517,12 +522,20 @@ impl RelayerImpl {
.collect();
slot_leaders = leader_schedule_cache.leaders_for_slots(&slots);

is_my_validator_leader = slot_leaders.contains(&my_validator_pubkey);
if is_my_validator_leader {
forward_transaction_signal.send(()).unwrap(); // here we trigger delay_packet_receiver from tx_relayer
}

let _ = relayer_metrics.crossbeam_slot_receiver_processing_us.increment(start.elapsed().as_micros() as u64);
},
recv(delay_packet_receiver) -> maybe_packet_batches => {
let start = Instant::now();
let failed_forwards = Self::forward_packets(maybe_packet_batches, packet_subscriptions, &slot_leaders, &mut relayer_metrics, &ofac_addresses, &address_lookup_table_cache, validator_packet_batch_size, forward_all)?;
Self::drop_connections(failed_forwards, packet_subscriptions, &mut relayer_metrics);
// only forward packets if this validator is the leader
if is_my_validator_leader {
let failed_forwards = Self::forward_packets(maybe_packet_batches, packet_subscriptions, &slot_leaders, &mut relayer_metrics, &ofac_addresses, &address_lookup_table_cache, validator_packet_batch_size, forward_all)?;
Self::drop_connections(failed_forwards, packet_subscriptions, &mut relayer_metrics);
}
let _ = relayer_metrics.crossbeam_delay_packet_receiver_processing_us.increment(start.elapsed().as_micros() as u64);
},
recv(subscription_receiver) -> maybe_subscription => {
Expand Down
43 changes: 15 additions & 28 deletions transaction-relayer/src/forwarder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -26,6 +26,7 @@ pub fn start_forward_and_delay_thread(
block_engine_sender: tokio::sync::mpsc::Sender<BlockEnginePackets>,
num_threads: u64,
disable_mempool: bool,
forward_transaction_signal: Receiver<()>,
exit: &Arc<AtomicBool>,
) -> Vec<JoinHandle<()>> {
const SLEEP_DURATION: Duration = Duration::from_millis(5);
Expand All @@ -36,17 +37,17 @@ pub fn start_forward_and_delay_thread(
let verified_receiver = verified_receiver.clone();
let delay_packet_sender = delay_packet_sender.clone();
let block_engine_sender = block_engine_sender.clone();

let forward_transaction_signal = forward_transaction_signal.clone();
let exit = exit.clone();
Builder::new()
.name(format!("forwarder_thread_{thread_id}"))
.spawn(move || {
let mut buffered_packet_batches: VecDeque<RelayerPacketBatches> =
let mut transaction_queue: VecDeque<RelayerPacketBatches> =
VecDeque::with_capacity(100_000);

let metrics_interval = Duration::from_secs(1);
let mut forwarder_metrics = ForwarderMetrics::new(
buffered_packet_batches.capacity(),
transaction_queue.capacity(),
verified_receiver.capacity().unwrap_or_default(), // TODO (LB): unbounded channel now, remove metric
block_engine_sender.capacity(),
);
Expand All @@ -57,7 +58,7 @@ pub fn start_forward_and_delay_thread(
forwarder_metrics.report(thread_id, packet_delay_ms);

forwarder_metrics = ForwarderMetrics::new(
buffered_packet_batches.capacity(),
transaction_queue.capacity(),
verified_receiver.capacity().unwrap_or_default(), // TODO (LB): unbounded channel now, remove metric
block_engine_sender.capacity(),
);
Expand Down Expand Up @@ -100,7 +101,8 @@ pub fn start_forward_and_delay_thread(
}
}
}
buffered_packet_batches.push_back(RelayerPacketBatches {

transaction_queue.push_back(RelayerPacketBatches {
stamp: instant,
banking_packet_batch,
});
Expand All @@ -111,31 +113,16 @@ pub fn start_forward_and_delay_thread(
}
}

while let Some(packet_batches) = buffered_packet_batches.front() {
if packet_batches.stamp.elapsed() < packet_delay {
break;
// Wait for the signal that it's time to forward transactions
if forward_transaction_signal.try_recv().is_ok() {
// Forward all queued transactions
while let Some(packet_batches) = transaction_queue.pop_front() {
// Send the transaction
delay_packet_sender
.send(packet_batches)
.expect("exiting forwarding delayed packets");
}
let batch = buffered_packet_batches.pop_front().unwrap();

let num_packets = batch
.banking_packet_batch
.0
.iter()
.map(|b| b.len() as u64)
.sum::<u64>();

forwarder_metrics.num_relayer_packets_forwarded += num_packets;
delay_packet_sender
.send(batch)
.expect("exiting forwarding delayed packets");
}

forwarder_metrics.update_queue_lengths(
buffered_packet_batches.len(),
buffered_packet_batches.capacity(),
verified_receiver.len(),
BLOCK_ENGINE_FORWARDER_QUEUE_CAPACITY - block_engine_sender.capacity(),
);
}
})
.unwrap()
Expand Down
7 changes: 6 additions & 1 deletion transaction-relayer/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -460,13 +460,17 @@ fn main() {
let (block_engine_sender, block_engine_receiver) =
channel(jito_transaction_relayer::forwarder::BLOCK_ENGINE_FORWARDER_QUEUE_CAPACITY);

let (forward_transaction_signal_tx, forward_transaction_signal_rx) =
crossbeam_channel::bounded(1);

let forward_and_delay_threads = start_forward_and_delay_thread(
verified_receiver,
delay_packet_sender,
args.packet_delay_ms,
block_engine_sender,
1,
args.disable_mempool,
forward_transaction_signal_rx.clone(),
&exit,
);

Expand Down Expand Up @@ -519,7 +523,8 @@ fn main() {
ofac_addresses,
address_lookup_table_cache,
args.validator_packet_batch_size,
args.forward_all,
forward_transaction_signal_tx,
false, // do not forward to all validators
);

let priv_key = fs::read(&args.signing_key_pem_path).unwrap_or_else(|_| {
Expand Down