use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::models::ExtraFields;
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct RawStreamMessage {
#[serde(rename = "type")]
pub kind: String,
#[serde(default)]
pub id: Option<u64>,
#[serde(default)]
pub sid: Option<u64>,
#[serde(default)]
pub seq: Option<u64>,
#[serde(default)]
pub msg: Value,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Envelope<T> {
pub sid: Option<u64>,
pub seq: Option<u64>,
pub msg: T,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct CommandEnvelope<T> {
pub id: Option<u64>,
pub sid: Option<u64>,
pub seq: Option<u64>,
pub msg: T,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct SubscribedMessage {
pub channel: String,
pub sid: u64,
#[serde(flatten)]
pub extra: ExtraFields,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct WsErrorMessage {
pub code: i64,
pub msg: String,
#[serde(flatten)]
pub extra: ExtraFields,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct TickerMessage {
pub market_ticker: String,
#[serde(default)]
pub market_id: Option<String>,
#[serde(default)]
pub price_dollars: Option<String>,
#[serde(default)]
pub yes_bid_dollars: Option<String>,
#[serde(default)]
pub yes_ask_dollars: Option<String>,
#[serde(default)]
pub volume_fp: Option<String>,
#[serde(default)]
pub open_interest_fp: Option<String>,
#[serde(default)]
pub yes_bid_size_fp: Option<String>,
#[serde(default)]
pub yes_ask_size_fp: Option<String>,
#[serde(default)]
pub last_trade_size_fp: Option<String>,
#[serde(default)]
pub ts: Option<i64>,
#[serde(default)]
pub ts_ms: Option<i64>,
#[serde(default)]
pub time: Option<String>,
#[serde(flatten)]
pub extra: ExtraFields,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct TradeMessage {
pub trade_id: String,
pub market_ticker: String,
#[serde(default)]
pub yes_price_dollars: Option<String>,
#[serde(default)]
pub no_price_dollars: Option<String>,
#[serde(default)]
pub count_fp: Option<String>,
#[serde(default)]
pub taker_side: Option<String>,
#[serde(default)]
pub ts: Option<i64>,
#[serde(default)]
pub ts_ms: Option<i64>,
#[serde(flatten)]
pub extra: ExtraFields,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct FillMessage {
pub trade_id: String,
pub order_id: String,
pub market_ticker: String,
#[serde(default)]
pub is_taker: Option<bool>,
#[serde(default)]
pub side: Option<String>,
#[serde(default)]
pub yes_price_dollars: Option<String>,
#[serde(default)]
pub count_fp: Option<String>,
#[serde(default)]
pub action: Option<String>,
#[serde(default)]
pub post_position_fp: Option<String>,
#[serde(default)]
pub purchased_side: Option<String>,
#[serde(default)]
pub subaccount: Option<u64>,
#[serde(default)]
pub ts: Option<i64>,
#[serde(default)]
pub ts_ms: Option<i64>,
#[serde(flatten)]
pub extra: ExtraFields,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct UserOrderMessage {
pub order_id: String,
#[serde(default)]
pub user_id: Option<String>,
#[serde(alias = "market_ticker")]
pub ticker: String,
pub status: String,
#[serde(default)]
pub side: Option<String>,
#[serde(default)]
pub is_yes: Option<bool>,
#[serde(default)]
pub yes_price_dollars: Option<String>,
#[serde(default)]
pub fill_count_fp: Option<String>,
#[serde(default)]
pub remaining_count_fp: Option<String>,
#[serde(default)]
pub initial_count_fp: Option<String>,
#[serde(default)]
pub client_order_id: Option<String>,
#[serde(default)]
pub order_group_id: Option<String>,
#[serde(default)]
pub created_time: Option<String>,
#[serde(default)]
pub created_ts_ms: Option<i64>,
#[serde(default)]
pub expiration_time: Option<String>,
#[serde(default)]
pub expiration_ts_ms: Option<i64>,
#[serde(default)]
pub subaccount_number: Option<u64>,
#[serde(flatten)]
pub extra: ExtraFields,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct OrderbookSnapshot {
pub market_ticker: String,
#[serde(default)]
pub market_id: Option<String>,
#[serde(default)]
pub yes_dollars_fp: Vec<(String, String)>,
#[serde(default)]
pub no_dollars_fp: Vec<(String, String)>,
#[serde(flatten)]
pub extra: ExtraFields,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct OrderbookDelta {
pub market_ticker: String,
#[serde(default)]
pub market_id: Option<String>,
pub price_dollars: String,
pub delta_fp: String,
pub side: String,
#[serde(default)]
pub ts: Option<String>,
#[serde(default)]
pub ts_ms: Option<i64>,
#[serde(flatten)]
pub extra: ExtraFields,
}
#[derive(Debug, Clone, PartialEq)]
pub enum StreamMessage {
Subscribed(CommandEnvelope<SubscribedMessage>),
Unsubscribed(CommandEnvelope<Value>),
Ok(CommandEnvelope<Value>),
ListSubscriptions(CommandEnvelope<Value>),
Error(CommandEnvelope<WsErrorMessage>),
OrderbookSnapshot(Envelope<OrderbookSnapshot>),
OrderbookDelta(Envelope<OrderbookDelta>),
Ticker(Envelope<TickerMessage>),
Trade(Envelope<TradeMessage>),
Fill(Envelope<FillMessage>),
MarketPosition(Envelope<Value>),
MarketLifecycleV2(Envelope<Value>),
EventLifecycle(Envelope<Value>),
EventFeeUpdate(Envelope<Value>),
MultivariateMarketLifecycle(Envelope<Value>),
MultivariateLookup(Envelope<Value>),
Communications(Envelope<Value>),
OrderGroupUpdates(Envelope<Value>),
UserOrder(Envelope<UserOrderMessage>),
CfBenchmarksValue(Envelope<Value>),
CfBenchmarksIndexList(Envelope<Value>),
Raw(RawStreamMessage),
}
impl StreamMessage {
pub fn parse_text(text: &str) -> serde_json::Result<Self> {
let raw: RawStreamMessage = serde_json::from_str(text)?;
Self::try_from_raw(raw)
}
pub fn try_from_raw(raw: RawStreamMessage) -> serde_json::Result<Self> {
Ok(match raw.kind.as_str() {
"subscribed" => Self::Subscribed(raw.command_envelope()?),
"unsubscribed" => Self::Unsubscribed(raw.command_envelope()?),
"ok" => Self::Ok(raw.command_envelope()?),
"list_subscriptions" => Self::ListSubscriptions(raw.command_envelope()?),
"error" => Self::Error(raw.command_envelope()?),
"orderbook_snapshot" => Self::OrderbookSnapshot(raw.envelope()?),
"orderbook_delta" => Self::OrderbookDelta(raw.envelope()?),
"ticker" => Self::Ticker(raw.envelope()?),
"trade" => Self::Trade(raw.envelope()?),
"fill" => Self::Fill(raw.envelope()?),
"market_position" => Self::MarketPosition(raw.envelope()?),
"market_lifecycle_v2" => Self::MarketLifecycleV2(raw.envelope()?),
"event_lifecycle" => Self::EventLifecycle(raw.envelope()?),
"event_fee_update" => Self::EventFeeUpdate(raw.envelope()?),
"multivariate_market_lifecycle" => Self::MultivariateMarketLifecycle(raw.envelope()?),
"multivariate_lookup" => Self::MultivariateLookup(raw.envelope()?),
"rfq_created" | "rfq_deleted" | "quote_created" | "quote_accepted"
| "quote_executed" => Self::Communications(raw.envelope()?),
"order_group_updates" => Self::OrderGroupUpdates(raw.envelope()?),
"user_order" => Self::UserOrder(raw.envelope()?),
"cfbenchmarks_value" => Self::CfBenchmarksValue(raw.envelope()?),
"cfbenchmarks_value_indexlist" => Self::CfBenchmarksIndexList(raw.envelope()?),
_ => Self::Raw(raw),
})
}
}
impl RawStreamMessage {
fn envelope<T>(self) -> serde_json::Result<Envelope<T>>
where
T: for<'de> Deserialize<'de>,
{
Ok(Envelope {
sid: self.sid,
seq: self.seq,
msg: serde_json::from_value(self.msg)?,
})
}
fn command_envelope<T>(self) -> serde_json::Result<CommandEnvelope<T>>
where
T: for<'de> Deserialize<'de>,
{
Ok(CommandEnvelope {
id: self.id,
sid: self.sid,
seq: self.seq,
msg: serde_json::from_value(self.msg)?,
})
}
}