Oregami
Repositories/oxedyne/fe2o3

oxedyne/fe2o3/fe2o3_steel/tests/ws_route.rs

19.9 KiB, 48 runs

created by r1870400018:20230, 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//! A `ws_route` in front of a real WebSocket server.
2//!
3//! The upstream in these tests is not a stub answering a fixed string: it is `fe2o3_net`'s own
4//! WebSocket machinery, handshaking from the request the relay forwarded and framing its replies.
5//! That is the point of the exercise -- a relay that satisfies a mock of itself has proved nothing,
6//! whereas a browser's decoder and this upstream agree on the same RFC.
7//!
8//! The relay is driven directly rather than through a TLS listener, because what is under test is
9//! the forwarding, and a certificate would only stand between the test and the bytes.
10//!
11//! [Written with AI entirely](https://need2know.ai/entirely-ai/code)\
12//! Anthropic Claude
13
14use oxedyne_fe2o3_core::prelude::*;
15use oxedyne_fe2o3_jdat::prelude::*;
16use oxedyne_fe2o3_net::{
17 http::{
18 fwd::ForwardedPolicy,
19 msg::HttpMessage,
20 },
21 ws::{
22 accept_key,
23 accept_response,
24 connect_request,
25 encode_message,
26 read_message,
27 WebSocketLimits,
28 WebSocketMessage,
29 },
30};
31use oxedyne_fe2o3_steel::srv::{
32 cfg::{
33 VhostConfig,
34 WsRoute,
35 },
36 wsproxy::tunnel_upgrade,
37};
38
39use std::{
40 net::SocketAddr,
41 pin::Pin,
42 sync::{
43 Arc,
44 Mutex,
45 },
46 time::Duration,
47};
48
49use tokio::{
50 io::{
51 AsyncReadExt,
52 AsyncWriteExt,
53 },
54 net::{
55 TcpListener,
56 TcpStream,
57 },
58};
59
60
61/// Frames a client-side message, as a browser would: masked.
62fn client_frame(msg: &WebSocketMessage) -> Outcome<Vec<u8>> {
63 encode_message(msg, true, 1024, 4096)
64}
65
66/// What the upstream saw of the request the relay forwarded to it.
67#[derive(Clone, Debug, Default)]
68struct UpstreamSaw {
69 path: String, // as the relay asked for it
70 ws_key: Option<String>, // the `Sec-WebSocket-Key` the client chose, if it survived
71 forwarded: Option<String>, // the `X-Forwarded-For` the relay added, if it did
72 host: Option<String>, // the `Host` header, which the relay owns
73}
74
75/// Start an echo WebSocket server on loopback.
76///
77/// It handshakes with `fe2o3_net`'s own machinery, echoes every text message back, and answers the
78/// text `bye` with a close frame before dropping the connection -- which is how the close-
79/// propagation test gets an upstream that goes away.
80async fn spawn_echo_upstream(saw: Arc<Mutex<UpstreamSaw>>) -> Outcome<SocketAddr> {
81 let listener = match TcpListener::bind("127.0.0.1:0").await {
82 Ok(l) => l,
83 Err(e) => return Err(err!(e, "Could not bind the upstream listener."; IO, Network, Init)),
84 };
85 let addr = match listener.local_addr() {
86 Ok(a) => a,
87 Err(e) => return Err(err!(e, "Upstream listener has no address."; IO, Network)),
88 };
89 tokio::spawn(async move {
90 let (mut stream, _peer) = match listener.accept().await {
91 Ok(p) => p,
92 Err(e) => {
93 error!(err!(e, "Upstream accept failed."; IO, Network));
94 return;
95 }
96 };
97 // Read the upgrade request the relay forwarded.
98 let read = HttpMessage::read::<1024, 1024, _>(
99 Pin::new(&mut stream),
100 &Vec::new(),
101 Some(true),
102 None,
103 ).await;
104 let req = match read {
105 Ok((Some(m), _)) => m,
106 Ok((None, _)) => {
107 error!(err!("Upstream saw no request at all."; IO, Network, Missing));
108 return;
109 }
110 Err(e) => {
111 error!(err!(e, "Upstream could not read the request."; IO, Network, Read));
112 return;
113 }
114 };
115 // Record what arrived, then answer it.
116 {
117 let mut guard = match saw.lock() {
118 Ok(g) => g,
119 Err(_) => return,
120 };
121 guard.path = match &req.header.headline {
122 oxedyne_fe2o3_net::http::header::HttpHeadline::Request { loc, .. } =>
123 loc.path.as_string().to_string(),
124 _ => String::new(),
125 };
126 guard.ws_key = header_of(&req, "Sec-WebSocket-Key");
127 guard.forwarded = header_of(&req, "X-Forwarded-For");
128 guard.host = header_of(&req, "Host");
129 }
130 let response = match accept_response(&req) {
131 Ok(r) => r,
132 Err(e) => {
133 error!(e);
134 return;
135 }
136 };
137 if let Err(e) = stream.write_all(response.as_bytes()).await {
138 error!(err!(e, "Upstream could not answer the upgrade."; IO, Network, Write));
139 return;
140 }
141 // Echo until told otherwise.
142 let mut buffer = Vec::new();
143 loop {
144 let msg = match read_message(&mut stream, &mut buffer, 1024, WebSocketLimits::default()).await {
145 Ok(Some(m)) => m,
146 Ok(None) => return,
147 Err(e) => {
148 error!(e);
149 return;
150 }
151 };
152 let reply = match msg {
153 WebSocketMessage::Text(txt) if txt == "bye" => {
154 let close = WebSocketMessage::Close(None, None);
155 match encode_message(&close, false, 1024, 4096) {
156 Ok(b) => {
157 let _ = stream.write_all(&b).await;
158 let _ = stream.shutdown().await;
159 return;
160 }
161 Err(e) => {
162 error!(e);
163 return;
164 }
165 }
166 }
167 WebSocketMessage::Text(txt) => WebSocketMessage::Text(txt),
168 WebSocketMessage::Binary(byts) => WebSocketMessage::Binary(byts),
169 WebSocketMessage::Close(_, _) => return,
170 other => {
171 error!(err!("Upstream got an unexpected {:?}.", other; Test, Unexpected));
172 return;
173 }
174 };
175 match encode_message(&reply, false, 1024, 4096) {
176 Ok(b) => {
177 if let Err(e) = stream.write_all(&b).await {
178 error!(err!(e, "Upstream could not echo."; IO, Network, Write));
179 return;
180 }
181 }
182 Err(e) => {
183 error!(e);
184 return;
185 }
186 }
187 }
188 });
189 Ok(addr)
190}
191
192fn header_of(msg: &HttpMessage, name: &str) -> Option<String> {
193 for (field_name, values) in msg.header.fields.iter() {
194 if fmt!("{}", field_name).eq_ignore_ascii_case(name) {
195 return values.first().map(|v| fmt!("{}", v));
196 }
197 }
198 None
199}
200
201async fn read_header_block(stream: &mut tokio::io::DuplexStream) -> Outcome<String> {
202 let mut accum: Vec<u8> = Vec::new();
203 let mut buf = [0u8; 256];
204 loop {
205 let n = match stream.read(&mut buf).await {
206 Ok(0) => return Err(err!(
207 "The relay closed before sending a response header block.";
208 IO, Network, Read, Missing)),
209 Ok(n) => n,
210 Err(e) => return Err(err!(e, "Reading the relayed response."; IO, Network, Read)),
211 };
212 accum.extend_from_slice(&buf[..n]);
213 if accum.windows(4).any(|w| w == b"\r\n\r\n") {
214 break;
215 }
216 if accum.len() > 65536 {
217 return Err(err!("Response header block never ended."; IO, Network, TooBig));
218 }
219 }
220 Ok(String::from_utf8_lossy(&accum).to_string())
221}
222
223/// The whole hop: a browser's handshake reaches a real WebSocket server through the relay, its
224/// accept value checks out, and text goes both ways.
225#[tokio::test]
226async fn test_ws_route_relays_handshake_and_bytes_00() -> Outcome<()> {
227 let saw = Arc::new(Mutex::new(UpstreamSaw::default()));
228 let upstream = res!(spawn_echo_upstream(saw.clone()).await);
229
230 let (mut browser, mut steel) = tokio::io::duplex(8192);
231 let (request, key) = res!(connect_request("oxegen.example", "/ws", None));
232 let src: SocketAddr = res!("203.0.113.9:51000".parse::<SocketAddr>(), Test);
233
234 let relay = tokio::spawn(async move {
235 tunnel_upgrade(
236 &mut steel,
237 &request,
238 "127.0.0.1",
239 upstream.port(),
240 "/ws",
241 src,
242 &ForwardedPolicy::none(),
243 "Test|ws",
244 ).await
245 });
246
247 // The 101 arrives, carrying the accept value derived from the key the browser chose. Nothing
248 // in the relay computes it: the upstream did.
249 let head = res!(read_header_block(&mut browser).await);
250 assert!(head.starts_with("HTTP/1.1 101 "),
251 "expected a 101, got: {}", head.lines().next().unwrap_or(""));
252 let expected = accept_key(&key);
253 assert!(head.contains(&fmt!("Sec-WebSocket-Accept: {}", expected)),
254 "the accept value must be the one derived from this handshake's key; got: {}", head);
255
256 // A masked client frame goes up, and the echo comes back.
257 let out = res!(client_frame(&WebSocketMessage::Text("oxe1:hello".to_string())));
258 res!(browser.write_all(&out).await, IO, Network, Write);
259 let mut buffer = Vec::new();
260 match res!(read_message(&mut browser, &mut buffer, 1024, WebSocketLimits::default()).await) {
261 Some(WebSocketMessage::Text(txt)) => assert_eq!(txt, "oxe1:hello"),
262 other => return Err(err!("Expected the echo, got {:?}.", other; Test, Mismatch)),
263 }
264
265 // And a second one, to show the tunnel is still open rather than one-shot.
266 let out = res!(client_frame(&WebSocketMessage::Text("second".to_string())));
267 res!(browser.write_all(&out).await, IO, Network, Write);
268 match res!(read_message(&mut browser, &mut buffer, 1024, WebSocketLimits::default()).await) {
269 Some(WebSocketMessage::Text(txt)) => assert_eq!(txt, "second"),
270 other => return Err(err!("Expected the second echo, got {:?}.", other; Test, Mismatch)),
271 }
272
273 // What the upstream saw: the client's key survived the hop (or the accept value above could
274 // not have matched), the path is the configured upstream one, and the relay named itself.
275 let seen = match saw.lock() {
276 Ok(g) => g.clone(),
277 Err(_) => return Err(err!("The upstream's record is poisoned."; Test, Poisoned)),
278 };
279 assert_eq!(seen.path, "/ws");
280 assert_eq!(seen.ws_key.as_deref(), Some(key.as_str()));
281 assert_eq!(seen.forwarded.as_deref(), Some("203.0.113.9:51000"));
282 assert_eq!(seen.host.as_deref(), Some("127.0.0.1"),
283 "the Host header belongs to this hop, not to the browser's original");
284
285 // Closing from the far end ends the relay.
286 let out = res!(client_frame(&WebSocketMessage::Text("bye".to_string())));
287 res!(browser.write_all(&out).await, IO, Network, Write);
288 match res!(read_message(&mut browser, &mut buffer, 1024, WebSocketLimits::default()).await) {
289 Some(WebSocketMessage::Close(_, _)) => (),
290 other => return Err(err!(
291 "Expected the upstream's close frame, got {:?}.", other; Test, Mismatch)),
292 }
293 // And with the upstream gone, the client's side ends too.
294 match res!(read_message(&mut browser, &mut buffer, 1024, WebSocketLimits::default()).await) {
295 None => (),
296 other => return Err(err!(
297 "Expected the relay to close after the upstream did, got {:?}.", other;
298 Test, Mismatch)),
299 }
300 match tokio::time::timeout(Duration::from_secs(5), relay).await {
301 Ok(Ok(outcome)) => res!(outcome),
302 Ok(Err(e)) => return Err(err!(e, "The relay task panicked."; Test)),
303 Err(_) => return Err(err!(
304 "The relay did not return after both ends closed."; Test, Timeout)),
305 }
306 Ok(())
307}
308
309/// A browser that goes away takes the tunnel with it, rather than leaving the relay holding a
310/// socket nobody is reading.
311#[tokio::test]
312async fn test_ws_route_client_close_ends_the_tunnel_00() -> Outcome<()> {
313 let saw = Arc::new(Mutex::new(UpstreamSaw::default()));
314 let upstream = res!(spawn_echo_upstream(saw).await);
315
316 let (mut browser, mut steel) = tokio::io::duplex(8192);
317 let (request, _key) = res!(connect_request("oxegen.example", "/ws", None));
318 let src: SocketAddr = res!("203.0.113.9:51001".parse::<SocketAddr>(), Test);
319
320 let relay = tokio::spawn(async move {
321 tunnel_upgrade(
322 &mut steel, &request, "127.0.0.1", upstream.port(), "/ws", src,
323 &ForwardedPolicy::none(), "Test|ws",
324 ).await
325 });
326
327 let head = res!(read_header_block(&mut browser).await);
328 assert!(head.starts_with("HTTP/1.1 101 "), "expected a 101, got: {}", head);
329
330 // Drop the browser end.
331 drop(browser);
332
333 match tokio::time::timeout(Duration::from_secs(5), relay).await {
334 Ok(Ok(outcome)) => res!(outcome),
335 Ok(Err(e)) => return Err(err!(e, "The relay task panicked."; Test)),
336 Err(_) => return Err(err!(
337 "The relay did not return after the client went away."; Test, Timeout)),
338 }
339 Ok(())
340}
341
342/// An upstream that refuses the upgrade has its refusal relayed verbatim. A client is entitled to
343/// read the answer the server gave, not one this hop invented.
344#[tokio::test]
345async fn test_ws_route_relays_a_refusal_00() -> Outcome<()> {
346 let listener = match TcpListener::bind("127.0.0.1:0").await {
347 Ok(l) => l,
348 Err(e) => return Err(err!(e, "Could not bind the grumpy upstream."; IO, Network, Init)),
349 };
350 let addr = res!(listener.local_addr(), IO, Network);
351 tokio::spawn(async move {
352 if let Ok((mut stream, _)) = listener.accept().await {
353 let mut buf = [0u8; 2048];
354 let _ = stream.read(&mut buf).await;
355 let _ = stream.write_all(
356 b"HTTP/1.1 403 Forbidden\r\nContent-Length: 0\r\n\r\n").await;
357 let _ = stream.shutdown().await;
358 }
359 });
360
361 let (mut browser, mut steel) = tokio::io::duplex(8192);
362 let (request, _key) = res!(connect_request("oxegen.example", "/ws", None));
363 let src: SocketAddr = res!("203.0.113.9:51002".parse::<SocketAddr>(), Test);
364 let relay = tokio::spawn(async move {
365 tunnel_upgrade(
366 &mut steel, &request, "127.0.0.1", addr.port(), "/ws", src,
367 &ForwardedPolicy::none(), "Test|ws",
368 ).await
369 });
370
371 let head = res!(read_header_block(&mut browser).await);
372 assert!(head.starts_with("HTTP/1.1 403 "),
373 "the upstream's refusal must reach the client unchanged; got: {}", head);
374 match tokio::time::timeout(Duration::from_secs(5), relay).await {
375 Ok(Ok(outcome)) => res!(outcome),
376 Ok(Err(e)) => return Err(err!(e, "The relay task panicked."; Test)),
377 Err(_) => return Err(err!("The relay did not return."; Test, Timeout)),
378 }
379 Ok(())
380}
381
382/// An upstream that is not listening is an error, not a hang and not a silent 101.
383#[tokio::test]
384async fn test_ws_route_errors_when_upstream_is_absent_00() -> Outcome<()> {
385 // Bind and drop, so the port is one nothing is listening on.
386 let listener = match TcpListener::bind("127.0.0.1:0").await {
387 Ok(l) => l,
388 Err(e) => return Err(err!(e, "Could not bind to find a free port."; IO, Network, Init)),
389 };
390 let addr = res!(listener.local_addr(), IO, Network);
391 drop(listener);
392
393 let (_browser, mut steel) = tokio::io::duplex(8192);
394 let (request, _key) = res!(connect_request("oxegen.example", "/ws", None));
395 let src: SocketAddr = res!("203.0.113.9:51003".parse::<SocketAddr>(), Test);
396 let result = tunnel_upgrade(
397 &mut steel, &request, "127.0.0.1", addr.port(), "/ws", src,
398 &ForwardedPolicy::none(), "Test|ws",
399 ).await;
400 assert!(result.is_err(), "an absent upstream must be reported");
401 Ok(())
402}
403
404/// The configuration a `ws_route` is written in: the URL's parts, and the defaults.
405#[test]
406fn test_ws_route_parses_its_upstream_00() -> Outcome<()> {
407 let route = res!(WsRoute::from_datmap(&mapdat!{
408 "path" => "/ws",
409 "upstream" => "ws://127.0.0.1:9080/ws",
410 }.get_map().unwrap_or_default()));
411 assert_eq!(route.path, "/ws");
412 assert_eq!(route.upstream_host, "127.0.0.1");
413 assert_eq!(route.upstream_port, 9080);
414 assert_eq!(route.upstream_path, "/ws");
415 assert!(route.matches("/ws"));
416 assert!(!route.matches("/ws/"), "the match is exact, not a prefix");
417 assert!(!route.matches("/wsx"));
418
419 // No port and no path: the HTTP default port, and the root.
420 let route = res!(WsRoute::from_datmap(&mapdat!{
421 "path" => "/socket",
422 "upstream" => "ws://gateway.internal",
423 }.get_map().unwrap_or_default()));
424 assert_eq!(route.upstream_host, "gateway.internal");
425 assert_eq!(route.upstream_port, 80);
426 assert_eq!(route.upstream_path, "/");
427
428 // The local path and the upstream path need not agree.
429 let route = res!(WsRoute::from_datmap(&mapdat!{
430 "path" => "/ws",
431 "upstream" => "ws://127.0.0.1:9080/gateway/v0",
432 }.get_map().unwrap_or_default()));
433 assert_eq!(route.upstream_path, "/gateway/v0");
434 Ok(())
435}
436
437/// What a `ws_route` refuses to be configured as. Each of these would otherwise become a surprise
438/// at runtime, on a connection an operator is watching.
439#[test]
440fn test_ws_route_refuses_bad_configuration_00() -> Outcome<()> {
441 let missing_upstream = mapdat!{ "path" => "/ws" };
442 assert!(WsRoute::from_datmap(&missing_upstream.get_map().unwrap_or_default()).is_err(),
443 "a route with no upstream is not a route");
444
445 let missing_path = mapdat!{ "upstream" => "ws://127.0.0.1:9080/ws" };
446 assert!(WsRoute::from_datmap(&missing_path.get_map().unwrap_or_default()).is_err(),
447 "a route with no local path claims nothing");
448
449 let tls = mapdat!{ "path" => "/ws", "upstream" => "wss://example.com/ws" };
450 assert!(WsRoute::from_datmap(&tls.get_map().unwrap_or_default()).is_err(),
451 "wss:// must be refused rather than quietly spoken as plaintext");
452
453 let http = mapdat!{ "path" => "/ws", "upstream" => "http://127.0.0.1:9080/ws" };
454 assert!(WsRoute::from_datmap(&http.get_map().unwrap_or_default()).is_err(),
455 "a scheme other than ws:// must be refused");
456
457 let bad_port = mapdat!{ "path" => "/ws", "upstream" => "ws://127.0.0.1:notaport/ws" };
458 assert!(WsRoute::from_datmap(&bad_port.get_map().unwrap_or_default()).is_err(),
459 "a port that is not a number must be refused");
460
461 let no_host = mapdat!{ "path" => "/ws", "upstream" => "ws:///ws" };
462 assert!(WsRoute::from_datmap(&no_host.get_map().unwrap_or_default()).is_err(),
463 "an upstream with no host must be refused");
464 Ok(())
465}
466
467/// A vhost written before `ws_routes` existed must still parse, and must have none. This is the
468/// property that keeps every deployed config working: a new field that is required breaks them all.
469#[test]
470fn test_vhost_without_ws_routes_still_parses_00() -> Outcome<()> {
471 let vhost = mapdat!{
472 "hostnames" => listdat!["example.com"],
473 "public_dir_rel" => "./www",
474 };
475 let cfg = res!(VhostConfig::from_datmap(&vhost.get_map().unwrap_or_default()));
476 assert!(cfg.ws_routes.is_empty(), "a config that names no ws_routes has none");
477
478 let with_routes = mapdat!{
479 "hostnames" => listdat!["example.com"],
480 "public_dir_rel" => "./www",
481 "ws_routes" => listdat![mapdat!{
482 "path" => "/ws",
483 "upstream" => "ws://127.0.0.1:9080/ws",
484 }],
485 };
486 let cfg = res!(VhostConfig::from_datmap(&with_routes.get_map().unwrap_or_default()));
487 assert_eq!(cfg.ws_routes.len(), 1);
488 assert_eq!(cfg.ws_routes[0].upstream_port, 9080);
489
490 // A malformed entry is a start-up failure, not a route silently dropped.
491 let broken = mapdat!{
492 "hostnames" => listdat!["example.com"],
493 "ws_routes" => listdat!["/ws"],
494 };
495 assert!(VhostConfig::from_datmap(&broken.get_map().unwrap_or_default()).is_err(),
496 "a ws_routes entry that is not a map must fail the parse");
497 Ok(())
498}
499
500/// The relay never has to be told the client's address twice: `TcpStream` is only mentioned here so
501/// the test file's imports match what a caller uses.
502#[allow(dead_code)]
503fn _unused(_: TcpStream) {}