1#![allow(unused_imports)]
15use async_trait::async_trait;
16use derive_builder::Builder;
17use reqwest;
18use rust_decimal::prelude::*;
19use serde::{Deserialize, Serialize};
20use serde_json::{Value, json};
21use std::collections::BTreeMap;
22
23use crate::common::{
24 config::ConfigurationRestApi,
25 models::{ParamBuildError, RestApiResponse},
26 utils::send_request,
27};
28use crate::stocks::rest_api::models;
29
30const HAS_TIME_UNIT: bool = false;
31
32#[async_trait]
33pub trait MarketDataApi: Send + Sync {
34 async fn exchange_info(
35 &self,
36 params: ExchangeInfoParams,
37 ) -> anyhow::Result<RestApiResponse<models::ExchangeInfoResponse>>;
38 async fn latest_quote(
39 &self,
40 params: LatestQuoteParams,
41 ) -> anyhow::Result<RestApiResponse<models::LatestQuoteResponse>>;
42 async fn tokenized_assets(
43 &self,
44 ) -> anyhow::Result<RestApiResponse<Vec<models::TokenizedAssetsResponseInner>>>;
45}
46
47#[derive(Debug, Clone)]
48pub struct MarketDataApiClient {
49 configuration: ConfigurationRestApi,
50}
51
52impl MarketDataApiClient {
53 pub fn new(configuration: ConfigurationRestApi) -> Self {
54 Self { configuration }
55 }
56}
57
58#[derive(Clone, Debug, Builder, Deserialize, Default)]
63#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
64pub struct ExchangeInfoParams {
65 #[builder(setter(into), default)]
69 #[serde(rename = "symbol", default)]
70 pub symbol: Option<String>,
71}
72
73impl ExchangeInfoParams {
74 #[must_use]
77 pub fn builder() -> ExchangeInfoParamsBuilder {
78 ExchangeInfoParamsBuilder::default()
79 }
80}
81#[derive(Clone, Debug, Builder, Deserialize)]
86#[builder(pattern = "owned", build_fn(error = "ParamBuildError"))]
87pub struct LatestQuoteParams {
88 #[builder(setter(into))]
92 #[serde(rename = "symbol")]
93 pub symbol: String,
94}
95
96impl LatestQuoteParams {
97 #[must_use]
104 pub fn builder(symbol: String) -> LatestQuoteParamsBuilder {
105 LatestQuoteParamsBuilder::default().symbol(symbol)
106 }
107}
108
109#[async_trait]
110impl MarketDataApi for MarketDataApiClient {
111 async fn exchange_info(
112 &self,
113 params: ExchangeInfoParams,
114 ) -> anyhow::Result<RestApiResponse<models::ExchangeInfoResponse>> {
115 let ExchangeInfoParams { symbol } = params;
116
117 let mut query_params = BTreeMap::new();
118 let body_params = BTreeMap::new();
119
120 if let Some(rw) = symbol {
121 query_params.insert("symbol".to_string(), json!(rw));
122 }
123
124 send_request::<models::ExchangeInfoResponse>(
125 &self.configuration,
126 "/sapi/v1/equity/market/exchangeInfo",
127 reqwest::Method::GET,
128 query_params,
129 body_params,
130 if HAS_TIME_UNIT {
131 self.configuration.time_unit
132 } else {
133 None
134 },
135 false,
136 )
137 .await
138 }
139
140 async fn latest_quote(
141 &self,
142 params: LatestQuoteParams,
143 ) -> anyhow::Result<RestApiResponse<models::LatestQuoteResponse>> {
144 let LatestQuoteParams { symbol } = params;
145
146 let mut query_params = BTreeMap::new();
147 let body_params = BTreeMap::new();
148
149 query_params.insert("symbol".to_string(), json!(symbol));
150
151 send_request::<models::LatestQuoteResponse>(
152 &self.configuration,
153 "/sapi/v1/equity/market/quote",
154 reqwest::Method::GET,
155 query_params,
156 body_params,
157 if HAS_TIME_UNIT {
158 self.configuration.time_unit
159 } else {
160 None
161 },
162 false,
163 )
164 .await
165 }
166
167 async fn tokenized_assets(
168 &self,
169 ) -> anyhow::Result<RestApiResponse<Vec<models::TokenizedAssetsResponseInner>>> {
170 let query_params = BTreeMap::new();
171 let body_params = BTreeMap::new();
172
173 send_request::<Vec<models::TokenizedAssetsResponseInner>>(
174 &self.configuration,
175 "/sapi/v1/equity/market/tokenized-assets",
176 reqwest::Method::GET,
177 query_params,
178 body_params,
179 if HAS_TIME_UNIT {
180 self.configuration.time_unit
181 } else {
182 None
183 },
184 false,
185 )
186 .await
187 }
188}
189
190#[cfg(all(test, feature = "stocks"))]
191mod tests {
192 use super::*;
193 use crate::TOKIO_SHARED_RT;
194 use crate::{errors::ConnectorError, models::DataFuture, models::RestApiRateLimit};
195 use async_trait::async_trait;
196 use std::collections::HashMap;
197
198 struct DummyRestApiResponse<T> {
199 inner: Box<dyn FnOnce() -> DataFuture<Result<T, ConnectorError>> + Send + Sync>,
200 status: u16,
201 headers: HashMap<String, String>,
202 rate_limits: Option<Vec<RestApiRateLimit>>,
203 }
204
205 impl<T> From<DummyRestApiResponse<T>> for RestApiResponse<T> {
206 fn from(dummy: DummyRestApiResponse<T>) -> Self {
207 Self {
208 data_fn: dummy.inner,
209 status: dummy.status,
210 headers: dummy.headers,
211 rate_limits: dummy.rate_limits,
212 }
213 }
214 }
215
216 struct MockMarketDataApiClient {
217 force_error: bool,
218 }
219
220 #[async_trait]
221 impl MarketDataApi for MockMarketDataApiClient {
222 async fn exchange_info(
223 &self,
224 _params: ExchangeInfoParams,
225 ) -> anyhow::Result<RestApiResponse<models::ExchangeInfoResponse>> {
226 if self.force_error {
227 return Err(ConnectorError::ConnectorClientError {
228 msg: "ResponseError".to_string(),
229 code: None,
230 }
231 .into());
232 }
233
234 let resp_json: Value = serde_json::from_str(r#"{"timezone":"UTC","symbols":[{"symbol":"AAPL","tradability":"BUY_SELL","tradabilityUpdateTime":1735900000000,"overnightSupported":true,"fractionable":true,"fractionableEh":false,"extendedSession":true,"maxNumOrders":200,"stepSize":"0.000000001","multiplierUp":"1.10","multiplierDown":"0.90","minQty":"0.000000001","maxQty":"100000.00000000","minNotional":"1.00","maxNotional":"1000000.00","listingTime":1700000000000,"delistingTime":1}]}"#).unwrap_or_else(|_| serde_json::json!({}));
235 let dummy_response: models::ExchangeInfoResponse =
236 serde_json::from_value(resp_json.clone())
237 .expect("should parse into models::ExchangeInfoResponse");
238
239 let dummy = DummyRestApiResponse {
240 inner: Box::new(move || Box::pin(async move { Ok(dummy_response) })),
241 status: 200,
242 headers: HashMap::new(),
243 rate_limits: None,
244 };
245
246 Ok(dummy.into())
247 }
248
249 async fn latest_quote(
250 &self,
251 _params: LatestQuoteParams,
252 ) -> anyhow::Result<RestApiResponse<models::LatestQuoteResponse>> {
253 if self.force_error {
254 return Err(ConnectorError::ConnectorClientError {
255 msg: "ResponseError".to_string(),
256 code: None,
257 }
258 .into());
259 }
260
261 let resp_json: Value = serde_json::from_str(r#"{"symbol":"AAPL","bidPrice":"180.50","askPrice":"180.52","bidSize":100,"askSize":200}"#).unwrap_or_else(|_| serde_json::json!({}));
262 let dummy_response: models::LatestQuoteResponse =
263 serde_json::from_value(resp_json.clone())
264 .expect("should parse into models::LatestQuoteResponse");
265
266 let dummy = DummyRestApiResponse {
267 inner: Box::new(move || Box::pin(async move { Ok(dummy_response) })),
268 status: 200,
269 headers: HashMap::new(),
270 rate_limits: None,
271 };
272
273 Ok(dummy.into())
274 }
275
276 async fn tokenized_assets(
277 &self,
278 ) -> anyhow::Result<RestApiResponse<Vec<models::TokenizedAssetsResponseInner>>> {
279 if self.force_error {
280 return Err(ConnectorError::ConnectorClientError {
281 msg: "ResponseError".to_string(),
282 code: None,
283 }
284 .into());
285 }
286
287 let resp_json: Value = serde_json::from_str(r#"[{"assetCode":"AAPLB","assetName":"Apple Inc. Tokenized Stock","underlyingEquitySymbol":"AAPL","multiplier":"1","multiplierValid":true}]"#).unwrap_or_else(|_| serde_json::json!({}));
288 let dummy_response: Vec<models::TokenizedAssetsResponseInner> =
289 serde_json::from_value(resp_json.clone())
290 .expect("should parse into Vec<models::TokenizedAssetsResponseInner>");
291
292 let dummy = DummyRestApiResponse {
293 inner: Box::new(move || Box::pin(async move { Ok(dummy_response) })),
294 status: 200,
295 headers: HashMap::new(),
296 rate_limits: None,
297 };
298
299 Ok(dummy.into())
300 }
301 }
302
303 #[test]
304 fn exchange_info_required_params_success() {
305 TOKIO_SHARED_RT.block_on(async {
306 let client = MockMarketDataApiClient { force_error: false };
307
308 let params = ExchangeInfoParams::builder().build().unwrap();
309
310 let resp_json: Value = serde_json::from_str(r#"{"timezone":"UTC","symbols":[{"symbol":"AAPL","tradability":"BUY_SELL","tradabilityUpdateTime":1735900000000,"overnightSupported":true,"fractionable":true,"fractionableEh":false,"extendedSession":true,"maxNumOrders":200,"stepSize":"0.000000001","multiplierUp":"1.10","multiplierDown":"0.90","minQty":"0.000000001","maxQty":"100000.00000000","minNotional":"1.00","maxNotional":"1000000.00","listingTime":1700000000000,"delistingTime":1}]}"#).unwrap_or_else(|_| serde_json::json!({}));
311 let expected_response : models::ExchangeInfoResponse = serde_json::from_value(resp_json.clone()).expect("should parse into models::ExchangeInfoResponse");
312
313 let resp = client.exchange_info(params).await.expect("Expected a response");
314 let data_future = resp.data();
315 let actual_response = data_future.await.unwrap();
316 assert_eq!(actual_response, expected_response);
317 });
318 }
319
320 #[test]
321 fn exchange_info_optional_params_success() {
322 TOKIO_SHARED_RT.block_on(async {
323 let client = MockMarketDataApiClient { force_error: false };
324
325 let params = ExchangeInfoParams::builder().symbol("AAPL".to_string()).build().unwrap();
326
327 let resp_json: Value = serde_json::from_str(r#"{"timezone":"UTC","symbols":[{"symbol":"AAPL","tradability":"BUY_SELL","tradabilityUpdateTime":1735900000000,"overnightSupported":true,"fractionable":true,"fractionableEh":false,"extendedSession":true,"maxNumOrders":200,"stepSize":"0.000000001","multiplierUp":"1.10","multiplierDown":"0.90","minQty":"0.000000001","maxQty":"100000.00000000","minNotional":"1.00","maxNotional":"1000000.00","listingTime":1700000000000,"delistingTime":1}]}"#).unwrap_or_else(|_| serde_json::json!({}));
328 let expected_response : models::ExchangeInfoResponse = serde_json::from_value(resp_json.clone()).expect("should parse into models::ExchangeInfoResponse");
329
330 let resp = client.exchange_info(params).await.expect("Expected a response");
331 let data_future = resp.data();
332 let actual_response = data_future.await.unwrap();
333 assert_eq!(actual_response, expected_response);
334 });
335 }
336
337 #[test]
338 fn exchange_info_response_error() {
339 TOKIO_SHARED_RT.block_on(async {
340 let client = MockMarketDataApiClient { force_error: true };
341
342 let params = ExchangeInfoParams::builder().build().unwrap();
343
344 match client.exchange_info(params).await {
345 Ok(_) => panic!("Expected an error"),
346 Err(err) => {
347 assert_eq!(err.to_string(), "Connector client error: ResponseError");
348 }
349 }
350 });
351 }
352
353 #[test]
354 fn latest_quote_required_params_success() {
355 TOKIO_SHARED_RT.block_on(async {
356 let client = MockMarketDataApiClient { force_error: false };
357
358 let params = LatestQuoteParams::builder("AAPL".to_string()).build().unwrap();
359
360 let resp_json: Value = serde_json::from_str(r#"{"symbol":"AAPL","bidPrice":"180.50","askPrice":"180.52","bidSize":100,"askSize":200}"#).unwrap_or_else(|_| serde_json::json!({}));
361 let expected_response : models::LatestQuoteResponse = serde_json::from_value(resp_json.clone()).expect("should parse into models::LatestQuoteResponse");
362
363 let resp = client.latest_quote(params).await.expect("Expected a response");
364 let data_future = resp.data();
365 let actual_response = data_future.await.unwrap();
366 assert_eq!(actual_response, expected_response);
367 });
368 }
369
370 #[test]
371 fn latest_quote_optional_params_success() {
372 TOKIO_SHARED_RT.block_on(async {
373 let client = MockMarketDataApiClient { force_error: false };
374
375 let params = LatestQuoteParams::builder("AAPL".to_string()).build().unwrap();
376
377 let resp_json: Value = serde_json::from_str(r#"{"symbol":"AAPL","bidPrice":"180.50","askPrice":"180.52","bidSize":100,"askSize":200}"#).unwrap_or_else(|_| serde_json::json!({}));
378 let expected_response : models::LatestQuoteResponse = serde_json::from_value(resp_json.clone()).expect("should parse into models::LatestQuoteResponse");
379
380 let resp = client.latest_quote(params).await.expect("Expected a response");
381 let data_future = resp.data();
382 let actual_response = data_future.await.unwrap();
383 assert_eq!(actual_response, expected_response);
384 });
385 }
386
387 #[test]
388 fn latest_quote_response_error() {
389 TOKIO_SHARED_RT.block_on(async {
390 let client = MockMarketDataApiClient { force_error: true };
391
392 let params = LatestQuoteParams::builder("AAPL".to_string())
393 .build()
394 .unwrap();
395
396 match client.latest_quote(params).await {
397 Ok(_) => panic!("Expected an error"),
398 Err(err) => {
399 assert_eq!(err.to_string(), "Connector client error: ResponseError");
400 }
401 }
402 });
403 }
404
405 #[test]
406 fn tokenized_assets_required_params_success() {
407 TOKIO_SHARED_RT.block_on(async {
408 let client = MockMarketDataApiClient { force_error: false };
409
410
411 let resp_json: Value = serde_json::from_str(r#"[{"assetCode":"AAPLB","assetName":"Apple Inc. Tokenized Stock","underlyingEquitySymbol":"AAPL","multiplier":"1","multiplierValid":true}]"#).unwrap_or_else(|_| serde_json::json!({}));
412 let expected_response : Vec<models::TokenizedAssetsResponseInner> = serde_json::from_value(resp_json.clone()).expect("should parse into Vec<models::TokenizedAssetsResponseInner>");
413
414 let resp = client.tokenized_assets().await.expect("Expected a response");
415 let data_future = resp.data();
416 let actual_response = data_future.await.unwrap();
417 assert_eq!(actual_response, expected_response);
418 });
419 }
420
421 #[test]
422 fn tokenized_assets_optional_params_success() {
423 TOKIO_SHARED_RT.block_on(async {
424 let client = MockMarketDataApiClient { force_error: false };
425
426
427 let resp_json: Value = serde_json::from_str(r#"[{"assetCode":"AAPLB","assetName":"Apple Inc. Tokenized Stock","underlyingEquitySymbol":"AAPL","multiplier":"1","multiplierValid":true}]"#).unwrap_or_else(|_| serde_json::json!({}));
428 let expected_response : Vec<models::TokenizedAssetsResponseInner> = serde_json::from_value(resp_json.clone()).expect("should parse into Vec<models::TokenizedAssetsResponseInner>");
429
430 let resp = client.tokenized_assets().await.expect("Expected a response");
431 let data_future = resp.data();
432 let actual_response = data_future.await.unwrap();
433 assert_eq!(actual_response, expected_response);
434 });
435 }
436
437 #[test]
438 fn tokenized_assets_response_error() {
439 TOKIO_SHARED_RT.block_on(async {
440 let client = MockMarketDataApiClient { force_error: true };
441
442 match client.tokenized_assets().await {
443 Ok(_) => panic!("Expected an error"),
444 Err(err) => {
445 assert_eq!(err.to_string(), "Connector client error: ResponseError");
446 }
447 }
448 });
449 }
450}