use crate::error::AppError;
use crate::presentation::account::AccountData;
use crate::presentation::chart::ChartData;
use crate::presentation::market::{MarketFields, PresentationMarketData};
use crate::presentation::price::PriceData;
use crate::presentation::trade::TradeData;
use lightstreamer_rs::subscription::ItemUpdate;
use std::collections::HashMap;
#[must_use]
fn changed_fields_as_options(changed: &HashMap<String, String>) -> HashMap<String, Option<String>> {
changed
.iter()
.map(|(key, value)| (key.clone(), Some(value.clone())))
.collect()
}
#[must_use = "the parse result must be handled"]
pub fn price_data_from_item_update(item_update: &ItemUpdate) -> Result<PriceData, AppError> {
let changed_fields = changed_fields_as_options(&item_update.changed_fields);
PriceData::from_fields(
item_update.item_name.as_deref(),
item_update.item_pos,
item_update.is_snapshot,
&item_update.fields,
&changed_fields,
)
.map_err(AppError::Deserialization)
}
impl From<&ItemUpdate> for PriceData {
fn from(item_update: &ItemUpdate) -> Self {
price_data_from_item_update(item_update).unwrap_or_else(|e| {
tracing::warn!(error = %e, "failed to convert ItemUpdate to PriceData, returning default");
PriceData::default()
})
}
}
#[must_use = "the parse result must be handled"]
pub fn market_data_from_item_update(
item_update: &ItemUpdate,
) -> Result<PresentationMarketData, AppError> {
let changed_fields = changed_fields_as_options(&item_update.changed_fields);
PresentationMarketData::from_fields(
item_update.item_name.as_deref(),
item_update.item_pos,
item_update.is_snapshot,
&item_update.fields,
&changed_fields,
)
.map_err(AppError::Deserialization)
}
impl From<&ItemUpdate> for PresentationMarketData {
fn from(item_update: &ItemUpdate) -> Self {
market_data_from_item_update(item_update).unwrap_or_else(|_| PresentationMarketData {
item_name: String::new(),
item_pos: 0,
fields: MarketFields::default(),
changed_fields: MarketFields::default(),
is_snapshot: false,
})
}
}
#[must_use = "the parse result must be handled"]
pub fn chart_data_from_item_update(item_update: &ItemUpdate) -> Result<ChartData, AppError> {
let changed_fields = changed_fields_as_options(&item_update.changed_fields);
ChartData::from_fields(
item_update.item_name.as_deref(),
item_update.item_pos,
item_update.is_snapshot,
&item_update.fields,
&changed_fields,
)
.map_err(AppError::Deserialization)
}
impl From<&ItemUpdate> for ChartData {
fn from(item_update: &ItemUpdate) -> Self {
chart_data_from_item_update(item_update).unwrap_or_default()
}
}
#[must_use = "the parse result must be handled"]
pub fn trade_data_from_item_update(item_update: &ItemUpdate) -> Result<TradeData, AppError> {
let changed_fields = changed_fields_as_options(&item_update.changed_fields);
TradeData::from_fields(
item_update.item_name.as_deref(),
item_update.item_pos,
item_update.is_snapshot,
&item_update.fields,
&changed_fields,
)
.map_err(AppError::Deserialization)
}
impl From<&ItemUpdate> for TradeData {
fn from(item_update: &ItemUpdate) -> Self {
trade_data_from_item_update(item_update).unwrap_or_default()
}
}
#[must_use = "the parse result must be handled"]
pub fn account_data_from_item_update(item_update: &ItemUpdate) -> Result<AccountData, AppError> {
let changed_fields = changed_fields_as_options(&item_update.changed_fields);
AccountData::from_fields(
item_update.item_name.as_deref(),
item_update.item_pos,
item_update.is_snapshot,
&item_update.fields,
&changed_fields,
)
.map_err(AppError::Deserialization)
}
impl From<&ItemUpdate> for AccountData {
fn from(item_update: &ItemUpdate) -> Self {
account_data_from_item_update(item_update).unwrap_or_else(|_| AccountData::default())
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::presentation::order::{Direction, OrderType, Status, TimeInForce};
fn trade_item_update(field_key: &str, field_value: &str) -> ItemUpdate {
let mut fields = HashMap::new();
fields.insert(field_key.to_string(), Some(field_value.to_string()));
ItemUpdate {
item_name: Some("TRADE:SYNTHETIC".to_string()),
item_pos: 2,
fields,
changed_fields: HashMap::new(),
is_snapshot: true,
}
}
const SANITIZED_OPU_JSON: &str = r#"{
"dealReference": "REF_SYNTHETIC_0001",
"dealId": "DIAAAAF00SYNTH01",
"dealIdOrigin": "DIAAAAF00SYNTH01",
"direction": "BUY",
"epic": "IX.D.DAX.DAILY.IP",
"status": "OPEN",
"dealStatus": "ACCEPTED",
"level": "18000.5",
"size": "1.0",
"currency": "EUR",
"timestamp": "2024-01-15T10:30:00.000",
"channel": "PublicRestOTC",
"expiry": "DFB"
}"#;
const SANITIZED_WOU_JSON: &str = r#"{
"dealReference": "REF_SYNTHETIC_0002",
"dealId": "DIAAAAF00SYNTH02",
"direction": "SELL",
"epic": "CS.D.EURUSD.CFD.IP",
"status": "OPEN",
"dealStatus": "ACCEPTED",
"level": "1.1000",
"size": "10000",
"currency": "USD",
"timestamp": "2024-01-15T10:31:00.000",
"channel": "PublicRestOTC",
"expiry": "-",
"stopDistance": "10.5",
"limitDistance": "20.0",
"guaranteedStop": false,
"orderType": "LIMIT",
"timeInForce": "GOOD_TILL_CANCELLED",
"goodTillDate": ""
}"#;
#[test]
fn test_trade_data_from_item_update_parses_open_position_update() {
let item = trade_item_update("OPU", SANITIZED_OPU_JSON);
let result = trade_data_from_item_update(&item);
assert!(result.is_ok(), "OPU payload should parse: {result:?}");
let data = result.expect("OPU parse checked ok above");
assert_eq!(data.item_name, "TRADE:SYNTHETIC");
assert_eq!(data.item_pos, 2);
assert!(data.is_snapshot);
let opu = data.fields.opu.expect("OPU should be populated");
assert_eq!(opu.deal_reference.as_deref(), Some("REF_SYNTHETIC_0001"));
assert_eq!(opu.deal_id.as_deref(), Some("DIAAAAF00SYNTH01"));
assert_eq!(opu.deal_id_origin.as_deref(), Some("DIAAAAF00SYNTH01"));
assert_eq!(opu.direction, Some(Direction::Buy));
assert_eq!(opu.epic.as_deref(), Some("IX.D.DAX.DAILY.IP"));
assert_eq!(opu.status, Some(Status::Open));
assert_eq!(opu.deal_status, Some(Status::Accepted));
assert_eq!(opu.level, Some(18000.5));
assert_eq!(opu.size, Some(1.0));
assert_eq!(opu.currency.as_deref(), Some("EUR"));
assert_eq!(opu.expiry.as_deref(), Some("DFB"));
assert!(data.fields.wou.is_none());
assert!(data.fields.confirms.is_none());
}
#[test]
fn test_trade_data_from_item_update_parses_working_order_update() {
let item = trade_item_update("WOU", SANITIZED_WOU_JSON);
let result = trade_data_from_item_update(&item);
assert!(result.is_ok(), "WOU payload should parse: {result:?}");
let data = result.expect("WOU parse checked ok above");
let wou = data.fields.wou.expect("WOU should be populated");
assert_eq!(wou.deal_reference.as_deref(), Some("REF_SYNTHETIC_0002"));
assert_eq!(wou.deal_id.as_deref(), Some("DIAAAAF00SYNTH02"));
assert_eq!(wou.direction, Some(Direction::Sell));
assert_eq!(wou.epic.as_deref(), Some("CS.D.EURUSD.CFD.IP"));
assert_eq!(wou.status, Some(Status::Open));
assert_eq!(wou.deal_status, Some(Status::Accepted));
assert_eq!(wou.level, Some(1.1000));
assert_eq!(wou.size, Some(10000.0));
assert_eq!(wou.currency.as_deref(), Some("USD"));
assert_eq!(wou.stop_distance, Some(10.5));
assert_eq!(wou.limit_distance, Some(20.0));
assert_eq!(wou.guaranteed_stop, Some(false));
assert_eq!(wou.order_type, Some(OrderType::Limit));
assert_eq!(wou.time_in_force, Some(TimeInForce::GoodTillCancelled));
assert!(wou.good_till_date.is_none());
assert!(data.fields.opu.is_none());
assert!(data.fields.confirms.is_none());
}
#[test]
fn test_trade_data_from_item_update_malformed_opu_is_error_on_direct_path() {
let item = trade_item_update("OPU", "{ this is not valid json ");
let result = trade_data_from_item_update(&item);
assert!(
result.is_err(),
"malformed OPU JSON must surface an error on the direct parse path"
);
}
#[test]
fn test_trade_data_from_impl_swallows_malformed_opu_into_default() {
let item = trade_item_update("OPU", "{ this is not valid json ");
let data = TradeData::from(&item);
assert!(
data.item_name.is_empty(),
"silent-default path discards item_name"
);
assert_eq!(data.item_pos, 0, "silent-default path discards item_pos");
assert!(
!data.is_snapshot,
"silent-default path discards is_snapshot"
);
assert!(data.fields.opu.is_none());
assert!(data.fields.wou.is_none());
assert!(data.fields.confirms.is_none());
}
#[test]
fn test_price_data_from_item_update_defaults_on_unknown_dealing_flag() {
let mut fields = HashMap::new();
fields.insert("DLG_FLAG".to_string(), Some("NOT_A_FLAG".to_string()));
let item = ItemUpdate {
item_name: Some("PRICE:CS.D.EURUSD.MINI.IP".to_string()),
item_pos: 1,
fields,
changed_fields: HashMap::new(),
is_snapshot: true,
};
assert!(price_data_from_item_update(&item).is_err());
let data = PriceData::from(&item);
assert!(data.item_name.is_empty());
assert_eq!(data.item_pos, 0);
assert!(!data.is_snapshot);
assert!(data.fields.dealing_flag.is_none());
}
}