Skip to main content

eggress_protocol_http/
lib.rs

1//! HTTP proxy protocol implementation.
2//!
3//! This crate provides HTTP/1.1 proxy protocol handlers including
4//! CONNECT tunneling and ordinary HTTP forwarding.
5
6pub 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, H2ConnectError, H2PoolGuard,
25    H2PoolKey, H2PoolRegistry, H2PoolStats, H2ProtocolMetrics, H2StreamRead, H2StreamWrite,
26    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    // ===== Integration tests for HTTP CONNECT =====
56
57    #[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            // Connect to the target
76            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        // Verify stream works after CONNECT
100        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        // Send a GET request instead of CONNECT
303        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        // Send a CONNECT with very long headers
325        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    // ===== Integration tests for HTTP Forwarding =====
338
339    #[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            // Proxy-Authorization should be removed
378            assert!(!request
379                .headers
380                .iter()
381                .any(|(name, _)| name.eq_ignore_ascii_case("Proxy-Authorization")));
382            // Other headers should remain
383            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            // The path should be origin-form (just the path, not the full URI)
414            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}