Implement initialization for packages/logging

- Add initialization logic for the `rustfs-logging` crate.
- Provide examples for logging utilities.
- Include tests for logging and telemetry functionalities.
- Ensure proper configuration of dependencies in `Cargo.toml`.
This commit is contained in:
houseme
2025-02-21 13:46:35 +08:00
parent 391eb5b6f9
commit 03589201fb
9 changed files with 587 additions and 7 deletions
+35
View File
@@ -0,0 +1,35 @@
[package]
name = "rustfs-logging"
edition.workspace = true
license.workspace = true
repository.workspace = true
rust-version.workspace = true
version.workspace = true
[lints]
workspace = true
[dependencies]
opentelemetry = { workspace = true }
opentelemetry_sdk = { workspace = true, features = ["rt-tokio"] }
opentelemetry-stdout = { workspace = true }
opentelemetry-otlp = { workspace = true, features = ["grpc-tonic", "metrics"] }
opentelemetry-semantic-conventions = { workspace = true, features = ["semconv_experimental"] }
serde = { workspace = true }
tracing = { workspace = true }
tracing-opentelemetry = { workspace = true }
tracing-subscriber = { workspace = true, features = ["fmt", "env-filter", "tracing-log", "time", "local-time", "json"] }
tokio = { workspace = true, features = ["full"] }
[dev-dependencies]
chrono = { workspace = true }
opentelemetry = { workspace = true, features = ["trace", "metrics"] }
opentelemetry_sdk = { workspace = true, features = ["trace", "rt-tokio"] }
opentelemetry-stdout = { workspace = true, features = ["trace", "metrics"] }
opentelemetry-otlp = { workspace = true, features = ["metrics", "grpc-tonic"] }
opentelemetry-semantic-conventions = { workspace = true, features = ["semconv_experimental"] }
tokio = { workspace = true, features = ["full"] }
tracing = { workspace = true, features = ["std", "attributes"] }
tracing-subscriber = { workspace = true, features = ["registry", "std", "fmt"] }
+153
View File
@@ -0,0 +1,153 @@
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_logging::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_logging::{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);
}
}
/// AuditLogger is a logger that logs audit entries
/// to multiple targets
///
/// # Example
/// ```
/// use rustfs_logging::{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)]
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_logging::{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_logging::{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) {
let _ = self.tx.send(entry).await;
}
}
+83
View File
@@ -0,0 +1,83 @@
//! Logging utilities
///
/// This crate provides utilities for logging.
///
/// # Examples
/// ```
/// use rustfs_logging::{log_info, log_error};
///
/// log_info("This is an informational message");
/// log_error("This is an error message");
/// ```
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 super::*;
use opentelemetry::global;
use opentelemetry::trace::TraceContextExt;
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: chrono::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]
async fn test_main() {
let telemetry = Telemetry::init();
// Initialize multiple audit objectives
let audit_targets: Vec<Box<dyn AuditTarget>> = vec![Box::new(FileAuditTarget)];
let audit_logger = AuditLogger::new(audit_targets);
// Test the PUT operation
put_object(&audit_logger, "my-bucket".to_string(), "my-object.txt".to_string(), "user123".to_string()).await;
// Wait for the export to complete
tokio::time::sleep(Duration::from_secs(2)).await;
drop(telemetry); // Make sure to clean up
}
}
+46
View File
@@ -0,0 +1,46 @@
use tracing::{debug, error, info};
/// Log an info message
///
/// # Arguments
/// msg: &str - The message to log
///
/// # Example
/// ```
/// use rustfs_logging::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_logging::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_logging::log_debug;
///
/// log_debug("This is a debug message");
/// ```
pub fn log_debug(msg: &str) {
debug!("{}", msg);
}
+113
View File
@@ -0,0 +1,113 @@
use opentelemetry::trace::TracerProvider;
use opentelemetry::{global, KeyValue};
use opentelemetry_otlp::{self, WithExportConfig};
use opentelemetry_sdk::{
metrics::{MeterProviderBuilder, PeriodicReader, SdkMeterProvider},
trace::{RandomIdGenerator, Sampler, SdkTracerProvider},
Resource,
};
use opentelemetry_semantic_conventions::{
attribute::{DEPLOYMENT_ENVIRONMENT_NAME, SERVICE_NAME, SERVICE_VERSION},
SCHEMA_URL,
};
use tracing::Level;
use tracing_opentelemetry::{MetricsLayer, OpenTelemetryLayer};
use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt};
/// 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_logging::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_schema_url(
[
KeyValue::new(SERVICE_NAME, env!("CARGO_PKG_NAME")),
KeyValue::new(SERVICE_VERSION, env!("CARGO_PKG_VERSION")),
KeyValue::new(DEPLOYMENT_ENVIRONMENT_NAME, "develop"),
],
SCHEMA_URL,
)
.build();
let tracer_exporter = opentelemetry_otlp::SpanExporter::builder()
.with_tonic()
.with_endpoint("http://localhost:4317")
.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_temporality(opentelemetry_sdk::metrics::Temporality::default())
.build()
.unwrap();
let meter_reader = PeriodicReader::builder(meter_exporter)
.with_interval(std::time::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)
.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("tracing-otel-subscriber-rustfs-service");
// 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))
.init();
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:?}");
}
}
}