#![expect(
clippy::unwrap_used,
clippy::expect_used,
reason = "example/test/bench: panic-on-error and print-for-output are the standard patterns for demos and harnesses"
)]
use rama::{
Layer,
bytes::Bytes,
extensions::ExtensionsRef,
http::{
Request, header,
layer::{
compression::CompressionLayer,
error_handling::ErrorHandlerLayer,
sensitive_headers::{
SetSensitiveRequestHeadersLayer, SetSensitiveResponseHeadersLayer,
},
trace::{DefaultMakeSpan, DefaultOnResponse, TraceLayer},
},
server::HttpServer,
service::web::{IntoEndpointService, response::Html},
},
layer::ArcLayer,
layer::{TimeoutLayer, TraceErrLayer},
net::stream::{
SocketInfo,
layer::{BytesRWTrackerHandle, IncomingBytesTrackerLayer},
},
rt::Executor,
tcp::server::TcpListener,
telemetry::tracing::{
self,
level_filters::LevelFilter,
subscriber::{EnvFilter, fmt, layer::SubscriberExt, util::SubscriberInitExt},
},
utils::latency::LatencyUnit,
};
use std::{sync::Arc, time::Duration};
#[tokio::main]
async fn main() {
tracing::subscriber::registry()
.with(fmt::layer())
.with(
EnvFilter::builder()
.with_default_directive(LevelFilter::DEBUG.into())
.from_env_lossy(),
)
.init();
let graceful = rama::graceful::Shutdown::default();
let sensitive_headers: Arc<[_]> = vec![header::AUTHORIZATION, header::COOKIE].into();
graceful.spawn_task_fn(async |guard| {
let exec = Executor::graceful(guard.clone());
let http_service = (
ArcLayer::new(),
CompressionLayer::new(),
SetSensitiveRequestHeadersLayer::from_shared(sensitive_headers.clone()),
TraceLayer::new_for_http()
.on_body_chunk(|chunk: &Bytes, latency: Duration, _: &tracing::Span| {
tracing::trace!(
http.request.body.chunk_size = chunk.len(),
http.request.body.chunk_read_ms = latency.as_millis(),
"sending body chunk"
)
})
.make_span_with(DefaultMakeSpan::new().with_include_headers(true))
.on_response(
DefaultOnResponse::new()
.with_include_headers(true)
.with_latency_unit(LatencyUnit::Micros),
),
SetSensitiveResponseHeadersLayer::from_shared(sensitive_headers),
ErrorHandlerLayer::new(),
)
.into_layer(
(|req: Request| {
let uri = req.uri();
let socket_info = req.extensions().get_ref::<SocketInfo>().unwrap();
let tracker = req.extensions().get_ref::<BytesRWTrackerHandle>().unwrap();
std::future::ready(Html(format!(
r##"
<html>
<head>
<title>Rama — Http Service Hello</title>
</head>
<body>
<h1>Hello</h1>
<p>Peer: {}</p>
<p>Path: {}</p>
<p>Stats (bytes):</p>
<ul>
<li>Read: {}</li>
<li>Written: {}</li>
</ul>
</body>
</html>"##,
socket_info.peer_addr(),
uri.path_or_root(),
tracker.read(),
tracker.written(),
)))
})
.into_endpoint_service(),
);
let tcp_http_service = HttpServer::auto(exec.clone()).service(http_service);
TcpListener::bind_address("127.0.0.1:62010", exec)
.await
.expect("bind TCP Listener")
.serve(
(
TraceErrLayer::new(),
TimeoutLayer::new(Duration::from_secs(8)),
IncomingBytesTrackerLayer::new(),
)
.into_layer(tcp_http_service),
)
.await;
});
graceful
.shutdown_with_limit(Duration::from_secs(30))
.await
.expect("graceful shutdown");
}