1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
use rsolace::solclient::SolClient;
use rsolace::SessionProps;
use rsolace::solmsg::SolMsgBuilder;
use rsolace::types::{SolClientLogLevel, SolClientSubscribeFlags};
use tracing_subscriber;
fn main() {
tracing_subscriber::fmt()
.with_max_level(tracing::Level::DEBUG)
.init();
let solclient = SolClient::new(SolClientLogLevel::Notice);
match solclient {
Ok(mut solclient) => {
#[cfg(feature = "raw")]
{
solclient.set_rx_event_callback(|_, event| {
tracing::info!("{:?}", event);
});
solclient.set_rx_msg_callback(|_, msg| {
tracing::info!(
"{} {} {:?}",
msg.get_topic().unwrap(),
msg.get_sender_time()
.unwrap_or(chrono::prelude::Utc::now())
//.format("%Y-%m-%d %H:%M:%S%.3f")
// .to_string(),
.to_rfc3339(),
msg.get_binary_attachment().unwrap()
);
});
}
#[cfg(feature = "channel")]
{
let event_recv = solclient.get_event_receiver();
std::thread::spawn(move || loop {
let event = event_recv.recv().unwrap();
tracing::info!("{:?}", event);
});
let msg_recv = solclient.get_msg_receiver();
// let solclient1 = SolClient::new(SolClientLogLevel::Notice).unwrap();
std::thread::spawn(move || loop {
// solclient1.get_event_receiver();
match msg_recv.recv() {
Ok(msg) => {
tracing::info!(
"{} {} {:?}",
msg.get_topic().unwrap(),
msg.get_sender_dt()
.unwrap_or(chrono::prelude::Utc::now())
.to_rfc3339(),
msg.get_binary_attachment().unwrap()
);
}
Err(e) => {
tracing::error!("recv msg error: {}", e);
}
}
});
}
let props = SessionProps::default()
.host("218.32.76.102:80")
.vpn("sinopac")
.username("shioaji")
.password("shioaji111")
.reapply_subscriptions(true)
.connect_retries(1)
.connect_timeout_ms(3000)
.compression_level(5);
let r = solclient.connect(props);
tracing::info!("connect: {}", r);
// solclient.set_rx_msg_callback(func)
solclient.subscribe_ext(
"TIC/v1/STK/*/TSE/2230",
SolClientSubscribeFlags::RequestConfirm,
);
solclient.subscribe_ext(
"QUO/v1/STK/*/TSE/2330",
SolClientSubscribeFlags::RequestConfirm,
);
std::thread::sleep(std::time::Duration::from_secs(5));
let msg = SolMsgBuilder::new()
.with_topic("api/v1/test")
.as_delivery_to_one(true)
.build()
.unwrap();
let rt = solclient.send_msg(&msg);
tracing::info!("send msg: {:?}", rt);
// let mut msgs = vec![SolMsg::new().unwrap(), SolMsg::new().unwrap()];
// for (i, msg) in msgs.iter_mut().enumerate() {
// msg.set_topic(format!("api/v1/test/{}", i).as_str());
// }
let msgs = vec![
SolMsgBuilder::new().with_topic("api/v1/test/0").build().unwrap(),
SolMsgBuilder::new().with_topic("api/v1/test/1").build().unwrap(),
];
let rt = solclient.send_multiple_msg(&msgs.iter().collect::<Vec<_>>());
tracing::info!("send multiple msg: {:?}", rt);
let msg = SolMsgBuilder::new().with_topic("api/v1/test").build().unwrap();
let res = solclient.send_request(&msg, 0);
tracing::info!("send request msg: {:?}", res);
tracing::info!("done");
}
Err(e) => {
println!("error: {}", e)
}
}
// let solclient2 = SolClient::new(SolClientLogLevel::Notice);
// match solclient2 {
// Ok(mut solclient) => {
// let r = solclient.connect(
// "218.32.76.102:80",
// "sinopac",
// "shioaji",
// "shioaji111",
// Some("c2"),
// None,
// None,
// );
// println!("connect: {}", r);
// }
// Err(e) => {
// println!("error: {}", e)
// }
// }
// std::thread::sleep(std::time::Duration::from_secs(5));
}