Skip to main content

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;