Compare commits
2 Commits
7f0ddd3ac5
...
05f42c7352
Author | SHA1 | Date |
---|---|---|
Riordan Panayides | 05f42c7352 | |
Riordan Panayides | 9f528603b0 |
|
@ -3,9 +3,9 @@ use anchor_lang::prelude::Pubkey;
|
||||||
use bytemuck::cast_slice;
|
use bytemuck::cast_slice;
|
||||||
use client::{Client, MangoGroupContext};
|
use client::{Client, MangoGroupContext};
|
||||||
use futures_channel::mpsc::{unbounded, UnboundedSender};
|
use futures_channel::mpsc::{unbounded, UnboundedSender};
|
||||||
use futures_util::{pin_mut, SinkExt, StreamExt};
|
use futures_util::{pin_mut, SinkExt, StreamExt, future::{self, Ready}, TryStreamExt};
|
||||||
use log::*;
|
use log::*;
|
||||||
use std::{collections::HashMap, fs::File, io::Read, net::SocketAddr, sync::Arc, sync::Mutex, time::Duration, convert::identity, str::FromStr};
|
use std::{collections::{HashMap, HashSet}, fs::File, io::Read, net::SocketAddr, sync::Arc, sync::Mutex, time::Duration, convert::identity, str::FromStr};
|
||||||
use tokio::{
|
use tokio::{
|
||||||
net::{TcpListener, TcpStream},
|
net::{TcpListener, TcpStream},
|
||||||
pin,
|
pin,
|
||||||
|
@ -13,7 +13,7 @@ use tokio::{
|
||||||
use tokio_tungstenite::tungstenite::{protocol::Message, Error};
|
use tokio_tungstenite::tungstenite::{protocol::Message, Error};
|
||||||
|
|
||||||
use serde::Deserialize;
|
use serde::Deserialize;
|
||||||
use solana_geyser_connector_lib::{metrics::{MetricType, MetricU64}, FilterConfig, fill_event_filter::SerumFillCheckpoint};
|
use solana_geyser_connector_lib::{metrics::{MetricType, MetricU64}, FilterConfig, fill_event_filter::SerumFillCheckpoint, StatusResponse};
|
||||||
use solana_geyser_connector_lib::{
|
use solana_geyser_connector_lib::{
|
||||||
fill_event_filter::{self, FillCheckpoint, FillEventFilterMessage},
|
fill_event_filter::{self, FillCheckpoint, FillEventFilterMessage},
|
||||||
grpc_plugin_source, metrics, websocket_source, MetricsConfig, SourceConfig,
|
grpc_plugin_source, metrics, websocket_source, MetricsConfig, SourceConfig,
|
||||||
|
@ -21,12 +21,42 @@ use solana_geyser_connector_lib::{
|
||||||
|
|
||||||
type CheckpointMap = Arc<Mutex<HashMap<String, FillCheckpoint>>>;
|
type CheckpointMap = Arc<Mutex<HashMap<String, FillCheckpoint>>>;
|
||||||
type SerumCheckpointMap = Arc<Mutex<HashMap<String, SerumFillCheckpoint>>>;
|
type SerumCheckpointMap = Arc<Mutex<HashMap<String, SerumFillCheckpoint>>>;
|
||||||
type PeerMap = Arc<Mutex<HashMap<SocketAddr, UnboundedSender<Message>>>>;
|
type PeerMap = Arc<Mutex<HashMap<SocketAddr, Peer>>>;
|
||||||
|
|
||||||
|
#[derive(Clone, Debug, Deserialize)]
|
||||||
|
#[serde(tag = "command")]
|
||||||
|
pub enum Command {
|
||||||
|
#[serde(rename = "subscribe")]
|
||||||
|
Subscribe(SubscribeCommand),
|
||||||
|
#[serde(rename = "unsubscribe")]
|
||||||
|
Unsubscribe(UnsubscribeCommand),
|
||||||
|
#[serde(rename = "getMarkets")]
|
||||||
|
GetMarkets,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Clone, Debug, Deserialize)]
|
||||||
|
#[serde(rename_all = "camelCase")]
|
||||||
|
pub struct SubscribeCommand {
|
||||||
|
pub market_id: String,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Clone, Debug, Deserialize)]
|
||||||
|
#[serde(rename_all = "camelCase")]
|
||||||
|
pub struct UnsubscribeCommand {
|
||||||
|
pub market_id: String,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Clone, Debug)]
|
||||||
|
pub struct Peer {
|
||||||
|
pub sender: UnboundedSender<Message>,
|
||||||
|
pub subscriptions: HashSet<String>,
|
||||||
|
}
|
||||||
|
|
||||||
async fn handle_connection_error(
|
async fn handle_connection_error(
|
||||||
checkpoint_map: CheckpointMap,
|
checkpoint_map: CheckpointMap,
|
||||||
serum_checkpoint_map: SerumCheckpointMap,
|
serum_checkpoint_map: SerumCheckpointMap,
|
||||||
peer_map: PeerMap,
|
peer_map: PeerMap,
|
||||||
|
market_ids: HashMap<String, String>,
|
||||||
raw_stream: TcpStream,
|
raw_stream: TcpStream,
|
||||||
addr: SocketAddr,
|
addr: SocketAddr,
|
||||||
metrics_opened_connections: MetricU64,
|
metrics_opened_connections: MetricU64,
|
||||||
|
@ -34,7 +64,7 @@ async fn handle_connection_error(
|
||||||
) {
|
) {
|
||||||
metrics_opened_connections.clone().increment();
|
metrics_opened_connections.clone().increment();
|
||||||
|
|
||||||
let result = handle_connection(checkpoint_map, serum_checkpoint_map, peer_map.clone(), raw_stream, addr).await;
|
let result = handle_connection(checkpoint_map, serum_checkpoint_map, peer_map.clone(), market_ids, raw_stream, addr).await;
|
||||||
if result.is_err() {
|
if result.is_err() {
|
||||||
error!("connection {} error {}", addr, result.unwrap_err());
|
error!("connection {} error {}", addr, result.unwrap_err());
|
||||||
};
|
};
|
||||||
|
@ -48,42 +78,156 @@ async fn handle_connection(
|
||||||
checkpoint_map: CheckpointMap,
|
checkpoint_map: CheckpointMap,
|
||||||
serum_checkpoint_map: SerumCheckpointMap,
|
serum_checkpoint_map: SerumCheckpointMap,
|
||||||
peer_map: PeerMap,
|
peer_map: PeerMap,
|
||||||
|
market_ids: HashMap<String, String>,
|
||||||
raw_stream: TcpStream,
|
raw_stream: TcpStream,
|
||||||
addr: SocketAddr,
|
addr: SocketAddr,
|
||||||
) -> Result<(), Error> {
|
) -> Result<(), Error> {
|
||||||
info!("ws connected: {}", addr);
|
info!("ws connected: {}", addr);
|
||||||
let ws_stream = tokio_tungstenite::accept_async(raw_stream).await?;
|
let ws_stream = tokio_tungstenite::accept_async(raw_stream).await?;
|
||||||
let (mut ws_tx, _ws_rx) = ws_stream.split();
|
let (ws_tx, ws_rx) = ws_stream.split();
|
||||||
|
|
||||||
// 1: publish channel in peer map
|
// 1: publish channel in peer map
|
||||||
let (chan_tx, chan_rx) = unbounded();
|
let (chan_tx, chan_rx) = unbounded();
|
||||||
{
|
{
|
||||||
peer_map.lock().unwrap().insert(addr, chan_tx);
|
peer_map.lock().unwrap().insert(
|
||||||
info!("ws published: {}", addr);
|
addr,
|
||||||
|
Peer {
|
||||||
|
sender: chan_tx,
|
||||||
|
subscriptions: HashSet::<String>::new(),
|
||||||
|
},
|
||||||
|
);
|
||||||
}
|
}
|
||||||
|
|
||||||
// todo: add subscribe logic
|
let receive_commands = ws_rx.try_for_each(|msg| {
|
||||||
// 2: send initial checkpoint
|
handle_commands(
|
||||||
{
|
addr,
|
||||||
let checkpoint_map_copy = checkpoint_map.lock().unwrap().clone();
|
msg,
|
||||||
for (_, ckpt) in checkpoint_map_copy.iter() {
|
peer_map.clone(),
|
||||||
ws_tx
|
checkpoint_map.clone(),
|
||||||
.feed(Message::Text(serde_json::to_string(ckpt).unwrap()))
|
serum_checkpoint_map.clone(),
|
||||||
.await?;
|
market_ids.clone(),
|
||||||
}
|
)
|
||||||
}
|
});
|
||||||
info!("ws ckpt sent: {}", addr);
|
|
||||||
ws_tx.flush().await?;
|
|
||||||
|
|
||||||
// 3: forward all events from channel to peer socket
|
|
||||||
let forward_updates = chan_rx.map(Ok).forward(ws_tx);
|
let forward_updates = chan_rx.map(Ok).forward(ws_tx);
|
||||||
pin_mut!(forward_updates);
|
|
||||||
forward_updates.await?;
|
|
||||||
|
|
||||||
|
pin_mut!(receive_commands, forward_updates);
|
||||||
|
future::select(receive_commands, forward_updates).await;
|
||||||
|
|
||||||
|
peer_map.lock().unwrap().remove(&addr);
|
||||||
info!("ws disconnected: {}", &addr);
|
info!("ws disconnected: {}", &addr);
|
||||||
Ok(())
|
Ok(())
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn handle_commands(
|
||||||
|
addr: SocketAddr,
|
||||||
|
msg: Message,
|
||||||
|
peer_map: PeerMap,
|
||||||
|
checkpoint_map: CheckpointMap,
|
||||||
|
serum_checkpoint_map: SerumCheckpointMap,
|
||||||
|
market_ids: HashMap<String, String>,
|
||||||
|
) -> Ready<Result<(), Error>> {
|
||||||
|
let msg_str = msg.clone().into_text().unwrap();
|
||||||
|
let command: Result<Command, serde_json::Error> = serde_json::from_str(&msg_str);
|
||||||
|
let mut peers = peer_map.lock().unwrap();
|
||||||
|
let peer = peers.get_mut(&addr).expect("peer should be in map");
|
||||||
|
match command {
|
||||||
|
Ok(Command::Subscribe(cmd)) => {
|
||||||
|
let market_id = cmd.clone().market_id;
|
||||||
|
match market_ids.get(&market_id) {
|
||||||
|
None => {
|
||||||
|
let res = StatusResponse {
|
||||||
|
success: false,
|
||||||
|
message: "market not found",
|
||||||
|
};
|
||||||
|
peer.sender
|
||||||
|
.unbounded_send(Message::Text(serde_json::to_string(&res).unwrap()))
|
||||||
|
.unwrap();
|
||||||
|
return future::ok(());
|
||||||
|
}
|
||||||
|
_ => {}
|
||||||
|
}
|
||||||
|
let subscribed = peer.subscriptions.insert(market_id.clone());
|
||||||
|
|
||||||
|
let res = if subscribed {
|
||||||
|
StatusResponse {
|
||||||
|
success: true,
|
||||||
|
message: "subscribed",
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
StatusResponse {
|
||||||
|
success: false,
|
||||||
|
message: "already subscribed",
|
||||||
|
}
|
||||||
|
};
|
||||||
|
peer.sender
|
||||||
|
.unbounded_send(Message::Text(serde_json::to_string(&res).unwrap()))
|
||||||
|
.unwrap();
|
||||||
|
|
||||||
|
if subscribed {
|
||||||
|
// todo: this is janky af
|
||||||
|
let checkpoint_map = checkpoint_map.lock().unwrap();
|
||||||
|
let serum_checkpoint_map = serum_checkpoint_map.lock().unwrap();
|
||||||
|
let checkpoint = checkpoint_map.get(&market_id);
|
||||||
|
match checkpoint {
|
||||||
|
Some(checkpoint) => {
|
||||||
|
peer.sender
|
||||||
|
.unbounded_send(Message::Text(
|
||||||
|
serde_json::to_string(&checkpoint).unwrap(),
|
||||||
|
))
|
||||||
|
.unwrap();
|
||||||
|
}
|
||||||
|
None => match serum_checkpoint_map.get(&market_id) {
|
||||||
|
Some(checkpoint) => {
|
||||||
|
peer.sender
|
||||||
|
.unbounded_send(Message::Text(
|
||||||
|
serde_json::to_string(&checkpoint).unwrap(),
|
||||||
|
))
|
||||||
|
.unwrap();
|
||||||
|
}
|
||||||
|
None => info!("no checkpoint available on client subscription"), // todo: what to do here?
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
Ok(Command::Unsubscribe(cmd)) => {
|
||||||
|
info!("unsubscribe {}", cmd.market_id);
|
||||||
|
let unsubscribed = peer.subscriptions.remove(&cmd.market_id);
|
||||||
|
let res = if unsubscribed {
|
||||||
|
StatusResponse {
|
||||||
|
success: true,
|
||||||
|
message: "unsubscribed",
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
StatusResponse {
|
||||||
|
success: false,
|
||||||
|
message: "not subscribed",
|
||||||
|
}
|
||||||
|
};
|
||||||
|
peer.sender
|
||||||
|
.unbounded_send(Message::Text(serde_json::to_string(&res).unwrap()))
|
||||||
|
.unwrap();
|
||||||
|
}
|
||||||
|
Ok(Command::GetMarkets) => {
|
||||||
|
info!("getMarkets");
|
||||||
|
peer.sender
|
||||||
|
.unbounded_send(Message::Text(serde_json::to_string(&market_ids).unwrap()))
|
||||||
|
.unwrap();
|
||||||
|
}
|
||||||
|
Err(err) => {
|
||||||
|
info!("error deserializing user input {:?}", err);
|
||||||
|
let res = StatusResponse {
|
||||||
|
success: false,
|
||||||
|
message: "invalid input",
|
||||||
|
};
|
||||||
|
peer.sender
|
||||||
|
.unbounded_send(Message::Text(serde_json::to_string(&res).unwrap()))
|
||||||
|
.unwrap();
|
||||||
|
}
|
||||||
|
};
|
||||||
|
|
||||||
|
future::ok(())
|
||||||
|
}
|
||||||
|
|
||||||
#[derive(Clone, Debug, Deserialize)]
|
#[derive(Clone, Debug, Deserialize)]
|
||||||
pub struct Config {
|
pub struct Config {
|
||||||
pub source: SourceConfig,
|
pub source: SourceConfig,
|
||||||
|
@ -170,6 +314,19 @@ async fn main() -> anyhow::Result<()> {
|
||||||
})
|
})
|
||||||
.collect();
|
.collect();
|
||||||
|
|
||||||
|
let a: Vec<(String, String)> = group_context
|
||||||
|
.serum3_markets
|
||||||
|
.iter()
|
||||||
|
.map(|(_, context)| (context.market.name().to_owned(), context.market.serum_market_external.to_string())).collect();
|
||||||
|
let b: Vec<(String, String)> = group_context
|
||||||
|
.perp_markets
|
||||||
|
.iter()
|
||||||
|
.map(|(_, context)| (context.market.name().to_owned(), context.address.to_string())).collect();
|
||||||
|
let market_pubkey_strings: HashMap<String, String> = [a, b]
|
||||||
|
.concat()
|
||||||
|
.into_iter()
|
||||||
|
.collect();
|
||||||
|
|
||||||
let (account_write_queue_sender, slot_queue_sender, fill_receiver) =
|
let (account_write_queue_sender, slot_queue_sender, fill_receiver) =
|
||||||
fill_event_filter::init(perp_queue_pks.clone(), serum_queue_pks.clone(), metrics_tx.clone()).await?;
|
fill_event_filter::init(perp_queue_pks.clone(), serum_queue_pks.clone(), metrics_tx.clone()).await?;
|
||||||
|
|
||||||
|
@ -190,18 +347,21 @@ async fn main() -> anyhow::Result<()> {
|
||||||
FillEventFilterMessage::Update(update) => {
|
FillEventFilterMessage::Update(update) => {
|
||||||
debug!("ws update {} {:?} fill", update.market, update.status);
|
debug!("ws update {} {:?} fill", update.market, update.status);
|
||||||
let mut peer_copy = peers_ref_thread.lock().unwrap().clone();
|
let mut peer_copy = peers_ref_thread.lock().unwrap().clone();
|
||||||
for (k, v) in peer_copy.iter_mut() {
|
for (addr, peer) in peer_copy.iter_mut() {
|
||||||
trace!(" > {}", k);
|
let json = serde_json::to_string(&update).unwrap();
|
||||||
let json = serde_json::to_string(&update);
|
|
||||||
let result = v.send(Message::Text(json.unwrap())).await;
|
// only send updates if the peer is subscribed
|
||||||
|
if peer.subscriptions.contains(&update.market) {
|
||||||
|
let result = peer.sender.send(Message::Text(json)).await;
|
||||||
if result.is_err() {
|
if result.is_err() {
|
||||||
error!(
|
error!(
|
||||||
"ws update {} {:?} fill could not reach {}",
|
"ws update {} fill could not reach {}",
|
||||||
update.market, update.status, k
|
update.market, addr
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
}
|
||||||
FillEventFilterMessage::Checkpoint(checkpoint) => {
|
FillEventFilterMessage::Checkpoint(checkpoint) => {
|
||||||
checkpoints_ref_thread
|
checkpoints_ref_thread
|
||||||
.lock()
|
.lock()
|
||||||
|
@ -211,18 +371,21 @@ async fn main() -> anyhow::Result<()> {
|
||||||
FillEventFilterMessage::SerumUpdate(update) => {
|
FillEventFilterMessage::SerumUpdate(update) => {
|
||||||
debug!("ws update {} {:?} serum fill", update.market, update.status);
|
debug!("ws update {} {:?} serum fill", update.market, update.status);
|
||||||
let mut peer_copy = peers_ref_thread.lock().unwrap().clone();
|
let mut peer_copy = peers_ref_thread.lock().unwrap().clone();
|
||||||
for (k, v) in peer_copy.iter_mut() {
|
for (addr, peer) in peer_copy.iter_mut() {
|
||||||
trace!(" > {}", k);
|
let json = serde_json::to_string(&update).unwrap();
|
||||||
let json = serde_json::to_string(&update);
|
|
||||||
let result = v.send(Message::Text(json.unwrap())).await;
|
// only send updates if the peer is subscribed
|
||||||
|
if peer.subscriptions.contains(&update.market) {
|
||||||
|
let result = peer.sender.send(Message::Text(json)).await;
|
||||||
if result.is_err() {
|
if result.is_err() {
|
||||||
error!(
|
error!(
|
||||||
"ws update {} {:?} serum fill could not reach {}",
|
"ws update {} fill could not reach {}",
|
||||||
update.market, update.status, k
|
update.market, addr
|
||||||
);
|
);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
}
|
||||||
FillEventFilterMessage::SerumCheckpoint(checkpoint) => {
|
FillEventFilterMessage::SerumCheckpoint(checkpoint) => {
|
||||||
serum_checkpoints_ref_thread
|
serum_checkpoints_ref_thread
|
||||||
.lock()
|
.lock()
|
||||||
|
@ -243,6 +406,7 @@ async fn main() -> anyhow::Result<()> {
|
||||||
checkpoints.clone(),
|
checkpoints.clone(),
|
||||||
serum_checkpoints.clone(),
|
serum_checkpoints.clone(),
|
||||||
peers.clone(),
|
peers.clone(),
|
||||||
|
market_pubkey_strings.clone(),
|
||||||
stream,
|
stream,
|
||||||
addr,
|
addr,
|
||||||
metrics_opened_connections.clone(),
|
metrics_opened_connections.clone(),
|
||||||
|
|
Loading…
Reference in New Issue