1use core::fmt;
2use std::vec;
3
4use serde::Deserialize;
5use serde_json::{from_str, Value};
6use tokio_tungstenite::tungstenite::{Bytes, Message};
7use tracing::warn;
8
9use crate::{
10 general::traits::MessageTransfer,
11 pocketoption::{
12 error::{PocketOptionError, PocketResult},
13 types::{
14 base::{ChangeSymbol, SubscribeSymbol},
15 info::MessageInfo,
16 order::{
17 Deal, FailOpenOrder, FailOpenPendingOrder, OpenOrder, OpenPendingOrder,
18 PocketMessageFail, SuccessCloseOrder, SuccessOpenPendingOrder, UpdateClosedDeals,
19 UpdateOpenedDeals,
20 },
21 success::SuccessAuth,
22 update::{
23 LoadHistoryPeriodResult, UpdateAssets, UpdateBalance, UpdateHistoryNew,
24 UpdateStream,
25 },
26 user::PocketUser,
27 },
28 ws::ssid::Ssid,
29 },
30};
31
32use super::basic::LoadHistoryPeriod;
33
34#[derive(Debug, Deserialize, Clone)]
35#[serde(untagged)]
36pub enum WebSocketMessage {
37 OpenOrder(OpenOrder),
38 ChangeSymbol(ChangeSymbol),
39 Auth(Ssid),
40 GetCandles(LoadHistoryPeriod),
41
42 LoadHistoryPeriod(LoadHistoryPeriodResult),
43 UpdateStream(UpdateStream),
44 UpdateHistoryNew(UpdateHistoryNew),
45 SubscribeSymbol(SubscribeSymbol),
46 UpdateAssets(UpdateAssets),
47 UpdateBalance(UpdateBalance),
48 SuccessAuth(SuccessAuth),
49 UpdateClosedDeals(UpdateClosedDeals),
50 SuccesscloseOrder(SuccessCloseOrder),
51 SuccessopenOrder(Deal),
52 SuccessupdateBalance(UpdateBalance),
53 UpdateOpenedDeals(UpdateOpenedDeals),
54 FailOpenOrder(FailOpenOrder),
55 FailOpenPendingOrder(FailOpenPendingOrder),
56 SuccessupdatePending(Value),
57 OpenPendingOrder(OpenPendingOrder),
58 SuccessOpenPendingOrder(SuccessOpenPendingOrder),
59 UserRequest(Box<PocketUser>),
60 None,
61}
62
63impl WebSocketMessage {
64 pub fn parse(data: impl ToString) -> PocketResult<Self> {
65 let data = data.to_string();
66 let message: Result<Self, serde_json::Error> = from_str(&data);
67 match message {
68 Ok(message) => Ok(message),
69 Err(e) => {
70 if let Ok(assets) = from_str::<UpdateAssets>(&data) {
71 return Ok(Self::UpdateAssets(assets));
72 }
73 if let Ok(history) = from_str::<UpdateHistoryNew>(&data) {
74 return Ok(Self::UpdateHistoryNew(history));
75 }
76 if let Ok(stream) = from_str::<UpdateStream>(&data) {
77 return Ok(Self::UpdateStream(stream));
78 }
79 if let Ok(balance) = from_str::<UpdateBalance>(&data) {
80 return Ok(Self::UpdateBalance(balance));
81 }
82 Err(e.into())
83 }
84 }
85 }
86
87 pub fn parse_with_context(data: impl ToString, previous: &MessageInfo) -> PocketResult<Self> {
88 let data = data.to_string();
89 match previous {
90 MessageInfo::OpenOrder => {
91 if let Ok(order) = from_str::<OpenOrder>(&data) {
92 return Ok(Self::OpenOrder(order));
93 }
94 }
95 MessageInfo::UpdateStream => {
96 if let Ok(stream) = from_str::<UpdateStream>(&data) {
97 return Ok(Self::UpdateStream(stream));
98 }
99 }
100 MessageInfo::UpdateHistoryNew => {
101 if let Ok(history) = from_str::<UpdateHistoryNew>(&data) {
102 return Ok(Self::UpdateHistoryNew(history));
103 }
104 }
105 MessageInfo::UpdateAssets => {
106 if let Ok(assets) = from_str::<UpdateAssets>(&data) {
107 return Ok(Self::UpdateAssets(assets));
108 }
109 }
110 MessageInfo::UpdateBalance => {
111 if let Ok(balance) = from_str::<UpdateBalance>(&data) {
112 return Ok(Self::UpdateBalance(balance));
113 }
114 }
115 MessageInfo::SuccesscloseOrder => {
116 if let Ok(order) = from_str::<SuccessCloseOrder>(&data) {
117 return Ok(Self::SuccesscloseOrder(order));
118 }
119 }
120 MessageInfo::Auth => {
121 if let Ok(auth) = from_str::<Ssid>(&data) {
122 return Ok(Self::Auth(auth));
123 }
124 }
125 MessageInfo::ChangeSymbol => {
126 if let Ok(symbol) = from_str::<ChangeSymbol>(&data) {
127 return Ok(Self::ChangeSymbol(symbol));
128 }
129 }
130 MessageInfo::SuccessupdateBalance => {
131 if let Ok(balance) = from_str::<UpdateBalance>(&data) {
132 return Ok(Self::SuccessupdateBalance(balance));
133 }
134 }
135 MessageInfo::SuccessupdatePending => {
136 if let Ok(pending) = from_str::<Value>(&data) {
137 return Ok(Self::SuccessupdatePending(pending));
138 }
139 }
140 MessageInfo::SubscribeSymbol => {
141 if let Ok(symbol) = from_str::<SubscribeSymbol>(&data) {
142 return Ok(Self::SubscribeSymbol(symbol));
143 }
144 }
145 MessageInfo::Successauth => {
146 if let Ok(auth) = from_str::<SuccessAuth>(&data) {
147 return Ok(Self::SuccessAuth(auth));
148 }
149 }
150 MessageInfo::UpdateOpenedDeals => {
151 if let Ok(deals) = from_str::<UpdateOpenedDeals>(&data) {
152 return Ok(Self::UpdateOpenedDeals(deals));
153 }
154 }
155 MessageInfo::UpdateClosedDeals => {
156 if let Ok(deals) = from_str::<UpdateClosedDeals>(&data) {
157 return Ok(Self::UpdateClosedDeals(deals));
158 }
159 }
160 MessageInfo::SuccessopenOrder => {
161 if let Ok(order) = from_str::<Deal>(&data) {
162 return Ok(Self::SuccessopenOrder(order));
163 }
164 }
165 MessageInfo::LoadHistoryPeriod => {
166 if let Ok(history) = from_str::<LoadHistoryPeriodResult>(&data) {
167 return Ok(Self::LoadHistoryPeriod(history));
168 }
169 }
170 MessageInfo::UpdateCharts => {
171 return Err(PocketOptionError::GeneralParsingError(
172 "This is expected, there is no parser for the 'updateCharts' message"
173 .to_string(),
174 ));
175 }
177 MessageInfo::GetCandles => {
178 if let Ok(candles) = from_str::<LoadHistoryPeriod>(&data) {
179 return Ok(Self::GetCandles(candles));
180 }
181 }
182 MessageInfo::FailopenOrder => {
183 if let Ok(fail) = from_str::<FailOpenOrder>(&data) {
184 return Ok(Self::FailOpenOrder(fail));
185 }
186 }
187 MessageInfo::FailopenPendingOrder => {
188 if let Ok(fail) = from_str::<FailOpenPendingOrder>(&data) {
189 return Ok(Self::FailOpenPendingOrder(fail));
190 }
191 }
192 MessageInfo::OpenPendingOrder => {
193 if let Ok(order) = from_str::<OpenPendingOrder>(&data) {
194 return Ok(Self::OpenPendingOrder(order));
195 }
196 }
197 MessageInfo::SuccessopenPendingOrder => {
198 if let Ok(order) = from_str::<SuccessOpenPendingOrder>(&data) {
199 return Ok(Self::SuccessOpenPendingOrder(order));
200 }
201 }
202 MessageInfo::None => return WebSocketMessage::parse(data.clone()),
203 }
204 warn!("Failed to parse message of type '{previous}':\n {data}");
205 Err(PocketOptionError::GeneralParsingError(format!(
206 "Error parsing message for message type '{}'",
207 previous
208 )))
209 }
210
211 pub fn info(&self) -> MessageInfo {
212 match self {
213 Self::UpdateStream(_) => MessageInfo::UpdateStream,
214 Self::UpdateHistoryNew(_) => MessageInfo::UpdateHistoryNew,
215 Self::UpdateAssets(_) => MessageInfo::UpdateAssets,
216 Self::UpdateBalance(_) => MessageInfo::UpdateBalance,
217 Self::OpenOrder(_) => MessageInfo::OpenOrder,
218 Self::SuccessAuth(_) => MessageInfo::Successauth,
219 Self::UpdateClosedDeals(_) => MessageInfo::UpdateClosedDeals,
220 Self::SuccesscloseOrder(_) => MessageInfo::SuccesscloseOrder,
221 Self::SuccessopenOrder(_) => MessageInfo::SuccessopenOrder,
222 Self::ChangeSymbol(_) => MessageInfo::ChangeSymbol,
223 Self::Auth(_) => MessageInfo::Auth,
224 Self::SuccessupdateBalance(_) => MessageInfo::SuccessupdateBalance,
225 Self::UpdateOpenedDeals(_) => MessageInfo::UpdateOpenedDeals,
226 Self::SubscribeSymbol(_) => MessageInfo::SubscribeSymbol,
227 Self::LoadHistoryPeriod(_) => MessageInfo::LoadHistoryPeriod,
228 Self::GetCandles(_) => MessageInfo::GetCandles,
229 Self::UserRequest(_) => MessageInfo::None,
230 Self::FailOpenOrder(_) => MessageInfo::FailopenOrder,
231 Self::SuccessupdatePending(_) => MessageInfo::SuccessupdatePending,
232 Self::FailOpenPendingOrder(_) => MessageInfo::FailopenPendingOrder,
233 Self::SuccessOpenPendingOrder(_) => MessageInfo::SuccessopenPendingOrder,
234 Self::OpenPendingOrder(_) => MessageInfo::OpenPendingOrder,
235 Self::None => MessageInfo::None,
236 }
237 }
238}
239
240impl fmt::Display for WebSocketMessage {
241 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
242 match self {
243 WebSocketMessage::ChangeSymbol(change_symbol) => {
244 write!(
245 f,
246 "42[{},{}]",
247 serde_json::to_string(&MessageInfo::ChangeSymbol).map_err(|_| fmt::Error)?,
248 serde_json::to_string(&change_symbol).map_err(|_| fmt::Error)?
249 )
250 }
251 WebSocketMessage::Auth(auth) => auth.fmt(f),
252 WebSocketMessage::GetCandles(candles) => {
253 write!(
254 f,
255 "42[{},{}]",
256 serde_json::to_string(&MessageInfo::LoadHistoryPeriod)
257 .map_err(|_| fmt::Error)?,
258 serde_json::to_string(candles).map_err(|_| fmt::Error)?
259 )
260 }
261 WebSocketMessage::OpenOrder(open_order) => {
262 write!(
263 f,
264 "42[{},{}]",
265 serde_json::to_string(&MessageInfo::OpenOrder).map_err(|_| fmt::Error)?,
266 serde_json::to_string(open_order).map_err(|_| fmt::Error)?
267 )
268 }
269 WebSocketMessage::SubscribeSymbol(subscribe_symbol) => {
270 write!(f, "{:?}", subscribe_symbol)
271 }
272
273 WebSocketMessage::UpdateStream(update_stream) => write!(f, "{:?}", update_stream),
274 WebSocketMessage::UpdateHistoryNew(update_history_new) => {
275 write!(f, "{:?}", update_history_new)
276 }
277 WebSocketMessage::UpdateAssets(update_assets) => write!(f, "{:?}", update_assets),
278 WebSocketMessage::UpdateBalance(update_balance) => write!(f, "{:?}", update_balance),
279 WebSocketMessage::SuccessAuth(success_auth) => write!(f, "{:?}", success_auth),
280 WebSocketMessage::UpdateClosedDeals(update_closed_deals) => {
281 write!(f, "{:?}", update_closed_deals)
282 }
283 WebSocketMessage::SuccesscloseOrder(success_close_order) => {
284 write!(f, "{:?}", success_close_order)
285 }
286 WebSocketMessage::SuccessopenOrder(success_open_order) => {
287 write!(f, "{:?}", success_open_order)
288 }
289 WebSocketMessage::SuccessupdateBalance(update_balance) => {
290 write!(f, "{:?}", update_balance)
291 }
292 WebSocketMessage::UpdateOpenedDeals(update_opened_deals) => {
293 write!(f, "{:?}", update_opened_deals)
294 }
295 WebSocketMessage::SuccessOpenPendingOrder(order) => write!(f, "{:?}", order),
296 WebSocketMessage::FailOpenPendingOrder(order) => write!(f, "{:?}", order),
297 WebSocketMessage::OpenPendingOrder(order) => write!(f, "{:?}", order),
298
299 WebSocketMessage::None => write!(f, "None"),
300 WebSocketMessage::LoadHistoryPeriod(period) => {
302 write!(
303 f,
304 "42[{}, {}]",
305 serde_json::to_string(&MessageInfo::LoadHistoryPeriod)
306 .map_err(|_| fmt::Error)?,
307 serde_json::to_string(&period).map_err(|_| fmt::Error)?
308 )
309 }
310 WebSocketMessage::UserRequest(user) => {
311 write!(f, "Request of type {:?}", user.info)
312 }
313 WebSocketMessage::FailOpenOrder(order) => order.fmt(f),
314 WebSocketMessage::SuccessupdatePending(pending) => pending.fmt(f),
315 }
316 }
317}
318
319impl From<WebSocketMessage> for Message {
320 fn from(value: WebSocketMessage) -> Self {
321 Box::new(value).into()
322 }
323}
324
325impl From<Box<WebSocketMessage>> for Message {
326 fn from(value: Box<WebSocketMessage>) -> Self {
327 if value.info() == MessageInfo::None {
328 return Message::Ping(Bytes::new());
329 }
330 Message::text(value.to_string())
331 }
332}
333
334impl MessageTransfer for WebSocketMessage {
335 type Error = PocketMessageFail;
336
337 type TransferError = PocketMessageFail;
338
339 type Info = MessageInfo;
340
341 fn info(&self) -> MessageInfo {
342 self.info()
343 }
344
345 fn error(&self) -> Option<Self::Error> {
346 if let Self::FailOpenOrder(fail) = self {
347 return Some(PocketMessageFail::Order(fail.to_owned()));
348 }
349 None
350 }
351
352 fn to_error(&self) -> Self::TransferError {
353 if let Self::FailOpenOrder(fail) = self {
354 PocketMessageFail::Order(fail.to_owned())
355 } else {
356 PocketMessageFail::Order(FailOpenOrder::new(
357 "This is unexpected and should never happend",
358 1.0,
359 "None",
360 ))
361 }
362 }
363
364 fn user_request(&self) -> Option<PocketUser> {
365 if let Self::UserRequest(user) = self {
366 let user = *user.to_owned();
367 return Some(user);
368 }
369 None
370 }
371
372 fn new_user(request: PocketUser) -> Self {
373 Self::UserRequest(Box::new(request))
374 }
375
376 fn error_info(&self) -> Option<Vec<Self::Info>> {
377 if let Self::FailOpenOrder(_) = self {
378 Some(vec![MessageInfo::SuccessopenOrder])
379 } else {
380 None
381 }
382 }
383}
384
385#[cfg(test)]
386mod tests {
387 use super::*;
388
389 use std::{
390 error::Error,
391 fs::File,
392 io::{BufReader, Read, Write},
393 };
394
395 use std::fs;
396 use std::path::Path;
397
398 fn get_files_in_directory(path: &str) -> Result<Vec<String>, std::io::Error> {
399 let dir_path = Path::new(path);
400
401 match fs::read_dir(dir_path) {
402 Ok(entries) => {
403 let mut file_names = Vec::new();
404
405 for entry in entries {
406 let file_name = entry?.file_name().to_string_lossy().to_string();
407 file_names.push(format!("{path}/{file_name}"));
408 }
409
410 Ok(file_names)
411 }
412 Err(e) => Err(e),
413 }
414 }
415
416 #[test]
417 fn test_descerialize_message() -> Result<(), Box<dyn Error>> {
418 let tests = [
419 r#"[["AUS200_otc",1732830010,6436.06]]"#,
420 r#"[["AUS200_otc",1732830108.205,6435.96]]"#,
421 r#"[["AEDCNY_otc",1732829668.352,1.89817]]"#,
422 r#"[["CADJPY_otc",1732830170.793,109.442]]"#,
423 ];
424 for item in tests.iter() {
425 let val = WebSocketMessage::parse(item)?;
426 dbg!(&val);
427 }
428 let mut history_raw = File::open("tests/update_history_new.txt")?;
429 let mut content = String::new();
430 history_raw.read_to_string(&mut content)?;
431 let history_new: WebSocketMessage = from_str(&content)?;
432 dbg!(&history_new);
433
434 let mut assets_raw = File::open("tests/data.json")?;
435 let mut content = String::new();
436 assets_raw.read_to_string(&mut content)?;
437 let assets_raw: WebSocketMessage = from_str(&content)?;
438 dbg!(&assets_raw);
439 Ok(())
440 }
441
442 #[test]
443 fn deep_test_descerialize_message() -> anyhow::Result<()> {
444 let dirs = get_files_in_directory("tests")?;
445 for dir in dirs {
446 dbg!(&dir);
447 if !dir.ends_with(".json") {
448 continue;
449 }
450 let file = File::open(dir)?;
451
452 let reader = BufReader::new(file);
453 let _: WebSocketMessage = serde_json::from_reader(reader)?;
454 }
455
456 Ok(())
457 }
458
459 #[test]
460 fn test_write_assets() -> anyhow::Result<()> {
461 let raw: UpdateAssets = serde_json::from_str(include_str!("../../../tests/data.json"))?;
462 let mut file = File::create("tests/assets.txt")?;
463 let data = raw.0.iter().fold(String::new(), |mut s, a| {
464 s.push_str(&format!("{}\n", a.symbol));
465 s
466 });
467 file.write_all(data.as_bytes())?;
468 Ok(())
469 }
470}