mirror of
https://github.com/rustfs/rustfs.git
synced 2026-08-05 21:07:43 +00:00
rename rustfs-logging to rustfs-obs
This commit is contained in:
@@ -0,0 +1,319 @@
|
||||
#[cfg(feature = "audit-kafka")]
|
||||
use rdkafka::{
|
||||
producer::{FutureProducer, FutureRecord},
|
||||
ClientConfig,
|
||||
};
|
||||
|
||||
#[cfg(feature = "audit-webhook")]
|
||||
use reqwest::Client;
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
use tokio::sync::mpsc;
|
||||
|
||||
/// AuditEntry is a struct that represents an audit entry
|
||||
/// that can be logged
|
||||
/// # Fields
|
||||
/// * `version` - The version of the audit entry
|
||||
/// * `event_type` - The type of event that occurred
|
||||
/// * `bucket` - The bucket that was accessed
|
||||
/// * `object` - The object that was accessed
|
||||
/// * `user` - The user that accessed the object
|
||||
/// * `time` - The time the event occurred
|
||||
/// * `user_agent` - The user agent that accessed the object
|
||||
/// * `span_id` - The span ID of the event
|
||||
/// # Example
|
||||
/// ```
|
||||
/// use rustfs_obs::AuditEntry;
|
||||
/// let entry = AuditEntry {
|
||||
/// version: "1.0".to_string(),
|
||||
/// event_type: "read".to_string(),
|
||||
/// bucket: "bucket".to_string(),
|
||||
/// object: "object".to_string(),
|
||||
/// user: "user".to_string(),
|
||||
/// time: "time".to_string(),
|
||||
/// user_agent: "user_agent".to_string(),
|
||||
/// span_id: "span_id".to_string(),
|
||||
/// };
|
||||
/// ```
|
||||
#[derive(Serialize, Deserialize, Debug, Clone)]
|
||||
pub struct AuditEntry {
|
||||
pub version: String,
|
||||
pub event_type: String,
|
||||
pub bucket: String,
|
||||
pub object: String,
|
||||
pub user: String,
|
||||
pub time: String,
|
||||
pub user_agent: String,
|
||||
pub span_id: String,
|
||||
}
|
||||
|
||||
/// AuditTarget is a trait that defines the interface for audit targets
|
||||
/// that can receive audit entries
|
||||
pub trait AuditTarget: Send + Sync {
|
||||
fn send(&self, entry: AuditEntry);
|
||||
}
|
||||
|
||||
/// FileAuditTarget is an audit target that logs audit entries to a file
|
||||
pub struct FileAuditTarget;
|
||||
|
||||
impl AuditTarget for FileAuditTarget {
|
||||
/// Send an audit entry to a file
|
||||
/// # Arguments
|
||||
/// * `entry` - The audit entry to send
|
||||
///
|
||||
/// # Example
|
||||
/// ```
|
||||
/// use rustfs_obs::{AuditEntry, AuditTarget, FileAuditTarget};
|
||||
/// let entry = AuditEntry {
|
||||
/// version: "1.0".to_string(),
|
||||
/// event_type: "read".to_string(),
|
||||
/// bucket: "bucket".to_string(),
|
||||
/// object: "object".to_string(),
|
||||
/// user: "user".to_string(),
|
||||
/// time: "time".to_string(),
|
||||
/// user_agent: "user_agent".to_string(),
|
||||
/// span_id: "span_id".to_string(),
|
||||
/// };
|
||||
/// FileAuditTarget.send(entry);
|
||||
/// ```
|
||||
fn send(&self, entry: AuditEntry) {
|
||||
println!("File audit: {:?}", entry);
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "audit-webhook")]
|
||||
/// Webhook audit objectives
|
||||
/// #Arguments
|
||||
/// * `client` - The reqwest client
|
||||
/// * `url` - The URL of the webhook
|
||||
/// # Example
|
||||
/// ```
|
||||
/// use rustfs_obs::WebhookAuditTarget;
|
||||
/// let target = WebhookAuditTarget::new("http://localhost:8080");
|
||||
/// ```
|
||||
pub struct WebhookAuditTarget {
|
||||
client: Client,
|
||||
url: String,
|
||||
}
|
||||
|
||||
#[cfg(feature = "audit-webhook")]
|
||||
impl WebhookAuditTarget {
|
||||
pub fn new(url: &str) -> Self {
|
||||
Self {
|
||||
client: Client::new(),
|
||||
url: url.to_string(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "audit-webhook")]
|
||||
impl AuditTarget for WebhookAuditTarget {
|
||||
fn send(&self, entry: AuditEntry) {
|
||||
let client = self.client.clone();
|
||||
let url = self.url.clone();
|
||||
tokio::spawn(async move {
|
||||
if let Err(e) = client.post(&url).json(&entry).send().await {
|
||||
eprintln!("Failed to send to Webhook: {:?}", e);
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "audit-kafka")]
|
||||
/// Kafka audit objectives
|
||||
/// # Arguments
|
||||
/// * `producer` - The Kafka producer
|
||||
/// * `topic` - The Kafka topic
|
||||
/// # Example
|
||||
/// ```
|
||||
/// use rustfs_obs::KafkaAuditTarget;
|
||||
/// let target = KafkaAuditTarget::new("localhost:9092", "rustfs-audit");
|
||||
/// ```
|
||||
/// # Note
|
||||
/// This feature requires the `rdkafka` crate
|
||||
/// # Example
|
||||
/// ```toml
|
||||
/// [dependencies]
|
||||
/// rdkafka = "0.26.0"
|
||||
/// rustfs_obs = { version = "0.1.0", features = ["audit-kafka"] }
|
||||
/// ```
|
||||
/// # Note
|
||||
/// The `rdkafka` crate requires the `librdkafka` library to be installed
|
||||
/// # Example
|
||||
/// ```sh
|
||||
/// sudo apt-get install librdkafka-dev
|
||||
/// ```
|
||||
/// # Note
|
||||
/// The `rdkafka` crate requires the `libssl-dev` and `pkg-config` packages to be installed
|
||||
/// # Example
|
||||
/// ```sh
|
||||
/// sudo apt-get install libssl-dev pkg-config
|
||||
/// ```
|
||||
/// # Note
|
||||
/// The `rdkafka` crate requires the `zlib1g-dev` package to be installed
|
||||
/// # Example
|
||||
/// ```sh
|
||||
/// sudo apt-get install zlib1g-dev
|
||||
/// ```
|
||||
/// # Note
|
||||
/// The `rdkafka` crate requires the `zstd` package to be installed
|
||||
/// # Example
|
||||
/// ```sh
|
||||
/// sudo apt-get install zstd
|
||||
/// ```
|
||||
/// # Note
|
||||
/// The `rdkafka` crate requires the `lz4` package to be installed
|
||||
/// # Example
|
||||
/// ```sh
|
||||
/// sudo apt-get install lz4
|
||||
/// ```
|
||||
pub struct KafkaAuditTarget {
|
||||
producer: FutureProducer,
|
||||
topic: String,
|
||||
}
|
||||
|
||||
#[cfg(feature = "audit-kafka")]
|
||||
impl KafkaAuditTarget {
|
||||
pub fn new(brokers: &str, topic: &str) -> Self {
|
||||
let producer: FutureProducer = ClientConfig::new()
|
||||
.set("bootstrap.servers", brokers)
|
||||
.set("message.timeout.ms", "5000")
|
||||
.create()
|
||||
.expect("Kafka producer creation failed");
|
||||
Self {
|
||||
producer,
|
||||
topic: topic.to_string(),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(feature = "audit-kafka")]
|
||||
impl AuditTarget for KafkaAuditTarget {
|
||||
fn send(&self, entry: AuditEntry) {
|
||||
let topic = self.topic.clone();
|
||||
let span_id = entry.span_id.clone();
|
||||
let payload = serde_json::to_string(&entry).unwrap();
|
||||
// let record = FutureRecord::to(&topic).payload(&payload).key(&span_id);
|
||||
tokio::spawn({
|
||||
// 在异步闭包内部创建 record
|
||||
let topic = topic;
|
||||
let payload = payload;
|
||||
let span_id = span_id;
|
||||
let producer = self.producer.clone();
|
||||
async move {
|
||||
let record = FutureRecord::to(&topic).payload(&payload).key(&span_id);
|
||||
if let Err(e) = producer.send(record, std::time::Duration::from_secs(0)).await {
|
||||
eprintln!("Failed to send to Kafka: {:?}", e);
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
/// AuditLogger is a logger that logs audit entries
|
||||
/// to multiple targets
|
||||
///
|
||||
/// # Example
|
||||
/// ```
|
||||
/// use rustfs_obs::{AuditEntry, AuditLogger, FileAuditTarget};
|
||||
///
|
||||
/// #[tokio::main]
|
||||
/// async fn main() {
|
||||
/// let logger = AuditLogger::new(vec![Box::new(FileAuditTarget)]);
|
||||
/// let entry = AuditEntry {
|
||||
/// version: "1.0".to_string(),
|
||||
/// event_type: "read".to_string(),
|
||||
/// bucket: "bucket".to_string(),
|
||||
/// object: "object".to_string(),
|
||||
/// user: "user".to_string(),
|
||||
/// time: "time".to_string(),
|
||||
/// user_agent: "user_agent".to_string(),
|
||||
/// span_id: "span_id".to_string(),
|
||||
/// };
|
||||
/// logger.log(entry).await;
|
||||
/// }
|
||||
/// ```
|
||||
|
||||
#[derive(Debug)]
|
||||
/// AuditLogger is a logger that logs audit entries
|
||||
/// to multiple targets
|
||||
/// # Example
|
||||
/// ```
|
||||
/// use rustfs_obs::{AuditEntry, AuditLogger, FileAuditTarget};
|
||||
/// let logger = AuditLogger::new(vec![Box::new(FileAuditTarget)]);
|
||||
/// ```
|
||||
/// # Note
|
||||
/// This feature requires the `tokio` crate
|
||||
/// # Example
|
||||
/// ```toml
|
||||
/// [dependencies]
|
||||
/// tokio = { version = "1", features = ["full"] }
|
||||
/// rustfs_obs = { version = "0.1.0"}
|
||||
/// ```
|
||||
/// # Note
|
||||
/// This feature requires the `serde` crate
|
||||
/// # Example
|
||||
/// ```toml
|
||||
/// [dependencies]
|
||||
/// serde = { version = "1", features = ["derive"] }
|
||||
/// rustfs_obs = { version = "0.1.0"}
|
||||
/// ```
|
||||
pub struct AuditLogger {
|
||||
tx: mpsc::Sender<AuditEntry>,
|
||||
}
|
||||
|
||||
impl AuditLogger {
|
||||
/// Create a new AuditLogger with the given targets
|
||||
/// that will receive audit entries
|
||||
/// # Arguments
|
||||
/// * `targets` - A vector of audit targets
|
||||
/// # Returns
|
||||
/// * An AuditLogger
|
||||
/// # Example
|
||||
/// ```
|
||||
/// use rustfs_obs::{AuditLogger, AuditEntry, FileAuditTarget};
|
||||
///
|
||||
/// let logger = AuditLogger::new(vec![Box::new(FileAuditTarget)]);
|
||||
/// ```
|
||||
pub fn new(targets: Vec<Box<dyn AuditTarget>>) -> Self {
|
||||
let (tx, mut rx) = mpsc::channel::<AuditEntry>(1000);
|
||||
tokio::spawn(async move {
|
||||
while let Some(entry) = rx.recv().await {
|
||||
for target in &targets {
|
||||
target.send(entry.clone());
|
||||
}
|
||||
}
|
||||
});
|
||||
Self { tx }
|
||||
}
|
||||
|
||||
/// Log an audit entry
|
||||
/// # Arguments
|
||||
/// * `entry` - The audit entry to log
|
||||
/// # Example
|
||||
/// ```
|
||||
/// use rustfs_obs::{AuditEntry, AuditLogger, FileAuditTarget};
|
||||
///
|
||||
/// #[tokio::main]
|
||||
/// async fn main() {
|
||||
/// let logger = AuditLogger::new(vec![Box::new(FileAuditTarget)]);
|
||||
/// let entry = AuditEntry {
|
||||
/// version: "1.0".to_string(),
|
||||
/// event_type: "read".to_string(),
|
||||
/// bucket: "bucket".to_string(),
|
||||
/// object: "object".to_string(),
|
||||
/// user: "user".to_string(),
|
||||
/// time: "time".to_string(),
|
||||
/// user_agent: "user_agent".to_string(),
|
||||
/// span_id: "span_id".to_string(),
|
||||
/// };
|
||||
/// logger.log(entry).await;
|
||||
/// }
|
||||
/// ```
|
||||
pub async fn log(&self, entry: AuditEntry) {
|
||||
// 将日志消息记录到当前 Span
|
||||
tracing::Span::current()
|
||||
.record("log_message", &entry.bucket)
|
||||
.record("source", &entry.event_type);
|
||||
let _ = self.tx.send(entry).await;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,190 @@
|
||||
//! Logging utilities
|
||||
|
||||
///
|
||||
/// This crate provides utilities for logging.
|
||||
///
|
||||
/// # Examples
|
||||
/// ```
|
||||
/// use rustfs_obs::{log_info, log_error};
|
||||
///
|
||||
/// log_info("This is an informational message");
|
||||
/// log_error("This is an error message");
|
||||
/// ```
|
||||
#[cfg(feature = "audit-kafka")]
|
||||
pub use audit::KafkaAuditTarget;
|
||||
#[cfg(feature = "audit-webhook")]
|
||||
pub use audit::WebhookAuditTarget;
|
||||
pub use audit::{AuditEntry, AuditLogger, AuditTarget, FileAuditTarget};
|
||||
pub use logger::{log_debug, log_error, log_info};
|
||||
pub use telemetry::Telemetry;
|
||||
|
||||
mod audit;
|
||||
mod logger;
|
||||
mod telemetry;
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use crate::{log_info, AuditEntry, AuditLogger, AuditTarget, FileAuditTarget, Telemetry};
|
||||
use chrono::Utc;
|
||||
use opentelemetry::global;
|
||||
use opentelemetry::trace::{TraceContextExt, Tracer};
|
||||
use std::time::{Duration, SystemTime};
|
||||
use tracing::{instrument, Span};
|
||||
use tracing_opentelemetry::OpenTelemetrySpanExt;
|
||||
|
||||
#[instrument(fields(bucket, object, user))]
|
||||
async fn put_object(audit_logger: &AuditLogger, bucket: String, object: String, user: String) {
|
||||
let start_time = SystemTime::now();
|
||||
log_info("Starting PUT operation");
|
||||
|
||||
// Simulate the operation
|
||||
tokio::time::sleep(Duration::from_millis(100)).await;
|
||||
|
||||
// Record Metrics
|
||||
let meter = global::meter("rustfs.rs");
|
||||
let request_duration = meter.f64_histogram("s3_request_duration_seconds").build();
|
||||
request_duration.record(
|
||||
start_time.elapsed().unwrap().as_secs_f64(),
|
||||
&[opentelemetry::KeyValue::new("operation", "put_object")],
|
||||
);
|
||||
|
||||
// Gets the current span
|
||||
let span = Span::current();
|
||||
|
||||
// Use 'OpenTelemetrySpanExt' to get 'SpanContext'
|
||||
let span_context = span.context(); // Get context via OpenTelemetrySpanExt
|
||||
let span_id = span_context.span().span_context().span_id().to_string(); // Get the SpanId
|
||||
|
||||
// Audit events are logged
|
||||
let audit_entry = AuditEntry {
|
||||
version: "1.0".to_string(),
|
||||
event_type: "s3_put_object".to_string(),
|
||||
bucket,
|
||||
object,
|
||||
user,
|
||||
time: Utc::now().to_rfc3339(),
|
||||
user_agent: "rustfs.rs-client".to_string(),
|
||||
span_id,
|
||||
};
|
||||
audit_logger.log(audit_entry).await;
|
||||
|
||||
log_info("PUT operation completed");
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
// #[cfg(feature = "audit-webhook")]
|
||||
// #[cfg(feature = "audit-kafka")]
|
||||
async fn test_main() {
|
||||
let telemetry = Telemetry::init();
|
||||
|
||||
// Initialize multiple audit objectives
|
||||
let audit_targets: Vec<Box<dyn AuditTarget>> = vec![
|
||||
Box::new(FileAuditTarget),
|
||||
// Box::new(KafkaAuditTarget::new("localhost:9092", "rustfs-audit")),
|
||||
// Box::new(WebhookAuditTarget::new("http://localhost:8080/audit")),
|
||||
];
|
||||
let audit_logger = AuditLogger::new(audit_targets);
|
||||
|
||||
// 创建根 Span 并执行操作
|
||||
// let tracer = global::tracer("main");
|
||||
// tracer.in_span("main_operation", |cx| {
|
||||
// Span::current().set_parent(cx);
|
||||
// log_info("Starting test async");
|
||||
// tokio::runtime::Runtime::new().unwrap().block_on(async {
|
||||
log_info("Starting test");
|
||||
// Test the PUT operation
|
||||
put_object(&audit_logger, "my-bucket".to_string(), "my-object.txt".to_string(), "user123".to_string()).await;
|
||||
tokio::time::sleep(Duration::from_millis(100)).await;
|
||||
query_object(&audit_logger, "my-bucket".to_string(), "my-object.txt".to_string(), "user123".to_string()).await;
|
||||
tokio::time::sleep(Duration::from_millis(100)).await;
|
||||
for i in 0..100 {
|
||||
put_object(
|
||||
&audit_logger,
|
||||
format!("my-bucket-{}", i),
|
||||
format!("my-object-{}", i),
|
||||
"user123".to_string(),
|
||||
)
|
||||
.await;
|
||||
tokio::time::sleep(Duration::from_millis(100)).await;
|
||||
query_object(
|
||||
&audit_logger,
|
||||
format!("my-bucket-{}", i),
|
||||
format!("my-object-{}", i),
|
||||
"user123".to_string(),
|
||||
)
|
||||
.await;
|
||||
tokio::time::sleep(Duration::from_millis(100)).await;
|
||||
}
|
||||
|
||||
// Wait for the export to complete
|
||||
tokio::time::sleep(Duration::from_secs(2)).await;
|
||||
log_info("Test completed");
|
||||
// });
|
||||
// });
|
||||
drop(telemetry); // Make sure to clean up
|
||||
}
|
||||
|
||||
#[instrument(fields(bucket, object, user))]
|
||||
async fn query_object(audit_logger: &AuditLogger, bucket: String, object: String, user: String) {
|
||||
let start_time = SystemTime::now();
|
||||
log_info("Starting query operation");
|
||||
|
||||
// Simulate the operation
|
||||
tokio::time::sleep(Duration::from_millis(100)).await;
|
||||
|
||||
// Record Metrics
|
||||
let meter = global::meter("rustfs.rs");
|
||||
let request_duration = meter.f64_histogram("s3_request_duration_seconds").build();
|
||||
request_duration.record(
|
||||
start_time.elapsed().unwrap().as_secs_f64(),
|
||||
&[opentelemetry::KeyValue::new("operation", "query_object")],
|
||||
);
|
||||
|
||||
// Gets the current span
|
||||
let span = Span::current();
|
||||
// Use 'OpenTelemetrySpanExt' to get 'SpanContext'
|
||||
let span_context = span.context(); // Get context via OpenTelemetrySpanExt
|
||||
let span_id = span_context.span().span_context().span_id().to_string(); // Get the SpanId
|
||||
query_one(user.clone());
|
||||
query_two(user.clone());
|
||||
query_three(user.clone());
|
||||
// Audit events are logged
|
||||
let audit_entry = AuditEntry {
|
||||
version: "1.0".to_string(),
|
||||
event_type: "s3_query_object".to_string(),
|
||||
bucket,
|
||||
object,
|
||||
user,
|
||||
time: Utc::now().to_rfc3339(),
|
||||
user_agent: "rustfs.rs-client".to_string(),
|
||||
span_id,
|
||||
};
|
||||
audit_logger.log(audit_entry).await;
|
||||
|
||||
log_info("query operation completed");
|
||||
}
|
||||
#[instrument(fields(user))]
|
||||
fn query_one(user: String) {
|
||||
// 初始化 OpenTelemetry Tracer
|
||||
let tracer = global::tracer("query_one");
|
||||
tracer.in_span("doing_work", |cx| {
|
||||
// Traced app logic here...
|
||||
Span::current().set_parent(cx);
|
||||
log_info("Doing work...");
|
||||
let current_span = Span::current();
|
||||
let span_context = current_span.context();
|
||||
let trace_id = span_context.clone().span().span_context().trace_id().to_string();
|
||||
let span_id = span_context.clone().span().span_context().span_id().to_string();
|
||||
log_info(format!("trace_id: {}, span_id: {}", trace_id, span_id).as_str());
|
||||
});
|
||||
log_info(format!("Starting query_one operation user:{}", user).as_str());
|
||||
}
|
||||
#[instrument(fields(user))]
|
||||
fn query_two(user: String) {
|
||||
log_info(format!("Starting query_two operation user:{}", user).as_str());
|
||||
}
|
||||
#[instrument(fields(user))]
|
||||
fn query_three(user: String) {
|
||||
log_info(format!("Starting query_three operation user: {}", user).as_str());
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,46 @@
|
||||
use tracing::{debug, error, info};
|
||||
|
||||
/// Log an info message
|
||||
///
|
||||
/// # Arguments
|
||||
/// msg: &str - The message to log
|
||||
///
|
||||
/// # Example
|
||||
/// ```
|
||||
/// use rustfs_obs::log_info;
|
||||
///
|
||||
/// log_info("This is an info message");
|
||||
/// ```
|
||||
pub fn log_info(msg: &str) {
|
||||
info!("{}", msg);
|
||||
}
|
||||
|
||||
/// Log an error message
|
||||
///
|
||||
/// # Arguments
|
||||
/// msg: &str - The message to log
|
||||
///
|
||||
/// # Example
|
||||
/// ```
|
||||
/// use rustfs_obs::log_error;
|
||||
///
|
||||
/// log_error("This is an error message");
|
||||
/// ```
|
||||
pub fn log_error(msg: &str) {
|
||||
error!("{}", msg);
|
||||
}
|
||||
|
||||
/// Log a debug message
|
||||
///
|
||||
/// # Arguments
|
||||
/// msg: &str - The message to log
|
||||
///
|
||||
/// # Example
|
||||
/// ```
|
||||
/// use rustfs_obs::log_debug;
|
||||
///
|
||||
/// log_debug("This is a debug message");
|
||||
/// ```
|
||||
pub fn log_debug(msg: &str) {
|
||||
debug!("{}", msg);
|
||||
}
|
||||
@@ -0,0 +1,149 @@
|
||||
use opentelemetry::trace::TracerProvider;
|
||||
use opentelemetry::{global, KeyValue};
|
||||
use opentelemetry_appender_tracing::layer;
|
||||
use opentelemetry_otlp::{self, WithExportConfig};
|
||||
use opentelemetry_sdk::logs::SdkLoggerProvider;
|
||||
use opentelemetry_sdk::{
|
||||
metrics::{MeterProviderBuilder, PeriodicReader, SdkMeterProvider},
|
||||
trace::{RandomIdGenerator, Sampler, SdkTracerProvider},
|
||||
Resource,
|
||||
};
|
||||
use opentelemetry_semantic_conventions::attribute::NETWORK_LOCAL_ADDRESS;
|
||||
use opentelemetry_semantic_conventions::{
|
||||
attribute::{DEPLOYMENT_ENVIRONMENT_NAME, SERVICE_NAME, SERVICE_VERSION},
|
||||
SCHEMA_URL,
|
||||
};
|
||||
use std::time::Duration;
|
||||
use tracing::{info, Level};
|
||||
use tracing_opentelemetry::{MetricsLayer, OpenTelemetryLayer};
|
||||
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt, EnvFilter, Layer};
|
||||
|
||||
/// Telemetry is a wrapper around the OpenTelemetry SDK Tracer and Meter providers.
|
||||
/// It initializes the global Tracer and Meter providers, and sets up the tracing subscriber.
|
||||
/// The Tracer and Meter providers are shut down when the Telemetry instance is dropped.
|
||||
/// This is a convenience struct to ensure that the global providers are properly initialized and shut down.
|
||||
///
|
||||
/// # Example
|
||||
/// ```
|
||||
/// use rustfs_obs::Telemetry;
|
||||
///
|
||||
/// let _telemetry = Telemetry::init();
|
||||
/// ```
|
||||
pub struct Telemetry {
|
||||
tracer_provider: SdkTracerProvider,
|
||||
meter_provider: SdkMeterProvider,
|
||||
}
|
||||
|
||||
impl Telemetry {
|
||||
pub fn init() -> Self {
|
||||
// Define service resource information
|
||||
let resource = Resource::builder()
|
||||
.with_service_name("rustfs-service")
|
||||
.with_schema_url(
|
||||
[
|
||||
KeyValue::new(SERVICE_NAME, "rustfs-service"),
|
||||
KeyValue::new(SERVICE_VERSION, "0.1.0"),
|
||||
KeyValue::new(DEPLOYMENT_ENVIRONMENT_NAME, "develop"),
|
||||
KeyValue::new(NETWORK_LOCAL_ADDRESS, "127.0.0.1"),
|
||||
],
|
||||
SCHEMA_URL,
|
||||
)
|
||||
.build();
|
||||
|
||||
let tracer_exporter = opentelemetry_otlp::SpanExporter::builder()
|
||||
.with_tonic()
|
||||
.with_endpoint("http://localhost:4317")
|
||||
.with_protocol(opentelemetry_otlp::Protocol::Grpc)
|
||||
.with_timeout(Duration::from_secs(3))
|
||||
.build()
|
||||
.unwrap();
|
||||
|
||||
// Configure Tracer Provider
|
||||
let tracer_provider = SdkTracerProvider::builder()
|
||||
.with_sampler(Sampler::AlwaysOn)
|
||||
.with_id_generator(RandomIdGenerator::default())
|
||||
.with_resource(resource.clone())
|
||||
.with_batch_exporter(tracer_exporter)
|
||||
// .with_simple_exporter(opentelemetry_stdout::SpanExporter::default())
|
||||
.build();
|
||||
|
||||
let meter_exporter = opentelemetry_otlp::MetricExporter::builder()
|
||||
.with_tonic()
|
||||
.with_endpoint("http://localhost:4317")
|
||||
.with_protocol(opentelemetry_otlp::Protocol::Grpc)
|
||||
.with_timeout(Duration::from_secs(3))
|
||||
.with_temporality(opentelemetry_sdk::metrics::Temporality::default())
|
||||
.build()
|
||||
.unwrap();
|
||||
|
||||
let meter_reader = PeriodicReader::builder(meter_exporter)
|
||||
.with_interval(Duration::from_secs(30))
|
||||
.build();
|
||||
|
||||
// For debugging in development
|
||||
// let meter_stdout_reader = PeriodicReader::builder(opentelemetry_stdout::MetricExporter::default()).build();
|
||||
|
||||
// Configure Meter Provider
|
||||
let meter_provider = MeterProviderBuilder::default()
|
||||
.with_resource(resource.clone())
|
||||
.with_reader(meter_reader)
|
||||
// .with_reader(meter_stdout_reader)
|
||||
.build();
|
||||
|
||||
// Set global Tracer and Meter providers
|
||||
global::set_tracer_provider(tracer_provider.clone());
|
||||
global::set_meter_provider(meter_provider.clone());
|
||||
|
||||
let tracer = tracer_provider.tracer("rustfs-service");
|
||||
|
||||
// // let _stdout_exporter = opentelemetry_stdout::LogExporter::default();
|
||||
// let otlp_exporter = opentelemetry_otlp::LogExporter::builder()
|
||||
// .with_tonic()
|
||||
// .with_endpoint("http://localhost:4317")
|
||||
// // .with_timeout(Duration::from_secs(3))
|
||||
// // .with_protocol(opentelemetry_otlp::Protocol::Grpc)
|
||||
// .build()
|
||||
// .unwrap();
|
||||
// let provider: SdkLoggerProvider = SdkLoggerProvider::builder()
|
||||
// .with_resource(resource.clone())
|
||||
// .with_simple_exporter(otlp_exporter)
|
||||
// .build();
|
||||
// let filter_otel = EnvFilter::new("debug")
|
||||
// // .add_directive("hyper=off".parse().unwrap())
|
||||
// // .add_directive("opentelemetry=off".parse().unwrap())
|
||||
// // .add_directive("tonic=off".parse().unwrap())
|
||||
// // .add_directive("h2=off".parse().unwrap())
|
||||
// .add_directive("reqwest=off".parse().unwrap());
|
||||
// let otel_layer = layer::OpenTelemetryTracingBridge::new(&provider).with_filter(filter_otel);
|
||||
let filter_fmt = EnvFilter::new("info").add_directive("opentelemetry=debug".parse().unwrap());
|
||||
let fmt_layer = tracing_subscriber::fmt::layer()
|
||||
.with_thread_names(true)
|
||||
.with_filter(filter_fmt);
|
||||
|
||||
// Configure `tracing subscriber`
|
||||
tracing_subscriber::registry()
|
||||
.with(tracing_subscriber::filter::LevelFilter::from_level(Level::DEBUG))
|
||||
.with(tracing_subscriber::fmt::layer().with_ansi(true))
|
||||
.with(MetricsLayer::new(meter_provider.clone()))
|
||||
.with(OpenTelemetryLayer::new(tracer))
|
||||
// .with(otel_layer)
|
||||
.with(fmt_layer)
|
||||
.init();
|
||||
info!(name: "my-event-name", target: "my-system", event_id = 20, user_name = "otel", user_email = "otel@opentelemetry.io", message = "This is an example message");
|
||||
Self {
|
||||
tracer_provider,
|
||||
meter_provider,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl Drop for Telemetry {
|
||||
fn drop(&mut self) {
|
||||
if let Err(err) = self.tracer_provider.shutdown() {
|
||||
eprintln!("{err:?}");
|
||||
}
|
||||
if let Err(err) = self.meter_provider.shutdown() {
|
||||
eprintln!("{err:?}");
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user