2022-03-28 10:58:14 -07:00
|
|
|
use futures_channel::mpsc::{unbounded, UnboundedSender};
|
|
|
|
use futures_util::{pin_mut, SinkExt, StreamExt};
|
|
|
|
use log::*;
|
2022-09-15 09:36:52 -07:00
|
|
|
use std::{collections::HashMap, fs::File, io::Read, net::SocketAddr, sync::Arc, sync::Mutex};
|
2022-03-28 10:58:14 -07:00
|
|
|
use tokio::{
|
|
|
|
net::{TcpListener, TcpStream},
|
|
|
|
pin,
|
|
|
|
};
|
2022-09-10 13:29:17 -07:00
|
|
|
use tokio_tungstenite::tungstenite::{protocol::Message, Error};
|
2022-03-28 10:58:14 -07:00
|
|
|
|
|
|
|
use serde::Deserialize;
|
|
|
|
use solana_geyser_connector_lib::{
|
|
|
|
fill_event_filter::{self, FillCheckpoint, FillEventFilterMessage, MarketConfig},
|
|
|
|
grpc_plugin_source, metrics, websocket_source, SourceConfig,
|
|
|
|
};
|
|
|
|
|
|
|
|
type CheckpointMap = Arc<Mutex<HashMap<String, FillCheckpoint>>>;
|
|
|
|
type PeerMap = Arc<Mutex<HashMap<SocketAddr, UnboundedSender<Message>>>>;
|
|
|
|
|
2022-09-10 13:29:17 -07:00
|
|
|
async fn handle_connection_error(
|
2022-03-28 10:58:14 -07:00
|
|
|
checkpoint_map: CheckpointMap,
|
|
|
|
peer_map: PeerMap,
|
|
|
|
raw_stream: TcpStream,
|
|
|
|
addr: SocketAddr,
|
|
|
|
) {
|
2022-09-10 13:29:17 -07:00
|
|
|
let result = handle_connection(checkpoint_map, peer_map.clone(), raw_stream, addr).await;
|
|
|
|
if result.is_err() {
|
|
|
|
error!("connection {} error {}", addr, result.unwrap_err());
|
|
|
|
};
|
|
|
|
peer_map.lock().unwrap().remove(&addr);
|
|
|
|
}
|
|
|
|
|
|
|
|
async fn handle_connection(
|
|
|
|
checkpoint_map: CheckpointMap,
|
|
|
|
peer_map: PeerMap,
|
|
|
|
raw_stream: TcpStream,
|
|
|
|
addr: SocketAddr,
|
|
|
|
) -> Result<(), Error> {
|
2022-03-28 10:58:14 -07:00
|
|
|
info!("ws connected: {}", addr);
|
2022-09-10 13:29:17 -07:00
|
|
|
let ws_stream = tokio_tungstenite::accept_async(raw_stream).await?;
|
2022-03-28 10:58:14 -07:00
|
|
|
let (mut ws_tx, _ws_rx) = ws_stream.split();
|
|
|
|
|
|
|
|
// 1: publish channel in peer map
|
|
|
|
let (chan_tx, chan_rx) = unbounded();
|
|
|
|
{
|
|
|
|
peer_map.lock().unwrap().insert(addr, chan_tx);
|
|
|
|
info!("ws published: {}", addr);
|
|
|
|
}
|
|
|
|
|
|
|
|
// 2: send initial checkpoint
|
|
|
|
{
|
|
|
|
let checkpoint_map_copy = checkpoint_map.lock().unwrap().clone();
|
|
|
|
for (_, ckpt) in checkpoint_map_copy.iter() {
|
|
|
|
ws_tx
|
|
|
|
.feed(Message::Text(serde_json::to_string(ckpt).unwrap()))
|
2022-09-10 13:29:17 -07:00
|
|
|
.await?;
|
2022-03-28 10:58:14 -07:00
|
|
|
}
|
|
|
|
}
|
2022-09-10 13:29:17 -07:00
|
|
|
info!("ws ckpt sent: {}", addr);
|
|
|
|
ws_tx.flush().await?;
|
2022-03-28 10:58:14 -07:00
|
|
|
|
|
|
|
// 3: forward all events from channel to peer socket
|
|
|
|
let forward_updates = chan_rx.map(Ok).forward(ws_tx);
|
|
|
|
pin_mut!(forward_updates);
|
2022-09-10 13:29:17 -07:00
|
|
|
forward_updates.await?;
|
2022-03-28 10:58:14 -07:00
|
|
|
|
2022-09-10 13:29:17 -07:00
|
|
|
info!("ws disconnected: {}", &addr);
|
|
|
|
Ok(())
|
2022-03-28 10:58:14 -07:00
|
|
|
}
|
|
|
|
|
|
|
|
#[derive(Clone, Debug, Deserialize)]
|
|
|
|
pub struct Config {
|
|
|
|
pub source: SourceConfig,
|
|
|
|
pub markets: Vec<MarketConfig>,
|
|
|
|
pub bind_ws_addr: String,
|
|
|
|
}
|
|
|
|
|
|
|
|
#[tokio::main]
|
|
|
|
async fn main() -> anyhow::Result<()> {
|
|
|
|
let args: Vec<String> = std::env::args().collect();
|
|
|
|
if args.len() < 2 {
|
|
|
|
error!("requires a config file argument");
|
|
|
|
return Ok(());
|
|
|
|
}
|
|
|
|
|
|
|
|
let config: Config = {
|
|
|
|
let mut file = File::open(&args[1])?;
|
|
|
|
let mut contents = String::new();
|
|
|
|
file.read_to_string(&mut contents)?;
|
|
|
|
toml::from_str(&contents).unwrap()
|
|
|
|
};
|
|
|
|
|
|
|
|
solana_logger::setup_with_default("info");
|
|
|
|
|
|
|
|
let metrics_tx = metrics::start();
|
|
|
|
|
|
|
|
let (account_write_queue_sender, slot_queue_sender, fill_receiver) =
|
|
|
|
fill_event_filter::init(config.markets.clone()).await?;
|
|
|
|
|
|
|
|
let checkpoints = CheckpointMap::new(Mutex::new(HashMap::new()));
|
|
|
|
let peers = PeerMap::new(Mutex::new(HashMap::new()));
|
|
|
|
|
|
|
|
let checkpoints_ref_thread = checkpoints.clone();
|
|
|
|
let peers_ref_thread = peers.clone();
|
|
|
|
|
|
|
|
// filleventfilter websocket sink
|
|
|
|
tokio::spawn(async move {
|
|
|
|
pin!(fill_receiver);
|
|
|
|
loop {
|
|
|
|
let message = fill_receiver.recv().await.unwrap();
|
|
|
|
match message {
|
|
|
|
FillEventFilterMessage::Update(update) => {
|
2022-09-15 09:35:16 -07:00
|
|
|
debug!("ws update {} {:?} fill", update.market, update.status);
|
2022-03-28 10:58:14 -07:00
|
|
|
let mut peer_copy = peers_ref_thread.lock().unwrap().clone();
|
|
|
|
for (k, v) in peer_copy.iter_mut() {
|
2022-09-15 09:35:16 -07:00
|
|
|
trace!(" > {}", k);
|
2022-03-28 10:58:14 -07:00
|
|
|
let json = serde_json::to_string(&update);
|
2022-09-15 09:35:16 -07:00
|
|
|
let result = v.send(Message::Text(json.unwrap())).await;
|
|
|
|
if result.is_err() {
|
|
|
|
error!(
|
|
|
|
"ws update {} {:?} fill could not reach {}",
|
|
|
|
update.market, update.status, k
|
|
|
|
);
|
|
|
|
}
|
2022-03-28 10:58:14 -07:00
|
|
|
}
|
|
|
|
}
|
|
|
|
FillEventFilterMessage::Checkpoint(checkpoint) => {
|
|
|
|
checkpoints_ref_thread
|
|
|
|
.lock()
|
|
|
|
.unwrap()
|
|
|
|
.insert(checkpoint.queue.clone(), checkpoint);
|
|
|
|
}
|
|
|
|
}
|
|
|
|
}
|
|
|
|
});
|
|
|
|
|
|
|
|
info!("ws listen: {}", config.bind_ws_addr);
|
|
|
|
let try_socket = TcpListener::bind(&config.bind_ws_addr).await;
|
|
|
|
let listener = try_socket.expect("Failed to bind");
|
|
|
|
tokio::spawn(async move {
|
|
|
|
// Let's spawn the handling of each connection in a separate task.
|
|
|
|
while let Ok((stream, addr)) = listener.accept().await {
|
2022-09-10 13:29:17 -07:00
|
|
|
tokio::spawn(handle_connection_error(
|
2022-03-28 10:58:14 -07:00
|
|
|
checkpoints.clone(),
|
|
|
|
peers.clone(),
|
|
|
|
stream,
|
|
|
|
addr,
|
|
|
|
));
|
|
|
|
}
|
|
|
|
});
|
|
|
|
|
|
|
|
info!(
|
|
|
|
"rpc connect: {}",
|
|
|
|
config
|
|
|
|
.source
|
|
|
|
.grpc_sources
|
|
|
|
.iter()
|
|
|
|
.map(|c| c.connection_string.clone())
|
|
|
|
.collect::<String>()
|
|
|
|
);
|
|
|
|
let use_geyser = true;
|
|
|
|
if use_geyser {
|
|
|
|
grpc_plugin_source::process_events(
|
|
|
|
&config.source,
|
|
|
|
account_write_queue_sender,
|
|
|
|
slot_queue_sender,
|
|
|
|
metrics_tx,
|
|
|
|
)
|
|
|
|
.await;
|
|
|
|
} else {
|
|
|
|
websocket_source::process_events(
|
|
|
|
&config.source,
|
|
|
|
account_write_queue_sender,
|
|
|
|
slot_queue_sender,
|
|
|
|
)
|
|
|
|
.await;
|
|
|
|
}
|
|
|
|
|
|
|
|
Ok(())
|
|
|
|
}
|