use super::ServiceContext;
use super::correlation_id::{
optional_client_order_id, optional_request_id, require_client_style_id,
};
use super::scope;
use super::unary;
use crate::codecs::decode::{
batch_cancel_from_proto, batch_create_from_proto, batch_replace_from_proto,
batch_replace_status_from_proto, cancel_all_after_from_proto, cancel_all_from_proto,
get_order_from_proto, modify_order_from_proto, order_mutation_from_cancel,
order_mutation_from_create, orders_list_from_history, orders_list_from_open,
preview_order_from_proto, user_trades_list_from_proto,
};
use crate::codecs::scalars::id_to_u64;
use crate::connect::orders::v1::{OrdersReadServiceClient, OrdersServiceClient};
use crate::errors::{Error, Result};
use crate::models::{
AttachedRisk, BatchCancelItem, BatchCancelOrdersResult, BatchCreateOrdersResult,
BatchReplaceItem, BatchReplaceOrdersResult, BatchReplaceStatusResult, CancelAllAfterResult,
CancelAllOpts, CancelAllOrdersResult, CancelOrderParams, CreateOrderParams, CreateOrderType,
CreateSide, CreateTimeInForce, FeeAsset, GetOrderOpts, GetOrderResult, ListOpenOrdersOpts,
ListOrderHistoryOpts, MaxSlippage, ModifyOrderParams, ModifyOrderResult, Order, OrderKey,
OrderMutationResult, OrderSelfTradePrevention, OrdersList, PreviewOrderParams,
PreviewOrderResult, RiskLeg, TrailingDistance, TrailingStop, UserTrade, UserTradesList,
};
use crate::proto::orders::v1::{
BatchCancelItem as ProtoBatchCancelItem, BatchCancelOrdersRequest, BatchCreateOrdersRequest,
BatchReplaceOrderItem as ProtoBatchReplaceOrderItem, BatchReplaceOrdersRequest,
CancelAllAfterRequest, CancelAllOrdersRequest, CancelOrderRequest, CreateOrderRequest,
FeeAsset as ProtoFeeAsset, GetBatchReplaceStatusRequest, GetOpenOrdersRequest,
GetOrderHistoryRequest, GetOrderRequest, GetUserTradesRequest, LimitFok, LimitGtc, LimitIoc,
MarketIoc, ModifyBehavior, ModifyOrderRequest, OrderIntent, PreviewOrderRequest, RiskExecution,
RiskLimitGtc, RiskPolicy, SelfTradePreventionMode, Side, StopLossPolicy, TakeProfitPolicy,
TrailingStopPolicy, batch_replace_order_item, cancel_order_request, get_order_request,
market_ioc, modify_order_request, order_intent, risk_execution, risk_policy,
trailing_stop_policy,
};
use crate::types::{
Price, Quantity, resolve_price_ticks, resolve_qty_scaled, resolve_quote_qty_scaled,
};
use rand_core::{OsRng, RngCore};
use std::time::Duration;
#[derive(Clone)]
pub struct OrdersService {
ctx: ServiceContext,
}
impl OrdersService {
const MAX_BATCH_ITEMS: usize = 20;
pub fn new(ctx: ServiceContext) -> Self {
Self { ctx }
}
fn write_client(&self) -> OrdersServiceClient<crate::transport::SharedTransport> {
OrdersServiceClient::new(
self.ctx.factory.transport(),
self.ctx.factory.connect_config(),
)
}
fn read_client(&self) -> OrdersReadServiceClient<crate::transport::SharedTransport> {
OrdersReadServiceClient::new(
self.ctx.factory.transport(),
self.ctx.factory.connect_config(),
)
}
pub async fn list_open(&self, subaccount_id: Option<u64>) -> Result<OrdersList> {
self.list_open_with(ListOpenOrdersOpts {
subaccount_id,
..Default::default()
})
.await
}
pub async fn list_open_with(&self, opts: ListOpenOrdersOpts) -> Result<OrdersList> {
let req = GetOpenOrdersRequest {
subaccount_id: scope::optional_subaccount(&self.ctx, opts.subaccount_id)?,
page_token: opts.page_token.unwrap_or_default(),
limit: opts.limit,
include_attached_risk: Some(opts.include_attached_risk),
include_attached_risk_state: Some(opts.include_attached_risk_state),
..Default::default()
};
let client = self.read_client();
let resp = unary::await_auth(
&self.ctx.factory,
"/orders.v1.OrdersReadService/GetOpenOrders",
req,
|req, opts| client.get_open_orders_with_options(req, opts),
)
.await?
.into_owned();
Ok(orders_list_from_open(&resp))
}
pub async fn list_history(
&self,
subaccount_id: Option<u64>,
limit: Option<u32>,
) -> Result<OrdersList> {
self.list_history_with(ListOrderHistoryOpts {
subaccount_id,
limit,
..Default::default()
})
.await
}
pub async fn list_history_with(&self, opts: ListOrderHistoryOpts) -> Result<OrdersList> {
let mut symbol_ids = Vec::new();
if let Some(sid) = opts.symbol_id {
if sid == 0 {
return Err(Error::validation(
"symbol_id must be non-zero when explicitly supplied",
));
}
symbol_ids.push(sid);
} else if let Some(ref symbol) = opts.symbol {
let resolved = self
.ctx
.catalogs
.symbol_id_for_symbol(symbol)
.ok_or_else(|| {
Error::validation(format!(
"unknown symbol {symbol}; call hydrate_catalogs / get_spot_config first"
))
})?;
symbol_ids.push(resolved);
}
let req = GetOrderHistoryRequest {
subaccount_id: scope::optional_subaccount(&self.ctx, opts.subaccount_id)?,
symbol_id: symbol_ids,
page_token: opts.page_token.unwrap_or_default(),
limit: opts.limit,
include_attached_risk: Some(opts.include_attached_risk),
include_attached_risk_state: Some(opts.include_attached_risk_state),
..Default::default()
};
let client = self.read_client();
let resp = unary::await_auth(
&self.ctx.factory,
"/orders.v1.OrdersReadService/GetOrderHistory",
req,
|req, opts| client.get_order_history_with_options(req, opts),
)
.await?
.into_owned();
Ok(orders_list_from_history(&resp))
}
pub async fn get(&self, key: OrderKey, subaccount_id: Option<u64>) -> Result<GetOrderResult> {
self.get_with(GetOrderOpts {
key,
subaccount_id,
include_attached_risk: false,
include_attached_risk_state: false,
})
.await
}
pub async fn get_with(&self, opts: GetOrderOpts) -> Result<GetOrderResult> {
let key = Some(Self::encode_get_order_key(&opts.key)?);
let req = GetOrderRequest {
subaccount_id: scope::optional_subaccount(&self.ctx, opts.subaccount_id)?,
key,
include_attached_risk: Some(opts.include_attached_risk),
include_attached_risk_state: Some(opts.include_attached_risk_state),
..Default::default()
};
let client = self.read_client();
let resp = unary::await_auth(
&self.ctx.factory,
"/orders.v1.OrdersReadService/GetOrder",
req,
|req, opts| client.get_order_with_options(req, opts),
)
.await?
.into_owned();
Ok(get_order_from_proto(&resp))
}
pub async fn wait_for_order_trades_complete(
&self,
key: OrderKey,
timeout: Duration,
) -> Result<GetOrderResult> {
let timeout = if timeout.is_zero() {
Duration::from_secs(15)
} else {
timeout
};
let deadline = tokio::time::Instant::now() + timeout;
loop {
let last = tokio::time::timeout_at(deadline, self.get(key.clone(), None))
.await
.map_err(|_| {
Error::transport(format!(
"timed out waiting for order trades to match cum_qty (key={key:?})"
))
})??;
if order_trades_projection_complete(&last) {
return Ok(last);
}
if tokio::time::Instant::now() >= deadline {
return Err(Error::transport(format!(
"timed out waiting for order trades to match cum_qty (key={key:?})"
)));
}
tokio::time::sleep_until(
deadline.min(tokio::time::Instant::now() + Duration::from_millis(100)),
)
.await;
}
}
fn encode_get_order_key(key: &OrderKey) -> Result<get_order_request::Key> {
match key {
OrderKey::OrderId(oid) => {
Ok(get_order_request::Key::OrderId(id_to_u64(oid, "order_id")?))
}
OrderKey::ClientOrderId(cid) => Ok(get_order_request::Key::ClientOrderId(
require_client_style_id(cid, "client_order_id")?,
)),
}
}
fn encode_cancel_order_key(key: &OrderKey) -> Result<cancel_order_request::Key> {
match key {
OrderKey::OrderId(oid) => Ok(cancel_order_request::Key::OrderId(id_to_u64(
oid, "order_id",
)?)),
OrderKey::ClientOrderId(cid) => Ok(cancel_order_request::Key::ClientOrderId(
require_client_style_id(cid, "client_order_id")?,
)),
}
}
fn encode_modify_order_key(key: &OrderKey) -> Result<modify_order_request::Key> {
match key {
OrderKey::OrderId(oid) => Ok(modify_order_request::Key::OrderId(id_to_u64(
oid, "order_id",
)?)),
OrderKey::ClientOrderId(cid) => Ok(modify_order_request::Key::ClientOrderId(
require_client_style_id(cid, "client_order_id")?,
)),
}
}
fn encode_batch_replace_key(key: &OrderKey) -> Result<batch_replace_order_item::Key> {
match key {
OrderKey::OrderId(oid) => Ok(batch_replace_order_item::Key::OrderId(id_to_u64(
oid, "order_id",
)?)),
OrderKey::ClientOrderId(cid) => Ok(batch_replace_order_item::Key::ClientOrderId(
require_client_style_id(cid, "client_order_id")?,
)),
}
}
fn require_quantity_scale(&self, symbol: &str, qty_scale: Option<u32>) -> Result<u32> {
if let Some(scale) = self.ctx.catalogs.base_quantity_scale_for_symbol(symbol) {
return Ok(scale);
}
if let Some(scale) = qty_scale {
return Ok(scale);
}
Err(Error::validation(format!(
"quantity scale for {symbol:?} is unavailable; await client.wait_for_catalogs() before placing orders, or pass a scaled Quantity"
)))
}
fn require_quote_quantity_scale(&self, symbol: &str) -> Result<u32> {
self.ctx
.catalogs
.quote_quantity_scale_for_symbol(symbol)
.ok_or_else(|| {
Error::validation(format!(
"quote quantity scale for {symbol:?} is unavailable; await client.wait_for_catalogs() before using a quote-debit budget"
))
})
}
fn validate_batch_size(operation: &str, len: usize) -> Result<()> {
if len == 0 {
return Err(Error::validation(format!(
"{operation} requires at least one item"
)));
}
if len > Self::MAX_BATCH_ITEMS {
return Err(Error::validation(format!(
"{operation} accepts at most {} items; received {len}",
Self::MAX_BATCH_ITEMS
)));
}
Ok(())
}
fn order_intent_from_params(&self, params: &CreateOrderParams) -> Result<OrderIntent> {
let mut intent = OrderIntent {
symbol: params.symbol.clone(),
side: match params.side {
CreateSide::Buy => Side::Buy.into(),
CreateSide::Sell => Side::Sell.into(),
},
..Default::default()
};
if let Some(client_order_id) = optional_client_order_id(params.client_order_id.as_deref())?
{
intent.client_order_id = client_order_id;
}
intent.sizing = Some(match (¶ms.quantity, ¶ms.max_quote_debit_scaled) {
(Some(quantity), None) => {
let scale = self.require_quantity_scale(¶ms.symbol, quantity.scale())?;
order_intent::Sizing::BaseQtyScaled(resolve_qty_scaled(
quantity,
scale,
Some(¶ms.symbol),
self.ctx.catalogs.symbol_id_for_symbol(¶ms.symbol),
)?)
}
(None, Some(max_quote_debit)) => {
let scale = self.require_quote_quantity_scale(¶ms.symbol)?;
order_intent::Sizing::MaxQuoteDebitScaled(resolve_quote_qty_scaled(
max_quote_debit,
scale,
Some(¶ms.symbol),
self.ctx.catalogs.symbol_id_for_symbol(¶ms.symbol),
)?)
}
(Some(_), Some(_)) | (None, None) => {
return Err(Error::validation(
"set exactly one of quantity or max_quote_debit_scaled",
));
}
});
intent.fee_asset = match params.fee_asset.unwrap_or(FeeAsset::Quote) {
FeeAsset::Quote => ProtoFeeAsset::Quote.into(),
FeeAsset::Base if matches!(params.side, CreateSide::Buy) => ProtoFeeAsset::Base.into(),
FeeAsset::Base => {
return Err(Error::validation(
"fee_asset=base is only valid for BUY orders",
));
}
};
intent.self_trade_prevention_mode = match params
.self_trade_prevention
.unwrap_or(OrderSelfTradePrevention::ExpireMaker)
{
OrderSelfTradePrevention::ExpireTaker => SelfTradePreventionMode::ExpireTaker.into(),
OrderSelfTradePrevention::ExpireMaker => SelfTradePreventionMode::ExpireMaker.into(),
OrderSelfTradePrevention::ExpireBoth => SelfTradePreventionMode::ExpireBoth.into(),
};
let post_only = params.post_only.unwrap_or(false);
intent.execution = Some(match params.order_type {
CreateOrderType::Market => {
if post_only {
return Err(Error::validation(
"post_only is not supported for market orders",
));
}
if params.price.is_some() {
return Err(Error::validation(
"price is not valid for market orders; use market_client_ref_price for a reservation reference",
));
}
let mut market = MarketIoc::default();
if let Some(ref_price) = params.market_client_ref_price.as_ref() {
market.client_ref_price_ticks =
resolve_price_ticks(ref_price, Some(¶ms.symbol))?;
}
market.max_slippage = match params.market_max_slippage {
Some(MaxSlippage::Ticks(value)) if value > 0 => {
Some(market_ioc::MaxSlippage::MaxSlippageTicks(value))
}
Some(MaxSlippage::Bps(value)) if value > 0 => {
Some(market_ioc::MaxSlippage::MaxSlippageBps(value))
}
Some(_) => {
return Err(Error::validation("market_max_slippage must be positive"));
}
None => None,
};
order_intent::Execution::MarketIoc(Box::new(market))
}
CreateOrderType::Limit => {
let price = params.price.as_ref().ok_or_else(|| {
Error::validation(
"price is required for limit orders (use Price::from_decimal or Price::from_ticks)",
)
})?;
let price_ticks = resolve_price_ticks(price, Some(¶ms.symbol))?;
match params.time_in_force {
Some(CreateTimeInForce::Ioc) => {
if post_only {
return Err(Error::validation(
"post_only is not supported for ioc limit orders",
));
}
order_intent::Execution::LimitIoc(Box::new(LimitIoc {
price_ticks,
..Default::default()
}))
}
Some(CreateTimeInForce::Fok) => {
if post_only {
return Err(Error::validation(
"post_only is not supported for fok limit orders",
));
}
order_intent::Execution::LimitFok(Box::new(LimitFok {
price_ticks,
..Default::default()
}))
}
_ => order_intent::Execution::LimitGtc(Box::new(LimitGtc {
price_ticks,
post_only,
..Default::default()
})),
}
}
});
if let Some(risk) = params.attached_risk.as_ref() {
*intent.attached_risk.get_or_insert_default() =
Self::encode_attached_risk(risk, Some(¶ms.symbol))?;
}
Ok(intent)
}
fn encode_create_params(&self, params: &CreateOrderParams) -> Result<CreateOrderRequest> {
let order = self.order_intent_from_params(params)?;
let mut req = CreateOrderRequest {
subaccount_id: scope::optional_subaccount(&self.ctx, params.subaccount_id)?,
..Default::default()
};
*req.order.get_or_insert_default() = order;
Ok(req)
}
#[allow(deprecated)]
fn encode_risk_child(leg: &RiskLeg, symbol: Option<&str>) -> Result<RiskExecution> {
if leg.trigger_price_source.is_some() {
return Err(Error::validation(
"attached risk always uses last trade; trigger_price_source cannot be supplied",
));
}
let child_ty = leg.order_type.unwrap_or(CreateOrderType::Market);
let execution = match (child_ty, leg.limit_price.as_ref()) {
(CreateOrderType::Market, None) => risk_execution::Execution::MarketIoc(Box::default()),
(CreateOrderType::Market, Some(_)) => {
return Err(Error::validation(
"attached_risk MARKET child must not set limit_price",
));
}
(CreateOrderType::Limit, Some(price)) => {
risk_execution::Execution::LimitGtc(Box::new(RiskLimitGtc {
price_ticks: resolve_price_ticks(price, symbol)?,
..Default::default()
}))
}
(CreateOrderType::Limit, None) => {
return Err(Error::validation(
"attached_risk LIMIT child requires limit_price",
));
}
};
Ok(RiskExecution {
execution: Some(execution),
..Default::default()
})
}
fn encode_take_profit(leg: &RiskLeg, symbol: Option<&str>) -> Result<TakeProfitPolicy> {
let mut policy = TakeProfitPolicy {
trigger_price_ticks: resolve_price_ticks(&leg.trigger_price, symbol)?,
..Default::default()
};
*policy.child.get_or_insert_default() = Self::encode_risk_child(leg, symbol)?;
Ok(policy)
}
fn encode_stop_loss(leg: &RiskLeg, symbol: Option<&str>) -> Result<StopLossPolicy> {
let mut policy = StopLossPolicy {
trigger_price_ticks: resolve_price_ticks(&leg.trigger_price, symbol)?,
..Default::default()
};
*policy.child.get_or_insert_default() = Self::encode_risk_child(leg, symbol)?;
Ok(policy)
}
#[allow(deprecated)]
fn encode_trailing_stop(
stop: &TrailingStop,
symbol: Option<&str>,
) -> Result<TrailingStopPolicy> {
if stop.trigger_price_source.is_some() {
return Err(Error::validation(
"attached risk always uses last trade; trigger_price_source cannot be supplied",
));
}
if stop.order_type.is_some() {
return Err(Error::validation(
"attached trailing_stop child is always market; order_type cannot be supplied",
));
}
let mut proto = TrailingStopPolicy::default();
if let Some(activation) = stop.activation_price.as_ref() {
proto.activation_price_ticks = resolve_price_ticks(activation, symbol)?;
}
proto.trailing_distance = Some(match stop.distance {
TrailingDistance::Ticks(v) => {
if v <= 0 {
return Err(Error::validation(
"trailing_distance_ticks must be positive",
));
}
trailing_stop_policy::TrailingDistance::TrailingDistanceTicks(v)
}
TrailingDistance::Bps(v) => {
if v <= 0 {
return Err(Error::validation("trailing_distance_bps must be positive"));
}
trailing_stop_policy::TrailingDistance::TrailingDistanceBps(v)
}
});
if let Some(slip) = stop.max_slippage {
proto.max_slippage = Some(match slip {
MaxSlippage::Ticks(v) => {
if v <= 0 {
return Err(Error::validation("max_slippage_ticks must be positive"));
}
trailing_stop_policy::MaxSlippage::MaxSlippageTicks(v)
}
MaxSlippage::Bps(v) => {
if v <= 0 {
return Err(Error::validation("max_slippage_bps must be positive"));
}
trailing_stop_policy::MaxSlippage::MaxSlippageBps(v)
}
});
}
Ok(proto)
}
fn encode_attached_risk(risk: &AttachedRisk, symbol: Option<&str>) -> Result<RiskPolicy> {
if risk.stop_loss.is_some() && risk.trailing_stop.is_some() {
return Err(Error::validation(
"attached_risk allows at most one of stop_loss or trailing_stop",
));
}
if risk.take_profit.is_none() && risk.stop_loss.is_none() && risk.trailing_stop.is_none() {
return Err(Error::validation(
"attached_risk requires take_profit and/or a stop leg",
));
}
let mut proto = RiskPolicy {
oco: risk.oco,
..Default::default()
};
if let Some(tp) = risk.take_profit.as_ref() {
*proto.take_profit.get_or_insert_default() = Self::encode_take_profit(tp, symbol)?;
}
if let Some(sl) = risk.stop_loss.as_ref() {
proto.stop_leg = Some(risk_policy::StopLeg::StopLoss(Box::new(
Self::encode_stop_loss(sl, symbol)?,
)));
} else if let Some(ts) = risk.trailing_stop.as_ref() {
proto.stop_leg = Some(risk_policy::StopLeg::TrailingStop(Box::new(
Self::encode_trailing_stop(ts, symbol)?,
)));
}
Ok(proto)
}
fn new_mutation_request_id(prefix: &str) -> Result<String> {
let mut random = [0_u8; 6];
OsRng
.try_fill_bytes(&mut random)
.map_err(|err| Error::transport(format!("secure randomness unavailable: {err}")))?;
Ok(format!("{prefix}-{}", hex::encode(random)))
}
fn coalesce_request_id(value: Option<String>, prefix: &str) -> Result<String> {
if let Some(id) = optional_request_id(value.as_deref())? {
return Ok(id);
}
Self::new_mutation_request_id(prefix)
}
fn modify_behavior(label: &str) -> Result<ModifyBehavior> {
match label.to_ascii_lowercase().as_str() {
"amend_or_replace" => Ok(ModifyBehavior::AmendOrReplace),
"amend_only" => Ok(ModifyBehavior::AmendOnly),
"replace_only" => Ok(ModifyBehavior::ReplaceOnly),
_ => Err(Error::validation(
"behavior must be amend_or_replace, amend_only, or replace_only",
)),
}
}
fn encode_modify_params(&self, params: ModifyOrderParams) -> Result<ModifyOrderRequest> {
if params.new_price.is_none()
&& params.new_qty.is_none()
&& params.new_attached_risk.is_none()
{
return Err(Error::validation(
"modify requires new_price, new_qty, and/or new_attached_risk",
));
}
let scale = self.require_quantity_scale(
¶ms.symbol,
params.new_qty.as_ref().and_then(Quantity::scale),
)?;
let mut req = ModifyOrderRequest {
subaccount_id: scope::optional_subaccount(&self.ctx, params.subaccount_id)?,
request_id: Self::coalesce_request_id(params.request_id, "mod")?,
key: Some(Self::encode_modify_order_key(¶ms.key)?),
..Default::default()
};
if let Some(price) = params.new_price.as_ref() {
req.new_price_ticks = Some(resolve_price_ticks(price, Some(¶ms.symbol))?);
}
if let Some(qty) = params.new_qty.as_ref() {
req.new_qty_scaled = Some(resolve_qty_scaled(
qty,
scale,
Some(¶ms.symbol),
self.ctx.catalogs.symbol_id_for_symbol(¶ms.symbol),
)?);
}
if let Some(risk) = params.new_attached_risk.as_ref() {
*req.new_attached_risk.get_or_insert_default() =
Self::encode_attached_risk(risk, Some(¶ms.symbol))?;
}
if let Some(behavior) = params.behavior.as_deref() {
req.behavior = Self::modify_behavior(behavior)?.into();
}
if let Some(ncid) = optional_client_order_id(params.new_client_order_id.as_deref())? {
req.new_client_order_id = ncid;
}
Ok(req)
}
pub async fn create(&self, params: CreateOrderParams) -> Result<OrderMutationResult> {
self.ctx.wait_for_catalogs().await?;
let req = self.encode_create_params(¶ms)?;
let client = self.write_client();
let resp = unary::await_auth(
&self.ctx.factory,
"/orders.v1.OrdersService/CreateOrder",
req,
|req, opts| client.create_order_with_options(req, opts),
)
.await?
.into_owned();
order_mutation_from_create(&resp)
}
pub async fn preview(&self, params: PreviewOrderParams) -> Result<PreviewOrderResult> {
self.ctx.wait_for_catalogs().await?;
let req = self.encode_preview_params(¶ms)?;
let client = self.write_client();
let resp = unary::await_auth(
&self.ctx.factory,
"/orders.v1.OrdersService/PreviewOrder",
req,
|req, opts| client.preview_order_with_options(req, opts),
)
.await?
.into_owned();
let base_scale = self.require_quantity_scale(¶ms.symbol, None)?;
preview_order_from_proto(
&resp,
base_scale,
¶ms.symbol,
self.ctx.catalogs.symbol_id_for_symbol(¶ms.symbol),
)
}
fn encode_preview_params(&self, params: &PreviewOrderParams) -> Result<PreviewOrderRequest> {
let create = CreateOrderParams {
symbol: params.symbol.clone(),
side: params.side,
order_type: params.order_type,
quantity: params.quantity.clone(),
max_quote_debit_scaled: params.max_quote_debit_scaled.clone(),
price: params.price.clone(),
time_in_force: params.time_in_force,
client_order_id: params.client_order_id.clone(),
subaccount_id: params.subaccount_id,
post_only: params.post_only,
market_client_ref_price: params.market_client_ref_price.clone(),
fee_asset: params.fee_asset,
self_trade_prevention: params.self_trade_prevention,
market_max_slippage: params.market_max_slippage,
attached_risk: params.attached_risk.clone(),
};
let order = self.order_intent_from_params(&create)?;
let mut req = PreviewOrderRequest {
subaccount_id: scope::optional_subaccount(&self.ctx, params.subaccount_id)?,
..Default::default()
};
*req.order.get_or_insert_default() = order;
Ok(req)
}
pub async fn batch_create(
&self,
items: Vec<CreateOrderParams>,
subaccount_id: Option<u64>,
request_id: Option<String>,
) -> Result<BatchCreateOrdersResult> {
Self::validate_batch_size("batch_create", items.len())?;
self.ctx.wait_for_catalogs().await?;
let mut encoded = Vec::with_capacity(items.len());
for item in &items {
encoded.push(self.order_intent_from_params(item)?);
}
let req = BatchCreateOrdersRequest {
subaccount_id: scope::optional_subaccount(&self.ctx, subaccount_id)?,
request_id: Self::coalesce_request_id(request_id, "batch-create")?,
items: encoded,
..Default::default()
};
let client = self.write_client();
let resp = unary::await_auth(
&self.ctx.factory,
"/orders.v1.OrdersService/BatchCreateOrders",
req,
|req, opts| client.batch_create_orders_with_options(req, opts),
)
.await?
.into_owned();
batch_create_from_proto(&resp)
}
pub async fn batch_cancel(
&self,
items: Vec<BatchCancelItem>,
subaccount_id: Option<u64>,
request_id: Option<String>,
) -> Result<BatchCancelOrdersResult> {
Self::validate_batch_size("batch_cancel", items.len())?;
let mut proto_items = Vec::with_capacity(items.len());
for item in items {
let mut proto = ProtoBatchCancelItem::default();
match &item.key {
OrderKey::OrderId(oid) => {
proto.order_id = id_to_u64(oid, "order_id")?;
}
OrderKey::ClientOrderId(cid) => {
proto.client_order_id = require_client_style_id(cid, "client_order_id")?;
}
}
if let Some(sid) = item.symbol_id {
proto.symbol_id = sid;
}
proto_items.push(proto);
}
let req = BatchCancelOrdersRequest {
subaccount_id: scope::optional_subaccount(&self.ctx, subaccount_id)?,
request_id: Self::coalesce_request_id(request_id, "batch-cancel")?,
items: proto_items,
..Default::default()
};
let client = self.write_client();
let resp = unary::await_auth(
&self.ctx.factory,
"/orders.v1.OrdersService/BatchCancelOrders",
req,
|req, opts| client.batch_cancel_orders_with_options(req, opts),
)
.await?
.into_owned();
batch_cancel_from_proto(&resp)
}
pub async fn batch_replace(
&self,
items: Vec<BatchReplaceItem>,
symbol: &str,
subaccount_id: Option<u64>,
request_id: Option<String>,
) -> Result<BatchReplaceOrdersResult> {
Self::validate_batch_size("batch_replace", items.len())?;
self.ctx.wait_for_catalogs().await?;
let symbol_id = self
.ctx
.catalogs
.symbol_id_for_symbol(symbol)
.ok_or_else(|| {
Error::validation(format!(
"unknown symbol {symbol}; call hydrate_catalogs / get_spot_config first"
))
})?;
let scale = Self::resolve_batch_replace_scale(&self.ctx.catalogs, symbol)?;
let mut proto_items = Vec::with_capacity(items.len());
for item in items {
if item.new_price.is_none()
&& item.new_qty.is_none()
&& item.new_attached_risk.is_none()
{
return Err(Error::validation(
"each batch item requires new_price, new_qty, and/or new_attached_risk",
));
}
let mut proto = ProtoBatchReplaceOrderItem {
key: Some(Self::encode_batch_replace_key(&item.key)?),
..Default::default()
};
if let Some(price) = item.new_price.as_ref() {
proto.new_price_ticks = Some(resolve_price_ticks(price, Some(symbol))?);
}
if let Some(qty) = item.new_qty.as_ref() {
proto.new_qty_scaled = Some(resolve_qty_scaled(
qty,
scale,
Some(symbol),
Some(symbol_id),
)?);
}
if let Some(risk) = item.new_attached_risk.as_ref() {
*proto.new_attached_risk.get_or_insert_default() =
Self::encode_attached_risk(risk, Some(symbol))?;
}
if let Some(ncid) = optional_client_order_id(item.new_client_order_id.as_deref())? {
proto.new_client_order_id = ncid;
}
proto_items.push(proto);
}
let req = BatchReplaceOrdersRequest {
subaccount_id: scope::optional_subaccount(&self.ctx, subaccount_id)?,
symbol_id,
request_id: Self::coalesce_request_id(request_id, "batch-replace")?,
items: proto_items,
..Default::default()
};
let client = self.write_client();
let resp = unary::await_auth(
&self.ctx.factory,
"/orders.v1.OrdersService/BatchReplaceOrders",
req,
|req, opts| client.batch_replace_orders_with_options(req, opts),
)
.await?
.into_owned();
batch_replace_from_proto(&resp)
}
pub async fn get_batch_replace_status(
&self,
batch_request_id: &str,
subaccount_id: Option<u64>,
) -> Result<BatchReplaceStatusResult> {
let req = GetBatchReplaceStatusRequest {
subaccount_id: scope::optional_subaccount(&self.ctx, subaccount_id)?,
batch_request_id: id_to_u64(batch_request_id, "batch_request_id")?,
..Default::default()
};
let client = self.read_client();
let resp = unary::await_auth(
&self.ctx.factory,
"/orders.v1.OrdersReadService/GetBatchReplaceStatus",
req,
|req, opts| client.get_batch_replace_status_with_options(req, opts),
)
.await?
.into_owned();
batch_replace_status_from_proto(&resp)
}
pub async fn cancel_all_after(
&self,
timeout_sec: u32,
symbol: Option<&str>,
subaccount_id: Option<u64>,
request_id: Option<String>,
) -> Result<CancelAllAfterResult> {
let req = CancelAllAfterRequest {
subaccount_id: scope::optional_subaccount(&self.ctx, subaccount_id)?,
timeout_sec,
symbol: symbol.unwrap_or("").to_owned(),
request_id: Self::coalesce_request_id(request_id, "cancel-after")?,
..Default::default()
};
let client = self.write_client();
let resp = unary::await_auth(
&self.ctx.factory,
"/orders.v1.OrdersService/CancelAllAfter",
req,
|req, opts| client.cancel_all_after_with_options(req, opts),
)
.await?
.into_owned();
cancel_all_after_from_proto(&resp)
}
pub async fn cancel(&self, req: CancelOrderRequest) -> Result<OrderMutationResult> {
let client = self.write_client();
let resp = unary::await_auth(
&self.ctx.factory,
"/orders.v1.OrdersService/CancelOrder",
req,
|req, opts| client.cancel_order_with_options(req, opts),
)
.await?
.into_owned();
order_mutation_from_cancel(&resp)
}
pub async fn cancel_with(&self, params: CancelOrderParams) -> Result<OrderMutationResult> {
if params.symbol_id.is_none() && params.symbol.is_some() {
self.ctx.wait_for_catalogs().await?;
}
let symbol_id = Self::resolve_cancel_symbol_id(
&self.ctx.catalogs,
params.symbol.as_deref(),
params.symbol_id,
)?;
let req = CancelOrderRequest {
symbol_id,
subaccount_id: scope::optional_subaccount(&self.ctx, params.subaccount_id)?,
key: Some(Self::encode_cancel_order_key(¶ms.key)?),
..Default::default()
};
self.cancel(req).await
}
fn resolve_cancel_symbol_id(
catalogs: &crate::catalogs::Manager,
symbol: Option<&str>,
symbol_id: Option<u32>,
) -> Result<u32> {
match (symbol, symbol_id) {
(None, None) => Ok(0),
(_, Some(0)) => Err(Error::validation(
"symbol_id must be non-zero when explicitly supplied",
)),
(Some(_), Some(_)) => Err(Error::validation(
"cancel accepts symbol or symbol_id, not both",
)),
(None, Some(symbol_id)) => Ok(symbol_id),
(Some(symbol), None) => catalogs.symbol_id_for_symbol(symbol).ok_or_else(|| {
Error::validation(format!(
"unknown symbol {symbol}; call hydrate_catalogs / get_spot_config first"
))
}),
}
}
pub async fn cancel_by_client_order_id(
&self,
client_order_id: &str,
symbol: Option<&str>,
subaccount_id: Option<u64>,
) -> Result<OrderMutationResult> {
self.cancel_with(CancelOrderParams {
key: OrderKey::ClientOrderId(client_order_id.to_owned()),
symbol: symbol.map(|s| s.to_owned()),
symbol_id: None,
subaccount_id,
})
.await
}
pub async fn cancel_by_order_id(
&self,
order_id: &str,
subaccount_id: Option<u64>,
) -> Result<OrderMutationResult> {
self.cancel_with(CancelOrderParams {
key: OrderKey::OrderId(order_id.to_owned()),
symbol: None,
symbol_id: None,
subaccount_id,
})
.await
}
pub async fn cancel_all(
&self,
symbol: Option<&str>,
dry_run: bool,
subaccount_id: Option<u64>,
) -> Result<CancelAllOrdersResult> {
self.cancel_all_with(CancelAllOpts {
symbol: symbol.map(|s| s.to_owned()),
dry_run,
subaccount_id,
..Default::default()
})
.await
}
pub async fn cancel_all_with(&self, opts: CancelAllOpts) -> Result<CancelAllOrdersResult> {
let mut req = CancelAllOrdersRequest {
subaccount_id: scope::optional_subaccount(&self.ctx, opts.subaccount_id)?,
symbol: opts.symbol.unwrap_or_default(),
dry_run: opts.dry_run,
request_id: Self::coalesce_request_id(opts.request_id, "cancel-all")?,
..Default::default()
};
if let Some(side) = opts.side.as_deref() {
req.side = Self::parse_side(side)?.into();
}
let client = self.write_client();
let resp = unary::await_auth(
&self.ctx.factory,
"/orders.v1.OrdersService/CancelAllOrders",
req,
|req, opts| client.cancel_all_orders_with_options(req, opts),
)
.await?
.into_owned();
cancel_all_from_proto(&resp)
}
fn parse_side(side: &str) -> Result<Side> {
match side.to_ascii_lowercase().as_str() {
"buy" => Ok(Side::Buy),
"sell" => Ok(Side::Sell),
_ => Err(Error::validation("side must be buy or sell")),
}
}
pub async fn modify(&self, params: ModifyOrderParams) -> Result<ModifyOrderResult> {
self.ctx.wait_for_catalogs().await?;
let req = self.encode_modify_params(params)?;
let client = self.write_client();
let resp = unary::await_auth(
&self.ctx.factory,
"/orders.v1.OrdersService/ModifyOrder",
req,
|req, opts| client.modify_order_with_options(req, opts),
)
.await?
.into_owned();
modify_order_from_proto(&resp)
}
pub fn create_params(
symbol: impl Into<String>,
side: CreateSide,
order_type: CreateOrderType,
quantity: Quantity,
price: Option<Price>,
client_order_id: Option<&str>,
) -> CreateOrderParams {
let client_order_id = client_order_id
.map(str::trim)
.filter(|s| !s.is_empty())
.map(|s| s.to_owned());
CreateOrderParams {
symbol: symbol.into(),
side,
order_type,
quantity: Some(quantity),
max_quote_debit_scaled: None,
price,
time_in_force: None,
client_order_id,
subaccount_id: None,
post_only: None,
market_client_ref_price: None,
fee_asset: None,
self_trade_prevention: None,
market_max_slippage: None,
attached_risk: None,
}
}
pub(crate) fn resolve_batch_replace_scale(
catalogs: &crate::catalogs::Manager,
symbol: &str,
) -> Result<u32> {
catalogs.base_quantity_scale_for_symbol(symbol).ok_or_else(|| {
Error::validation(format!(
"quantity scale for {symbol:?} is unavailable; await client.wait_for_catalogs() before placing orders"
))
})
}
pub async fn subscribe(
&self,
account_id: Option<&str>,
) -> Result<crate::realtime::TypedSubscription<Order>> {
let account = scope::resolve_account_id(&self.ctx, account_id)?;
let channel = format!("private:spot:orders:{account}:proto");
self.ctx
.realtime
.subscribe_proto(&channel, crate::codecs::decode::order_from_bytes)
.await
}
}
#[derive(Clone)]
pub struct TradesService {
ctx: ServiceContext,
}
impl TradesService {
pub fn new(ctx: ServiceContext) -> Self {
Self { ctx }
}
pub async fn list(
&self,
subaccount_id: Option<u64>,
limit: Option<u32>,
) -> Result<UserTradesList> {
let req = GetUserTradesRequest {
subaccount_id: scope::optional_subaccount(&self.ctx, subaccount_id)?,
limit,
..Default::default()
};
let client = OrdersReadServiceClient::new(
self.ctx.factory.transport(),
self.ctx.factory.connect_config(),
);
let resp = unary::await_auth(
&self.ctx.factory,
"/orders.v1.OrdersReadService/GetUserTrades",
req,
|req, opts| client.get_user_trades_with_options(req, opts),
)
.await?
.into_owned();
Ok(user_trades_list_from_proto(&resp))
}
pub async fn subscribe(
&self,
account_id: Option<&str>,
) -> Result<crate::realtime::TypedSubscription<UserTrade>> {
let account = scope::resolve_account_id(&self.ctx, account_id)?;
let channel = format!("private:spot:trades:{account}:proto");
self.ctx
.realtime
.subscribe_proto(&channel, crate::codecs::decode::user_trade_from_bytes)
.await
}
}
fn order_trades_projection_complete(result: &GetOrderResult) -> bool {
let Some(order) = result.order.as_ref() else {
return false;
};
if !matches!(order.status.as_str(), "filled" | "canceled" | "rejected") {
return false;
}
let Some(cum) = order.cum_qty.as_ref() else {
return false;
};
let cum = cum.as_scaled();
if cum == 0 {
return true;
}
let mut trade_sum = 0_i64;
for trade in &result.trades {
let Some(qty) = trade.qty.as_ref() else {
return false;
};
let Some(sum) = trade_sum.checked_add(qty.as_scaled()) else {
return false;
};
trade_sum = sum;
}
trade_sum == cum
}
#[cfg(test)]
mod tests {
use super::*;
use crate::codecs::scalars::format_id;
use buffa::Message;
use serde_json::json;
fn client() -> crate::Client {
let client = crate::Client::new(crate::Config {
hydrate_catalogs: false,
..Default::default()
})
.unwrap();
client
.catalogs
.hydrate_spot_config_json(json!({
"pairs": [{
"symbol": "BTC-USDT",
"symbol_id": 7,
"base_quantity_scale": 8,
"quote_quantity_scale": 6
}]
}))
.expect("hydrate");
client
}
fn create_params(quantity: Quantity, price: Price) -> CreateOrderParams {
CreateOrderParams {
symbol: "BTC-USDT".into(),
side: CreateSide::Buy,
order_type: CreateOrderType::Limit,
quantity: Some(quantity),
max_quote_debit_scaled: None,
price: Some(price),
time_in_force: Some(CreateTimeInForce::Gtc),
client_order_id: Some("order-equivalence".into()),
subaccount_id: None,
post_only: Some(true),
market_client_ref_price: None,
fee_asset: None,
self_trade_prevention: None,
market_max_slippage: None,
attached_risk: None,
}
}
#[test]
fn decimal_and_scaled_create_encode_identically() {
let client = client();
let decimal = create_params(
Quantity::from_decimal_str("0.1", 8, Some("BTC-USDT".into()), Some(7)).unwrap(),
Price::from_decimal_str("50000", Some("BTC-USDT".into())).unwrap(),
);
let scaled = create_params(
Quantity::from_scaled(
10_000_000,
Some(8),
crate::QuantityDomain::OrderBase,
Some("BTC-USDT".into()),
Some(7),
)
.unwrap(),
Price::from_ticks(50_000_000_000, Some("BTC-USDT".into())).unwrap(),
);
let decimal_wire = client.orders.encode_create_params(&decimal).unwrap();
let scaled_wire = client.orders.encode_create_params(&scaled).unwrap();
assert_eq!(decimal_wire.encode_to_vec(), scaled_wire.encode_to_vec());
}
fn modify_params(new_price: Option<Price>, new_qty: Option<Quantity>) -> ModifyOrderParams {
ModifyOrderParams {
symbol: "BTC-USDT".into(),
key: OrderKey::OrderId("1".into()),
subaccount_id: None,
request_id: Some("modify-equivalence".into()),
new_price,
new_qty,
new_attached_risk: None,
behavior: Some("amend_or_replace".into()),
new_client_order_id: None,
}
}
#[test]
fn decimal_and_scaled_modify_encode_identically() {
let client = client();
let decimal = modify_params(
Some(Price::from_decimal_str("50001", Some("BTC-USDT".into())).unwrap()),
Some(Quantity::from_decimal_str("0.2", 8, Some("BTC-USDT".into()), Some(7)).unwrap()),
);
let scaled = modify_params(
Some(Price::from_ticks(50_001_000_000, Some("BTC-USDT".into())).unwrap()),
Some(
Quantity::from_scaled(
20_000_000,
Some(8),
crate::QuantityDomain::OrderBase,
Some("BTC-USDT".into()),
Some(7),
)
.unwrap(),
),
);
let decimal_wire = client.orders.encode_modify_params(decimal).unwrap();
let scaled_wire = client.orders.encode_modify_params(scaled).unwrap();
assert_eq!(decimal_wire.encode_to_vec(), scaled_wire.encode_to_vec());
}
#[test]
fn batch_replace_requires_catalog_quantity_scale() {
let catalogs = crate::catalogs::Manager::new();
let err = OrdersService::resolve_batch_replace_scale(&catalogs, "BTC-USDT").unwrap_err();
assert!(
err.to_string().contains("quantity scale"),
"unexpected error: {err}"
);
}
#[test]
fn batch_replace_uses_symbol_catalog_quantity_scale() {
let client = client();
assert_eq!(
OrdersService::resolve_batch_replace_scale(&client.catalogs, "BTC-USDT").unwrap(),
8
);
}
#[test]
fn modify_validates_key_and_patch() {
let client = client();
let empty_key = ModifyOrderParams {
key: OrderKey::ClientOrderId(String::new()),
..modify_params(Some(Price::from_ticks(1, None).unwrap()), None)
};
assert!(client.orders.encode_modify_params(empty_key).is_err());
let no_patch = modify_params(None, None);
assert!(client.orders.encode_modify_params(no_patch).is_err());
}
#[test]
#[allow(deprecated)]
fn attached_risk_encodes_on_create_and_modify() {
use crate::models::{AttachedRisk, RiskLeg, TriggerPriceSourceKind};
let client = client();
let risk = AttachedRisk {
take_profit: Some(RiskLeg {
trigger_price: Price::from_ticks(51_000_000_000, Some("BTC-USDT".into())).unwrap(),
trigger_price_source: None,
order_type: Some(CreateOrderType::Market),
limit_price: None,
}),
stop_loss: Some(RiskLeg {
trigger_price: Price::from_ticks(49_000_000_000, Some("BTC-USDT".into())).unwrap(),
trigger_price_source: None,
order_type: Some(CreateOrderType::Limit),
limit_price: Some(
Price::from_ticks(48_900_000_000, Some("BTC-USDT".into())).unwrap(),
),
}),
trailing_stop: None,
oco: true,
};
let mut create = create_params(
Quantity::from_scaled(
10_000_000,
Some(8),
crate::QuantityDomain::OrderBase,
Some("BTC-USDT".into()),
Some(7),
)
.unwrap(),
Price::from_ticks(50_000_000_000, Some("BTC-USDT".into())).unwrap(),
);
create.attached_risk = Some(risk.clone());
let create_wire = client.orders.encode_create_params(&create).unwrap();
let order = create_wire.order.as_option().unwrap();
assert!(order.attached_risk.is_set());
assert!(order.attached_risk.as_option().unwrap().oco);
let mut modify = modify_params(None, None);
modify.new_attached_risk = Some(risk);
let modify_wire = client.orders.encode_modify_params(modify).unwrap();
assert!(modify_wire.new_attached_risk.is_set());
let mut unsupported = create;
unsupported
.attached_risk
.as_mut()
.unwrap()
.take_profit
.as_mut()
.unwrap()
.trigger_price_source = Some(TriggerPriceSourceKind::IndexPrice);
let err = client
.orders
.encode_create_params(&unsupported)
.unwrap_err();
assert!(matches!(&err, Error::Validation(_)));
assert!(err.to_string().contains("always uses last trade"));
}
#[test]
#[allow(deprecated)]
fn attached_trailing_stop_validates_positive_fields_and_rejects_silent_compat() {
use crate::models::{
AttachedRisk, MaxSlippage, TrailingDistance, TrailingStop, TriggerPriceSourceKind,
};
let client = client();
let base = create_params(
Quantity::from_scaled(
10_000_000,
Some(8),
crate::QuantityDomain::OrderBase,
Some("BTC-USDT".into()),
Some(7),
)
.unwrap(),
Price::from_ticks(50_000_000_000, Some("BTC-USDT".into())).unwrap(),
);
let mut zero_distance = base.clone();
zero_distance.attached_risk = Some(AttachedRisk {
trailing_stop: Some(TrailingStop {
distance: TrailingDistance::Ticks(0),
activation_price: None,
trigger_price_source: None,
order_type: None,
max_slippage: None,
}),
..Default::default()
});
let err = client
.orders
.encode_create_params(&zero_distance)
.unwrap_err();
assert!(err.to_string().contains("trailing_distance_ticks must be positive"));
let mut zero_slip = base.clone();
zero_slip.attached_risk = Some(AttachedRisk {
trailing_stop: Some(TrailingStop {
distance: TrailingDistance::Bps(25),
activation_price: None,
trigger_price_source: None,
order_type: None,
max_slippage: Some(MaxSlippage::Ticks(0)),
}),
..Default::default()
});
let err = client.orders.encode_create_params(&zero_slip).unwrap_err();
assert!(err.to_string().contains("max_slippage_ticks must be positive"));
let mut with_source = base.clone();
with_source.attached_risk = Some(AttachedRisk {
trailing_stop: Some(TrailingStop {
distance: TrailingDistance::Bps(25),
activation_price: None,
trigger_price_source: Some(TriggerPriceSourceKind::IndexPrice),
order_type: None,
max_slippage: None,
}),
..Default::default()
});
let err = client
.orders
.encode_create_params(&with_source)
.unwrap_err();
assert!(err.to_string().contains("always uses last trade"));
let mut with_order_type = base;
with_order_type.attached_risk = Some(AttachedRisk {
trailing_stop: Some(TrailingStop {
distance: TrailingDistance::Bps(25),
activation_price: None,
trigger_price_source: None,
order_type: Some(CreateOrderType::Limit),
max_slippage: None,
}),
..Default::default()
});
let err = client
.orders
.encode_create_params(&with_order_type)
.unwrap_err();
assert!(err.to_string().contains("always market"));
}
#[test]
fn preview_encodes_full_order_intent() {
let client = client();
let create = create_params(
Quantity::from_scaled(
10_000_000,
Some(8),
crate::QuantityDomain::OrderBase,
Some("BTC-USDT".into()),
Some(7),
)
.unwrap(),
Price::from_ticks(50_000_000_000, Some("BTC-USDT".into())).unwrap(),
);
let preview = PreviewOrderParams {
symbol: create.symbol.clone(),
side: create.side,
order_type: create.order_type,
quantity: create.quantity.clone(),
max_quote_debit_scaled: None,
price: create.price.clone(),
time_in_force: create.time_in_force,
client_order_id: Some("preview-cid".into()),
subaccount_id: Some(9),
post_only: create.post_only,
market_client_ref_price: None,
fee_asset: Some(FeeAsset::Quote),
self_trade_prevention: Some(OrderSelfTradePrevention::ExpireTaker),
market_max_slippage: None,
attached_risk: None,
};
let wire = client.orders.encode_preview_params(&preview).unwrap();
assert_eq!(wire.subaccount_id, Some(9));
let intent = wire.order.as_option().expect("preview order intent");
assert_eq!(intent.symbol, "BTC-USDT");
assert_eq!(intent.side.as_known(), Some(Side::Buy));
assert_eq!(intent.client_order_id, "preview-cid");
assert!(matches!(
intent.sizing,
Some(order_intent::Sizing::BaseQtyScaled(10_000_000))
));
assert!(matches!(
intent.execution,
Some(order_intent::Execution::LimitGtc(_))
));
assert_eq!(
intent.self_trade_prevention_mode.as_known(),
Some(SelfTradePreventionMode::ExpireTaker)
);
}
#[test]
fn create_allows_omitted_client_order_id_and_encodes_market_maker_controls() {
let client = client();
let mut params = create_params(
Quantity::from_scaled(
10_000_000,
Some(8),
crate::QuantityDomain::OrderBase,
Some("BTC-USDT".into()),
Some(7),
)
.unwrap(),
Price::from_ticks(50_000_000_000, Some("BTC-USDT".into())).unwrap(),
);
params.client_order_id = None;
let omitted = client.orders.encode_create_params(¶ms).unwrap();
assert!(
omitted
.order
.as_option()
.unwrap()
.client_order_id
.is_empty()
);
params.client_order_id = Some(" ".into());
let whitespace = client.orders.encode_create_params(¶ms).unwrap();
assert!(
whitespace
.order
.as_option()
.unwrap()
.client_order_id
.is_empty()
);
params.client_order_id = Some("mm-create-1".into());
params.order_type = CreateOrderType::Market;
params.price = None;
params.post_only = None;
params.fee_asset = Some(FeeAsset::Base);
params.self_trade_prevention = Some(OrderSelfTradePrevention::ExpireBoth);
params.market_max_slippage = Some(MaxSlippage::Bps(25));
let wire = client.orders.encode_create_params(¶ms).unwrap();
let intent = wire.order.as_option().unwrap();
assert_eq!(intent.fee_asset.as_known(), Some(ProtoFeeAsset::Base));
assert_eq!(
intent.self_trade_prevention_mode.as_known(),
Some(SelfTradePreventionMode::ExpireBoth)
);
let Some(order_intent::Execution::MarketIoc(market)) = intent.execution.as_ref() else {
panic!("expected market execution");
};
assert!(matches!(
market.max_slippage,
Some(market_ioc::MaxSlippage::MaxSlippageBps(25))
));
params.price = Some(Price::from_ticks(1, None).unwrap());
let err = client.orders.encode_create_params(¶ms).unwrap_err();
assert!(
err.to_string().contains("price is not valid for market"),
"unexpected error: {err}"
);
params.price = None;
params.market_max_slippage = Some(MaxSlippage::Ticks(0));
assert!(client.orders.encode_create_params(¶ms).is_err());
}
#[test]
fn create_encodes_quote_budget_sizing_and_rejects_ambiguous_sizing() {
let client = client();
let mut params = create_params(
Quantity::from_scaled(
10_000_000,
Some(8),
crate::QuantityDomain::OrderBase,
Some("BTC-USDT".into()),
Some(7),
)
.unwrap(),
Price::from_ticks(50_000_000_000, Some("BTC-USDT".into())).unwrap(),
);
params.quantity = None;
params.max_quote_debit_scaled = Some(
Quantity::from_quote_scaled(5_000_000, 6, Some("BTC-USDT".into()), Some(7)).unwrap(),
);
let wire = client.orders.encode_create_params(¶ms).unwrap();
let intent = wire.order.as_option().unwrap();
assert!(matches!(
intent.sizing,
Some(order_intent::Sizing::MaxQuoteDebitScaled(5_000_000))
));
params.quantity = Some(
Quantity::from_scaled(
10_000_000,
Some(8),
crate::QuantityDomain::OrderBase,
Some("BTC-USDT".into()),
Some(7),
)
.unwrap(),
);
assert!(client.orders.encode_create_params(¶ms).is_err());
params.quantity = None;
params.max_quote_debit_scaled =
Some(Quantity::from_quote_scaled(5_000_000, 8, None, None).unwrap());
let err = client.orders.encode_create_params(¶ms).unwrap_err();
assert!(err.to_string().contains("scale mismatch"));
}
#[test]
fn batch_size_guard_rejects_empty_and_more_than_twenty() {
assert!(OrdersService::validate_batch_size("batch_create", 1).is_ok());
assert!(OrdersService::validate_batch_size("batch_create", 20).is_ok());
assert!(
OrdersService::validate_batch_size("batch_create", 0)
.unwrap_err()
.to_string()
.contains("at least one")
);
assert!(
OrdersService::validate_batch_size("batch_create", 21)
.unwrap_err()
.to_string()
.contains("at most 20")
);
}
#[test]
fn create_rejects_invalid_client_order_id_before_wire() {
let client = client();
let mut params = create_params(
Quantity::from_scaled(
10_000_000,
Some(8),
crate::QuantityDomain::OrderBase,
Some("BTC-USDT".into()),
Some(7),
)
.unwrap(),
Price::from_ticks(50_000_000_000, Some("BTC-USDT".into())).unwrap(),
);
params.client_order_id = Some("bad id".into());
let err = client.orders.encode_create_params(¶ms).unwrap_err();
assert!(err.to_string().contains("invalid characters"));
params.client_order_id = Some("a".repeat(37));
let err = client.orders.encode_create_params(¶ms).unwrap_err();
assert!(err.to_string().contains("1 to 36"));
params.client_order_id = Some("ok-id_1.2:3/4".into());
assert!(client.orders.encode_create_params(¶ms).is_ok());
let err = OrdersService::coalesce_request_id(Some("bad id".into()), "mod").unwrap_err();
assert!(err.to_string().contains("invalid characters"));
let err = OrdersService::coalesce_request_id(Some("r".repeat(65)), "mod").unwrap_err();
assert!(err.to_string().contains("1 to 64"));
}
#[tokio::test]
async fn singular_order_methods_reject_invalid_client_order_id_before_transport() {
let client = client();
let err = client
.orders
.cancel_by_client_order_id("bad id!", None, None)
.await
.unwrap_err();
assert!(matches!(err, Error::Validation(_)));
assert!(err.to_string().contains("invalid characters"));
let err = client
.orders
.get(OrderKey::ClientOrderId("bad id!".into()), None)
.await
.unwrap_err();
assert!(matches!(err, Error::Validation(_)));
assert!(err.to_string().contains("invalid characters"));
}
#[test]
fn cancel_symbol_routing_distinguishes_omitted_and_invalid_inputs() {
let client = client();
assert_eq!(
OrdersService::resolve_cancel_symbol_id(&client.catalogs, None, None).unwrap(),
0
);
assert_eq!(
OrdersService::resolve_cancel_symbol_id(&client.catalogs, Some("BTC-USDT"), None)
.unwrap(),
7
);
for err in [
OrdersService::resolve_cancel_symbol_id(&client.catalogs, Some("UNKNOWN-USDT"), None)
.unwrap_err(),
OrdersService::resolve_cancel_symbol_id(&client.catalogs, None, Some(0)).unwrap_err(),
OrdersService::resolve_cancel_symbol_id(&client.catalogs, Some("BTC-USDT"), Some(7))
.unwrap_err(),
] {
assert!(matches!(err, Error::Validation(_)));
}
}
#[tokio::test]
async fn cancel_rejects_unknown_supplied_symbol_before_transport() {
let client = client();
let err = client
.orders
.cancel_with(CancelOrderParams {
key: OrderKey::OrderId(format_id(9)),
symbol: Some("UNKNOWN-USDT".into()),
symbol_id: None,
subaccount_id: None,
})
.await
.unwrap_err();
assert!(matches!(&err, Error::Validation(_)));
assert!(err.to_string().contains("unknown symbol"));
}
#[test]
fn mutation_request_ids_are_generated_when_omitted_like_go_python_typescript() {
for prefix in [
"cancel-all",
"cancel-after",
"mod",
"batch-create",
"batch-cancel",
"batch-replace",
] {
let generated = OrdersService::coalesce_request_id(None, prefix).unwrap();
assert!(
generated.starts_with(&format!("{prefix}-")),
"unexpected generated id for {prefix}: {generated}"
);
assert_eq!(generated.len(), prefix.len() + 1 + 12);
let blank = OrdersService::coalesce_request_id(Some(" ".into()), prefix).unwrap();
assert!(blank.starts_with(&format!("{prefix}-")));
assert_ne!(generated, blank);
}
assert_eq!(
OrdersService::coalesce_request_id(Some(" retry-mod-1 ".into()), "mod").unwrap(),
"retry-mod-1"
);
assert_eq!(
OrdersService::coalesce_request_id(Some("same-retry".into()), "batch-create").unwrap(),
"same-retry"
);
}
#[test]
fn wait_helper_detects_trade_projection_complete() {
let incomplete = GetOrderResult {
order: Some(Order {
order_id: "1".into(),
symbol_id: 7,
client_order_id: "c".into(),
side: "buy".into(),
status: "filled".into(),
order_type: "market".into(),
tif: "ioc".into(),
orig_qty: None,
cum_qty: Some(
Quantity::from_scaled(
100,
Some(8),
crate::QuantityDomain::OrderBase,
None,
None,
)
.unwrap(),
),
leaves_qty: None,
price: None,
avg_px: None,
created_ts_ns: String::new(),
version: 1,
post_only: false,
fee_asset: "quote".into(),
submitted_max_quote_debit_scaled: None,
attached_risk: None,
}),
trades: vec![],
};
assert!(!order_trades_projection_complete(&incomplete));
let open_unfilled = GetOrderResult {
order: Some(Order {
status: "working".into(),
cum_qty: Some(
Quantity::from_scaled(0, Some(8), crate::QuantityDomain::OrderBase, None, None)
.unwrap(),
),
..incomplete.order.clone().unwrap()
}),
trades: vec![],
};
assert!(
!order_trades_projection_complete(&open_unfilled),
"an unfilled working order is not a stable projection"
);
let complete = GetOrderResult {
order: Some(Order {
status: "filled".into(),
..incomplete.order.clone().unwrap()
}),
trades: vec![UserTrade {
symbol_id: 7,
match_id: "m".into(),
order_id: "1".into(),
side: "buy".into(),
is_maker: false,
price: None,
qty: Some(
Quantity::from_scaled(
100,
Some(8),
crate::QuantityDomain::OrderBase,
None,
None,
)
.unwrap(),
),
fee_scaled: "0".into(),
fee_asset: "quote".into(),
referral_share_scaled: "0".into(),
ts_ns: String::new(),
}],
};
assert!(order_trades_projection_complete(&complete));
}
}