diff --git a/engine/Cargo.lock b/engine/Cargo.lock index 71f3bce0..17d692a9 100644 --- a/engine/Cargo.lock +++ b/engine/Cargo.lock @@ -1655,6 +1655,17 @@ dependencies = [ "tokio", ] +[[package]] +name = "delegate" +version = "0.13.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "297806318ef30ad066b15792a8372858020ae3ca2e414ee6c2133b1eb9e9e945" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.94", +] + [[package]] name = "der" version = "0.6.1" @@ -1791,6 +1802,7 @@ dependencies = [ "color-eyre", "config 0.15.4", "dashmap 6.1.0", + "delegate", "dotenvy", "futures", "hex", @@ -1798,13 +1810,16 @@ dependencies = [ "ipnetwork", "moka", "opentelemetry 0.28.0", + "opentelemetry-auto-span", "opentelemetry-http 0.28.0", "opentelemetry-otlp", + "opentelemetry-prometheus", "opentelemetry-semantic-conventions 0.28.0", "opentelemetry-stdout", "opentelemetry_sdk", "poem", "poem-openapi", + "prometheus", "reqwest", "rust-s3", "rustls 0.23.20", @@ -1815,6 +1830,8 @@ dependencies = [ "sqlx", "thiserror 2.0.9", "tracing", + "tracing-futures", + "tracing-log", "tracing-opentelemetry", "tracing-subscriber", "url", @@ -3346,6 +3363,20 @@ version = "0.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ff011a302c396a5197692431fc1948019154afc178baf7d8e37367442a4601cf" +[[package]] +name = "opentelemetry" +version = "0.26.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "570074cc999d1a58184080966e5bd3bf3a9a4af650c3b05047c2621e7405cd17" +dependencies = [ + "futures-core", + "futures-sink", + "js-sys", + "once_cell", + "pin-project-lite", + "thiserror 1.0.69", +] + [[package]] name = "opentelemetry" version = "0.27.1" @@ -3374,6 +3405,20 @@ dependencies = [ "tracing", ] +[[package]] +name = "opentelemetry-auto-span" +version = "0.4.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "829e46e3ae41aa8cf45f958abf33d591b7a3fa9e4f835cca0ed5821b98fd2a83" +dependencies = [ + "darling", + "opentelemetry 0.26.0", + "proc-macro2", + "quote", + "regex", + "syn 2.0.94", +] + [[package]] name = "opentelemetry-http" version = "0.27.0" @@ -3421,6 +3466,20 @@ dependencies = [ "tracing", ] +[[package]] +name = "opentelemetry-prometheus" +version = "0.28.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "765a76ba13ec77043903322f85dc5434d7d01a37e75536d0f871ed7b9b5bbf0d" +dependencies = [ + "once_cell", + "opentelemetry 0.28.0", + "opentelemetry_sdk", + "prometheus", + "protobuf", + "tracing", +] + [[package]] name = "opentelemetry-proto" version = "0.28.0" @@ -3914,6 +3973,21 @@ dependencies = [ "unicode-ident", ] +[[package]] +name = "prometheus" +version = "0.13.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "3d33c28a30771f7f96db69893f78b857f7450d7e0237e9c8fc6427a81bae7ed1" +dependencies = [ + "cfg-if", + "fnv", + "lazy_static", + "memchr", + "parking_lot", + "protobuf", + "thiserror 1.0.69", +] + [[package]] name = "prost" version = "0.13.5" @@ -3937,6 +4011,12 @@ dependencies = [ "syn 2.0.94", ] +[[package]] +name = "protobuf" +version = "2.28.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "106dd99e98437432fed6519dedecfade6a06a73bb7b2a1e019fdd2bee5778d94" + [[package]] name = "quick-xml" version = "0.32.0" @@ -5688,6 +5768,16 @@ dependencies = [ "tracing-subscriber", ] +[[package]] +name = "tracing-futures" +version = "0.2.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "97d095ae15e245a057c8e8451bab9b3ee1e1f68e9ba2b4fbc18d0ac5237835f2" +dependencies = [ + "pin-project", + "tracing", +] + [[package]] name = "tracing-log" version = "0.2.0" diff --git a/engine/Cargo.toml b/engine/Cargo.toml index 645f068c..c91e028d 100644 --- a/engine/Cargo.toml +++ b/engine/Cargo.toml @@ -25,7 +25,7 @@ poem = { version = "3.1.6", git = "https://github.com/poem-web/poem", branch = " "sse", "tempfile", "opentelemetry", - "requestid" + "requestid", ] } poem-openapi = { version = "5.1.5", git = "https://github.com/poem-web/poem", branch = "master", features = [ "chrono", @@ -61,16 +61,32 @@ rust-s3 = { version = "0.36.0-beta.2", default-features = false, features = [ ] } aws-sdk-s3 = "1.64.0" aws-config = "1.5.10" -async_zip = { version = "0.0.17", features = ["bzip2", "lzma", "zstd", "xz", "deflate"] } +async_zip = { version = "0.0.17", features = [ + "bzip2", + "lzma", + "zstd", + "xz", + "deflate", +] } dashmap = { version = "6.1.0", features = ["serde"] } serde_json = "1.0.138" moka = { version = "0.12.10", features = ["future"] } infer = "0.19.0" # tokio = { version = "1.39.0", features = [""], default-features = false} -tracing-opentelemetry = "0.29.0" -opentelemetry-otlp = { version = "0.28.0", features = ["trace", "metrics", "grpc-tonic"] } -opentelemetry_sdk = {version = "0.28.0", features = ["trace", "rt-async-std"]} -opentelemetry = { version = "0.28.0", features = ["trace"]} +tracing-opentelemetry = { version = "0.29.0" } +opentelemetry-otlp = { version = "0.28.0", features = [ + "trace", + "metrics", + "grpc-tonic", +] } +opentelemetry_sdk = { version = "0.28.0", features = ["trace", "rt-async-std"] } +opentelemetry = { version = "0.28.0", features = ["trace"] } opentelemetry-stdout = "0.28.0" opentelemetry-http = "0.28.0" opentelemetry-semantic-conventions = "0.28.0" +opentelemetry-prometheus = "0.28.0" +prometheus = "0.13.4" +tracing-log = "0.2.0" +tracing-futures = "0.2.5" +opentelemetry-auto-span = "0.4.0" +delegate = "0.13.2" diff --git a/engine/compose.yaml b/engine/compose.yaml index 14c6b620..abc3efa9 100644 --- a/engine/compose.yaml +++ b/engine/compose.yaml @@ -37,3 +37,70 @@ services: echo $0 "$@" exec minio $0 "$@" # start Minio in the foreground ' + +# Monitoring + + # Tempo runs as user 10001, and docker compose creates the volume as root. + # As such, we need to chown the volume in order for Tempo to start correctly. + init: + image: &tempoImage grafana/tempo:latest + user: root + entrypoint: + - "chown" + - "10001:10001" + - "/var/tempo" + volumes: + - ./tempo-data:/var/tempo + + memcached: + image: memcached:1.6.29 + container_name: memcached + ports: + - "11211:11211" + environment: + - MEMCACHED_MAX_MEMORY=64m # Set the maximum memory usage + - MEMCACHED_THREADS=4 # Number of threads to use + + tempo: + image: *tempoImage + command: [ "-config.file=/etc/tempo.yaml" ] + volumes: + - ../shared/tempo.yaml:/etc/tempo.yaml + - ./tempo-data:/var/tempo + ports: + - "14268:14268" # jaeger ingest + - "3200:3200" # tempo + - "9095:9095" # tempo grpc + - "4317:4317" # otlp grpc + - "4318:4318" # otlp http + - "9411:9411" # zipkin + depends_on: + - init + - memcached + + prometheus: + image: prom/prometheus:latest + command: + - --config.file=/etc/prometheus.yaml + - --web.enable-remote-write-receiver + - --enable-feature=exemplar-storage + - --enable-feature=native-histograms + volumes: + - ../shared/prometheus.yaml:/etc/prometheus.yaml + ports: + - "9090:9090" + extra_hosts: + - "host.docker.internal:host-gateway" + + grafana: + image: grafana/grafana:11.2.0 + volumes: + - ../shared/grafana-datasources.yaml:/etc/grafana/provisioning/datasources/datasources.yaml + environment: + - GF_AUTH_ANONYMOUS_ENABLED=true + - GF_AUTH_ANONYMOUS_ORG_ROLE=Admin + - GF_AUTH_DISABLE_LOGIN_FORM=true + - GF_FEATURE_TOGGLES_ENABLE=traceqlEditor metricsSummary + - GF_INSTALL_PLUGINS=https://storage.googleapis.com/integration-artifacts/grafana-exploretraces-app/grafana-exploretraces-app-latest.zip;grafana-traces-app + ports: + - "3001:3000" diff --git a/engine/src/database.rs b/engine/src/database.rs index 43c42245..4d3dd3c5 100644 --- a/engine/src/database.rs +++ b/engine/src/database.rs @@ -1,4 +1,7 @@ -use sqlx::{postgres::PgPoolOptions, PgPool}; +use sqlx::{ + postgres::{PgConnectOptions, PgPoolOptions}, + ConnectOptions, PgPool, +}; use tracing::info; #[derive(Debug)] @@ -8,7 +11,23 @@ pub struct Database { impl Database { pub async fn new(url: &str) -> Result { - let pool = PgPoolOptions::new().max_connections(5).connect(url).await?; + let mut options: PgConnectOptions = url.parse().unwrap(); + + options = options.log_statements(tracing_log::log::LevelFilter::Trace); + + let pool = PgPoolOptions::new() + .max_connections(5) + .before_acquire(|conn, _meta| Box::pin(async move { + info!("Before Acquire"); + Ok(true) + })) + .after_connect(|conn, _meta| Box::pin(async move { + info!("After Connect"); + + Ok(()) + })) + .connect_with(options) + .await?; // Initialization code here let s = Self { pool }; diff --git a/engine/src/main.rs b/engine/src/main.rs index 07c3e2a7..76f13e92 100644 --- a/engine/src/main.rs +++ b/engine/src/main.rs @@ -3,9 +3,10 @@ use std::env; use opentelemetry::KeyValue; use opentelemetry::{global, trace::TracerProvider}; use opentelemetry_otlp::WithExportConfig; +use opentelemetry_sdk::trace::SdkTracerProvider; use opentelemetry_sdk::Resource; use state::AppState; -use tracing::{error, info}; +use tracing::{error, info, warn, Subscriber}; pub mod assets; pub mod cache; @@ -17,24 +18,25 @@ pub mod state; pub mod storage; pub mod utils; +use tracing_subscriber::layer::Layered; use tracing_subscriber::{prelude::*, EnvFilter}; #[async_std::main] async fn main() { - dotenvy::dotenv().ok(); - - // let exporter = opentelemetry_otlp::SpanExporter::builder() - // .with_http() - // .with_endpoint("http://localhost:4317") - // .build() - // // .install_batch(opentelemetry_sdk::runtime::AsyncStd) - // .expect("Couldn't create OTLP tracer"); + // Initialize the log bridge early so that all log records are captured. + // tracing_log::LogTracer::init().expect("Failed to set logger"); - let otlp_endpoint = env::var("OTLP_ENDPOINT").ok(); + dotenvy::dotenv().ok(); - if let Some(endpoint) = otlp_endpoint { - info!("Starting Edgerouter with OTLP tracing"); + let otlp_endpoint = std::env::var("OTLP_ENDPOINT").ok(); + let fmt_layer = tracing_subscriber::fmt::layer(); + let env_filter = tracing_subscriber::EnvFilter::from_default_env(); + // Create a subscriber based on whether OTLP is configured. + let (subscriber, tracer_provider): ( + Box, + Option, + ) = if let Some(endpoint) = otlp_endpoint { let exporter = opentelemetry_otlp::SpanExporter::builder() .with_tonic() .with_endpoint(endpoint) @@ -42,38 +44,56 @@ async fn main() { .expect("Couldn't create OTLP tracer"); let hostname = std::env::var("HOSTNAME").unwrap_or_else(|_| "unknown".to_string()); - - let resource = Resource::builder() + let resource = opentelemetry_sdk::Resource::builder() .with_service_name("edgeserver") - .with_attribute(KeyValue::new("host.name", hostname)) + .with_attribute(opentelemetry::KeyValue::new("host.name", hostname)) .build(); - let trace_provider = opentelemetry_sdk::trace::SdkTracerProvider::builder() + // let sql_resource = opentelemetry_sdk::Resource::builder() + // .with_service_name("postgresql") + // .with_attribute(opentelemetry::KeyValue::new("host.name", hostname)) + // .build(); + + let tracer_provider = opentelemetry_sdk::trace::SdkTracerProvider::builder() .with_resource(resource) .with_batch_exporter(exporter) .build(); - global::set_tracer_provider(trace_provider.clone()); - - let tracer = trace_provider.tracer("edgeserver"); + opentelemetry::global::set_tracer_provider(tracer_provider.clone()); + let tracer = tracer_provider.tracer("edgeserver"); let telemetry_layer = tracing_opentelemetry::layer() + .with_level(true) .with_tracer(tracer.clone()) .with_error_fields_to_exceptions(true) .with_tracked_inactivity(true); - let fmt_layer = tracing_subscriber::fmt::layer(); - tracing_subscriber::registry() - .with(fmt_layer) - .with(telemetry_layer) - .with(EnvFilter::from_default_env()) - .init(); - // tracing_subscriber::fmt::init(); - + ( + Box::new( + tracing_subscriber::registry() + .with(fmt_layer) + .with(env_filter) + .with(telemetry_layer), + ), + Some(tracer_provider), + ) } else { - info!("Starting Edgerouter without OTLP tracing, provide OTLP_ENDPOINT to enable tracing"); - tracing_subscriber::fmt::init(); - } + ( + Box::new( + tracing_subscriber::registry() + .with(fmt_layer) + .with(env_filter), + ), + None, + ) + }; + + // Use try_init() to avoid a panic if a global subscriber is already set. + subscriber.try_init().unwrap_or_else(|err| { + eprintln!("Global subscriber already set: {}", err); + }); + + tracing::info!("Starting Edgerouter..."); let state = match AppState::new().await { Ok(state) => state, @@ -84,4 +104,10 @@ async fn main() { }; routes::serve(state).await; + + if let Some(tracer_provider) = tracer_provider { + warn!("Shutting down tracer provider"); + tracer_provider.force_flush(); + tracer_provider.shutdown(); + } } diff --git a/engine/src/middlewares/metrics.rs b/engine/src/middlewares/metrics.rs new file mode 100644 index 00000000..0c44bff9 --- /dev/null +++ b/engine/src/middlewares/metrics.rs @@ -0,0 +1,123 @@ +use std::time::Instant; + +use opentelemetry::{ + global, + metrics::{Counter, Histogram}, + Key, KeyValue, +}; +use opentelemetry_semantic_conventions::trace; + +use poem::{PathPattern, Endpoint, IntoResponse, Middleware, Request, Response, Result}; + +/// Middleware for metrics with OpenTelemetry. +#[cfg_attr(docsrs, doc(cfg(feature = "opentelemetry")))] +pub struct OpenTelemetryMetrics { + request_count: Counter, + error_count: Counter, + duration: Histogram, +} + +impl Default for OpenTelemetryMetrics { + fn default() -> Self { + Self::new() + } +} + +impl OpenTelemetryMetrics { + /// Create `OpenTelemetryMetrics` middleware with `meter`. + pub fn new() -> Self { + let meter = global::meter("poem"); + Self { + request_count: meter + .u64_counter("poem_requests_count") + .with_description("total request count (since start of service)") + .build(), + error_count: meter + .u64_counter("poem_errors_count") + .with_description("failed request count (since start of service)") + .build(), + duration: meter + .f64_histogram("poem_request_duration_ms") + .with_unit("milliseconds") + .with_description( + "request duration histogram (in milliseconds, since start of service)", + ) + .build(), + } + } +} + +impl Middleware for OpenTelemetryMetrics { + type Output = OpenTelemetryMetricsEndpoint; + + fn transform(&self, ep: E) -> Self::Output { + OpenTelemetryMetricsEndpoint { + request_count: self.request_count.clone(), + error_count: self.error_count.clone(), + duration: self.duration.clone(), + inner: ep, + } + } +} + +/// Endpoint for the OpenTelemetryMetrics middleware. +#[cfg_attr(docsrs, doc(cfg(feature = "opentelemetry")))] +pub struct OpenTelemetryMetricsEndpoint { + request_count: Counter, + error_count: Counter, + duration: Histogram, + inner: E, +} + +impl Endpoint for OpenTelemetryMetricsEndpoint { + type Output = Response; + + async fn call(&self, req: Request) -> Result { + let mut labels = Vec::with_capacity(3); + labels.push(KeyValue::new( + trace::HTTP_REQUEST_METHOD, + req.method().to_string(), + )); + labels.push(KeyValue::new( + trace::URL_FULL, + req.original_uri().to_string(), + )); + + let s = Instant::now(); + let res = self.inner.call(req).await.map(IntoResponse::into_response); + let elapsed = s.elapsed(); + + match &res { + Ok(resp) => { + if let Some(path_pattern) = resp.data::() { + const HTTP_PATH_PATTERN: Key = Key::from_static_str("http.path_pattern"); + labels.push(KeyValue::new(HTTP_PATH_PATTERN, path_pattern.0.to_string())); + } + + labels.push(KeyValue::new( + trace::HTTP_RESPONSE_STATUS_CODE, + resp.status().as_u16() as i64, + )); + } + Err(err) => { + if let Some(path_pattern) = err.data::() { + const HTTP_PATH_PATTERN: Key = Key::from_static_str("http.path_pattern"); + labels.push(KeyValue::new(HTTP_PATH_PATTERN, path_pattern.0.to_string())); + } + + labels.push(KeyValue::new( + trace::HTTP_RESPONSE_STATUS_CODE, + err.status().as_u16() as i64, + )); + self.error_count.add(1, &labels); + labels.push(KeyValue::new(trace::EXCEPTION_MESSAGE, err.to_string())); + } + } + + self.request_count.add(1, &labels); + self.duration + .record(elapsed.as_secs_f64() * 1000.0, &labels); + + res + } +} diff --git a/engine/src/middlewares/mod.rs b/engine/src/middlewares/mod.rs index b41280aa..08f09b88 100644 --- a/engine/src/middlewares/mod.rs +++ b/engine/src/middlewares/mod.rs @@ -1,2 +1,3 @@ pub mod auth; pub mod tracing; +pub mod metrics; diff --git a/engine/src/middlewares/tracing.rs b/engine/src/middlewares/tracing.rs index 8d8d53ff..04ff9290 100644 --- a/engine/src/middlewares/tracing.rs +++ b/engine/src/middlewares/tracing.rs @@ -2,7 +2,7 @@ use std::sync::Arc; use opentelemetry::{ global, - trace::{FutureExt, Span, SpanKind, TraceContextExt, Tracer}, + trace::{self, FutureExt, Span, SpanKind, TraceContextExt, Tracer}, Context, Key, KeyValue, }; use opentelemetry_http::HeaderExtractor; @@ -108,15 +108,14 @@ where span.add_event("request.started".to_string(), vec![]); - // let tracing_span = info_span!("request.started"); let parent_context = Context::current_with_span(span); - // tracing_span.set_parent(parent_context.clone()); - - // set the tracing_opentelemetry span default for new spans to the parent_context - let tracing_span = tracing::span!(tracing::Level::INFO, "tracing-request-started"); - tracing_span.set_parent(parent_context.clone()); async move { + let cx = Context::current(); + cx.attach(); + // let tracing_span = tracing::span!(tracing::Level::INFO, "test"); + // tracing_span.set_parent(cx); + let res = self.inner.call(req).await; let cx = Context::current(); let span = cx.span(); @@ -135,15 +134,26 @@ where } span.add_event("request.completed".to_string(), vec![]); + + // Set the http.response.status_code otlp attribute span.set_attribute(KeyValue::new( attribute::HTTP_RESPONSE_STATUS_CODE, resp.status().as_u16() as i64, )); - // Grafana specific override + // Grafana specific override (http.status_code) span.set_attribute(KeyValue::new( "http.status_code", resp.status().as_u16() as i64, )); + + // Span status update + if resp.status().is_success() { + span.set_status(trace::Status::Ok); + } else { + let status_string = resp.status().to_string(); + span.set_status(trace::Status::Error { description: status_string.into() }); + } + if let Some(content_length) = resp.headers().typed_get::() { @@ -193,8 +203,6 @@ where } } .with_context(parent_context) - .instrument(tracing_span) - // .instrument(tracing_span) .await } } diff --git a/engine/src/models/user/mod.rs b/engine/src/models/user/mod.rs index 560bee88..e092272e 100644 --- a/engine/src/models/user/mod.rs +++ b/engine/src/models/user/mod.rs @@ -1,10 +1,13 @@ use chrono::{DateTime, Utc}; +use opentelemetry::{global::ObjectSafeSpan, trace::{SpanKind, Tracer}, KeyValue}; use poem_openapi::Object; use serde::{Deserialize, Serialize}; -use sqlx::{query_as, query_scalar}; +use sqlx::{query_as, query_scalar, Acquire, Connection, ConnectOptions}; +use tracing::{info_span, instrument}; use crate::{ - database::Database, utils::id::{generate_id, IdType} + database::Database, + utils::id::{generate_id, IdType}, }; use super::team::Team; @@ -41,14 +44,11 @@ impl User { let password = password.as_ref(); // check if no user with this name already exists - let exists = query_scalar!( - "SELECT COUNT(*) FROM users WHERE name = $1", - name - ) - .fetch_one(&db.pool) - .await? - .unwrap_or(0) - > 0; + let exists = query_scalar!("SELECT COUNT(*) FROM users WHERE name = $1", name) + .fetch_one(&db.pool) + .await? + .unwrap_or(0) + > 0; if exists { // TODO: nicer `Username taken` error @@ -77,7 +77,49 @@ impl User { } } + #[opentelemetry_auto_span::auto_span] pub async fn get_by_id(db: &Database, user_id: impl AsRef) -> Result { + let query = "SELECT * FROM users WHERE user_id = $1"; + + let tracer = opentelemetry::global::tracer("sqlx"); + let attributes: Vec = vec![ + KeyValue::new("db.system", "postgresql"), + KeyValue::new("db.query.text", query), + KeyValue::new("db.statement", query), + KeyValue::new("db.system.name", "postgresql"), + KeyValue::new("otel.kind", "client"), + KeyValue::new("db.name", "edgeserver"), + KeyValue::new("db.collection", "users"), + KeyValue::new("db.operation", "SELECT"), + // + KeyValue::new("peer.service", "postgresql"), + KeyValue::new("server.address", "localhost:5432"), + // KeyValue::new("net.peer.name", "localhost"), + // KeyValue::new("net.peer.port", 5432), + ]; + let span = tracer.span_builder("SELECT * FROM users ...").with_kind(SpanKind::Client).with_attributes(attributes); + let span = tracer.build(span); + + let span = info_span!( + "SELECT * FROM users ...", + otel.kind = "client", + db.system = "postgresql", + db.system.name = "postgresql", + db.name = "edgeserver", + db.collection.name = "users", + db.collection = "users", + db.operation.name = "SELECT", + db.operation = "SELECT", + db.query.text = %query, + db.statement = %query, + peer.service = "postgresql", + server.address = "localhost:5432", + net.peer.name = "localhost", // Database host + net.peer.port = 5432, // Database port + service.name = "postgres", + otel.service.name = "postgres", + ); + query_as!( User, "SELECT * FROM users WHERE user_id = $1", @@ -103,13 +145,17 @@ impl User { } pub async fn get_all_minimal(db: &Database) -> Result, sqlx::Error> { - query_as!(UserMinimal, "SELECT user_id, name, avatar_url, admin FROM users") - .fetch_all(&db.pool) - .await + query_as!( + UserMinimal, + "SELECT user_id, name, avatar_url, admin FROM users" + ) + .fetch_all(&db.pool) + .await } + #[instrument(skip(db))] pub async fn can_bootstrap(db: &Database) -> Result { - Ok(query_scalar!("SELECT COUNT(*) FROM users") + Ok(query_scalar!("SELECT COUNT(*) FROM users LIMIT 1") .fetch_one(&db.pool) .await? .unwrap_or(0) diff --git a/engine/src/routes/mod.rs b/engine/src/routes/mod.rs index d2b445a3..2c37f271 100644 --- a/engine/src/routes/mod.rs +++ b/engine/src/routes/mod.rs @@ -4,18 +4,26 @@ use async_std::path::Path; use auth::AuthApi; use invite::InviteApi; use opentelemetry::global; -use poem::middleware::OpenTelemetryMetrics; +use opentelemetry::metrics::Counter; +use opentelemetry::KeyValue; +use opentelemetry_prometheus::PrometheusExporter; +use opentelemetry_prometheus::ExporterBuilder; +use opentelemetry_sdk::metrics::SdkMeterProvider; +use poem::web::Data; use poem::{ endpoint::StaticFilesEndpoint, get, handler, listener::TcpListener, middleware::Cors, web::Html, EndpointExt, Route, Server, }; use poem_openapi::payload::PlainText; use poem_openapi::{OpenApi, OpenApiService, Tags}; +use prometheus::Registry; +use prometheus::{Encoder, TextEncoder}; use site::SiteApi; use team::TeamApi; use tracing::info; use user::UserApi; +use crate::middlewares::metrics::OpenTelemetryMetrics; use crate::middlewares::tracing::TraceId; use crate::state::AppState; @@ -59,19 +67,28 @@ pub async fn serve(state: AppState) { .index_file("index.html") .fallback_to_index(); + let registry = prometheus::Registry::new(); + let exporter = opentelemetry_prometheus::exporter().with_registry(registry.clone()).build().unwrap(); + let provider = SdkMeterProvider::builder().with_reader(exporter).build(); + global::set_meter_provider(provider); + + let meter = global::meter("edgeserver"); + let counter = meter.u64_counter("test_counter").build(); + let app = Route::new() .nest("/api", api_service) .nest("/openapi.json", spec) .at("/docs", get(get_openapi_docs)) + .at("/metrics", get(metrics_handler).data(Arc::new(registry)).data(Arc::new(counter))) .nest("/", file_endpoint) .with(Cors::new()) - .with(TraceId::new(Arc::new(global::tracer("edgeserver")))) .with(OpenTelemetryMetrics::new()) + .with(TraceId::new(Arc::new(global::tracer("edgeserver")))) .data(state); let listener = TcpListener::bind("0.0.0.0:3000"); - Server::new(listener).run(app).await.unwrap() + Server::new(listener).run(app).await.unwrap(); } #[handler] @@ -84,3 +101,16 @@ async fn not_found() -> Html<&'static str> { // inline 404 template Html(include_str!("./404.html")) } + +#[handler] +async fn metrics_handler(registry: Data<&Arc>, counter: Data<&Arc>>) -> String { + counter.0.add(1, &[KeyValue::new("test_key", "test_value")]); + + // // Gather and format the metrics for Prometheus. + let encoder = TextEncoder::new(); + let mf = registry.0.gather(); + let mut buffer = Vec::new(); + encoder.encode(&mf, &mut buffer).unwrap(); + + String::from_utf8(buffer).unwrap() +} diff --git a/shared/grafana-datasources.yaml b/shared/grafana-datasources.yaml new file mode 100644 index 00000000..794d5a1d --- /dev/null +++ b/shared/grafana-datasources.yaml @@ -0,0 +1,32 @@ +apiVersion: 1 + +datasources: +- name: Prometheus + type: prometheus + uid: prometheus + access: proxy + orgId: 1 + url: http://prometheus:9090 + basicAuth: false + isDefault: false + version: 1 + editable: false + jsonData: + httpMethod: GET +- name: Tempo + type: tempo + access: proxy + orgId: 1 + url: http://tempo:3200 + basicAuth: false + isDefault: true + version: 1 + editable: false + apiVersion: 1 + uid: tempo + jsonData: + httpMethod: GET + serviceMap: + datasourceUid: prometheus + streamingEnabled: + search: true diff --git a/shared/prometheus.yaml b/shared/prometheus.yaml new file mode 100644 index 00000000..0d3afc11 --- /dev/null +++ b/shared/prometheus.yaml @@ -0,0 +1,15 @@ +global: + scrape_interval: 15s + evaluation_interval: 15s + +scrape_configs: + - job_name: 'prometheus' + static_configs: + - targets: [ 'localhost:9090' ] + - job_name: 'tempo' + static_configs: + - targets: [ 'tempo:3200' ] + - job_name: 'edgeserver' + metrics_path: /metrics + static_configs: + - targets: [ 'host.docker.internal:3000' ] diff --git a/shared/tempo.yaml b/shared/tempo.yaml new file mode 100644 index 00000000..e17868f5 --- /dev/null +++ b/shared/tempo.yaml @@ -0,0 +1,89 @@ +stream_over_http_enabled: true +server: + http_listen_port: 3200 + log_level: info + +cache: + background: + writeback_goroutines: 5 + caches: + - roles: + - frontend-search + memcached: + addresses: dns+memcached:11211 + +query_frontend: + search: + duration_slo: 5s + throughput_bytes_slo: 1.073741824e+09 + metadata_slo: + duration_slo: 5s + throughput_bytes_slo: 1.073741824e+09 + trace_by_id: + duration_slo: 100ms + metrics: + max_duration: 120h # maximum duration of a metrics query, increase for local setups + query_backend_after: 5m + duration_slo: 5s + throughput_bytes_slo: 1.073741824e+09 + +distributor: + receivers: # this configuration will listen on all ports and protocols that tempo is capable of. + jaeger: # the receives all come from the OpenTelemetry collector. more configuration information can + protocols: # be found there: https://github.com/open-telemetry/opentelemetry-collector/tree/main/receiver + thrift_http: # + endpoint: "tempo:14268" # for a production deployment you should only enable the receivers you need! + grpc: + endpoint: "tempo:14250" + thrift_binary: + endpoint: "tempo:6832" + thrift_compact: + endpoint: "tempo:6831" + zipkin: + endpoint: "tempo:9411" + otlp: + protocols: + grpc: + endpoint: "tempo:4317" + http: + endpoint: "tempo:4318" + opencensus: + endpoint: "tempo:55678" + +ingester: + max_block_duration: 5m # cut the headblock when this much time passes. this is being set for demo purposes and should probably be left alone normally + +compactor: + compaction: + block_retention: 24h # overall Tempo trace retention. set for demo purposes + +metrics_generator: + registry: + external_labels: + source: tempo + cluster: docker-compose + storage: + path: /var/tempo/generator/wal + remote_write: + - url: http://prometheus:9090/api/v1/write + send_exemplars: true + traces_storage: + path: /var/tempo/generator/traces + processor: + local_blocks: + filter_server_spans: false + flush_to_storage: true + +storage: + trace: + backend: local # backend configuration to use + wal: + path: /var/tempo/wal # where to store the wal locally + local: + path: /var/tempo/blocks + +overrides: + defaults: + metrics_generator: + processors: [service-graphs, span-metrics, local-blocks] # enables metrics generator + generate_native_histograms: both diff --git a/web/package.json b/web/package.json index 7335a2a6..0b63278f 100644 --- a/web/package.json +++ b/web/package.json @@ -26,6 +26,7 @@ "@esbuild-plugins/node-globals-polyfill": "^0.1.1", "@headlessui/react": "^1.6.5", "@hookform/resolvers": "^3.1.0", + "@opentelemetry/api": "^1.9.0", "@radix-ui/react-accordion": "^1.2.2", "@radix-ui/react-alert-dialog": "^1.1.6", "@radix-ui/react-avatar": "^1.1.2", diff --git a/web/pnpm-lock.yaml b/web/pnpm-lock.yaml index 85bd3ae7..f0a56da9 100644 --- a/web/pnpm-lock.yaml +++ b/web/pnpm-lock.yaml @@ -20,6 +20,9 @@ importers: '@hookform/resolvers': specifier: ^3.1.0 version: 3.9.1(react-hook-form@7.54.2(react@18.3.1)) + '@opentelemetry/api': + specifier: ^1.9.0 + version: 1.9.0 '@radix-ui/react-accordion': specifier: ^1.2.2 version: 1.2.2(@types/react-dom@18.3.5(@types/react@18.3.18))(@types/react@18.3.18)(react-dom@18.3.1(react@18.3.1))(react@18.3.1) @@ -944,6 +947,10 @@ packages: resolution: {integrity: sha512-oGB+UxlgWcgQkgwo8GcEGwemoTFt3FIO9ababBmaGwXIoBKZ+GTy0pP185beGg7Llih/NSHSV2XAs1lnznocSg==} engines: {node: '>= 8'} + '@opentelemetry/api@1.9.0': + resolution: {integrity: sha512-3giAOQvZiH5F9bMlMiv8+GSPMeqg0dbaeo58/0SlA9sxSqZhnUtxzX9/2FzyhS9sWQf5S0GJE0AKBrFqjpeYcg==} + engines: {node: '>=8.0.0'} + '@parcel/watcher-android-arm64@2.5.0': resolution: {integrity: sha512-qlX4eS28bUcQCdribHkg/herLe+0A9RyYC+mm2PXpncit8z5b3nSqGVzMNR3CmtAOgRutiZ02eIJJgP/b1iEFQ==} engines: {node: '>= 10.0.0'} @@ -5524,6 +5531,8 @@ snapshots: '@nodelib/fs.scandir': 2.1.5 fastq: 1.18.0 + '@opentelemetry/api@1.9.0': {} + '@parcel/watcher-android-arm64@2.5.0': optional: true