Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_steel/tests/client.rs

13.9 KiB, 86 runs

created by r1870400018:1009, which is this file's identity for as long as the history lasts, whatever it is later renamed to

download · who wrote it · its history

1use oxedyne_fe2o3_steel::srv::{
2 constant,
3 context,
4 id,
5 ws::syntax::WebSocketSyntax,
6};
7
8use oxedyne_fe2o3_core::{
9 prelude::*,
10};
11use oxedyne_fe2o3_crypto::enc::EncryptionScheme;
12use oxedyne_fe2o3_hash::{
13 csum::ChecksumScheme,
14 hash::HashScheme,
15};
16use oxedyne_fe2o3_jdat::version::SemVer;
17use oxedyne_fe2o3_net::{
18 conc::AsyncReadIterator,
19 http::msg::{
20 //AsyncReadIterator,
21 HttpMessage,
22 HttpMessageReader,
23 },
24 ws::{
25 self,
26 WebSocket,
27 WebSocketMessage,
28 handler::WebSocketSinkHandler,
29 status::WebSocketStatusCode,
30 },
31};
32use oxedyne_fe2o3_o3db_sync::O3db;
33
34use std::{
35 fs::File,
36 io::BufReader,
37 path::Path,
38 pin::Pin,
39 sync::{
40 Arc,
41 RwLock,
42 },
43 thread,
44 time::Duration,
45};
46
47//use rustls_pki_types;
48use tokio::{
49 self,
50 io::{
51 AsyncWriteExt,
52 },
53 net::TcpStream,
54};
55use tokio_rustls::{
56 client::TlsStream,
57 rustls::{
58 self,
59 client::danger::ServerCertVerifier,
60 ClientConfig,
61 RootCertStore,
62 },
63 TlsConnector,
64};
65
66fn load_certs() -> Outcome<RootCertStore> {
67 let mut root_store = RootCertStore::empty();
68 let home = res!(std::env::var("HOME"));
69 let path = Path::new(&home).join("usr/code/web/apps/test/tls/fullchain.pem");
70 let cert_file = res!(File::open(path));
71 let mut reader = BufReader::new(cert_file);
72 let certs = res!(rustls_pemfile::certs(&mut reader).collect::<Result<Vec<_>, _>>());
73 for cert in certs {
74 res!(root_store.add(cert));
75 }
76 Ok(root_store)
77}
78
79pub async fn new_stream(host: &str, port: u16) -> Outcome<TlsStream<TcpStream>> {
80 //let mut root_cert_store = RootCertStore::empty();
81 //root_cert_store.extend(webpki_roots::TLS_SERVER_ROOTS.iter().cloned());
82 let root_cert_store = res!(load_certs());
83 let config = ClientConfig::builder()
84 .with_root_certificates(root_cert_store)
85 .with_no_client_auth();
86
87 let connector = TlsConnector::from(Arc::new(config));
88 let host_clone = host.to_string();
89 let dnsname = res!(rustls::pki_types::ServerName::try_from(host_clone));
90
91 let result = TcpStream::connect((host, port)).await;
92 let stream = res!(result);
93 let result = connector.connect(dnsname, stream).await;
94 let stream = res!(result);
95 Ok(stream)
96}
97
98pub async fn test_client(filter: &'static str) -> Outcome<()> {
99
100 let host = "localhost";
101 let port = 8443;
102
103 match filter {
104 "all" | "firefox" | "get" => {
105 let result = new_stream(host, port).await;
106 let mut stream = res!(result);
107
108 //tokio::time::sleep(tokio::time::Duration::from_secs(10)).await;
109
110 let request_string = fmt!("GET / HTTP/1.1\r\n\
111 Host: {}\r\n\
112 User-Agent: Mozilla/5.0 (Windows NT 10.0; Win64; x64; rv:93.0) Gecko/20100101 Firefox/93.0\r\n\
113 Accept: text/html,application/xhtml+xml,application/xml;q=0.9,image/avif,image/webp,*/*;q=0.8\r\n\
114 Accept-Language: en-US,en;q=0.5\r\n\
115 Accept-Encoding: gzip, deflate, br\r\n\
116 Connection: keep-alive\r\n\
117 Upgrade-Insecure-Requests: 1\r\n\
118 Sec-Fetch-Dest: document\r\n\
119 Sec-Fetch-Mode: navigate\r\n\
120 Sec-Fetch-Site: none\r\n\
121 Sec-Fetch-User: ?1\r\n\
122 Cache-Control: max-age=0\r\n\r\n",
123 host);
124
125 debug!("Writing to stream now...");
126 for line in request_string.lines() {
127 debug!(" {}", line);
128 }
129 let result = stream.write_all(request_string.as_bytes()).await;
130 res!(result);
131
132 loop {
133 trace!("Entered response read loop");
134 let result = HttpMessage::read::<
135 { constant::HTTP_DEFAULT_HEADER_CHUNK_SIZE },
136 { constant::HTTP_DEFAULT_BODY_CHUNK_SIZE },
137 _,
138 >(Pin::new(&mut stream), &Vec::new(), Some(false), None).await;
139
140 match result {
141 Ok((Some(response), remnant)) => {
142 trace!("Incoming Response:");
143 response.log(log_get_level!());
144 trace!("##### Remnant is {} bytes", remnant.len());
145 }
146 Ok((None, _)) => break,
147 Err(e) => return Err(e),
148 }
149 }
150 }
151 _ => (),
152 }
153
154 match filter {
155 "all" | "firefox" | "get" | "iter" => {
156 let result = new_stream(host, port).await;
157 let mut stream = res!(result);
158
159 //tokio::time::sleep(tokio::time::Duration::from_secs(10)).await;
160
161 let request_string = fmt!("GET / HTTP/1.1\r\n\
162 Host: {}\r\n\
163 User-Agent: Mozilla/5.0 (Windows NT 10.0; Win64; x64; rv:93.0) Gecko/20100101 Firefox/93.0\r\n\
164 Accept: text/html,application/xhtml+xml,application/xml;q=0.9,image/avif,image/webp,*/*;q=0.8\r\n\
165 Accept-Language: en-US,en;q=0.5\r\n\
166 Accept-Encoding: gzip, deflate, br\r\n\
167 Connection: keep-alive\r\n\
168 Upgrade-Insecure-Requests: 1\r\n\
169 Sec-Fetch-Dest: document\r\n\
170 Sec-Fetch-Mode: navigate\r\n\
171 Sec-Fetch-Site: none\r\n\
172 Sec-Fetch-User: ?1\r\n\
173 Cache-Control: max-age=0\r\n\r\n",
174 host);
175
176 debug!("Writing to stream now...");
177 for line in request_string.lines() {
178 debug!(" {}", line);
179 }
180 let result = stream.write_all(request_string.as_bytes()).await;
181 res!(result);
182
183 let mut reader: HttpMessageReader<
184 '_,
185 { constant::HTTP_DEFAULT_HEADER_CHUNK_SIZE },
186 { constant::HTTP_DEFAULT_BODY_CHUNK_SIZE },
187 _,
188 > = HttpMessageReader::new(Pin::new(&mut stream));
189
190 while let Some(result) = reader.next().await {
191 match result {
192 Ok(response) => {
193 trace!("Incoming Response:");
194 response.log(log_get_level!());
195 }
196 Err(e) => {
197 return Err(e);
198 }
199 }
200 }
201
202 }
203 _ => (),
204 }
205
206 match filter {
207 "all" | "websocket" | "text" => {
208 let result = new_stream(host, port).await;
209 let mut stream = res!(result);
210
211 // Create the websocket and connect to the server, using a standard HTTP request
212 // message.
213 let result = context::new_ws_no_db(
214 &mut stream,
215 WebSocketSinkHandler,
216 );
217 let mut ws = res!(result);
218 let (request, key) = res!(ws::connect_request(host, "/ws", None));
219 let result = ws.connect(request, Some(key)).await;
220 res!(result);
221
222 // Send a websocket text message..
223 let txt = fmt!("echo (str|\"Hello, WebSocket!\")");
224 let msg = WebSocketMessage::Text(txt.clone());
225 let result = ws.send(&msg).await;
226 res!(result);
227 test!("Sent text message: '{}'", txt);
228 test!("Sent text message as bytes: {:02x?}", txt.clone().as_bytes().to_vec());
229
230 // ..and wait for the echo.
231 match ws.read().await {
232 Ok(Some(WebSocketMessage::Text(txt2))) => {
233 test!("Received text message: '{}'", txt2);
234 req!(txt2, txt);
235 }
236 Err(e) => return Err(err!(e, "Error receiving message."; IO, Network, Read, Wire)),
237 Ok(None) => return Err(err!(
238 "The server has closed the connection unexpectedly.";
239 IO, Network, Read, Wire)),
240 Ok(Some(msg)) => return Err(err!(
241 "Expecting text message from server, received: {:?}", msg;
242 IO, Network, Unexpected, Read)),
243 }
244
245 // Send a small websocket binary message..
246 let byts = vec![0xaa, 0x12, 0x34];
247 let msg = WebSocketMessage::Binary(byts.clone());
248 let result = ws.send(&msg).await;
249 res!(result);
250 test!("Sent binary message: {:02x?}", byts);
251
252 // ..and wait for the echo.
253 match ws.read().await {
254 Ok(Some(WebSocketMessage::Binary(byts2))) => {
255 test!("Received binary message: {:02x?}", byts2);
256 req!(byts2, byts);
257 }
258 Err(e) => return Err(err!(e, "Error receiving message."; IO, Network, Read, Wire)),
259 Ok(None) => return Err(err!(
260 "The server has closed the connection unexpectedly.";
261 IO, Network, Read, Wire)),
262 Ok(Some(msg)) => return Err(err!(
263 "Expecting binary message from server, received: {:?}", msg;
264 IO, Network, Unexpected, Read)),
265 }
266
267 // Send a large websocket binary message..
268 let byts = vec![
269 0xef, 0xd3, 0x1f, 0x85, 0xfe, 0xd2, 0x36, 0xd4, 0xbd, 0x90,
270 0xae, 0x09, 0x31, 0x59, 0x3f, 0xe6, 0x96, 0x6a, 0x84, 0x11,
271 0xac, 0x0f, 0x57, 0x5a, 0xbf, 0x3f, 0xd3, 0x1b, 0x06, 0x7e,
272 0xcd, 0x86, 0x1a, 0x84, 0x29, 0xbd, 0x24,
273 ];
274 let msg = WebSocketMessage::Binary(byts.clone());
275 let result = ws.send(&msg).await;
276 res!(result);
277 test!("Sent binary message: {:02x?}", byts);
278
279 // ..and wait for the echo.
280 match ws.read().await {
281 Ok(Some(WebSocketMessage::Binary(byts2))) => {
282 test!("Received binary message: {:02x?}", byts2);
283 req!(byts2, byts);
284 }
285 Err(e) => return Err(err!(e, "Error receiving message."; IO, Network, Read, Wire)),
286 Ok(None) => return Err(err!(
287 "The server has closed the connection unexpectedly.";
288 IO, Network, Read, Wire)),
289 Ok(Some(msg)) => return Err(err!(
290 "Expecting binary message from server, received: {:?}", msg;
291 IO, Network, Unexpected, Read)),
292 }
293
294 // Send some data as text to store..
295 let txt = r#"insert (t2|[(str|a/b/c),{(str|name):(str|jane),(str|age):(u8|21)}])"#;
296 let msg = WebSocketMessage::Text(txt.to_string());
297 let result = ws.send(&msg).await;
298 res!(result);
299 test!("Sent text message: '{}'", txt);
300 test!("Sent text message as bytes: {:02x?}", txt.clone().as_bytes().to_vec());
301
302 // ..and wait for the reply.
303 match ws.read().await {
304 Ok(Some(WebSocketMessage::Text(txt2))) => {
305 test!("Received text message: '{}'", txt2);
306 }
307 Err(e) => return Err(err!(e, "Error receiving message."; IO, Network, Read, Wire)),
308 Ok(None) => return Err(err!(
309 "The server has closed the connection unexpectedly.";
310 IO, Network, Read, Wire)),
311 Ok(Some(msg)) => return Err(err!(
312 "Expecting binary message from server, received: {:?}", msg;
313 IO, Network, Unexpected, Read)),
314 }
315
316 thread::sleep(Duration::from_secs(1));
317
318 // Retrieve data as text..
319 let txt = fmt!("get_data (str|a/b/c)");
320 let msg = WebSocketMessage::Text(txt.clone());
321 let result = ws.send(&msg).await;
322 res!(result);
323 test!("Sent text message: '{}'", txt);
324 test!("Sent text message as bytes: {:02x?}", txt.clone().as_bytes().to_vec());
325
326 // ..and wait for the reply.
327 match ws.read().await {
328 Ok(Some(WebSocketMessage::Text(txt2))) => {
329 test!("Received text message: '{}'", txt2);
330 }
331 Err(e) => return Err(err!(e, "Error receiving message."; IO, Network, Read, Wire)),
332 Ok(None) => return Err(err!(
333 "The server has closed the connection unexpectedly.";
334 IO, Network, Read, Wire)),
335 Ok(Some(msg)) => return Err(err!(
336 "Expecting binary message from server, received: {:?}", msg;
337 IO, Network, Unexpected, Read)),
338 }
339
340 // Now listen to the websocket..
341 let listen_time_limit = tokio::time::Duration::from_secs(60);
342
343 let ws_syntax = res!(WebSocketSyntax::new(
344 "steel_ws",
345 &SemVer::new(0, 1, 0),
346 "Steel Websocket Test Client",
347 ));
348
349 tokio::time::timeout(listen_time_limit, async {
350 // Client side: no database is needed for the listen loop.
351 ws.listen(
352 None,
353 ws_syntax,
354 Some(30),
355 0,
356 &fmt!("client"),
357 ).await
358 }).await;
359
360 // And finally, close it.
361 test!("Closing websocket now.");
362 let result = ws.close(
363 Some(WebSocketStatusCode::NormalClosure),
364 Some(fmt!("Closing the connection")),
365 ).await;
366 res!(result);
367 }
368 _ => (),
369 }
370
371 Ok(())
372}
373