Skip to main content

tradingview/live/
models.rs

1use core::fmt;
2
3use futures_util::stream::SplitStream;
4use serde::{Deserialize, Serialize};
5use serde_json::Value;
6use tokio::{net::TcpStream, sync::MutexGuard};
7use tokio_tungstenite::{
8    MaybeTlsStream, WebSocketStream,
9    tungstenite::{
10        http::{HeaderMap, HeaderValue},
11        protocol::Message,
12    },
13};
14use ustr::Ustr;
15
16use crate::{
17    Result,
18    error::{Error, TradingViewError},
19    utils::format_packet,
20};
21use std::sync::LazyLock;
22
23pub static WEBSOCKET_HEADERS: LazyLock<HeaderMap<HeaderValue>> = LazyLock::new(|| {
24    let mut headers = HeaderMap::new();
25    headers.insert("Origin", "https://www.tradingview.com/".parse().unwrap());
26    headers.insert(
27        "User-Agent",
28        "Mozilla/5.0 (Macintosh; Intel Mac OS X 10.15; rv:155.0) Gecko/20100101 Firefox/155.0"
29            .parse()
30            .unwrap(),
31    );
32    headers
33});
34
35/// WebSocket event types dispatched by TradingView's data server.
36///
37/// Maps TradingView's wire protocol event names (`"timescale_update"`,
38/// `"du"`, `"qsd"`, etc.) to Rust enum variants.
39#[derive(Debug, Clone, PartialEq, Eq, Hash, Copy)]
40pub enum TradingViewDataEvent {
41    OnChartData,
42    OnChartDataUpdate,
43    OnQuoteData,
44    OnQuoteCompleted,
45    OnSeriesLoading,
46    OnSeriesCompleted,
47    OnSymbolResolved,
48    OnReplayOk,
49    OnReplayPoint,
50    OnReplayInstanceId,
51    OnReplayResolutions,
52    OnReplayDataEnd,
53    OnStudyLoading,
54    OnStudyCompleted,
55    OnError(TradingViewError),
56    UnknownEvent(Ustr),
57}
58
59impl From<String> for TradingViewDataEvent {
60    fn from(s: String) -> Self {
61        match s.as_str() {
62            "timescale_update" => TradingViewDataEvent::OnChartData,
63            "du" => TradingViewDataEvent::OnChartDataUpdate,
64
65            "qsd" => TradingViewDataEvent::OnQuoteData,
66            "quote_completed" => TradingViewDataEvent::OnQuoteCompleted,
67
68            "series_loading" => TradingViewDataEvent::OnSeriesLoading,
69            "series_completed" => TradingViewDataEvent::OnSeriesCompleted,
70
71            "symbol_resolved" => TradingViewDataEvent::OnSymbolResolved,
72
73            "replay_ok" => TradingViewDataEvent::OnReplayOk,
74            "replay_point" => TradingViewDataEvent::OnReplayPoint,
75            "replay_instance_id" => TradingViewDataEvent::OnReplayInstanceId,
76            "replay_resolutions" => TradingViewDataEvent::OnReplayResolutions,
77            "replay_data_end" => TradingViewDataEvent::OnReplayDataEnd,
78
79            "study_loading" => TradingViewDataEvent::OnStudyLoading,
80            "study_completed" => TradingViewDataEvent::OnStudyCompleted,
81
82            "symbol_error" => TradingViewDataEvent::OnError(TradingViewError::SymbolError),
83            "series_error" => TradingViewDataEvent::OnError(TradingViewError::SeriesError),
84            "critical_error" => TradingViewDataEvent::OnError(TradingViewError::CriticalError),
85            "study_error" => TradingViewDataEvent::OnError(TradingViewError::StudyError),
86            "protocol_error" => TradingViewDataEvent::OnError(TradingViewError::ProtocolError),
87            "replay_error" => TradingViewDataEvent::OnError(TradingViewError::ReplayError),
88
89            s => TradingViewDataEvent::UnknownEvent(s.into()),
90        }
91    }
92}
93
94impl From<Ustr> for TradingViewDataEvent {
95    fn from(s: Ustr) -> Self {
96        TradingViewDataEvent::from(s.to_string())
97    }
98}
99
100/// A serialized WebSocket message ready for transmission.
101#[derive(Debug, Clone, PartialEq, Serialize)]
102pub struct SocketMessageSer {
103    pub m: Value,
104    pub p: Value,
105}
106
107/// A deserialized WebSocket message from TradingView.
108#[derive(Debug, Clone, PartialEq, Deserialize)]
109pub struct SocketMessageDe {
110    pub m: Ustr,
111    pub p: Vec<Value>,
112    #[serde(default)]
113    pub t: u64, // Timestamp in seconds (0 when absent, e.g. error messages)
114    #[serde(default)]
115    pub t_ms: u64, // Timestamp in milliseconds (0 when absent)
116}
117
118impl SocketMessageSer {
119    pub fn new<M, P>(m: M, p: P) -> Self
120    where
121        M: Serialize,
122        P: Serialize,
123    {
124        let m = serde_json::to_value(m).expect("Failed to serialize Socket Message");
125        let p = serde_json::to_value(p).expect("Failed to serialize Socket Message");
126        SocketMessageSer { m, p }
127    }
128
129    pub fn to_message(&self) -> Result<Message> {
130        let msg = format_packet(self)?;
131        Ok(msg)
132    }
133}
134
135/// Server metadata sent by TradingView after a successful WebSocket connection.
136///
137/// Contains the session ID, server timestamp, and base URL for chart data.
138#[derive(Default, Debug, Clone, PartialEq, Serialize, Deserialize)]
139#[serde(rename_all = "camelCase")]
140pub struct SocketServerInfo {
141    #[serde(rename = "session_id")]
142    pub session_id: Ustr,
143    pub timestamp: i64,
144    pub timestamp_ms: i64,
145    pub release: Ustr,
146    #[serde(rename = "studies_metadata_hash")]
147    pub studies_metadata_hash: Ustr,
148    #[serde(rename = "auth_scheme_vsn")]
149    pub auth_scheme_vsn: i64,
150    pub protocol: Ustr,
151    pub via: Ustr,
152    #[serde(rename = "javastudies")]
153    pub sjavastudies: Vec<Ustr>,
154}
155
156impl fmt::Display for SocketServerInfo {
157    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
158        write!(
159            f,
160            "SocketServerInfo {{ session_id: {}, timestamp: {}, timestamp_ms: {}, release: {}, studies_metadata_hash: {}, auth_scheme_vsn: {}, protocol: {}, via: {}, sjavastudies: {:?} }}",
161            self.session_id,
162            self.timestamp,
163            self.timestamp_ms,
164            self.release,
165            self.studies_metadata_hash,
166            self.auth_scheme_vsn,
167            self.protocol,
168            self.via,
169            self.sjavastudies
170        )
171    }
172}
173
174#[derive(Debug, Clone, PartialEq, Deserialize, Serialize)]
175#[serde(untagged)]
176pub enum SocketMessage<T> {
177    SocketServerInfo(SocketServerInfo),
178    SocketMessage(T),
179    Heartbeat(u64),
180    Other(Value),
181    Unknown(String),
182}
183
184impl<T> SocketMessage<T> {
185    pub fn heartbeat_echo(&self) -> Option<String> {
186        match self {
187            SocketMessage::Heartbeat(counter) => {
188                let payload = format!("~h~{counter}");
189                Some(format!("~m~{}~m~{payload}", payload.len()))
190            }
191            _ => None,
192        }
193    }
194}
195
196/// Which TradingView data server tier to connect to.
197///
198/// `ProData` is the default and recommended server. `Data` and
199/// `DataExtended` are alternatives with different capabilities.
200#[derive(Default, Clone, Debug, PartialEq, Serialize, Deserialize, Copy, Eq)]
201pub enum DataServer {
202    #[default]
203    Data,
204    ProData,
205    WidgetData,
206    MobileData,
207}
208
209impl std::fmt::Display for DataServer {
210    fn fmt(&self, f: &mut std::fmt::Formatter) -> std::fmt::Result {
211        match *self {
212            DataServer::Data => write!(f, "data"),
213            DataServer::ProData => write!(f, "prodata"),
214            DataServer::WidgetData => write!(f, "widgetdata"),
215            DataServer::MobileData => write!(f, "mobile-data"),
216        }
217    }
218}
219
220/// Trait for WebSocket message handling — serialize and deserialize from
221/// TradingView's wire format.
222///
223/// Implemented by [`SocketMessage`] for the two protocol variants.
224///
225/// [`SocketMessage`]: crate::live::models::SocketMessage
226pub trait Socket {
227    fn event_loop(
228        &self,
229        read: MutexGuard<SplitStream<WebSocketStream<MaybeTlsStream<TcpStream>>>>,
230    ) -> impl Future<Output = Result<()>> + Send;
231
232    fn handle_raw_messages(&self, raw: Message) -> impl Future<Output = Result<()>> + Send;
233
234    fn handle_parsed_messages(
235        &self,
236        messages: Vec<SocketMessage<SocketMessageDe>>,
237        raw: &Message,
238    ) -> impl Future<Output = Result<()>> + Send;
239
240    fn handle_message_data(
241        &self,
242        message: SocketMessageDe,
243    ) -> impl Future<Output = Result<()>> + Send;
244
245    fn handle_error(&self, error: Error, context: Ustr) -> impl Future<Output = Result<()>> + Send;
246}
247
248// ---------------------------------------------------------------------------
249// Tests
250// ---------------------------------------------------------------------------
251
252#[cfg(test)]
253mod tests {
254    use super::*;
255
256    #[test]
257    fn test_study_loading_maps_to_on_study_loading() {
258        let event = TradingViewDataEvent::from("study_loading".to_string());
259        assert_eq!(event, TradingViewDataEvent::OnStudyLoading);
260    }
261
262    #[test]
263    fn test_series_loading_is_not_study_loading() {
264        let event = TradingViewDataEvent::from("series_loading".to_string());
265        assert_eq!(event, TradingViewDataEvent::OnSeriesLoading);
266    }
267
268    #[test]
269    fn test_study_completed_maps_correctly() {
270        let event = TradingViewDataEvent::from("study_completed".to_string());
271        assert_eq!(event, TradingViewDataEvent::OnStudyCompleted);
272    }
273
274    #[test]
275    fn test_series_completed_maps_correctly() {
276        let event = TradingViewDataEvent::from("series_completed".to_string());
277        assert_eq!(event, TradingViewDataEvent::OnSeriesCompleted);
278    }
279
280    #[test]
281    fn test_all_study_and_series_events_are_distinct() {
282        let loading = TradingViewDataEvent::from("study_loading".to_string());
283        let completed = TradingViewDataEvent::from("study_completed".to_string());
284        let series_loading = TradingViewDataEvent::from("series_loading".to_string());
285        let series_completed = TradingViewDataEvent::from("series_completed".to_string());
286
287        // All four should be distinct
288        assert_ne!(loading, series_loading);
289        assert_ne!(completed, series_completed);
290        assert_ne!(loading, completed);
291        assert_ne!(series_loading, series_completed);
292
293        // Verify mapping expectations
294        assert_eq!(loading, TradingViewDataEvent::OnStudyLoading);
295        assert_eq!(completed, TradingViewDataEvent::OnStudyCompleted);
296        assert_eq!(series_loading, TradingViewDataEvent::OnSeriesLoading);
297        assert_eq!(series_completed, TradingViewDataEvent::OnSeriesCompleted);
298    }
299
300    #[test]
301    fn test_ustr_from_maps_correctly() {
302        let event: TradingViewDataEvent = ustr::ustr("study_loading").into();
303        assert_eq!(event, TradingViewDataEvent::OnStudyLoading);
304    }
305}