#![expect(
clippy::unwrap_used,
reason = "example/test/bench: panic-on-error and print-for-output are the standard patterns for demos and harnesses"
)]
use rama::{
Layer,
extensions::{Extension, Extensions},
http::{
client::EasyHttpWebClient,
layer::{opentelemetry::RequestMetricsLayer, trace::TraceLayer},
server::HttpServer,
service::web::{WebService, response::Html},
},
layer::AddInputExtensionLayer,
net::{stream::layer::opentelemetry::NetworkMetricsLayer, uri::Uri},
rt::Executor,
tcp::server::TcpListener,
telemetry::{
opentelemetry::{
self, InstrumentationScope, KeyValue,
collector::OtelExporter,
logs::{LogRecord, Logger, LoggerProvider},
metrics::UpDownCounter,
sdk::{
Resource,
logs::SdkLoggerProvider,
metrics::{PeriodicReader, SdkMeterProvider},
},
semantic_conventions::{
self,
resource::{HOST_ARCH, OS_NAME, SERVICE_NAME, SERVICE_VERSION},
},
},
tracing::{
self,
level_filters::LevelFilter,
subscriber::{EnvFilter, fmt, layer::SubscriberExt, util::SubscriberInitExt},
},
},
};
use std::{sync::Arc, time::Duration};
#[derive(Debug, Extension)]
struct AppState {
counter: UpDownCounter<i64>,
logger: <SdkLoggerProvider as LoggerProvider>::Logger,
}
impl AppState {
fn new(logger_provider: &SdkLoggerProvider) -> Self {
let scope = InstrumentationScope::builder("example.http_telemetry")
.with_version(env!("CARGO_PKG_VERSION"))
.with_schema_url(semantic_conventions::SCHEMA_URL)
.with_attributes(vec![
KeyValue::new(OS_NAME, std::env::consts::OS),
KeyValue::new(HOST_ARCH, std::env::consts::ARCH),
])
.build();
let meter = opentelemetry::global::meter_with_scope(scope.clone());
let counter = meter.i64_up_down_counter("visitor_counter").build();
let logger = logger_provider.logger_with_scope(scope);
Self { counter, logger }
}
}
#[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 exporter = OtelExporter::new_http(EasyHttpWebClient::default())
.with_endpoint(Uri::from_static("http://localhost:4318"))
.with_timeout(Duration::from_secs(10));
let meter_reader = PeriodicReader::builder(exporter.clone())
.with_interval(Duration::from_secs(3))
.build();
let resource = Resource::builder()
.with_attribute(KeyValue::new(SERVICE_NAME, "http_telemetry"))
.with_attribute(KeyValue::new(SERVICE_VERSION, rama::utils::info::VERSION))
.build();
let meter = SdkMeterProvider::builder()
.with_resource(resource.clone())
.with_reader(meter_reader)
.build();
opentelemetry::global::set_meter_provider(meter);
let logger_provider = SdkLoggerProvider::builder()
.with_resource(resource)
.with_batch_exporter(exporter)
.build();
let state = Arc::new(AppState::new(&logger_provider));
let graceful = rama::graceful::Shutdown::default();
graceful.spawn_task_fn(async |guard| {
let exec = Executor::graceful(guard);
let http_service = HttpServer::auto(exec.clone()).service(
(TraceLayer::new_for_http(), RequestMetricsLayer::default()).into_layer(
WebService::default().with_get("/", async |ext: Extensions| {
let state = ext.get_ref::<AppState>().unwrap();
state.counter.add(1, &[]);
let mut record = state.logger.create_log_record();
record.set_severity_text("INFO");
record.set_body("visitor".into());
state.logger.emit(record);
Html("<h1>Hello!</h1>")
}),
),
);
TcpListener::build(exec)
.bind_address("127.0.0.1:62012")
.await
.unwrap()
.serve(
(
AddInputExtensionLayer::new_arc(state),
NetworkMetricsLayer::default(),
)
.into_layer(http_service),
)
.await;
});
graceful
.shutdown_with_limit(Duration::from_secs(30))
.await
.unwrap();
}