use std::sync::atomic::{AtomicUsize, Ordering};
use std::{cell::RefCell, collections::HashMap, io, io::Read, io::Write, net, rc::Rc, sync::Arc};
use coo_kie::Cookie;
use flate2::{Compression, read::GzDecoder, write::GzEncoder, write::ZlibEncoder};
use rand::Rng;
use ntex::client::{Client, ClientConfig, error::ClientError};
use ntex::http::test::server as test_server;
use ntex::http::{HttpMessage, HttpService, Method, Response, header};
use ntex::io::IoConfig;
use ntex::service::{cfg::SharedCfg, fn_layer, service};
use ntex::web::middleware::Compress;
use ntex::web::{self, App, BodyEncoding, HttpRequest, HttpResponse, WebError, test};
use ntex::{client, time::Millis, time::Seconds, time::sleep, util::Bytes};
const STR: &str = "Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World \
Hello World Hello World Hello World Hello World Hello World";
#[ntex::test]
async fn test_simple() {
let srv = test::server(async |_| {
App::new().service(web::resource("/").route(web::to(async || HttpResponse::Ok().body(STR))))
});
let request = srv.get("/").header("x-test", "111").send();
let response = request.await.unwrap();
assert!(response.status().is_success());
let bytes = response.body().await.unwrap();
assert_eq!(bytes, Bytes::from_static(STR.as_ref()));
let response = srv.post("/").timeout(Seconds(30)).send().await.unwrap();
assert!(response.status().is_success());
let bytes = response.body().await.unwrap();
assert_eq!(bytes, Bytes::from_static(STR.as_ref()));
}
#[ntex::test]
async fn test_json() {
let srv = test::server(async |_| {
App::new().service(web::resource("/").route(web::to(
async |_: web::types::Json<String>| HttpResponse::Ok(),
)))
});
let response = srv
.get("/")
.header("x-test", "111")
.send_json(&"TEST".to_string())
.await
.unwrap();
assert!(response.status().is_success());
}
#[ntex::test]
async fn test_form() {
let srv = test::server(async |_| {
App::new().service(web::resource("/").route(web::to(
async |_: web::types::Form<HashMap<String, String>>| HttpResponse::Ok(),
)))
});
let mut data = HashMap::new();
let _ = data.insert("key".to_string(), "TEST".to_string());
let request = srv.get("/").header("x-test", "111").send_form(&data);
let response = request.await.unwrap();
assert!(response.status().is_success());
}
#[ntex::test]
async fn test_timeout() {
let srv = test::server(async |_| {
App::new().service(web::resource("/").route(web::to(async || {
sleep(Millis(5000)).await;
HttpResponse::Ok().body(STR)
})))
});
let client = Client::with_config(
SharedCfg::new("SVC")
.add(IoConfig::new().set_connect_timeout(2500))
.add(ClientConfig::new().set_response_timeout(Seconds(3))),
);
let err = client.get(srv.url("/")).send().await.err().unwrap();
assert!(matches!(err.into_error(), ClientError::Timeout));
}
#[ntex::test]
async fn test_timeout_override() {
let srv = test::server(async |_| {
App::new().service(web::resource("/").route(web::to(async || {
sleep(Millis(2000)).await;
HttpResponse::Ok().body(STR)
})))
});
let client = Client::with_config(ClientConfig::new().set_response_timeout(Seconds(50)));
let err = client
.get(srv.url("/"))
.timeout(Seconds(1))
.send()
.await
.err();
assert!(matches!(err.unwrap().into_error(), ClientError::Timeout));
}
#[ntex::test]
async fn test_connection_reuse() {
let num = Arc::new(AtomicUsize::new(0));
let num2 = num.clone();
let srv = test_server(async move |_| {
let num2 = num2.clone();
service(async move |io| {
num2.fetch_add(1, Ordering::Relaxed);
Ok(io)
})
.and_then(HttpService::new(
App::new().service(web::resource("/").route(web::to(async || HttpResponse::Ok()))),
))
});
let client = Client::with_config(ClientConfig::new().set_response_timeout(Seconds(30)));
let request = client.get(srv.url("/")).send();
let response = request.await.unwrap();
assert!(response.status().is_success());
let req = client.post(srv.url("/"));
let response = req.send().await.unwrap();
assert!(response.status().is_success());
assert_eq!(num.load(Ordering::Relaxed), 1);
}
#[ntex::test]
async fn test_connection_close() {
let srv = test_server(async move |_| {
HttpService::new(async |_| Ok::<_, io::Error>(Response::Ok().body(STR)))
});
let response = srv
.request(Method::GET, "/")
.force_close()
.send()
.await
.unwrap();
assert!(response.status().is_success());
}
#[ntex::test]
async fn test_connection_force_close() {
let num = Arc::new(AtomicUsize::new(0));
let num2 = num.clone();
let srv = test_server(async move |_| {
let num2 = num2.clone();
service(async move |io| {
num2.fetch_add(1, Ordering::Relaxed);
Ok(io)
})
.and_then(HttpService::new(
App::new().service(web::resource("/").route(web::to(async || HttpResponse::Ok()))),
))
});
let client = Client::with_config(ClientConfig::new().set_response_timeout(Seconds(30)));
let request = client.get(srv.url("/")).force_close().send();
let response = request.await.unwrap();
assert!(response.status().is_success());
let client = Client::with_config(ClientConfig::new().set_response_timeout(Seconds(30)));
let req = client.post(srv.url("/")).force_close();
let response = req.send().await.unwrap();
assert!(response.status().is_success());
assert_eq!(num.load(Ordering::Relaxed), 2);
}
#[ntex::test]
async fn test_connection_server_close() {
let num = Arc::new(AtomicUsize::new(0));
let num2 = num.clone();
let srv = test_server(async move |_| {
let num2 = num2.clone();
service(async move |io| {
num2.fetch_add(1, Ordering::Relaxed);
Ok(io)
})
.and_then(HttpService::new(App::new().service(
web::resource("/").route(web::to(async || HttpResponse::Ok().force_close().finish())),
)))
});
let client = Client::with_config(ClientConfig::new().set_response_timeout(Seconds(30)));
let request = client.get(srv.url("/")).send();
let response = request.await.unwrap();
assert!(response.status().is_success());
let req = client.post(srv.url("/"));
let response = req.send().await.unwrap();
assert!(response.status().is_success());
assert_eq!(num.load(Ordering::Relaxed), 2);
}
#[ntex::test]
async fn test_connection_wait_queue() {
let num = Arc::new(AtomicUsize::new(0));
let num2 = num.clone();
let srv = test_server(async move |_| {
let num2 = num2.clone();
service(async move |io| {
num2.fetch_add(1, Ordering::Relaxed);
Ok(io)
})
.and_then(HttpService::new(App::new().service(
web::resource("/").route(web::to(async || HttpResponse::Ok().body(STR))),
)))
});
let client = Client::with_config(
ClientConfig::new()
.set_response_timeout(Seconds(30))
.set_limit(1),
);
let request = client.get(srv.url("/")).send();
let response = request.await.unwrap();
assert!(response.status().is_success());
let req2 = client.post(srv.url("/"));
let req2_fut = req2.send();
let bytes = response.body().await.unwrap();
assert_eq!(bytes, Bytes::from_static(STR.as_ref()));
let response = req2_fut.await.unwrap();
assert!(response.status().is_success());
assert_eq!(num.load(Ordering::Relaxed), 1);
}
#[ntex::test]
async fn test_connection_wait_queue_force_close() {
let num = Arc::new(AtomicUsize::new(0));
let num2 = num.clone();
let srv = test_server(async move |_| {
let num2 = num2.clone();
service(async move |io| {
num2.fetch_add(1, Ordering::Relaxed);
Ok(io)
})
.and_then(HttpService::new(App::new().service(
web::resource("/").route(web::to(async || HttpResponse::Ok().force_close().body(STR))),
)))
});
let client = Client::with_config(
ClientConfig::new()
.set_limit(1)
.set_response_timeout(Seconds(30)),
);
let request = client.get(srv.url("/")).send();
let response = request.await.unwrap();
assert!(response.status().is_success());
let req2 = client.post(srv.url("/"));
let req2_fut = req2.send();
let bytes = response.body().await.unwrap();
assert_eq!(bytes, Bytes::from_static(STR.as_ref()));
let response = req2_fut.await.unwrap();
assert!(response.status().is_success());
assert_eq!(num.load(Ordering::Relaxed), 2);
}
#[ntex::test]
async fn test_with_query_parameter() {
let srv = test::server(async |_| {
App::new().service(web::resource("/").to(async move |req: HttpRequest| {
if req.query_string().contains("qp") {
HttpResponse::Ok()
} else {
HttpResponse::BadRequest()
}
}))
});
let res = srv.get("/?qp=5").send().await.unwrap();
assert!(res.status().is_success());
let request = srv.request(Method::GET, srv.url("/?qp=5"));
let response = request.send().await.unwrap();
assert!(response.status().is_success());
}
#[ntex::test]
async fn test_no_decompress() {
let srv = test::server(async |_| {
App::new()
.middleware(Compress::default())
.service(web::resource("/").route(web::to(async || {
let mut res = HttpResponse::Ok().body(STR);
res.encoding(header::ContentEncoding::Gzip);
res
})))
});
let client = Client::builder().build(ClientConfig::new().set_response_timeout(Seconds(30)));
let res = client
.get(srv.url("/"))
.no_decompress()
.send()
.await
.unwrap();
assert!(res.status().is_success());
let bytes = res.body().await.unwrap();
let mut e = GzDecoder::new(&bytes[..]);
let mut dec = Vec::new();
e.read_to_end(&mut dec).unwrap();
assert_eq!(Bytes::from(dec), Bytes::from_static(STR.as_ref()));
let res = client
.post(srv.url("/"))
.no_decompress()
.send()
.await
.unwrap();
assert!(res.status().is_success());
let bytes = res.body().await.unwrap();
let mut e = GzDecoder::new(&bytes[..]);
let mut dec = Vec::new();
e.read_to_end(&mut dec).unwrap();
assert_eq!(Bytes::from(dec), Bytes::from_static(STR.as_ref()));
}
#[ntex::test]
async fn test_client_gzip_encoding() {
let srv = test::server(async |_| {
App::new().service(web::resource("/").route(web::to(async || {
let mut e = GzEncoder::new(Vec::new(), Compression::default());
e.write_all(STR.as_ref()).unwrap();
let data = e.finish().unwrap();
HttpResponse::Ok()
.header("content-encoding", "gzip")
.body(data)
})))
});
let response = srv.post("/").send().await.unwrap();
assert!(response.status().is_success());
let bytes = response.body().await.unwrap();
assert_eq!(bytes, Bytes::from_static(STR.as_ref()));
}
#[ntex::test]
async fn test_client_gzip_encoding_large() {
let srv = test::server(async |_| {
App::new().service(web::resource("/").route(web::to(async || {
let mut e = GzEncoder::new(Vec::new(), Compression::default());
e.write_all(STR.repeat(10).as_ref()).unwrap();
let data = e.finish().unwrap();
HttpResponse::Ok()
.header("content-encoding", "gzip")
.body(data)
})))
});
let response = srv.post("/").send().await.unwrap();
assert!(response.status().is_success());
let bytes = response.body().await.unwrap();
assert_eq!(bytes, Bytes::from(STR.repeat(10)));
}
#[ntex::test]
async fn test_client_gzip_encoding_large_random() {
let data = rand::rng()
.sample_iter(&rand::distr::Alphanumeric)
.take(1_048_500)
.map(char::from)
.collect::<String>();
let srv = test::server_with(
test::config().server_cfg(
SharedCfg::new("SRV").add(
web::WebAppConfig::new()
.set_state(web::types::PayloadConfig::default().limit(1_048_576)),
),
),
async |_| {
App::new().service(web::resource("/").route(web::to(async move |data: Bytes| {
let mut e = GzEncoder::new(Vec::new(), Compression::default());
e.write_all(&data).unwrap();
let data = e.finish().unwrap();
HttpResponse::Ok()
.header("content-encoding", "gzip")
.body(data)
})))
},
);
let response = srv.post("/").send_body(data.clone()).await.unwrap();
assert!(response.status().is_success());
let bytes = response.body().limit(1_048_576).await.unwrap();
assert_eq!(bytes, Bytes::from(data));
}
#[ntex::test]
async fn test_client_deflate_encoding() {
let srv = test::server(async |_| {
App::new().service(web::resource("/").route(web::to(async move |data: Bytes| {
let mut e = ZlibEncoder::new(Vec::new(), flate2::Compression::fast());
e.write_all(&data).unwrap();
let data = e.finish().unwrap();
HttpResponse::Ok()
.header("content-encoding", "deflate")
.body(data)
})))
});
let response = srv.post("/").send_body(STR).await.unwrap();
assert!(response.status().is_success());
let bytes = response.body().await.unwrap();
assert_eq!(bytes, Bytes::from_static(STR.as_ref()));
}
#[ntex::test]
async fn test_client_deflate_encoding_large_random() {
let data = rand::rng()
.sample_iter(&rand::distr::Alphanumeric)
.take(70_000)
.map(char::from)
.collect::<String>();
let srv = test::server(async |_| {
App::new().service(web::resource("/").route(web::to(async move |data: Bytes| {
let mut e = ZlibEncoder::new(Vec::new(), flate2::Compression::fast());
e.write_all(&data).unwrap();
let data = e.finish().unwrap();
HttpResponse::Ok()
.header("content-encoding", "deflate")
.body(data)
})))
});
let response = srv.post("/").send_body(data.clone()).await.unwrap();
assert!(response.status().is_success());
let bytes = response.body().await.unwrap();
assert_eq!(bytes, Bytes::from(data));
}
#[ntex::test]
async fn test_client_cookie_handling() {
use std::io::{Error as IoError, ErrorKind};
let cookie1 = Cookie::build(("cookie1", "value1"));
let cookie2 = Cookie::build(("cookie2", "value2"))
.domain("www.example.org")
.path("/")
.secure(true)
.http_only(true);
let cookie1b = cookie1.clone();
let cookie2b = cookie2.clone();
let srv = test::server(async move |_| {
let cookie1 = cookie1b.clone();
let cookie2 = cookie2b.clone();
App::new().route(
"/",
web::to(web::dev::__assert_handler1(
async move |req: HttpRequest| {
let cookie1 = cookie1.clone();
let cookie2 = cookie2.clone();
let res: Result<(), WebError> = req
.cookie("cookie1")
.ok_or(())
.and_then(|c1| if c1.value() == "value1" { Ok(()) } else { Err(()) })
.and_then(|()| req.cookie("cookie2").ok_or(()))
.and_then(|c2| if c2.value() == "value2" { Ok(()) } else { Err(()) })
.map_err(|_| WebError::new(IoError::from(ErrorKind::NotFound)));
res?;
Ok::<_, WebError>(HttpResponse::Ok().cookie(cookie1).cookie(cookie2).finish())
},
)),
)
});
let request = srv.get("/").cookie(cookie1.clone()).cookie(cookie2.clone());
let response = request.send().await.unwrap();
assert!(response.status().is_success());
let c1 = response.cookie("cookie1").expect("Missing cookie1");
assert_eq!(c1, cookie1);
let c2 = response.cookie("cookie2").expect("Missing cookie2");
assert_eq!(c2, cookie2);
}
#[ntex::test]
async fn test_client_timeout() {
let srv = test_server(async move |_| {
HttpService::new(async |_| {
sleep(Seconds(10)).await;
Ok::<_, io::Error>(Response::Ok().body(STR))
})
})
.set_client_timeout(Seconds(1), Millis(30_000));
let err = srv
.request(Method::GET, "/")
.force_close()
.send()
.await
.err()
.unwrap();
assert!(matches!(err.into_error(), ClientError::Timeout));
}
#[ntex::test]
async fn client_read_until_eof() {
let addr = ntex::server::TestServer::unused_addr();
std::thread::spawn(move || {
let lst = std::net::TcpListener::bind(addr).unwrap();
for stream in lst.incoming() {
if let Ok(mut stream) = stream {
let mut b = [0; 1000];
log::debug!("Reading request");
let res = stream.read(&mut b).unwrap();
log::debug!("Read {res:?}");
const S: &[u8] = b"HTTP/1.0 200 OK\r\nconnection: close\r\n\r\nwelcome!";
let res = stream.write_all(S);
log::debug!("Sent {res:?} cnt:{}", S.len());
let _ = stream.shutdown(net::Shutdown::Both);
} else {
break;
}
}
});
sleep(Millis(250)).await;
let req = Client::builder()
.build(ClientConfig::new().set_response_timeout(Seconds(30)))
.get(format!("http://{addr}/").as_str());
let response = req.send().await.unwrap();
assert!(response.status().is_success());
let bytes = response.body().await.unwrap();
assert_eq!(bytes, Bytes::from_static(b"welcome!"));
}
#[ntex::test]
async fn client_basic_auth() {
let srv = test::server(async |_| {
App::new().route(
"/",
web::to(async move |req: HttpRequest| {
if req
.headers()
.get(header::AUTHORIZATION)
.unwrap()
.to_str()
.unwrap()
== "Basic dXNlcm5hbWU6cGFzc3dvcmQ="
{
HttpResponse::Ok()
} else {
HttpResponse::BadRequest()
}
}),
)
});
let request = srv.get("/").basic_auth("username", Some("password"));
let response = request.send().await.unwrap();
assert!(response.status().is_success());
}
#[ntex::test]
async fn client_bearer_auth() {
let srv = test::server(async |_| {
App::new().route(
"/",
web::to(async move |req: HttpRequest| {
if req
.headers()
.get(header::AUTHORIZATION)
.unwrap()
.to_str()
.unwrap()
== "Bearer someS3cr3tAutht0k3n"
{
HttpResponse::Ok()
} else {
HttpResponse::BadRequest()
}
}),
)
});
let request = srv.get("/").bearer_auth("someS3cr3tAutht0k3n");
let response = request.send().await.unwrap();
assert!(response.status().is_success());
}
#[ntex::test]
async fn middleware() {
let srv = test::server(async |_| {
App::new().service(web::resource("/").route(web::to(async || HttpResponse::Ok().body(STR))))
});
let data = Rc::new(RefCell::new(Vec::new()));
let data2 = data.clone();
let client = Client::builder()
.middleware(fn_layer(
async move |mut req: client::ServiceRequest, svc| {
assert!(req.headers().is_empty());
assert!(req.address().is_none());
let s = format!("{:?}", req.head().uri);
data2.borrow_mut().push(format!("1 -- {}", s));
let result = svc.call(req).await;
data2.borrow_mut().push(format!("2 -- {}", s));
result
},
))
.build(SharedCfg::default());
assert!(client.ready().await.is_ok());
let request = client.get(srv.url("/")).header("x-test", "111").send();
let response = request.await.unwrap();
assert!(response.status().is_success());
let bytes = response.body().await.unwrap();
assert_eq!(bytes, Bytes::from_static(STR.as_ref()));
let host = srv.url("/");
assert_eq!(
&data.borrow()[..],
&[format!("1 -- {}", host), format!("2 -- {}", host)]
);
}
#[ntex::test]
async fn test_h1_v2() {
let srv = test_server(async move |_| {
HttpService::new(async |_| Ok::<_, io::Error>(Response::Ok().body(STR)))
});
let response = srv.request(Method::GET, "/").send().await.unwrap();
assert!(response.status().is_success());
let request = srv.request(Method::GET, "/").header("x-test", "111").send();
let response = request.await.unwrap();
assert!(response.status().is_success());
let bytes = response.body().await.unwrap();
assert_eq!(bytes, Bytes::from_static(STR.as_ref()));
let response = srv.request(Method::POST, "/").send().await.unwrap();
assert!(response.status().is_success());
let bytes = response.body().await.unwrap();
assert_eq!(bytes, Bytes::from_static(STR.as_ref()));
}
#[ntex::test]
async fn untruncated_host_header() {
let srv = test::server(async |_| {
App::new().route(
"/host",
web::get().to(async move |req: HttpRequest| {
req.headers()
.get(ntex::http::header::HOST)
.and_then(|v| v.to_str().ok())
.unwrap_or("<missing>")
.to_string()
}),
)
});
let port = srv.addr().port();
let resp = srv
.request(ntex::http::Method::GET, srv.url("/host"))
.send()
.await
.unwrap();
let body = resp.body().await.unwrap();
let received = String::from_utf8(body.to_vec()).unwrap();
assert_eq!(received, format!("localhost:{port}"));
}
#[ntex::test]
async fn test_query_method() {
use ntex::util::BytesMut;
let srv = test::server(async |_| {
App::new().route(
"/query",
web::query().to(async move |mut payload: web::types::Payload| {
let mut bytes = BytesMut::new();
while let Some(item) = ntex::util::stream_recv(&mut payload).await {
bytes.extend_from_slice(&item.unwrap());
}
bytes
}),
)
});
let test_body = Bytes::from_static(b"SuperBody");
let resp = srv
.query("/query")
.send_body(test_body.clone())
.await
.unwrap();
let body = resp.body().await.unwrap();
assert_eq!(test_body, body);
}