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
| 1 | use oxedyne_fe2o3_steel::srv::{ |
| 2 | constant, |
| 3 | context, |
| 4 | id, |
| 5 | ws::syntax::WebSocketSyntax, |
| 6 | }; |
| 7 | |
| 8 | use oxedyne_fe2o3_core::{ |
| 9 | prelude::*, |
| 10 | }; |
| 11 | use oxedyne_fe2o3_crypto::enc::EncryptionScheme; |
| 12 | use oxedyne_fe2o3_hash::{ |
| 13 | csum::ChecksumScheme, |
| 14 | hash::HashScheme, |
| 15 | }; |
| 16 | use oxedyne_fe2o3_jdat::version::SemVer; |
| 17 | use 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 | }; |
| 32 | use oxedyne_fe2o3_o3db_sync::O3db; |
| 33 | |
| 34 | use 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; |
| 48 | use tokio::{ |
| 49 | self, |
| 50 | io::{ |
| 51 | AsyncWriteExt, |
| 52 | }, |
| 53 | net::TcpStream, |
| 54 | }; |
| 55 | use tokio_rustls::{ |
| 56 | client::TlsStream, |
| 57 | rustls::{ |
| 58 | self, |
| 59 | client::danger::ServerCertVerifier, |
| 60 | ClientConfig, |
| 61 | RootCertStore, |
| 62 | }, |
| 63 | TlsConnector, |
| 64 | }; |
| 65 | |
| 66 | fn 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 | |
| 79 | pub 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 | |
| 98 | pub 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 |