1pub mod connect;
7pub mod detect;
8pub mod error;
9pub mod forward;
10pub mod h2_connect;
11
12pub use connect::{
13 handle_connect, http_connect, validate_credentials, ConnectRequest, HttpConnectLimits,
14};
15pub use detect::HttpDetector;
16pub use error::HttpError;
17pub use forward::{
18 build_origin_request, copy_request_body, determine_request_body_kind, filter_hop_by_hop,
19 forward_request, forward_request_stream, forward_response, has_unsupported_expectation,
20 BodyCopyLimits, BodyCopyReport, ForwardRequest, ForwardResponse, ForwardResponseReport,
21 ForwardResult, RequestBodyKind,
22};
23pub use h2_connect::{
24 h2_connect_client, h2_connect_client_pooled, h2_connect_relay, handle_h2_connect,
25 H2ConnectError, H2PoolGuard, H2PoolKey, H2PoolRegistry, H2PoolStats, H2ProtocolMetrics,
26 H2StreamRead, H2StreamWrite, H2_POOL_REGISTRY, H2_PROTOCOL_METRICS,
27};
28
29#[cfg(test)]
30mod tests {
31 use super::*;
32 use eggress_core::detect::{DetectResult, ProtocolDetector};
33 use eggress_core::{BoxStream, TargetAddr, TargetHost};
34 use tokio::io::{AsyncReadExt, AsyncWriteExt};
35
36 #[test]
37 fn test_http_detector_identifies_http() {
38 let detector = HttpDetector;
39 assert_eq!(detector.id(), eggress_core::ProtocolId::Http);
40 assert_eq!(
41 detector.detect(b"GET / HTTP/1.1\r\n"),
42 DetectResult::Match { confidence: 100 }
43 );
44 }
45
46 #[test]
47 fn test_http_detector_identifies_connect() {
48 let detector = HttpDetector;
49 assert_eq!(
50 detector.detect(b"CONNECT example.com:443 HTTP/1.1\r\n"),
51 DetectResult::Match { confidence: 100 }
52 );
53 }
54
55 #[tokio::test]
58 async fn test_connect_to_echo_server() {
59 let (addr, jh) = eggress_testkit::start_echo_server().await;
60 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
61 let proxy_addr = listener.local_addr().unwrap();
62
63 let server_jh = tokio::spawn(async move {
64 let (stream, _) = listener.accept().await.unwrap();
65 let boxed: BoxStream = Box::new(stream);
66 let (request, stream) = connect::handle_connect(boxed, false, None).await.unwrap();
67 assert_eq!(
68 request.target,
69 TargetAddr {
70 host: TargetHost::Ip(addr.ip()),
71 port: addr.port(),
72 }
73 );
74
75 let target_stream = tokio::net::TcpStream::connect((addr.ip(), addr.port()))
77 .await
78 .unwrap();
79 let (mut cr, mut cw) = tokio::io::split(stream);
80 let (mut tr, mut tw) = tokio::io::split(target_stream);
81 tokio::spawn(async move {
82 let _ = tokio::io::copy(&mut cr, &mut tw).await;
83 });
84 tokio::spawn(async move {
85 let _ = tokio::io::copy(&mut tr, &mut cw).await;
86 });
87 });
88
89 let stream = tokio::net::TcpStream::connect(proxy_addr).await.unwrap();
90 let boxed: BoxStream = Box::new(stream);
91 let target = TargetAddr {
92 host: TargetHost::Ip(addr.ip()),
93 port: addr.port(),
94 };
95 let mut conn = connect::http_connect(boxed, &target, None, &Default::default())
96 .await
97 .unwrap();
98
99 conn.write_all(b"hello connect").await.unwrap();
101 conn.shutdown().await.unwrap();
102
103 let mut buf = [0u8; 13];
104 conn.read_exact(&mut buf).await.unwrap();
105 assert_eq!(&buf, b"hello connect");
106
107 server_jh.await.unwrap();
108 jh.abort();
109 }
110
111 #[tokio::test]
112 async fn test_connect_with_basic_auth() {
113 let (addr, jh) = eggress_testkit::start_echo_server().await;
114 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
115 let proxy_addr = listener.local_addr().unwrap();
116
117 let server_jh = tokio::spawn(async move {
118 let (stream, _) = listener.accept().await.unwrap();
119 let boxed: BoxStream = Box::new(stream);
120 let (request, stream) = connect::handle_connect(boxed, true, Some(("user", "pass")))
121 .await
122 .unwrap();
123 assert_eq!(request.proxy_auth, Some(("user".into(), "pass".into())));
124
125 let target_stream = tokio::net::TcpStream::connect((addr.ip(), addr.port()))
126 .await
127 .unwrap();
128 let (mut cr, mut cw) = tokio::io::split(stream);
129 let (mut tr, mut tw) = tokio::io::split(target_stream);
130 tokio::spawn(async move {
131 let _ = tokio::io::copy(&mut cr, &mut tw).await;
132 });
133 tokio::spawn(async move {
134 let _ = tokio::io::copy(&mut tr, &mut cw).await;
135 });
136 });
137
138 let stream = tokio::net::TcpStream::connect(proxy_addr).await.unwrap();
139 let boxed: BoxStream = Box::new(stream);
140 let target = TargetAddr {
141 host: TargetHost::Ip(addr.ip()),
142 port: addr.port(),
143 };
144 let mut conn =
145 connect::http_connect(boxed, &target, Some(("user", "pass")), &Default::default())
146 .await
147 .unwrap();
148
149 conn.write_all(b"auth test").await.unwrap();
150 conn.shutdown().await.unwrap();
151
152 let mut buf = [0u8; 9];
153 conn.read_exact(&mut buf).await.unwrap();
154 assert_eq!(&buf, b"auth test");
155
156 server_jh.await.unwrap();
157 jh.abort();
158 }
159
160 #[tokio::test]
161 async fn test_connect_auth_required() {
162 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
163 let proxy_addr = listener.local_addr().unwrap();
164
165 let server_jh = tokio::spawn(async move {
166 let (stream, _) = listener.accept().await.unwrap();
167 let boxed: BoxStream = Box::new(stream);
168 let result = connect::handle_connect(boxed, true, Some(("user", "pass"))).await;
169 assert!(result.is_err());
170 });
171
172 let stream = tokio::net::TcpStream::connect(proxy_addr).await.unwrap();
173 let boxed: BoxStream = Box::new(stream);
174 let target = TargetAddr {
175 host: TargetHost::Ip("127.0.0.1".parse().unwrap()),
176 port: 80,
177 };
178 let result = connect::http_connect(boxed, &target, None, &Default::default()).await;
179 assert!(result.is_err());
180 match result {
181 Err(HttpError::AuthRequired) => {}
182 _ => panic!("expected AuthRequired error"),
183 }
184
185 server_jh.await.unwrap();
186 }
187
188 #[tokio::test]
189 async fn test_connect_with_domain_target() {
190 let (addr, jh) = eggress_testkit::start_echo_server().await;
191 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
192 let proxy_addr = listener.local_addr().unwrap();
193
194 let server_jh = tokio::spawn(async move {
195 let (stream, _) = listener.accept().await.unwrap();
196 let boxed: BoxStream = Box::new(stream);
197 let (request, stream) = connect::handle_connect(boxed, false, None).await.unwrap();
198 assert_eq!(
199 request.target,
200 TargetAddr {
201 host: TargetHost::Domain("example.com".to_string()),
202 port: 443,
203 }
204 );
205
206 let target_stream = tokio::net::TcpStream::connect((addr.ip(), addr.port()))
207 .await
208 .unwrap();
209 let (mut cr, mut cw) = tokio::io::split(stream);
210 let (mut tr, mut tw) = tokio::io::split(target_stream);
211 tokio::spawn(async move {
212 let _ = tokio::io::copy(&mut cr, &mut tw).await;
213 });
214 tokio::spawn(async move {
215 let _ = tokio::io::copy(&mut tr, &mut cw).await;
216 });
217 });
218
219 let stream = tokio::net::TcpStream::connect(proxy_addr).await.unwrap();
220 let boxed: BoxStream = Box::new(stream);
221 let target = TargetAddr {
222 host: TargetHost::Domain("example.com".to_string()),
223 port: 443,
224 };
225 let mut conn = connect::http_connect(boxed, &target, None, &Default::default())
226 .await
227 .unwrap();
228
229 conn.write_all(b"domain test").await.unwrap();
230 conn.shutdown().await.unwrap();
231
232 let mut buf = [0u8; 11];
233 conn.read_exact(&mut buf).await.unwrap();
234 assert_eq!(&buf, b"domain test");
235
236 server_jh.await.unwrap();
237 jh.abort();
238 }
239
240 #[tokio::test]
241 async fn test_connect_with_ipv6_target() {
242 let (addr, jh) = eggress_testkit::start_echo_server().await;
243 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
244 let proxy_addr = listener.local_addr().unwrap();
245
246 let server_jh = tokio::spawn(async move {
247 let (stream, _) = listener.accept().await.unwrap();
248 let boxed: BoxStream = Box::new(stream);
249 let (request, stream) = connect::handle_connect(boxed, false, None).await.unwrap();
250 assert!(matches!(
251 request.target.host,
252 TargetHost::Ip(std::net::IpAddr::V6(_))
253 ));
254
255 let target_stream = tokio::net::TcpStream::connect((addr.ip(), addr.port()))
256 .await
257 .unwrap();
258 let (mut cr, mut cw) = tokio::io::split(stream);
259 let (mut tr, mut tw) = tokio::io::split(target_stream);
260 tokio::spawn(async move {
261 let _ = tokio::io::copy(&mut cr, &mut tw).await;
262 });
263 tokio::spawn(async move {
264 let _ = tokio::io::copy(&mut tr, &mut cw).await;
265 });
266 });
267
268 let stream = tokio::net::TcpStream::connect(proxy_addr).await.unwrap();
269 let boxed: BoxStream = Box::new(stream);
270 let target = TargetAddr {
271 host: TargetHost::Ip("::1".parse().unwrap()),
272 port: addr.port(),
273 };
274 let mut conn = connect::http_connect(boxed, &target, None, &Default::default())
275 .await
276 .unwrap();
277
278 conn.write_all(b"ipv6 test").await.unwrap();
279 conn.shutdown().await.unwrap();
280
281 let mut buf = [0u8; 9];
282 conn.read_exact(&mut buf).await.unwrap();
283 assert_eq!(&buf, b"ipv6 test");
284
285 server_jh.await.unwrap();
286 jh.abort();
287 }
288
289 #[tokio::test]
290 async fn test_connect_invalid_method_rejection() {
291 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
292 let proxy_addr = listener.local_addr().unwrap();
293
294 let server_jh = tokio::spawn(async move {
295 let (stream, _) = listener.accept().await.unwrap();
296 let boxed: BoxStream = Box::new(stream);
297 let result = connect::handle_connect(boxed, false, None).await;
298 assert!(result.is_err());
299 });
300
301 let mut stream = tokio::net::TcpStream::connect(proxy_addr).await.unwrap();
302 stream
304 .write_all(b"GET / HTTP/1.1\r\nHost: example.com\r\n\r\n")
305 .await
306 .unwrap();
307
308 server_jh.await.unwrap();
309 }
310
311 #[tokio::test]
312 async fn test_connect_header_size_limit() {
313 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
314 let proxy_addr = listener.local_addr().unwrap();
315
316 let server_jh = tokio::spawn(async move {
317 let (stream, _) = listener.accept().await.unwrap();
318 let boxed: BoxStream = Box::new(stream);
319 let result = connect::handle_connect(boxed, false, None).await;
320 assert!(result.is_err());
321 });
322
323 let mut stream = tokio::net::TcpStream::connect(proxy_addr).await.unwrap();
324 let mut request = b"CONNECT example.com:443 HTTP/1.1\r\n".to_vec();
326 for _ in 0..1000 {
327 request.extend_from_slice(b"X-Long-Header: ");
328 request.extend(&[b'A'; 100]);
329 request.extend_from_slice(b"\r\n");
330 }
331 request.extend_from_slice(b"\r\n");
332 stream.write_all(&request).await.unwrap();
333
334 server_jh.await.unwrap();
335 }
336
337 #[tokio::test]
340 async fn test_forward_get_absolute_form() {
341 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
342 let proxy_addr = listener.local_addr().unwrap();
343
344 let server_jh = tokio::spawn(async move {
345 let (stream, _) = listener.accept().await.unwrap();
346 let boxed: BoxStream = Box::new(stream);
347 let (request, _stream) = forward::forward_request(boxed).await.unwrap();
348 assert_eq!(request.method, "GET");
349 assert_eq!(request.path, "/index.html");
350 assert_eq!(
351 request.target,
352 TargetAddr {
353 host: TargetHost::Domain("example.com".to_string()),
354 port: 80,
355 }
356 );
357 });
358
359 let mut stream = tokio::net::TcpStream::connect(proxy_addr).await.unwrap();
360 stream
361 .write_all(b"GET http://example.com/index.html HTTP/1.1\r\nHost: example.com\r\n\r\n")
362 .await
363 .unwrap();
364
365 server_jh.await.unwrap();
366 }
367
368 #[tokio::test]
369 async fn test_forward_removes_proxy_auth() {
370 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
371 let proxy_addr = listener.local_addr().unwrap();
372
373 let server_jh = tokio::spawn(async move {
374 let (stream, _) = listener.accept().await.unwrap();
375 let boxed: BoxStream = Box::new(stream);
376 let (request, _stream) = forward::forward_request(boxed).await.unwrap();
377 assert!(!request
379 .headers
380 .iter()
381 .any(|(name, _)| name.eq_ignore_ascii_case("Proxy-Authorization")));
382 assert!(request
384 .headers
385 .iter()
386 .any(|(name, _)| name.eq_ignore_ascii_case("Authorization")));
387 });
388
389 let mut stream = tokio::net::TcpStream::connect(proxy_addr).await.unwrap();
390 stream
391 .write_all(
392 b"GET http://example.com/ HTTP/1.1\r\n\
393 Host: example.com\r\n\
394 Authorization: Bearer token123\r\n\
395 Proxy-Authorization: Basic dXNlcjpwYXNz\r\n\
396 \r\n",
397 )
398 .await
399 .unwrap();
400
401 server_jh.await.unwrap();
402 }
403
404 #[tokio::test]
405 async fn test_forward_origin_form_conversion() {
406 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
407 let proxy_addr = listener.local_addr().unwrap();
408
409 let server_jh = tokio::spawn(async move {
410 let (stream, _) = listener.accept().await.unwrap();
411 let boxed: BoxStream = Box::new(stream);
412 let (request, _stream) = forward::forward_request(boxed).await.unwrap();
413 assert_eq!(request.path, "/api/data");
415 assert_eq!(
416 request.target.host,
417 TargetHost::Domain("api.example.com".to_string())
418 );
419 assert_eq!(request.target.port, 8080);
420 });
421
422 let mut stream = tokio::net::TcpStream::connect(proxy_addr).await.unwrap();
423 stream
424 .write_all(
425 b"POST http://api.example.com:8080/api/data HTTP/1.1\r\n\
426 Host: api.example.com:8080\r\n\
427 Content-Length: 11\r\n\
428 \r\n\
429 hello world",
430 )
431 .await
432 .unwrap();
433
434 server_jh.await.unwrap();
435 }
436
437 #[tokio::test]
438 async fn test_forward_post_with_body() {
439 let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
440 let proxy_addr = listener.local_addr().unwrap();
441
442 let server_jh = tokio::spawn(async move {
443 let (stream, _) = listener.accept().await.unwrap();
444 let boxed: BoxStream = Box::new(stream);
445 let (request, _stream) = forward::forward_request(boxed).await.unwrap();
446 assert_eq!(request.method, "POST");
447 assert!(request.has_body);
448 assert_eq!(request.content_length, Some(11));
449 });
450
451 let mut stream = tokio::net::TcpStream::connect(proxy_addr).await.unwrap();
452 stream
453 .write_all(
454 b"POST http://example.com/api HTTP/1.1\r\n\
455 Host: example.com\r\n\
456 Content-Length: 11\r\n\
457 \r\n\
458 hello world",
459 )
460 .await
461 .unwrap();
462
463 server_jh.await.unwrap();
464 }
465}