use std::sync::{Arc, OnceLock};
use tracing::{debug, error, info, warn};
use tokio::{
sync::{broadcast, mpsc, oneshot},
task::JoinHandle,
};
use crate::{
ConnectStrategy,
api::{
commands::{
RithmicBracketLevelAdjustment, RithmicBracketOrder, RithmicCancelAllOrders,
RithmicCancelOrder, RithmicExitPosition, RithmicLinkOrders, RithmicModifyOrder,
RithmicModifyOrderReferenceData, RithmicOcoOrder, RithmicOrder,
},
receiver_api::RithmicResponse,
sender_api::LoginScope,
},
config::{LoginConfig, RithmicAccount, RithmicConfig},
error::RithmicError,
plants::{
await_all_responses, await_first_response,
core::{PlantActor, PlantCore, SelectResult},
subscription::SubscriptionFilter,
trade_routes::TradeRouteCache,
},
rti::{TradeRoute, messages::RithmicMessage, request_login::SysInfraType},
types::{EasyToBorrowRequest, FillHistoryRange, RmsUpdateBits},
};
pub(crate) enum OrderPlantCommand {
Close,
Abort,
GetSystemInfo {
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
Login {
config: LoginConfig,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
SetLogin,
Logout {
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
UpdateHeartbeat {
seconds: u64,
},
AccountList {
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
SubscribeOrderUpdates {
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
SubscribeBracketUpdates {
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
PlaceBracketOrder {
bracket_order: Box<RithmicBracketOrder>,
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
ModifyOrder {
order: RithmicModifyOrder,
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
ModifyStop {
adjustment: RithmicBracketLevelAdjustment,
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
ModifyTarget {
adjustment: RithmicBracketLevelAdjustment,
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
CancelOrder {
order: RithmicCancelOrder,
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
ShowOrders {
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
CancelAllOrders {
command: RithmicCancelAllOrders,
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
GetAccountRmsInfo {
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
GetProductRmsInfo {
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
GetTradeRoutes {
subscribe_for_updates: bool,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
RecordTradeRoutes(Vec<RithmicResponse>),
RecordTradeRouteUpdate(Box<TradeRoute>),
TradeRouteFor {
exchange: String,
response_sender: oneshot::Sender<Result<String, RithmicError>>,
},
ShowOrderHistoryDates {
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
ShowOrderHistorySummary {
date: String,
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
ShowOrderHistoryDetail {
basket_id: String,
date: String,
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
ShowOrderHistory {
basket_id: Option<String>,
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
PlaceOrder {
order: RithmicOrder,
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
PlaceOcoOrder {
order: RithmicOcoOrder,
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
ShowBrackets {
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
ShowBracketStops {
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
ExitPosition {
command: RithmicExitPosition,
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
LinkOrders {
command: RithmicLinkOrders,
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
GetEasyToBorrowList {
request_type: EasyToBorrowRequest,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
ModifyOrderReferenceData {
command: RithmicModifyOrderReferenceData,
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
GetOrderSessionConfig {
should_defer_request: Option<bool>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
ReplayExecutions {
start_index_sec: i32,
finish_index_sec: i32,
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
GetUserInfo {
user: Option<String>,
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
ShowFillHistory {
range: FillHistoryRange,
max_record_count: Option<i32>,
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
SubscribeAccountRmsUpdates {
subscribe: bool,
update_bits: Vec<RmsUpdateBits>,
account: Arc<RithmicAccount>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
GetLoginInfo {
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
ListUnacceptedAgreements {
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
ListAcceptedAgreements {
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
AcceptAgreement {
agreement_id: String,
market_data_usage_capacity: Option<String>,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
ShowAgreement {
agreement_id: String,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
SetRithmicMrktDataSelfCertStatus {
agreement_id: String,
market_data_usage_capacity: String,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
ListExchangePermissions {
user: String,
response_sender: oneshot::Sender<Result<Vec<RithmicResponse>, RithmicError>>,
},
}
#[derive(Debug)]
pub struct RithmicOrderPlant {
pub(crate) connection_handle: JoinHandle<()>,
sender: mpsc::Sender<OrderPlantCommand>,
subscription_sender: broadcast::Sender<RithmicResponse>,
login_scope: Arc<OnceLock<LoginScope>>,
}
impl RithmicOrderPlant {
pub async fn connect(
config: &RithmicConfig,
strategy: ConnectStrategy,
) -> Result<RithmicOrderPlant, RithmicError> {
let (req_tx, req_rx) = mpsc::channel::<OrderPlantCommand>(64);
let (sub_tx, _sub_rx) = broadcast::channel(10_000);
let login_scope = Arc::new(OnceLock::new());
let mut order_plant = OrderPlant::new(
req_rx,
sub_tx.clone(),
config,
strategy,
Arc::clone(&login_scope),
)
.await?;
let connection_handle = tokio::spawn(async move {
order_plant.run().await;
});
Ok(RithmicOrderPlant {
connection_handle,
sender: req_tx,
subscription_sender: sub_tx,
login_scope,
})
}
}
impl RithmicOrderPlant {
pub async fn await_shutdown(self) -> Result<(), tokio::task::JoinError> {
self.connection_handle.await
}
pub fn get_handle(&self, account: &RithmicAccount) -> RithmicOrderPlantHandle {
let account = Arc::new(account.clone());
let account_for_filter = Arc::clone(&account);
RithmicOrderPlantHandle {
account,
login_scope: Arc::clone(&self.login_scope),
sender: self.sender.clone(),
subscription_receiver: SubscriptionFilter::new(
account_for_filter,
self.subscription_sender.subscribe(),
),
}
}
}
#[derive(Debug)]
struct OrderPlant {
core: PlantCore,
request_receiver: mpsc::Receiver<OrderPlantCommand>,
login_scope: Arc<OnceLock<LoginScope>>,
trade_routes: TradeRouteCache,
}
impl OrderPlant {
async fn new(
request_receiver: mpsc::Receiver<OrderPlantCommand>,
subscription_sender: broadcast::Sender<RithmicResponse>,
config: &RithmicConfig,
strategy: ConnectStrategy,
login_scope: Arc<OnceLock<LoginScope>>,
) -> Result<OrderPlant, RithmicError> {
let core = PlantCore::new(subscription_sender, config, strategy, "order_plant").await?;
Ok(OrderPlant {
core,
request_receiver,
login_scope,
trade_routes: TradeRouteCache::default(),
})
}
}
impl PlantActor for OrderPlant {
type Command = OrderPlantCommand;
async fn run(&mut self) {
loop {
let result = self.core.next_event(&mut self.request_receiver).await;
let stop = match result {
SelectResult::HeartbeatFired => self.core.send_heartbeat().await,
SelectResult::PingFired => self.core.send_ping().await,
SelectResult::PingTimeout => self.core.handle_ping_timeout(),
SelectResult::Command(cmd) => {
if matches!(cmd, OrderPlantCommand::Abort) {
self.core.handle_abort()
} else {
self.handle_command(cmd).await;
false
}
}
SelectResult::RithmicMessage(msg) => self.core.handle_rithmic_message(msg).await,
SelectResult::StreamClosed => self.core.handle_stream_closed(),
};
if stop {
break;
}
}
}
async fn handle_command(&mut self, command: OrderPlantCommand) {
if self.core.close_requested
&& !matches!(
command,
OrderPlantCommand::Close
| OrderPlantCommand::SetLogin
| OrderPlantCommand::UpdateHeartbeat { .. }
| OrderPlantCommand::Abort
)
{
debug!("order_plant: dropping a command queued after close was requested");
return;
}
match command {
OrderPlantCommand::Close => {
self.core.handle_close().await;
}
OrderPlantCommand::GetSystemInfo { response_sender } => {
self.core.handle_get_system_info(response_sender).await;
}
OrderPlantCommand::Login {
config,
response_sender,
} => {
self.core
.handle_login(config, SysInfraType::OrderPlant, response_sender)
.await;
}
OrderPlantCommand::SetLogin => {
self.core.handle_set_login();
}
OrderPlantCommand::Logout { response_sender } => {
self.core.handle_logout(response_sender).await;
}
OrderPlantCommand::UpdateHeartbeat { seconds } => {
self.core.handle_update_heartbeat(seconds);
}
OrderPlantCommand::AccountList { response_sender } => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_account_list(self.login_scope.get());
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::SubscribeOrderUpdates {
account,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_subscribe_for_order_updates(&account);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::SubscribeBracketUpdates {
account,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_subscribe_to_bracket_updates(&account);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::PlaceBracketOrder {
bracket_order,
account,
response_sender,
} => {
let trade_route = match self.trade_routes.resolve(
bracket_order.trade_route.as_deref(),
&bracket_order.exchange,
) {
Ok(trade_route) => trade_route,
Err(err) => {
let _ = response_sender.send(Err(err));
return;
}
};
let (req_buf, id) = self.core.rithmic_sender_api.request_bracket_order(
*bracket_order,
&account,
self.login_scope.get(),
&trade_route,
);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::ModifyOrder {
order,
account,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_modify_order(&order, &account);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::CancelOrder {
order,
account,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_cancel_order(&order, &account);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::ModifyStop {
adjustment,
account,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_update_stop_bracket_level(&adjustment, &account);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::ModifyTarget {
adjustment,
account,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_update_target_bracket_level(&adjustment, &account);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::ShowOrders {
account,
response_sender,
} => {
let (req_buf, id) = self.core.rithmic_sender_api.request_show_orders(&account);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::CancelAllOrders {
command,
account,
response_sender,
} => {
let (req_buf, id) = self.core.rithmic_sender_api.request_cancel_all_orders(
&command,
&account,
self.login_scope.get(),
);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::GetAccountRmsInfo {
account,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_account_rms_info(&account, self.login_scope.get());
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::GetProductRmsInfo {
account,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_product_rms_info(&account);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::GetTradeRoutes {
subscribe_for_updates,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_trade_routes(subscribe_for_updates);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::RecordTradeRoutes(responses) => {
let loaded = responses
.iter()
.filter(|response| self.trade_routes.record_response(response))
.count();
match loaded {
0 => error!(
"order_plant: no trade routes published, orders will fail with NoTradeRoute"
),
loaded => info!("order_plant: {} trade routes loaded", loaded),
}
}
OrderPlantCommand::RecordTradeRouteUpdate(update) => {
self.trade_routes.record_update(&update);
}
OrderPlantCommand::TradeRouteFor {
exchange,
response_sender,
} => {
let _ = response_sender.send(self.trade_routes.resolve(None, &exchange));
}
OrderPlantCommand::ShowOrderHistoryDates { response_sender } => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_show_order_history_dates();
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::ShowOrderHistorySummary {
date,
account,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_show_order_history_summary(&date, &account);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::ShowOrderHistoryDetail {
basket_id,
date,
account,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_show_order_history_detail(&basket_id, &date, &account);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::ShowOrderHistory {
basket_id,
account,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_show_order_history(basket_id.as_deref(), &account);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::PlaceOrder {
order,
account,
response_sender,
} => {
let trade_route = match self
.trade_routes
.resolve(order.trade_route.as_deref(), &order.exchange)
{
Ok(trade_route) => trade_route,
Err(err) => {
let _ = response_sender.send(Err(err));
return;
}
};
let (req_buf, id) =
self.core
.rithmic_sender_api
.request_order(&order, &account, &trade_route);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::PlaceOcoOrder {
order,
account,
response_sender,
} => {
let timing = order.cancel_timing();
let legs = match self.trade_routes.resolve_legs(order.legs) {
Ok(legs) => legs,
Err(err) => {
let _ = response_sender.send(Err(err));
return;
}
};
let (req_buf, id) = match self
.core
.rithmic_sender_api
.request_oco_order(legs, timing, &account)
{
Ok(request) => request,
Err(err) => {
let _ = response_sender.send(Err(err));
return;
}
};
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::ShowBrackets {
account,
response_sender,
} => {
let (req_buf, id) = self.core.rithmic_sender_api.request_show_brackets(&account);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::ShowBracketStops {
account,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_show_bracket_stops(&account);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::ExitPosition {
command,
account,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_exit_position(&command, &account);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::LinkOrders {
command,
account,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_link_orders(command, &account);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::GetEasyToBorrowList {
request_type,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_easy_to_borrow_list(request_type);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::ModifyOrderReferenceData {
command,
account,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_modify_order_reference_data(&command, &account);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::GetOrderSessionConfig {
should_defer_request,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_order_session_config(should_defer_request);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::ReplayExecutions {
start_index_sec,
finish_index_sec,
account,
response_sender,
} => {
let (req_buf, id) = self.core.rithmic_sender_api.request_replay_executions(
start_index_sec,
finish_index_sec,
&account,
);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::GetUserInfo {
user,
account,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_get_user_info(user.as_deref(), &account);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::ShowFillHistory {
range,
max_record_count,
account,
response_sender,
} => {
let (req_buf, id) = self.core.rithmic_sender_api.request_show_fill_history(
range,
max_record_count,
&account,
);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::SubscribeAccountRmsUpdates {
subscribe,
update_bits,
account,
response_sender,
} => {
let (req_buf, id) = self.core.rithmic_sender_api.request_account_rms_updates(
subscribe,
update_bits,
&account,
);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::GetLoginInfo { response_sender } => {
let (req_buf, id) = self.core.rithmic_sender_api.request_login_info();
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::ListUnacceptedAgreements { response_sender } => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_list_unaccepted_agreements();
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::ListAcceptedAgreements { response_sender } => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_list_accepted_agreements();
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::AcceptAgreement {
agreement_id,
market_data_usage_capacity,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_accept_agreement(&agreement_id, market_data_usage_capacity.as_deref());
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::ShowAgreement {
agreement_id,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_show_agreement(&agreement_id);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::SetRithmicMrktDataSelfCertStatus {
agreement_id,
market_data_usage_capacity,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_set_rithmic_mrkt_data_self_cert_status(
&agreement_id,
&market_data_usage_capacity,
);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::ListExchangePermissions {
user,
response_sender,
} => {
let (req_buf, id) = self
.core
.rithmic_sender_api
.request_list_exchange_permissions(&user);
self.core
.register_and_send(req_buf, id, response_sender)
.await;
}
OrderPlantCommand::Abort => {
unreachable!("Abort is handled in run() before handle_command");
}
}
}
}
pub struct RithmicOrderPlantHandle {
account: Arc<RithmicAccount>,
login_scope: Arc<OnceLock<LoginScope>>,
sender: mpsc::Sender<OrderPlantCommand>,
pub subscription_receiver: SubscriptionFilter,
}
impl std::fmt::Debug for RithmicOrderPlantHandle {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("RithmicOrderPlantHandle")
.field("account", &self.account)
.field("sender", &self.sender)
.finish_non_exhaustive()
}
}
impl RithmicOrderPlantHandle {
pub async fn get_system_info(&self) -> Result<RithmicResponse, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::GetSystemInfo {
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_first_response(rx).await
}
pub async fn login(&self) -> Result<RithmicResponse, RithmicError> {
self.login_with_config(LoginConfig::default()).await
}
pub async fn login_with_config(
&self,
config: LoginConfig,
) -> Result<RithmicResponse, RithmicError> {
info!("order_plant: logging in");
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let mut config = config;
config.aggregated_quotes = None;
let command = OrderPlantCommand::Login {
config,
response_sender: tx,
};
let _ = self.sender.send(command).await;
let response = await_first_response(rx).await?;
if let Some(err) = response.error.clone() {
error!("order_plant: login failed {:?}", err);
return Err(err);
}
let _ = self.sender.send(OrderPlantCommand::SetLogin).await;
if let RithmicMessage::ResponseLogin(resp) = &response.message {
if let Some(hb) = resp.heartbeat_interval {
let secs = hb as u64;
self.update_heartbeat(secs).await;
}
if let Some(session_id) = &resp.unique_user_id {
info!("order_plant: session id: {}", session_id);
}
}
match self.get_login_info().await {
Ok(response) => {
if let Some(err) = &response.error {
warn!(
"order_plant: login info rejected, account list will be unscoped: {:?}",
err
);
}
}
Err(err) => warn!(
"order_plant: login info unavailable, account list will be unscoped: {:?}",
err
),
}
self.prime_trade_routes().await;
info!("order_plant: logged in");
Ok(response)
}
async fn prime_trade_routes(&self) {
match self.get_trade_routes(true).await {
Ok(responses) => {
for rejection in responses.iter().filter_map(|resp| resp.error.as_ref()) {
error!(
"order_plant: trade route request rejected, orders will fail: {}",
rejection
);
}
let _ = self
.sender
.send(OrderPlantCommand::RecordTradeRoutes(responses))
.await;
}
Err(err) => error!(
"order_plant: trade routes unavailable, orders will fail: {}",
err
),
}
}
pub async fn disconnect(&self) -> Result<RithmicResponse, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::Logout {
response_sender: tx,
};
let _ = self.sender.send(command).await;
let outcome = rx.await.map_err(|_| RithmicError::ConnectionClosed);
let _ = self.sender.send(OrderPlantCommand::Close).await;
outcome??
.into_iter()
.next()
.ok_or(RithmicError::EmptyResponse)
}
pub fn abort(&self) {
let _ = self.sender.try_send(OrderPlantCommand::Abort);
}
pub async fn get_account_list(&self) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
if self.login_scope.get().is_none() {
warn!("order_plant: no login info retained, listing accounts unscoped");
}
let command = OrderPlantCommand::AccountList {
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn subscribe_order_updates(&self) -> Result<RithmicResponse, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::SubscribeOrderUpdates {
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_first_response(rx).await
}
pub async fn subscribe_bracket_updates(&self) -> Result<RithmicResponse, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::SubscribeBracketUpdates {
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_first_response(rx).await
}
pub async fn place_bracket_order(
&self,
bracket_order: RithmicBracketOrder,
) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::PlaceBracketOrder {
bracket_order: Box::new(bracket_order),
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn modify_order(
&self,
order: RithmicModifyOrder,
) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::ModifyOrder {
order,
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn cancel_order(
&self,
order: RithmicCancelOrder,
) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::CancelOrder {
order,
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn adjust_target(
&self,
adjustment: RithmicBracketLevelAdjustment,
) -> Result<RithmicResponse, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::ModifyTarget {
adjustment,
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_first_response(rx).await
}
pub async fn adjust_stop(
&self,
adjustment: RithmicBracketLevelAdjustment,
) -> Result<RithmicResponse, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::ModifyStop {
adjustment,
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_first_response(rx).await
}
pub async fn show_orders(&self) -> Result<RithmicResponse, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::ShowOrders {
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_first_response(rx).await
}
async fn update_heartbeat(&self, seconds: u64) {
let command = OrderPlantCommand::UpdateHeartbeat { seconds };
let _ = self.sender.send(command).await;
}
pub async fn cancel_all_orders(
&self,
command: RithmicCancelAllOrders,
) -> Result<RithmicResponse, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::CancelAllOrders {
command,
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_first_response(rx).await
}
pub async fn get_account_rms_info(&self) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::GetAccountRmsInfo {
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn get_product_rms_info(&self) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::GetProductRmsInfo {
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn get_trade_routes(
&self,
subscribe_for_updates: bool,
) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::GetTradeRoutes {
subscribe_for_updates,
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn record_trade_route(&self, update: &TradeRoute) -> Result<(), RithmicError> {
self.sender
.send(OrderPlantCommand::RecordTradeRouteUpdate(Box::new(
update.clone(),
)))
.await
.map_err(|_| RithmicError::ConnectionClosed)
}
pub async fn trade_route_for(&self, exchange: &str) -> Result<String, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<String, RithmicError>>();
let command = OrderPlantCommand::TradeRouteFor {
exchange: exchange.to_string(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
rx.await.map_err(|_| RithmicError::ConnectionClosed)?
}
pub async fn show_order_history_dates(&self) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::ShowOrderHistoryDates {
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn show_order_history_summary(
&self,
date: &str,
) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::ShowOrderHistorySummary {
date: date.to_string(),
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn show_order_history_detail(
&self,
basket_id: &str,
date: &str,
) -> Result<RithmicResponse, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::ShowOrderHistoryDetail {
basket_id: basket_id.to_string(),
date: date.to_string(),
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_first_response(rx).await
}
pub async fn show_order_history(
&self,
basket_id: Option<&str>,
) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::ShowOrderHistory {
basket_id: basket_id.map(|s| s.to_string()),
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn place_order(
&self,
order: RithmicOrder,
) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::PlaceOrder {
order,
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn place_oco_order(
&self,
order: RithmicOcoOrder,
) -> Result<Vec<RithmicResponse>, RithmicError> {
if order.legs.len() < 2 {
return Err(RithmicError::InvalidArgument(
"OCO order requires at least 2 legs".to_string(),
));
}
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::PlaceOcoOrder {
order,
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn show_brackets(&self) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::ShowBrackets {
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn show_bracket_stops(&self) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::ShowBracketStops {
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn exit_position(
&self,
command: RithmicExitPosition,
) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::ExitPosition {
command,
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn link_orders(
&self,
command: RithmicLinkOrders,
) -> Result<RithmicResponse, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::LinkOrders {
command,
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_first_response(rx).await
}
pub async fn get_easy_to_borrow_list(
&self,
request_type: EasyToBorrowRequest,
) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::GetEasyToBorrowList {
request_type,
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn modify_order_reference_data(
&self,
command: RithmicModifyOrderReferenceData,
) -> Result<RithmicResponse, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::ModifyOrderReferenceData {
command,
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_first_response(rx).await
}
pub async fn get_order_session_config(
&self,
should_defer_request: Option<bool>,
) -> Result<RithmicResponse, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::GetOrderSessionConfig {
should_defer_request,
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_first_response(rx).await
}
pub async fn replay_executions(
&self,
start_index_sec: i32,
finish_index_sec: i32,
) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::ReplayExecutions {
start_index_sec,
finish_index_sec,
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn get_user_info(
&self,
user: Option<&str>,
) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::GetUserInfo {
user: user.map(str::to_string),
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn show_fill_history(
&self,
range: FillHistoryRange,
max_record_count: Option<i32>,
) -> Result<Vec<RithmicResponse>, RithmicError> {
if let Some(count) = max_record_count.filter(|count| !(0..=10_000).contains(count)) {
return Err(RithmicError::InvalidArgument(format!(
"max_record_count must be between 0 and 10,000, got {count}"
)));
}
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::ShowFillHistory {
range,
max_record_count,
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn subscribe_account_rms_updates(
&self,
subscribe: bool,
update_bits: Vec<RmsUpdateBits>,
) -> Result<RithmicResponse, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::SubscribeAccountRmsUpdates {
subscribe,
update_bits,
account: self.account.clone(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_first_response(rx).await
}
pub async fn get_login_info(&self) -> Result<RithmicResponse, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::GetLoginInfo {
response_sender: tx,
};
let _ = self.sender.send(command).await;
let response = await_first_response(rx).await?;
let scope = match &response.message {
RithmicMessage::ResponseLoginInfo(info) if response.error.is_none() => {
LoginScope::from_login_info(info)
}
_ => None,
};
if let Some(scope) = scope {
let _ = self.login_scope.set(scope);
}
Ok(response)
}
pub async fn list_unaccepted_agreements(&self) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::ListUnacceptedAgreements {
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn list_accepted_agreements(&self) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::ListAcceptedAgreements {
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn accept_agreement(
&self,
agreement_id: &str,
market_data_usage_capacity: Option<&str>,
) -> Result<RithmicResponse, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::AcceptAgreement {
agreement_id: agreement_id.to_string(),
market_data_usage_capacity: market_data_usage_capacity.map(|s| s.to_string()),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_first_response(rx).await
}
pub async fn show_agreement(
&self,
agreement_id: &str,
) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::ShowAgreement {
agreement_id: agreement_id.to_string(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
pub async fn set_rithmic_mrkt_data_self_cert_status(
&self,
agreement_id: &str,
market_data_usage_capacity: &str,
) -> Result<RithmicResponse, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::SetRithmicMrktDataSelfCertStatus {
agreement_id: agreement_id.to_string(),
market_data_usage_capacity: market_data_usage_capacity.to_string(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_first_response(rx).await
}
pub async fn list_exchange_permissions(
&self,
user: &str,
) -> Result<Vec<RithmicResponse>, RithmicError> {
let (tx, rx) = oneshot::channel::<Result<Vec<RithmicResponse>, RithmicError>>();
let command = OrderPlantCommand::ListExchangePermissions {
user: user.to_string(),
response_sender: tx,
};
let _ = self.sender.send(command).await;
await_all_responses(rx).await
}
}
impl Clone for RithmicOrderPlantHandle {
fn clone(&self) -> Self {
RithmicOrderPlantHandle {
account: Arc::clone(&self.account),
login_scope: Arc::clone(&self.login_scope),
sender: self.sender.clone(),
subscription_receiver: self.subscription_receiver.resubscribe(),
}
}
}
#[cfg(test)]
mod tests;