diff --git a/Cargo.lock b/Cargo.lock index 0db20be5..9eabbf95 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -129,7 +129,7 @@ version = "1.1.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "40c48f72fd53cd289104fc64099abca73db4166ad86ea0b4341abe65af83dadc" dependencies = [ - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -140,7 +140,7 @@ checksum = "291e6a250ff86cd4a820112fb8898808a366d8f9f58ce16d1f538353ad55747d" dependencies = [ "anstyle", "once_cell_polyfill", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -487,6 +487,10 @@ dependencies = [ "humantime", "humantime-serde", "itoa", + "opentelemetry", + "opentelemetry-appender-tracing", + "opentelemetry-otlp", + "opentelemetry_sdk", "rand 0.10.1", "regex", "rsa 0.10.0-rc.18", @@ -503,6 +507,7 @@ dependencies = [ "tower-http 0.7.0", "tracing", "tracing-error", + "tracing-opentelemetry", "tracing-subscriber", "uuid", "xdg", @@ -2184,7 +2189,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "39cab71617ae0d63f51a36d69f866391735b51691dbda63cf6f96d042b63efeb" dependencies = [ "libc", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -2857,7 +2862,7 @@ dependencies = [ "libc", "percent-encoding", "pin-project-lite", - "socket2 0.5.10", + "socket2 0.6.4", "tokio", "tower-service", "tracing", @@ -3549,7 +3554,7 @@ version = "0.50.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7957b9740744892f114936ab4a57b3f487491bbeafaf8083688b16841a4240e5" dependencies = [ - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -3647,6 +3652,93 @@ version = "0.2.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7c87def4c32ab89d880effc9e097653c8da5d6ef28e6b539d313baaacfbafcbe" +[[package]] +name = "opentelemetry" +version = "0.32.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b0142c63252a9e054e68a4c61a5778f7b14f576274d593f8ce883d191a099682" +dependencies = [ + "futures-core", + "futures-sink", + "js-sys", + "pin-project-lite", + "thiserror 2.0.18", + "tracing", +] + +[[package]] +name = "opentelemetry-appender-tracing" +version = "0.32.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "2c0080f0dc1d7c786f467cd85a4e395fcab11ee852004f39a29a18ab7c25d837" +dependencies = [ + "opentelemetry", + "tracing", + "tracing-core", + "tracing-subscriber", +] + +[[package]] +name = "opentelemetry-http" +version = "0.32.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5683015d09e2df236ef005b17f6f196f0d5f6313c4fa43a7b6a53b52776e4331" +dependencies = [ + "async-trait", + "bytes", + "http 1.4.2", + "opentelemetry", + "reqwest", +] + +[[package]] +name = "opentelemetry-otlp" +version = "0.32.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9966929966d17620d7c316c643ba62631826e10021409357772d5eea84f62c35" +dependencies = [ + "http 1.4.2", + "opentelemetry", + "opentelemetry-http", + "opentelemetry-proto", + "opentelemetry_sdk", + "prost", + "reqwest", + "thiserror 2.0.18", + "tokio", + "tonic", + "tonic-types", +] + +[[package]] +name = "opentelemetry-proto" +version = "0.32.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "56d658ba1faf63f7b9c492cfbe6e0ec365440a16132d3270c1065f7b33f1b638" +dependencies = [ + "opentelemetry", + "opentelemetry_sdk", + "prost", + "tonic", + "tonic-prost", +] + +[[package]] +name = "opentelemetry_sdk" +version = "0.32.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9b59f80e1ac4d5ff7a2db8fb6c80badb7f0f3f858211fba08dd9aaec750894f9" +dependencies = [ + "futures-channel", + "futures-executor", + "futures-util", + "opentelemetry", + "percent-encoding", + "portable-atomic", + "rand 0.9.4", + "thiserror 2.0.18", +] + [[package]] name = "ordered-float" version = "4.6.0" @@ -4028,7 +4120,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "b570b25f7617e43d59005d0990ccb79e950a423952cea19671b7a876da390adf" dependencies = [ "anyhow", - "itertools 0.13.0", + "itertools 0.14.0", "proc-macro2", "quote", "syn 2.0.118", @@ -4076,7 +4168,7 @@ dependencies = [ "quinn-udp", "rustc-hash", "rustls 0.23.41", - "socket2 0.5.10", + "socket2 0.6.4", "thiserror 2.0.18", "tokio", "tracing", @@ -4114,7 +4206,7 @@ dependencies = [ "cfg_aliases", "libc", "once_cell", - "socket2 0.5.10", + "socket2 0.6.4", "tracing", "windows-sys 0.60.2", ] @@ -4489,7 +4581,7 @@ dependencies = [ "errno", "libc", "linux-raw-sys", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -4511,6 +4603,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6b92b125634d9b795e7beca796cc790df15a7fb38323bf3196fda83292d06b1f" dependencies = [ "aws-lc-rs", + "log", "once_cell", "ring", "rustls-pki-types", @@ -4559,7 +4652,7 @@ dependencies = [ "security-framework", "security-framework-sys", "webpki-root-certs", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -5164,7 +5257,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "52d1cfed4120b4d927bf7c0f86d2087a4a7d6027c906d9f9d525a80573b9be51" dependencies = [ "libc", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -5517,7 +5610,7 @@ dependencies = [ "getrandom 0.4.3", "once_cell", "rustix", - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] @@ -5797,11 +5890,13 @@ dependencies = [ "socket2 0.6.4", "sync_wrapper", "tokio", + "tokio-rustls 0.26.4", "tokio-stream", "tower", "tower-layer", "tower-service", "tracing", + "webpki-roots", ] [[package]] @@ -5815,6 +5910,17 @@ dependencies = [ "tonic", ] +[[package]] +name = "tonic-types" +version = "0.14.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "73ab1b02061f83d519bba3caa167f88f261ef05720ab8ebc954ade70de3348e8" +dependencies = [ + "prost", + "prost-types", + "tonic", +] + [[package]] name = "tower" version = "0.5.3" @@ -5937,6 +6043,22 @@ dependencies = [ "tracing-core", ] +[[package]] +name = "tracing-opentelemetry" +version = "0.33.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "adbc64cba7137545b8044cb1fe9814f7aacf3c6b5f9b45be8bb5db538befdb26" +dependencies = [ + "js-sys", + "opentelemetry", + "smallvec", + "tracing", + "tracing-core", + "tracing-log", + "tracing-subscriber", + "web-time", +] + [[package]] name = "tracing-serde" version = "0.2.0" @@ -6287,7 +6409,7 @@ version = "0.1.11" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" dependencies = [ - "windows-sys 0.60.2", + "windows-sys 0.61.2", ] [[package]] diff --git a/book/src/SUMMARY.md b/book/src/SUMMARY.md index bef7afc9..ed97014c 100644 --- a/book/src/SUMMARY.md +++ b/book/src/SUMMARY.md @@ -4,10 +4,11 @@ - [Tutorial](./tutorial.md) - [User Guide](./user-guide/README.md) - [Admin Guide](./admin-guide/README.md) - - [Deploying to NixOS](./admin-guide/deployment/nixos.md) - - [Chunking](./admin-guide/chunking.md) + - [Deploying to NixOS](./admin-guide/deployment/nixos.md) + - [OpenTelemetry](./admin-guide/opentelemetry.md) + - [Chunking](./admin-guide/chunking.md) - [FAQs](./faqs.md) - [Reference](./reference/README.md) - - [attic](./reference/attic-cli.md) - - [atticd](./reference/atticd-cli.md) - - [atticadm](./reference/atticadm-cli.md) + - [attic](./reference/attic-cli.md) + - [atticd](./reference/atticd-cli.md) + - [atticadm](./reference/atticadm-cli.md) diff --git a/book/src/admin-guide/README.md b/book/src/admin-guide/README.md index 85cf559e..b09504d3 100644 --- a/book/src/admin-guide/README.md +++ b/book/src/admin-guide/README.md @@ -6,4 +6,5 @@ This section describes how to set up and administer an Attic Server. For a quick start, read the [Tutorial](../tutorial.md). - **[Deploying to NixOS](./deployment/nixos.md)** - Deploying to a NixOS machine +- **[OpenTelemetry](./opentelemetry.md)** - Exporting traces, logs, and metrics with OTLP - **[Chunking](./chunking.md)** - Configuring Content-Defined Chunking data deduplication in Attic diff --git a/book/src/admin-guide/opentelemetry.md b/book/src/admin-guide/opentelemetry.md new file mode 100644 index 00000000..25ae96e5 --- /dev/null +++ b/book/src/admin-guide/opentelemetry.md @@ -0,0 +1,78 @@ +# OpenTelemetry + +`atticd` can export traces, structured logs, and metrics to an OpenTelemetry collector using OTLP. Export is configured with the standard OpenTelemetry environment variables. + +OTLP export is disabled unless a common or signal-specific OTLP endpoint is configured. Set `OTEL_SDK_DISABLED=true` to explicitly disable export even when endpoints are present. Local formatted logs remain enabled in either case. + +## Transport + +Attic supports both OTLP over HTTP with protobuf and OTLP over gRPC. HTTP/protobuf is the default. + +```sh +# Configure an endpoint to enable export. HTTP/protobuf is the default transport. +export OTEL_EXPORTER_OTLP_ENDPOINT=http://localhost:4318 + +# Optional: gRPC +export OTEL_EXPORTER_OTLP_PROTOCOL=grpc +export OTEL_EXPORTER_OTLP_ENDPOINT=http://localhost:4317 +``` + +`OTEL_EXPORTER_OTLP_PROTOCOL` accepts `http/protobuf` or `grpc`. Signal-specific protocol variables are also supported: + +- `OTEL_EXPORTER_OTLP_TRACES_PROTOCOL` +- `OTEL_EXPORTER_OTLP_LOGS_PROTOCOL` +- `OTEL_EXPORTER_OTLP_METRICS_PROTOCOL` + +Signal-specific settings take precedence over `OTEL_EXPORTER_OTLP_PROTOCOL`. + +Export is enabled when any of these endpoint variables has a non-empty value: + +- `OTEL_EXPORTER_OTLP_ENDPOINT` +- `OTEL_EXPORTER_OTLP_TRACES_ENDPOINT` +- `OTEL_EXPORTER_OTLP_LOGS_ENDPOINT` +- `OTEL_EXPORTER_OTLP_METRICS_ENDPOINT` + +When none are configured, Attic does not create exporter workers or attempt connections to the OpenTelemetry SDK's default localhost endpoint. + +## Endpoints and authentication + +Standard common and signal-specific OTLP settings are honored, including: + +- `OTEL_EXPORTER_OTLP_ENDPOINT` +- `OTEL_EXPORTER_OTLP_TRACES_ENDPOINT` +- `OTEL_EXPORTER_OTLP_LOGS_ENDPOINT` +- `OTEL_EXPORTER_OTLP_METRICS_ENDPOINT` +- `OTEL_EXPORTER_OTLP_HEADERS` +- `OTEL_EXPORTER_OTLP_TRACES_HEADERS` +- `OTEL_EXPORTER_OTLP_LOGS_HEADERS` +- `OTEL_EXPORTER_OTLP_METRICS_HEADERS` +- `OTEL_EXPORTER_OTLP_TIMEOUT` + +For HTTP, the exporter appends `/v1/traces`, `/v1/logs`, or `/v1/metrics` to the common endpoint. Signal-specific HTTP endpoints should include the complete signal path. + +## Resource attributes + +The service name defaults to `atticd`. Standard OpenTelemetry resource configuration can add deployment-specific metadata: + +```sh +export OTEL_RESOURCE_ATTRIBUTES='service.namespace=cache,deployment.environment.name=production' +``` + +## Local logging + +OTLP export runs alongside the existing formatted log output. `RUST_LOG` controls local output. OTLP logs and traces use an `info` filter and suppress exporter transport targets to avoid telemetry feedback loops. + +## Exported metrics + +Attic currently emits: + +- `http.server.request.count` +- `http.server.active_requests` +- `http.server.request.duration` +- `attic.operation.count` +- `attic.operation.duration` +- `attic.operation.bytes` + +Operation attributes distinguish uploads, downloads, database connection and migration work, garbage collection runs, and garbage-collected object types. HTTP metrics use method and status attributes; request paths and cache names are intentionally excluded from metric attributes to prevent unbounded cardinality. + +Telemetry providers are flushed during orderly `atticd` shutdown. diff --git a/server/Cargo.toml b/server/Cargo.toml index 5de0b91c..f8b38c33 100644 --- a/server/Cargo.toml +++ b/server/Cargo.toml @@ -44,6 +44,10 @@ http-body-util = "0.1.3" humantime = "2.2.0" humantime-serde = "1.1.1" itoa = "1.0.15" +opentelemetry = { version = "0.32.0", features = ["logs", "metrics", "trace"] } +opentelemetry-appender-tracing = "0.32.0" +opentelemetry-otlp = { version = "0.32.0", default-features = false, features = ["grpc-tonic", "http-proto", "logs", "metrics", "reqwest-client", "reqwest-rustls", "tls-webpki-roots", "trace"] } +opentelemetry_sdk = { version = "0.32.0", features = ["logs", "metrics", "trace"] } rand = "0.10" regex = "1.11.1" ryu = "1.0.20" @@ -56,6 +60,7 @@ toml = "1.1.2" tower-http = { version = "0.7.0", features = [ "catch-panic", "trace" ] } tracing = "0.1.41" tracing-error = "0.2.1" +tracing-opentelemetry = "0.33.0" tracing-subscriber = { version = "0.3.19", features = [ "json" ] } uuid = { version = "1.17.0", features = ["v4"] } console-subscriber = "0.5.0" diff --git a/server/src/api/binary_cache.rs b/server/src/api/binary_cache.rs index 475e18eb..cb45a3f9 100644 --- a/server/src/api/binary_cache.rs +++ b/server/src/api/binary_cache.rs @@ -8,6 +8,7 @@ use std::collections::VecDeque; use std::io::Error as IoError; use std::path::PathBuf; use std::sync::Arc; +use std::time::Instant; use axum::http; use axum::{ @@ -27,6 +28,7 @@ use tracing::instrument; use crate::database::AtticDatabase; use crate::database::entity::chunk::ChunkModel; use crate::error::{ErrorKind, ServerResult}; +use crate::metrics; use crate::narinfo::NarInfo; use crate::nix_manifest; use crate::storage::{Download, StorageBackend, StorageBackendImpl}; @@ -172,6 +174,7 @@ async fn get_nar( Extension(req_state): Extension, Path((cache_name, path)): Path<(CacheName, String)>, ) -> ServerResult { + let started_at = Instant::now(); let components: Vec<&str> = path.splitn(2, '.').collect(); if components.len() != 2 { @@ -192,7 +195,7 @@ async fn get_nar( let database = state.database().await?; - let (object, cache, _nar, chunks) = database + let (object, cache, nar, chunks) = database .find_object_and_chunks_by_store_path_hash(&cache_name, &store_path_hash, true) .await?; @@ -211,6 +214,15 @@ async fn get_nar( database.bump_object_last_accessed(object.id).await?; + metrics::record_operation("download", "success", started_at.elapsed()); + metrics::add_bytes("download.nar", nar.nar_size as u64, &[]); + tracing::info!( + nar.size = nar.nar_size, + chunk.count = chunks.len(), + duration_ms = started_at.elapsed().as_secs_f64() * 1000.0, + "NAR download prepared" + ); + if chunks.len() == 1 { // single chunk let chunk = chunks[0].as_ref().unwrap(); diff --git a/server/src/api/v1/upload_path.rs b/server/src/api/v1/upload_path.rs index 245f4d19..f1fa92bc 100644 --- a/server/src/api/v1/upload_path.rs +++ b/server/src/api/v1/upload_path.rs @@ -3,6 +3,7 @@ use std::io; use std::io::Cursor; use std::marker::Unpin; use std::sync::Arc; +use std::time::Instant; use anyhow::anyhow; use async_compression::Level as CompressionLevel; @@ -30,6 +31,7 @@ use uuid::Uuid; use crate::compression::{CompressionStream, CompressorFn}; use crate::config::CompressionType; use crate::error::{ErrorKind, ServerError, ServerResult}; +use crate::metrics; use crate::narinfo::Compression; use crate::storage::StorageBackend; use crate::{RequestState, State}; @@ -89,6 +91,7 @@ pub(crate) async fn upload_path( headers: HeaderMap, body: Body, ) -> ServerResult> { + let started_at = Instant::now(); let stream = body.into_data_stream(); let mut stream = StreamReader::new(stream.map(|r| r.map_err(|e| io::Error::other(e.to_string())))); @@ -149,6 +152,8 @@ pub(crate) async fn upload_path( let username = req_state.auth.username().map(str::to_string); + let nar_size = upload_info.nar_size; + // Try to acquire a lock on an existing NAR if let Some(existing_nar) = database.find_and_lock_nar(&upload_info.nar_hash).await? { // Deduplicate? @@ -163,7 +168,7 @@ pub(crate) async fn upload_path( if missing_chunk.is_none() { // Can actually be deduplicated - return upload_path_dedup( + let result = upload_path_dedup( username, cache, upload_info, @@ -173,10 +178,44 @@ pub(crate) async fn upload_path( existing_nar, ) .await; + record_upload(&result, nar_size, "deduplicated", started_at); + return result; } } - upload_path_new(username, cache, upload_info, stream, database, &state).await + let result = upload_path_new(username, cache, upload_info, stream, database, &state).await; + let kind = result + .as_ref() + .map(|result| match result.kind { + UploadPathResultKind::Uploaded => "uploaded", + UploadPathResultKind::Deduplicated => "deduplicated", + _ => "unknown", + }) + .unwrap_or("error"); + record_upload(&result, nar_size, kind, started_at); + result +} + +fn record_upload( + result: &ServerResult>, + nar_size: usize, + kind: &'static str, + started_at: Instant, +) { + let outcome = if result.is_ok() { "success" } else { "error" }; + metrics::record_operation("upload", outcome, started_at.elapsed()); + metrics::add_bytes( + "upload.nar", + nar_size as u64, + &[opentelemetry::KeyValue::new("upload.kind", kind)], + ); + tracing::info!( + upload.kind = kind, + nar.size = nar_size, + operation.outcome = outcome, + duration_ms = started_at.elapsed().as_secs_f64() * 1000.0, + "Upload completed" + ); } /// Uploads a path when there is already a matching NAR in the global cache. diff --git a/server/src/gc.rs b/server/src/gc.rs index 2708274c..5d50ac58 100644 --- a/server/src/gc.rs +++ b/server/src/gc.rs @@ -1,7 +1,7 @@ //! Garbage collection. use std::sync::Arc; -use std::time::Duration; +use std::time::{Duration, Instant}; use anyhow::{Result, anyhow}; use chrono::{Duration as ChronoDuration, Utc}; @@ -22,6 +22,7 @@ use crate::database::entity::chunk::{self, ChunkState, Entity as Chunk}; use crate::database::entity::chunkref::{self, Entity as ChunkRef}; use crate::database::entity::nar::{self, Entity as Nar, NarState}; use crate::database::entity::object::{self, Entity as Object}; +use crate::metrics; use crate::storage::StorageBackend; #[derive(Debug, FromQueryResult)] @@ -67,14 +68,21 @@ pub async fn run_garbage_collection(config: Config, shutdown: CancellationToken) /// Runs garbage collection once. #[instrument(skip_all)] pub async fn run_garbage_collection_once(config: Config) -> Result<()> { - tracing::info!("Running garbage collection..."); + tracing::info!("Running garbage collection"); + let started_at = Instant::now(); let state = StateInner::new(config).await; - run_time_based_garbage_collection(&state).await?; - run_reap_orphan_nars(&state).await?; - run_reap_orphan_chunks(&state).await?; + let result = async { + run_time_based_garbage_collection(&state).await?; + run_reap_orphan_nars(&state).await?; + run_reap_orphan_chunks(&state).await?; + Result::<_, anyhow::Error>::Ok(()) + } + .await; - Ok(()) + let outcome = if result.is_ok() { "success" } else { "error" }; + metrics::record_operation("garbage_collection", outcome, started_at.elapsed()); + result } #[instrument(skip_all)] @@ -133,7 +141,8 @@ async fn run_time_based_garbage_collection(state: &State) -> Result<()> { objects_deleted += deletion.rows_affected; } - tracing::info!("Deleted {} objects in total", objects_deleted); + tracing::info!(objects_deleted, "Deleted expired objects"); + metrics::add_count("garbage_collection.objects_deleted", objects_deleted, &[]); Ok(()) } @@ -164,7 +173,12 @@ async fn run_reap_orphan_nars(state: &State) -> Result<()> { .exec(db) .await?; - tracing::info!("Deleted {} orphan NARs", deletion.rows_affected,); + tracing::info!(nars_deleted = deletion.rows_affected, "Deleted orphan NARs"); + metrics::add_count( + "garbage_collection.nars_deleted", + deletion.rows_affected, + &[], + ); Ok(()) } @@ -269,7 +283,15 @@ async fn run_reap_orphan_chunks(state: &State) -> Result<()> { .exec(db) .await?; - tracing::info!("Deleted {} orphan chunks", deletion.rows_affected); + tracing::info!( + chunks_deleted = deletion.rows_affected, + "Deleted orphan chunks" + ); + metrics::add_count( + "garbage_collection.chunks_deleted", + deletion.rows_affected, + &[], + ); Ok(()) } diff --git a/server/src/lib.rs b/server/src/lib.rs index a6fac7df..6e208200 100644 --- a/server/src/lib.rs +++ b/server/src/lib.rs @@ -20,16 +20,18 @@ pub mod config; pub mod database; pub mod error; pub mod gc; +mod metrics; mod middleware; mod narinfo; pub mod nix_manifest; pub mod oobe; mod storage; +pub mod telemetry; use std::net::SocketAddr; use std::sync::Arc; use std::sync::atomic::{AtomicBool, Ordering}; -use std::time::Duration; +use std::time::{Duration, Instant}; use anyhow::Result; use axum::{ @@ -45,7 +47,8 @@ use tokio::sync::OnceCell; use tokio::time; use tokio_util::sync::CancellationToken; use tower_http::catch_panic::CatchPanicLayer; -use tower_http::trace::TraceLayer; +use tower_http::trace::{DefaultOnRequest, TraceLayer}; +use tracing::Level; use access::http::{AuthState, apply_auth}; use attic::cache::CacheName; @@ -109,6 +112,7 @@ impl StateInner { async fn database(&self) -> ServerResult<&DatabaseConnection> { self.database .get_or_try_init(|| async { + let started_at = Instant::now(); let db = Database::connect(&self.config.database.url) .await .map_err(ServerError::database_error); @@ -132,6 +136,11 @@ impl StateInner { .await; } + let outcome = if db.is_ok() { "success" } else { "error" }; + metrics::record_operation("database.connect", outcome, started_at.elapsed()); + if let Err(error) = &db { + tracing::error!(%error, "Database connection failed"); + } db }) .await @@ -243,7 +252,33 @@ pub async fn run_api_server( .layer(axum::middleware::from_fn(init_request_state)) .layer(axum::middleware::from_fn(restrict_host)) .layer(Extension(state.clone())) - .layer(TraceLayer::new_for_http()) + .layer( + TraceLayer::new_for_http() + .make_span_with(|request: &axum::http::Request<_>| { + tracing::info_span!( + "http.request", + otel.name = %format!("{} {}", request.method(), request.uri().path()), + http.request.method = %request.method(), + url.path = %request.uri().path(), + http.response.status_code = tracing::field::Empty, + ) + }) + .on_request(DefaultOnRequest::new().level(Level::DEBUG)) + .on_response( + |response: &axum::http::Response<_>, + latency: Duration, + span: &tracing::Span| { + span.record("http.response.status_code", response.status().as_u16()); + tracing::info!( + parent: span, + http.response.status_code = response.status().as_u16(), + duration_ms = latency.as_secs_f64() * 1000.0, + "HTTP request completed" + ); + }, + ), + ) + .layer(axum::middleware::from_fn(track_http_metrics)) .layer(CatchPanicLayer::new()); eprintln!("Listening on {:?}...", listen); @@ -275,13 +310,33 @@ pub async fn run_api_server( Ok(()) } +async fn track_http_metrics( + request: axum::extract::Request, + next: axum::middleware::Next, +) -> axum::response::Response { + let request_metrics = metrics::start_http_request(request.method().as_str()); + let response = next.run(request).await; + request_metrics.finish(response.status()); + response +} + /// Runs database migrations. pub async fn run_migrations(config: Config) -> Result<()> { - eprintln!("Running migrations..."); - - let state = StateInner::new(config).await; - let db = state.database().await?; - Migrator::up(db, None).await?; + tracing::info!("Running database migrations"); + let started_at = Instant::now(); + + let result = async { + let state = StateInner::new(config).await; + let db = state.database().await?; + Migrator::up(db, None).await?; + Result::<_, anyhow::Error>::Ok(()) + } + .await; - Ok(()) + let outcome = if result.is_ok() { "success" } else { "error" }; + metrics::record_operation("database.migrate", outcome, started_at.elapsed()); + if let Err(error) = &result { + tracing::error!(%error, "Database migration failed"); + } + result } diff --git a/server/src/main.rs b/server/src/main.rs index 1a841bf4..04f7f9dd 100644 --- a/server/src/main.rs +++ b/server/src/main.rs @@ -3,15 +3,12 @@ use std::net::SocketAddr; use std::path::PathBuf; use anyhow::Result; +use attic_server::config; +use attic_server::telemetry; use clap::{Parser, ValueEnum}; use tokio::signal::unix::{SignalKind, signal}; use tokio::task::spawn; use tokio_util::sync::CancellationToken; -use tracing_error::ErrorLayer; -use tracing_subscriber::EnvFilter; -use tracing_subscriber::prelude::*; - -use attic_server::config; /// Nix binary cache server. #[derive(Debug, Parser)] @@ -64,50 +61,59 @@ enum ServerMode { async fn main() -> Result<()> { let opts = Opts::parse(); - init_logging(opts.tokio_console); + let telemetry = telemetry::init(opts.tokio_console)?; dump_version(); - let config = - config::load_config(opts.config.as_deref(), opts.mode == ServerMode::Monolithic).await?; + let result = async { + let config = + config::load_config(opts.config.as_deref(), opts.mode == ServerMode::Monolithic) + .await?; - match opts.mode { - ServerMode::Monolithic => { - attic_server::run_migrations(config.clone()).await?; + match opts.mode { + ServerMode::Monolithic => { + attic_server::run_migrations(config.clone()).await?; - let shutdown = run_shutdown_handler(); - let gc_handle = spawn(attic_server::gc::run_garbage_collection( - config.clone(), - shutdown.clone(), - )); + let shutdown = run_shutdown_handler(); + let gc_handle = spawn(attic_server::gc::run_garbage_collection( + config.clone(), + shutdown.clone(), + )); - let api_server = - attic_server::run_api_server(opts.listen, config.clone(), shutdown.clone()).await; + let api_server = + attic_server::run_api_server(opts.listen, config.clone(), shutdown.clone()) + .await; - shutdown.cancel(); - let _ = gc_handle.await; + shutdown.cancel(); + let _ = gc_handle.await; - api_server?; - } - ServerMode::ApiServer => { - let shutdown = run_shutdown_handler(); - attic_server::run_api_server(opts.listen, config, shutdown).await?; - } - ServerMode::GarbageCollector => { - let shutdown = run_shutdown_handler(); - attic_server::gc::run_garbage_collection(config.clone(), shutdown).await; - } - ServerMode::DbMigrations => { - attic_server::run_migrations(config).await?; - } - ServerMode::GarbageCollectorOnce => { - attic_server::gc::run_garbage_collection_once(config).await?; - } - ServerMode::CheckConfig => { - // config is valid, let's just exit :) + api_server?; + } + ServerMode::ApiServer => { + let shutdown = run_shutdown_handler(); + attic_server::run_api_server(opts.listen, config, shutdown).await?; + } + ServerMode::GarbageCollector => { + let shutdown = run_shutdown_handler(); + attic_server::gc::run_garbage_collection(config.clone(), shutdown).await; + } + ServerMode::DbMigrations => { + attic_server::run_migrations(config).await?; + } + ServerMode::GarbageCollectorOnce => { + attic_server::gc::run_garbage_collection_once(config).await?; + } + ServerMode::CheckConfig => { + // config is valid, let's just exit :) + } } + + Result::<_, anyhow::Error>::Ok(()) } + .await; - Ok(()) + let shutdown_result = telemetry.shutdown(); + result?; + shutdown_result } fn run_shutdown_handler() -> CancellationToken { @@ -142,31 +148,6 @@ fn run_shutdown_handler() -> CancellationToken { shutdown } -fn init_logging(tokio_console: bool) { - let env_filter = EnvFilter::from_default_env(); - let fmt_layer = tracing_subscriber::fmt::layer().with_filter(env_filter); - - let error_layer = ErrorLayer::default(); - - let console_layer = if tokio_console { - let (layer, server) = console_subscriber::ConsoleLayer::new(); - spawn(server.serve()); - Some(layer) - } else { - None - }; - - tracing_subscriber::registry() - .with(fmt_layer) - .with(error_layer) - .with(console_layer) - .init(); - - if tokio_console { - eprintln!("Note: tokio-console is enabled"); - } -} - fn dump_version() { #[cfg(debug_assertions)] eprintln!("Attic Server {} (debug)", env!("CARGO_PKG_VERSION")); diff --git a/server/src/metrics.rs b/server/src/metrics.rs new file mode 100644 index 00000000..aaa0a0e9 --- /dev/null +++ b/server/src/metrics.rs @@ -0,0 +1,150 @@ +use std::sync::OnceLock; +use std::time::Duration; + +use axum::http::StatusCode; +use opentelemetry::KeyValue; +use opentelemetry::global; +use opentelemetry::metrics::{Counter, Histogram, UpDownCounter}; + +struct HttpMetrics { + requests: Counter, + active_requests: UpDownCounter, + request_duration: Histogram, +} + +struct OperationMetrics { + operations: Counter, + duration: Histogram, + bytes: Counter, +} + +pub struct ActiveHttpRequest { + method: String, + started_at: std::time::Instant, +} + +impl ActiveHttpRequest { + pub fn finish(self, status: StatusCode) { + let metrics = http_metrics(); + let attributes = [ + KeyValue::new("http.request.method", self.method.clone()), + KeyValue::new("http.response.status_code", i64::from(status.as_u16())), + KeyValue::new("http.response.status_class", status_class(status)), + ]; + + metrics.active_requests.add( + -1, + &[KeyValue::new("http.request.method", self.method.clone())], + ); + metrics.requests.add(1, &attributes); + metrics + .request_duration + .record(self.started_at.elapsed().as_secs_f64(), &attributes); + } +} + +pub fn start_http_request(method: &str) -> ActiveHttpRequest { + http_metrics().active_requests.add( + 1, + &[KeyValue::new("http.request.method", method.to_owned())], + ); + + ActiveHttpRequest { + method: method.to_owned(), + started_at: std::time::Instant::now(), + } +} + +pub fn record_operation(name: &'static str, outcome: &'static str, duration: Duration) { + let attributes = [ + KeyValue::new("operation.name", name), + KeyValue::new("operation.outcome", outcome), + ]; + let metrics = operation_metrics(); + metrics.operations.add(1, &attributes); + metrics.duration.record(duration.as_secs_f64(), &attributes); +} + +pub fn add_bytes(name: &'static str, bytes: u64, attributes: &[KeyValue]) { + let mut all_attributes = Vec::with_capacity(attributes.len() + 1); + all_attributes.push(KeyValue::new("operation.name", name)); + all_attributes.extend_from_slice(attributes); + operation_metrics().bytes.add(bytes, &all_attributes); +} + +pub fn add_count(name: &'static str, count: u64, attributes: &[KeyValue]) { + let mut all_attributes = Vec::with_capacity(attributes.len() + 1); + all_attributes.push(KeyValue::new("operation.name", name)); + all_attributes.extend_from_slice(attributes); + operation_metrics().operations.add(count, &all_attributes); +} + +fn http_metrics() -> &'static HttpMetrics { + static METRICS: OnceLock = OnceLock::new(); + METRICS.get_or_init(|| { + let meter = global::meter("atticd.http"); + HttpMetrics { + requests: meter + .u64_counter("http.server.request.count") + .with_description("HTTP requests received") + .build(), + active_requests: meter + .i64_up_down_counter("http.server.active_requests") + .with_description("HTTP requests currently being processed") + .build(), + request_duration: meter + .f64_histogram("http.server.request.duration") + .with_description("HTTP request duration") + .with_unit("s") + .build(), + } + }) +} + +fn operation_metrics() -> &'static OperationMetrics { + static METRICS: OnceLock = OnceLock::new(); + METRICS.get_or_init(|| { + let meter = global::meter("atticd.operations"); + OperationMetrics { + operations: meter + .u64_counter("attic.operation.count") + .with_description("Attic operations and processed objects") + .build(), + duration: meter + .f64_histogram("attic.operation.duration") + .with_description("Attic operation duration") + .with_unit("s") + .build(), + bytes: meter + .u64_counter("attic.operation.bytes") + .with_description("Bytes processed by Attic operations") + .with_unit("By") + .build(), + } + }) +} + +pub(crate) fn status_class(status: StatusCode) -> &'static str { + match status.as_u16() / 100 { + 1 => "1xx", + 2 => "2xx", + 3 => "3xx", + 4 => "4xx", + 5 => "5xx", + _ => "unknown", + } +} + +#[cfg(test)] +mod tests { + use axum::http::StatusCode; + + use super::status_class; + + #[test] + fn groups_http_status_by_class() { + assert_eq!(status_class(StatusCode::OK), "2xx"); + assert_eq!(status_class(StatusCode::NOT_FOUND), "4xx"); + assert_eq!(status_class(StatusCode::INTERNAL_SERVER_ERROR), "5xx"); + } +} diff --git a/server/src/telemetry.rs b/server/src/telemetry.rs new file mode 100644 index 00000000..2bd236b9 --- /dev/null +++ b/server/src/telemetry.rs @@ -0,0 +1,327 @@ +use std::env; + +use anyhow::{Context, Result, bail}; +use opentelemetry::global; +use opentelemetry::trace::TracerProvider as _; +use opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge; +use opentelemetry_otlp::{LogExporter, MetricExporter, Protocol, SpanExporter, WithExportConfig}; +use opentelemetry_sdk::Resource; +use opentelemetry_sdk::logs::SdkLoggerProvider; +use opentelemetry_sdk::metrics::SdkMeterProvider; +use opentelemetry_sdk::trace::SdkTracerProvider; +use tracing_error::ErrorLayer; +use tracing_subscriber::layer::SubscriberExt; +use tracing_subscriber::util::SubscriberInitExt; +use tracing_subscriber::{EnvFilter, Layer}; + +const DEFAULT_PROTOCOL: &str = "http/protobuf"; +const OTEL_SDK_DISABLED: &str = "OTEL_SDK_DISABLED"; +const OTLP_ENDPOINT: &str = "OTEL_EXPORTER_OTLP_ENDPOINT"; +const OTLP_PROTOCOL: &str = "OTEL_EXPORTER_OTLP_PROTOCOL"; + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum OtlpProtocol { + HttpProtobuf, + Grpc, +} + +pub struct TelemetryGuard { + tracer_provider: Option, + logger_provider: Option, + meter_provider: Option, +} + +impl TelemetryGuard { + pub fn shutdown(self) -> Result<()> { + let mut errors = Vec::new(); + + if let Some(provider) = self.tracer_provider + && let Err(error) = provider.shutdown() + { + errors.push(format!("tracer provider: {error}")); + } + if let Some(provider) = self.logger_provider + && let Err(error) = provider.shutdown() + { + errors.push(format!("logger provider: {error}")); + } + if let Some(provider) = self.meter_provider + && let Err(error) = provider.shutdown() + { + errors.push(format!("meter provider: {error}")); + } + + if errors.is_empty() { + Ok(()) + } else { + bail!("failed to shut down OpenTelemetry: {}", errors.join("; ")) + } + } +} + +pub fn init(tokio_console: bool) -> Result { + let fmt_layer = tracing_subscriber::fmt::layer().with_filter(EnvFilter::from_default_env()); + let error_layer = ErrorLayer::default(); + let console_layer = if tokio_console { + let (layer, server) = console_subscriber::ConsoleLayer::new(); + tokio::spawn(server.serve()); + Some(layer) + } else { + None + }; + + if !export_configured() { + tracing_subscriber::registry() + .with(fmt_layer) + .with(error_layer) + .with(console_layer) + .try_init() + .context("failed to initialize tracing subscriber")?; + + return Ok(TelemetryGuard { + tracer_provider: None, + logger_provider: None, + meter_provider: None, + }); + } + + let mut resource = Resource::builder().with_attribute(opentelemetry::KeyValue::new( + "service.version", + env!("CARGO_PKG_VERSION"), + )); + if !service_name_configured() { + resource = resource.with_service_name("atticd"); + } + let resource = resource.build(); + + let meter_provider = if signal_export_configured("METRICS") { + let metric_exporter = build_metric_exporter(protocol_for("METRICS")?)?; + let provider = SdkMeterProvider::builder() + .with_resource(resource.clone()) + .with_periodic_exporter(metric_exporter) + .build(); + global::set_meter_provider(provider.clone()); + Some(provider) + } else { + None + }; + + let logger_provider = if signal_export_configured("LOGS") { + Some( + SdkLoggerProvider::builder() + .with_resource(resource.clone()) + .with_batch_exporter(build_log_exporter(protocol_for("LOGS")?)?) + .build(), + ) + } else { + None + }; + let log_layer = logger_provider + .as_ref() + .map(|provider| OpenTelemetryTracingBridge::new(provider).with_filter(otel_filter())); + + let tracer_provider = if signal_export_configured("TRACES") { + let provider = SdkTracerProvider::builder() + .with_resource(resource) + .with_batch_exporter(build_span_exporter(protocol_for("TRACES")?)?) + .build(); + global::set_tracer_provider(provider.clone()); + Some(provider) + } else { + None + }; + let trace_layer = tracer_provider.as_ref().map(|provider| { + tracing_opentelemetry::layer() + .with_tracer(provider.tracer("atticd")) + .with_filter(otel_filter()) + }); + + tracing_subscriber::registry() + .with(fmt_layer) + .with(error_layer) + .with(console_layer) + .with(trace_layer) + .with(log_layer) + .try_init() + .context("failed to initialize tracing subscriber")?; + + Ok(TelemetryGuard { + tracer_provider, + logger_provider, + meter_provider, + }) +} + +fn export_configured() -> bool { + export_enabled( + env::var(OTEL_SDK_DISABLED).is_ok_and(|value| value.eq_ignore_ascii_case("true")), + non_empty_env(OTLP_ENDPOINT).as_deref(), + non_empty_env("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT").as_deref(), + non_empty_env("OTEL_EXPORTER_OTLP_LOGS_ENDPOINT").as_deref(), + non_empty_env("OTEL_EXPORTER_OTLP_METRICS_ENDPOINT").as_deref(), + ) +} + +fn signal_export_configured(signal: &str) -> bool { + non_empty_env(OTLP_ENDPOINT).is_some() + || non_empty_env(&format!("OTEL_EXPORTER_OTLP_{signal}_ENDPOINT")).is_some() +} + +fn export_enabled( + sdk_disabled: bool, + endpoint: Option<&str>, + traces_endpoint: Option<&str>, + logs_endpoint: Option<&str>, + metrics_endpoint: Option<&str>, +) -> bool { + !sdk_disabled + && [endpoint, traces_endpoint, logs_endpoint, metrics_endpoint] + .into_iter() + .flatten() + .any(|endpoint| !endpoint.trim().is_empty()) +} + +fn non_empty_env(key: &str) -> Option { + env::var(key).ok().filter(|value| !value.trim().is_empty()) +} + +fn service_name_configured() -> bool { + env::var("OTEL_SERVICE_NAME").is_ok_and(|value| !value.is_empty()) + || env::var("OTEL_RESOURCE_ATTRIBUTES").is_ok_and(|attributes| { + attributes + .split(',') + .filter_map(|attribute| attribute.split_once('=')) + .any(|(key, value)| key.trim() == "service.name" && !value.trim().is_empty()) + }) +} + +fn protocol_for(signal: &str) -> Result { + let signal_protocol = env::var(format!("OTEL_EXPORTER_OTLP_{signal}_PROTOCOL")).ok(); + let protocol = signal_protocol.or_else(|| env::var(OTLP_PROTOCOL).ok()); + select_otlp_protocol(protocol.as_deref()) +} + +pub fn select_otlp_protocol(protocol: Option<&str>) -> Result { + match protocol.unwrap_or(DEFAULT_PROTOCOL) { + "http/protobuf" => Ok(OtlpProtocol::HttpProtobuf), + "grpc" => Ok(OtlpProtocol::Grpc), + protocol => { + bail!("unsupported OTLP protocol {protocol:?}; expected \"http/protobuf\" or \"grpc\"") + } + } +} + +fn build_span_exporter(protocol: OtlpProtocol) -> Result { + match protocol { + OtlpProtocol::HttpProtobuf => SpanExporter::builder() + .with_http() + .with_protocol(Protocol::HttpBinary) + .build(), + OtlpProtocol::Grpc => SpanExporter::builder().with_tonic().build(), + } + .context("failed to create OTLP trace exporter") +} + +fn build_metric_exporter(protocol: OtlpProtocol) -> Result { + match protocol { + OtlpProtocol::HttpProtobuf => MetricExporter::builder() + .with_http() + .with_protocol(Protocol::HttpBinary) + .build(), + OtlpProtocol::Grpc => MetricExporter::builder().with_tonic().build(), + } + .context("failed to create OTLP metrics exporter") +} + +fn build_log_exporter(protocol: OtlpProtocol) -> Result { + match protocol { + OtlpProtocol::HttpProtobuf => LogExporter::builder() + .with_http() + .with_protocol(Protocol::HttpBinary) + .build(), + OtlpProtocol::Grpc => LogExporter::builder().with_tonic().build(), + } + .context("failed to create OTLP log exporter") +} + +fn otel_filter() -> EnvFilter { + EnvFilter::new("info") + .add_directive("hyper=off".parse().unwrap()) + .add_directive("tonic=off".parse().unwrap()) + .add_directive("h2=off".parse().unwrap()) + .add_directive("reqwest=off".parse().unwrap()) +} + +#[cfg(test)] +mod tests { + use super::{OtlpProtocol, export_enabled, select_otlp_protocol}; + + #[test] + fn export_is_disabled_without_an_endpoint() { + assert!(!export_enabled(false, None, None, None, None)); + } + + #[test] + fn common_endpoint_enables_export() { + assert!(export_enabled( + false, + Some("http://collector:4318"), + None, + None, + None + )); + } + + #[test] + fn signal_endpoint_enables_export() { + assert!(export_enabled( + false, + None, + Some("http://collector:4318/v1/traces"), + None, + None + )); + } + + #[test] + fn sdk_disabled_overrides_endpoints() { + assert!(!export_enabled( + true, + Some("http://collector:4318"), + None, + None, + None + )); + } + + #[test] + fn defaults_to_http_protobuf() { + assert_eq!( + select_otlp_protocol(None).unwrap(), + OtlpProtocol::HttpProtobuf + ); + } + + #[test] + fn supports_http_protobuf() { + assert_eq!( + select_otlp_protocol(Some("http/protobuf")).unwrap(), + OtlpProtocol::HttpProtobuf + ); + } + + #[test] + fn supports_grpc() { + assert_eq!( + select_otlp_protocol(Some("grpc")).unwrap(), + OtlpProtocol::Grpc + ); + } + + #[test] + fn rejects_unsupported_protocols() { + let error = select_otlp_protocol(Some("http/json")).unwrap_err(); + + assert!(error.to_string().contains("http/json")); + } +}