Skip to content

Commit b634cf3

Browse files
committed
add observability with otel
1 parent 0dec4ee commit b634cf3

78 files changed

Lines changed: 1241 additions & 304 deletions

Some content is hidden

Large Commits have some content hidden by default. Use the searchbox below for content that may be hidden.

Cargo.lock

Lines changed: 448 additions & 31 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

Cargo.toml

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,12 @@ readme = "README.md"
1414
[workspace.dependencies]
1515
sqlx = { version = "0.8.6", features = ["chrono", "ipnetwork", "json", "postgres", "runtime-tokio", "tls-rustls-ring-webpki"] }
1616
ipnetwork = { version = "=0.20.0", features = ["serde"] } # needs to be in sync with the version sqlx uses^
17+
opentelemetry = "0.31.0"
18+
opentelemetry_sdk = { version = "0.31.0", features = ["rt-tokio"] }
19+
opentelemetry-otlp = { version = "0.31.1", features = ["grpc-tonic"] }
20+
opentelemetry-semantic-conventions = "0.31.0"
21+
opentelemetry-appender-tracing = { version = "0.31.1", default-features = false, features = ["log"] }
22+
tracing-opentelemetry = "0.32.1"
1723

1824
[workspace.lints.clippy]
1925
pedantic = { level = "warn", priority = -1 }

docker-compose.observability.yml

Lines changed: 65 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,65 @@
1+
services:
2+
otel-collector:
3+
image: otel/opentelemetry-collector-contrib:0.150.0
4+
command: ["--config=/etc/otel-collector.yml"]
5+
volumes:
6+
- ./observability/otel-collector.yml:/etc/otel-collector.yml:ro
7+
ports:
8+
- "4317:4317" # gRPC
9+
- "4318:4318" # HTTP
10+
- "8889:8889" # Prometheus
11+
restart: unless-stopped
12+
13+
prometheus:
14+
image: prom/prometheus:v3.11.2
15+
command:
16+
- "--config.file=/etc/prometheus/prometheus.yml"
17+
- "--storage.tsdb.path=/prometheus"
18+
- "--web.enable-remote-write-receiver"
19+
volumes:
20+
- ./observability/prometheus.yml:/etc/prometheus/prometheus.yml:ro
21+
- prometheus_data:/prometheus
22+
expose:
23+
- "9090"
24+
restart: unless-stopped
25+
26+
jaeger:
27+
image: jaegertracing/all-in-one:1.76.0
28+
environment:
29+
- COLLECTOR_OTLP_ENABLED=true
30+
volumes:
31+
- jaeger_data:/badger
32+
ports:
33+
- "16686:16686"
34+
restart: unless-stopped
35+
36+
loki:
37+
image: grafana/loki:3.7.1
38+
command: -config.file=/etc/loki/local-config.yaml
39+
volumes:
40+
- loki_data:/loki
41+
expose:
42+
- "3100"
43+
restart: unless-stopped
44+
45+
grafana:
46+
image: grafana/grafana:13.0.1
47+
environment:
48+
- GF_SECURITY_ADMIN_PASSWORD=admin
49+
- GF_USERS_ALLOW_SIGN_UP=false
50+
volumes:
51+
- ./observability/grafana/provisioning:/etc/grafana/provisioning:ro
52+
- grafana_data:/var/lib/grafana
53+
ports:
54+
- "3000:3000"
55+
depends_on:
56+
- prometheus
57+
- jaeger
58+
- loki
59+
restart: unless-stopped
60+
61+
volumes:
62+
prometheus_data:
63+
grafana_data:
64+
jaeger_data:
65+
loki_data:

gitarena-common/Cargo.toml

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,11 @@ num-derive = "0.3.3"
2323
num-traits = "0.2.14"
2424
num_cpus = "1.13.1"
2525
once_cell = "1.9.0"
26+
opentelemetry = { workspace = true }
27+
opentelemetry_sdk = { workspace = true }
28+
opentelemetry-otlp = { workspace = true }
29+
opentelemetry-appender-tracing = { workspace = true }
30+
tracing-opentelemetry = { workspace = true }
2631
serde = { version = "1.0.133", features = ["derive"] }
2732
sqlx = { workspace = true }
2833
tokio = { version = "1.15.0", features = ["full", "tracing"] }

gitarena-common/src/database.rs

Lines changed: 5 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -4,10 +4,10 @@ use std::str::FromStr;
44
use std::time::Duration;
55

66
use anyhow::{Context, Result, anyhow, bail};
7-
use log::info;
87
use once_cell::sync::OnceCell;
98
use sqlx::{Executor, Postgres};
109
use tokio::fs;
10+
use tracing::{info, instrument};
1111

1212
pub mod models;
1313

@@ -18,10 +18,10 @@ pub type ConnectOptions = PgConnectOptions;
1818
pub type Database = Postgres;
1919
pub type DatabaseError = PgDatabaseError;
2020

21-
// These type aliases get their values from above so will not need to be #[cfg]'ed
22-
pub type Pool = sqlx::pool::Pool<Database>;
21+
pub type Pool = sqlx::Pool<Database>;
2322
pub type PoolOptions = sqlx::pool::PoolOptions<Database>;
2423

24+
#[instrument(err)]
2525
pub async fn create_postgres_pool(module: &'static str, max_conns: Option<u32>) -> Result<Pool> {
2626
static ONCE: OnceCell<String> = OnceCell::new();
2727

@@ -43,9 +43,11 @@ pub async fn create_postgres_pool(module: &'static str, max_conns: Option<u32>)
4343
.await?;
4444

4545
sqlx::migrate!("../migrations").run(&pool).await?;
46+
4647
Ok(pool)
4748
}
4849

50+
#[instrument(err)]
4951
async fn read_database_config() -> Result<ConnectOptions> {
5052
let mut options = match (env::var_os("DATABASE_URL"), env::var_os("DATABASE_URL_FILE")) {
5153
(Some(url), None) => {

gitarena-common/src/ipc.rs

Lines changed: 5 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -1,8 +1,8 @@
11
use std::fmt::{Debug, Display, Formatter};
22
use std::io::{Read, Write};
3-
use std::{fmt, fs, mem};
3+
use std::{fmt, mem};
44

5-
use anyhow::{Context, Result};
5+
use anyhow::Result;
66
use bincode::config::{AllowTrailing, Bounded, LittleEndian, VarintEncoding, WithOtherEndian, WithOtherIntEncoding, WithOtherLimit, WithOtherTrailing};
77
use bincode::{DefaultOptions, Options as _};
88
use serde::de::DeserializeOwned;
@@ -33,11 +33,12 @@ impl<T: Serialize + Sized + PacketId> IpcPacket<T> {
3333
}
3434

3535
impl<T: Sized> IpcPacket<T> {
36-
/// Maximum size that this struct can be serialized from (mem::size_of::<Self> + 1 MB)
36+
/// Maximum size that this struct can be serialized from (`mem::size_of::<Self>` + 1 MB)
3737
#[inline]
38+
#[must_use]
3839
pub const fn max_size() -> u64 {
3940
// Allow 1 MB additional limit
40-
mem::size_of::<T>() as u64 + 1_000_000
41+
size_of::<T>() as u64 + 1_000_000
4142
}
4243

4344
#[inline]
@@ -97,15 +98,10 @@ pub trait PacketId {
9798
}
9899

99100
/// Cross-platform way to get the socket/pipe path.
100-
///
101-
/// # Side effects
102-
///
103-
/// On *non-Windows systems* this function exhibits side effects (creation of directory `/run/gitarena`)
104101
pub fn ipc_path() -> Result<&'static str> {
105102
Ok(if cfg!(windows) {
106103
r"\\.\pipe\gitarena-workhorse"
107104
} else {
108-
fs::create_dir_all("/run/gitarena").context("Failed to create directory")?;
109105
"/run/gitarena/workhorse"
110106
})
111107
}

gitarena-common/src/lib.rs

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,3 +3,4 @@ pub mod ipc;
33
pub mod log;
44
pub mod packets;
55
pub mod prelude;
6+
pub mod telemetry;

gitarena-common/src/log.rs

Lines changed: 14 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -4,11 +4,14 @@ use std::path::Path;
44
use std::{env, fs, io};
55

66
use anyhow::{Context, Result};
7-
use log::debug;
7+
use opentelemetry::global;
8+
use opentelemetry_appender_tracing::layer::OpenTelemetryTracingBridge;
9+
use opentelemetry_sdk::logs::SdkLoggerProvider;
810
use tracing::Subscriber;
911
use tracing::metadata::LevelFilter;
1012
use tracing_appender::non_blocking::WorkerGuard;
1113
use tracing_appender::rolling;
14+
use tracing_opentelemetry::OpenTelemetryLayer;
1215
use tracing_subscriber::filter::FromEnvError;
1316
use tracing_subscriber::fmt::Layer;
1417
use tracing_subscriber::layer::SubscriberExt;
@@ -21,6 +24,7 @@ pub fn init_logger(
2124
module: &str,
2225
directives: &'static [&str],
2326
additional: Option<Box<dyn layer::Layer<Registry> + Send + Sync + 'static>>,
27+
logger_provider: Option<&SdkLoggerProvider>,
2428
) -> Result<Vec<WorkerGuard>> {
2529
let mut guards = Vec::new();
2630

@@ -38,21 +42,27 @@ pub fn init_logger(
3842

3943
let (env_filter, tokio_console_layer) = tokio_console(env_filter);
4044

45+
let otel_tracing_layer = OpenTelemetryLayer::new(global::tracer(module.to_owned()));
46+
let otel_log_bridge = logger_provider.map(OpenTelemetryTracingBridge::new);
47+
4148
// https://stackoverflow.com/a/66138267
4249
Registry::default()
4350
.with(additional)
4451
.with(env_filter)
4552
.with(stdout_layer)
4653
.with(file_layer)
4754
.with(tokio_console_layer)
55+
.with(otel_tracing_layer)
56+
.with(otel_log_bridge)
4857
.try_init()
4958
.context("Failed to initialize logger")?;
5059

51-
debug!("Successfully initialized logger for {}", module);
60+
tracing::debug!("Successfully initialized logger for {module}");
5261

5362
Ok(guards)
5463
}
5564

65+
#[must_use]
5666
pub fn stdout<S: Subscriber + for<'a> LookupSpan<'a>>() -> Option<(impl layer::Layer<S>, WorkerGuard)> {
5767
if env::var_os("NO_STDOUT_LOG").is_some() {
5868
return None;
@@ -98,11 +108,11 @@ pub fn tokio_console<S: Subscriber + for<'a> LookupSpan<'a>>(filter: EnvFilter)
98108
(filter, Some(layer))
99109
}
100110

111+
#[must_use]
101112
pub fn default_env(err: FromEnvError, directives: &[&str]) -> EnvFilter {
102113
let not_found = err
103114
.source()
104-
.map(|o| o.downcast_ref::<VarError>().map_or_else(|| false, |err| matches!(err, VarError::NotPresent)))
105-
.unwrap_or(false);
115+
.is_some_and(|o| o.downcast_ref::<VarError>().map_or_else(|| false, |err| matches!(err, VarError::NotPresent)));
106116

107117
if !not_found {
108118
eprintln!(

gitarena-common/src/telemetry.rs

Lines changed: 94 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,94 @@
1+
use std::env;
2+
3+
use anyhow::Result;
4+
use opentelemetry::global;
5+
use opentelemetry_otlp::{LogExporter, MetricExporter, SpanExporter};
6+
use opentelemetry_sdk::Resource;
7+
use opentelemetry_sdk::logs::BatchLogProcessor;
8+
use opentelemetry_sdk::logs::SdkLoggerProvider;
9+
use opentelemetry_sdk::metrics::{PeriodicReader, SdkMeterProvider};
10+
use opentelemetry_sdk::trace::{BatchSpanProcessor, SdkTracerProvider};
11+
12+
/// dropping this struct causes the observability stack to be flushed
13+
pub struct TelemetryGuards {
14+
tracer: Option<SdkTracerProvider>,
15+
meter: Option<SdkMeterProvider>,
16+
logger: SdkLoggerProvider,
17+
}
18+
19+
impl TelemetryGuards {
20+
/// Whether this guard even guards something lol
21+
pub fn is_guarding(&self) -> bool {
22+
self.tracer.is_some() && self.meter.is_some()
23+
}
24+
}
25+
26+
impl Drop for TelemetryGuards {
27+
fn drop(&mut self) {
28+
if let Some(ref p) = self.tracer {
29+
let _ = p.shutdown();
30+
}
31+
if let Some(ref p) = self.meter {
32+
let _ = p.shutdown();
33+
}
34+
let _ = self.logger.shutdown();
35+
}
36+
}
37+
38+
fn is_configured() -> bool {
39+
env::var_os("OTEL_EXPORTER_OTLP_ENDPOINT").is_some()
40+
|| env::var_os("OTEL_EXPORTER_OTLP_TRACES_ENDPOINT").is_some()
41+
|| env::var_os("OTEL_EXPORTER_OTLP_METRICS_ENDPOINT").is_some()
42+
|| env::var_os("OTEL_EXPORTER_OTLP_LOGS_ENDPOINT").is_some()
43+
}
44+
45+
/// needs to be called before `init_logger`
46+
pub fn init(service_name: &'static str) -> Result<(TelemetryGuards, SdkLoggerProvider)> {
47+
if !is_configured() {
48+
let noop_logger = SdkLoggerProvider::builder().build();
49+
let clone = noop_logger.clone();
50+
return Ok((
51+
TelemetryGuards {
52+
tracer: None,
53+
meter: None,
54+
logger: noop_logger,
55+
},
56+
clone,
57+
));
58+
}
59+
60+
let resource = Resource::builder()
61+
.with_service_name(env::var("OTEL_SERVICE_NAME").unwrap_or_else(|_| service_name.to_owned()))
62+
.build();
63+
64+
let span_exporter = SpanExporter::builder().with_tonic().build()?;
65+
let tracer_provider = SdkTracerProvider::builder()
66+
.with_resource(resource.clone())
67+
.with_span_processor(BatchSpanProcessor::builder(span_exporter).build())
68+
.build();
69+
global::set_tracer_provider(tracer_provider.clone());
70+
71+
let metric_exporter = MetricExporter::builder().with_tonic().build()?;
72+
let meter_provider = SdkMeterProvider::builder()
73+
.with_resource(resource.clone())
74+
.with_reader(PeriodicReader::builder(metric_exporter).build())
75+
.build();
76+
global::set_meter_provider(meter_provider.clone());
77+
78+
let log_exporter = LogExporter::builder().with_tonic().build()?;
79+
let logger_provider = SdkLoggerProvider::builder()
80+
.with_resource(resource)
81+
.with_log_processor(BatchLogProcessor::builder(log_exporter).build())
82+
.build();
83+
84+
let logger_for_bridge = logger_provider.clone();
85+
86+
Ok((
87+
TelemetryGuards {
88+
tracer: Some(tracer_provider),
89+
meter: Some(meter_provider),
90+
logger: logger_provider,
91+
},
92+
logger_for_bridge,
93+
))
94+
}

gitarena-workhorse/Cargo.toml

Lines changed: 0 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -15,7 +15,6 @@ console-subscriber = { version = "0.1.3", features = ["parking_lot"] }
1515
futures = "0.3.19"
1616
futures-locks = "0.7.0"
1717
gitarena-common = { version = "0.0.0", path = "../gitarena-common" }
18-
log = "0.4.14"
1918
parity-tokio-ipc = "0.9.0"
2019
tokio = { version = "1.15.0", features = ["full", "tracing"] }
2120
tracing = "0.1.29"

0 commit comments

Comments
 (0)