ibapi/wsh/async.rs
1//! Asynchronous implementation of Wall Street Horizon functionality
2
3use time::Date;
4
5use crate::{
6 common::request_helpers::{self, expect_proto},
7 protocol::{check_version, Features},
8 subscriptions::Subscription,
9 Client, Error,
10};
11
12use super::builder::{WshEventDataBuilder, WshEventFilterBuilder};
13use super::{common::decoders, encoders, AutoFill, WshEventData, WshMetadata};
14
15impl Client {
16 /// Fetch Wall Street Horizon metadata table with retry semantics.
17 ///
18 /// # Examples
19 ///
20 /// ```no_run
21 /// use ibapi::prelude::*;
22 ///
23 /// #[tokio::main]
24 /// async fn main() {
25 /// let client = Client::connect("127.0.0.1:4002", 100).await.expect("connection failed");
26 /// let metadata = client.wsh_metadata().await.expect("request wsh metadata failed");
27 /// println!("{metadata:?}");
28 /// }
29 /// ```
30 pub async fn wsh_metadata(&self) -> Result<WshMetadata, Error> {
31 check_version(self.server_version(), Features::WSHE_CALENDAR)?;
32
33 request_helpers::one_shot_by_request_id(
34 self,
35 encoders::encode_request_wsh_metadata,
36 expect_proto(decoders::decode_wsh_metadata_proto),
37 )
38 .await
39 }
40
41 /// Build a request for Wall Street Horizon events on one contract.
42 ///
43 /// Terminal: [`WshEventDataBuilder::fetch`]. Optional narrowing via
44 /// `.starting()` / `.ending()` / `.limit()` / `.auto_fill()`, each of which
45 /// carries its own server-version requirement.
46 ///
47 /// # Arguments
48 ///
49 /// * `contract_id` - Contract identifier for the event request.
50 ///
51 /// # Examples
52 ///
53 /// ```no_run
54 /// use ibapi::Client;
55 ///
56 /// #[tokio::main]
57 /// async fn main() {
58 /// let client = Client::connect("127.0.0.1:4002", 100).await.expect("connection failed");
59 ///
60 /// let contract_id = 76792991; // TSLA
61 /// let event_data = client
62 /// .wsh_event_data_by_contract(contract_id)
63 /// .fetch()
64 /// .await
65 /// .expect("request wsh event data failed");
66 /// println!("{event_data:?}");
67 /// }
68 /// ```
69 pub fn wsh_event_data_by_contract(&self, contract_id: i32) -> WshEventDataBuilder<'_, Self> {
70 WshEventDataBuilder::new(self, contract_id)
71 }
72
73 /// Build a request for Wall Street Horizon events matching a JSON filter.
74 ///
75 /// Terminal: [`WshEventFilterBuilder::subscribe`].
76 ///
77 /// # Arguments
78 ///
79 /// * `filter` - Json-formatted string containing all filter values.
80 ///
81 /// # Examples
82 ///
83 /// ```no_run
84 /// use ibapi::Client;
85 /// use ibapi::subscriptions::r#async::SubscriptionItemStreamExt;
86 /// use futures::StreamExt;
87 ///
88 /// #[tokio::main]
89 /// async fn main() {
90 /// let client = Client::connect("127.0.0.1:4002", 100).await.expect("connection failed");
91 ///
92 /// let filter = "{}"; // see https://www.interactivebrokers.com/campus/ibkr-api-page/twsapi-doc/#wsheventdata-object
93 /// let mut subscription = client
94 /// .wsh_event_data_by_filter(filter)
95 /// .subscribe()
96 /// .await
97 /// .expect("request wsh event data failed");
98 ///
99 /// let mut events = subscription.filter_data();
100 /// while let Some(event) = events.next().await {
101 /// println!("{:?}", event.expect("decode error"));
102 /// }
103 /// }
104 /// ```
105 pub fn wsh_event_data_by_filter<'a>(&'a self, filter: &'a str) -> WshEventFilterBuilder<'a, Self> {
106 WshEventFilterBuilder::new(self, filter)
107 }
108}
109
110/// Request events for one contract. Reached through
111/// [`WshEventDataBuilder::fetch`](super::builder::WshEventDataBuilder::fetch).
112pub(crate) async fn wsh_event_data_by_contract(
113 client: &Client,
114 contract_id: i32,
115 start_date: Option<Date>,
116 end_date: Option<Date>,
117 limit: Option<i32>,
118 auto_fill: Option<AutoFill>,
119) -> Result<WshEventData, Error> {
120 check_version(client.server_version(), Features::WSHE_CALENDAR)?;
121
122 if auto_fill.is_some() {
123 check_version(client.server_version(), Features::WSH_EVENT_DATA_FILTERS)?;
124 }
125
126 if start_date.is_some() || end_date.is_some() || limit.is_some() {
127 check_version(client.server_version(), Features::WSH_EVENT_DATA_FILTERS_DATE)?;
128 }
129
130 request_helpers::one_shot_by_request_id(
131 client,
132 |request_id| encoders::encode_request_wsh_event_data(request_id, Some(contract_id), None, start_date, end_date, limit, auto_fill),
133 expect_proto(decoders::decode_wsh_event_data_proto),
134 )
135 .await
136}
137
138/// Request events matching a JSON filter. Reached through
139/// [`WshEventFilterBuilder::subscribe`](super::builder::WshEventFilterBuilder::subscribe).
140pub(crate) async fn wsh_event_data_by_filter(
141 client: &Client,
142 filter: &str,
143 limit: Option<i32>,
144 auto_fill: Option<AutoFill>,
145) -> Result<Subscription<WshEventData>, Error> {
146 if limit.is_some() {
147 check_version(client.server_version(), Features::WSH_EVENT_DATA_FILTERS_DATE)?;
148 }
149
150 request_helpers::request_with_id(client, Features::WSH_EVENT_DATA_FILTERS, |request_id| {
151 encoders::encode_request_wsh_event_data(
152 request_id,
153 None,
154 Some(filter),
155 None, // start_date
156 None, // end_date
157 limit,
158 auto_fill,
159 )
160 })
161 .await
162}
163
164#[cfg(test)]
165#[path = "async_tests.rs"]
166mod tests;