From ef162396cfb66776753258a80e37286e9f6058da Mon Sep 17 00:00:00 2001 From: houseme Date: Sun, 16 Mar 2025 16:33:52 +0800 Subject: [PATCH] improve code --- Cargo.lock | 1 + packages/obs/Cargo.toml | 5 +- packages/obs/examples/config.toml | 4 +- packages/obs/examples/server.rs | 69 +-- packages/obs/src/config.rs | 12 +- packages/obs/src/entry/audit.rs | 88 ++++ packages/obs/src/entry/base.rs | 114 +++++ packages/obs/src/{entry.rs => entry/log.rs} | 3 + packages/obs/src/entry/mod.rs | 3 + packages/obs/src/lib.rs | 43 +- packages/obs/src/logger.rs | 464 +++++++++++++++----- packages/obs/src/sink.rs | 7 +- packages/obs/src/telemetry.rs | 107 +++-- packages/obs/src/worker.rs | 2 +- rustfs/src/main.rs | 8 +- 15 files changed, 724 insertions(+), 206 deletions(-) create mode 100644 packages/obs/src/entry/audit.rs create mode 100644 packages/obs/src/entry/base.rs rename packages/obs/src/{entry.rs => entry/log.rs} (95%) create mode 100644 packages/obs/src/entry/mod.rs diff --git a/Cargo.lock b/Cargo.lock index ca81fe744..f619b7a52 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5744,6 +5744,7 @@ dependencies = [ "tokio", "tracing", "tracing-core", + "tracing-error", "tracing-opentelemetry", "tracing-subscriber", ] diff --git a/packages/obs/Cargo.toml b/packages/obs/Cargo.toml index 216e4296a..86a9fb6c4 100644 --- a/packages/obs/Cargo.toml +++ b/packages/obs/Cargo.toml @@ -11,7 +11,7 @@ workspace = true [features] default = ["file"] -kafka = ["dep:rdkafka", "dep:serde_json"] +kafka = ["dep:rdkafka"] webhook = ["dep:reqwest"] file = [] @@ -28,12 +28,13 @@ opentelemetry-semantic-conventions = { workspace = true, features = ["semconv_ex serde = { workspace = true } tracing = { workspace = true } tracing-core = { workspace = true } +tracing-error = { workspace = true } tracing-opentelemetry = { workspace = true } tracing-subscriber = { workspace = true, features = ["fmt", "env-filter", "tracing-log", "time", "local-time", "json"] } tokio = { workspace = true, features = ["sync", "fs", "rt-multi-thread"] } rdkafka = { workspace = true, features = ["tokio"], optional = true } reqwest = { workspace = true, optional = true, default-features = false } -serde_json = { workspace = true, optional = true } +serde_json = { workspace = true } thiserror = { workspace = true } local-ip-address = { workspace = true } diff --git a/packages/obs/examples/config.toml b/packages/obs/examples/config.toml index 9c3c68a26..2ceac4397 100644 --- a/packages/obs/examples/config.toml +++ b/packages/obs/examples/config.toml @@ -3,7 +3,7 @@ endpoint = "http://localhost:4317" use_stdout = true sample_ratio = 0.5 meter_interval = 30 -service_name = "logging_service" +service_name = "rustfs_obs_service" service_version = "0.1.0" deployment_environment = "develop" @@ -23,7 +23,7 @@ batch_timeout_ms = 1000 # Default is 100ms if not specified [sinks.file] enabled = true -path = "app.log" +path = "logs/app.log" batch_size = 100 batch_timeout_ms = 1000 # Default is 8192 bytes if not specified diff --git a/packages/obs/examples/server.rs b/packages/obs/examples/server.rs index b2c104dd3..60b275de2 100644 --- a/packages/obs/examples/server.rs +++ b/packages/obs/examples/server.rs @@ -1,47 +1,53 @@ use opentelemetry::global; -use opentelemetry::trace::TraceContextExt; -use rustfs_obs::{init_logging, load_config, LogEntry}; +use rustfs_obs::{get_logger, init_obs, load_config, log_info, LogEntry}; use std::time::{Duration, SystemTime}; -use tracing::{info, instrument, Span}; +use tracing::{info, instrument}; use tracing_core::Level; -use tracing_opentelemetry::OpenTelemetrySpanExt; #[tokio::main] async fn main() { - let start_time = SystemTime::now(); let config = load_config(Some("packages/obs/examples/config".to_string())); - info!("Configuration file loading is complete {:?}", config.clone()); - let (logger, _guard) = init_logging(config); - info!("Log module initialization is completed"); + let (_logger, _guard) = init_obs(config.clone()).await; // Simulate the operation tokio::time::sleep(Duration::from_millis(100)).await; + run( + "service-demo".to_string(), + "object-demo".to_string(), + "user-demo".to_string(), + "service-demo".to_string(), + ) + .await; + info!("Program ends"); +} + +#[instrument(fields(bucket, object, user))] +async fn run(bucket: String, object: String, user: String, service_name: String) { + let start_time = SystemTime::now(); + info!("Log module initialization is completed service_name: {:?}", service_name); // 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")], + &[opentelemetry::KeyValue::new("operation", "run")], ); - // 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 - let trace_id = span_context.span().span_context().trace_id().to_string(); // Get the TraceId - let result = logger + let result = get_logger() + .lock() + .await .log(LogEntry::new( Level::INFO, "Process user requests".to_string(), "api_handler".to_string(), + Some("demo-audit".to_string()), Some("req-12345".to_string()), - Some("user-6789".to_string()), + Some(user), vec![ ("endpoint".to_string(), "/api/v1/data".to_string()), ("method".to_string(), "GET".to_string()), - ("span_id".to_string(), span_id), - ("trace_id".to_string(), trace_id), + ("bucket".to_string(), bucket), + ("object-length".to_string(), object.len().to_string()), ], )) .await; @@ -49,28 +55,31 @@ async fn main() { put_object("bucket".to_string(), "object".to_string(), "user".to_string()).await; info!("Logging is completed"); tokio::time::sleep(Duration::from_secs(2)).await; - info!("Program ends"); + info!("Program run ends"); } #[instrument(fields(bucket, object, user))] async fn put_object(bucket: String, object: String, user: String) { let start_time = SystemTime::now(); - info!("Starting PUT operation"); - // 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 - let trace_id = span_context.span().span_context().trace_id().to_string(); // Get the TraceId + info!("Starting put_object operation"); + + 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")], + ); + info!( - "Starting PUT operation content: bucket = {}, object = {}, user = {},span_id = {},trace_id = {},start_time = {}", + "Starting PUT operation content: bucket = {}, object = {}, user = {},start_time = {}", bucket, object, user, - span_id, - trace_id, start_time.elapsed().unwrap().as_secs_f64() ); + + let result = log_info("put_object logger info", "put_object").await; + info!("put_object is completed {:?}", result); // Simulate the operation tokio::time::sleep(Duration::from_millis(100)).await; diff --git a/packages/obs/src/config.rs b/packages/obs/src/config.rs index 8b14f7d0d..c6e57c421 100644 --- a/packages/obs/src/config.rs +++ b/packages/obs/src/config.rs @@ -15,7 +15,7 @@ pub struct OtelConfig { } /// Kafka Sink Configuration - Add batch parameters -#[derive(Debug, Deserialize, Clone)] +#[derive(Debug, Deserialize, Clone, Default)] pub struct KafkaSinkConfig { pub enabled: bool, pub bootstrap_servers: String, @@ -25,7 +25,7 @@ pub struct KafkaSinkConfig { } /// Webhook Sink Configuration - Add Retry Parameters -#[derive(Debug, Deserialize, Clone)] +#[derive(Debug, Deserialize, Clone, Default)] pub struct WebhookSinkConfig { pub enabled: bool, pub url: String, @@ -34,7 +34,7 @@ pub struct WebhookSinkConfig { } /// File Sink Configuration - Add buffering parameters -#[derive(Debug, Deserialize, Clone)] +#[derive(Debug, Deserialize, Clone, Default)] pub struct FileSinkConfig { pub enabled: bool, pub path: String, @@ -44,7 +44,7 @@ pub struct FileSinkConfig { } /// Sink configuration collection -#[derive(Debug, Deserialize, Clone)] +#[derive(Debug, Deserialize, Clone, Default)] pub struct SinkConfig { pub kafka: KafkaSinkConfig, pub webhook: WebhookSinkConfig, @@ -52,13 +52,13 @@ pub struct SinkConfig { } ///Logger Configuration -#[derive(Debug, Deserialize, Clone)] +#[derive(Debug, Deserialize, Clone, Default)] pub struct LoggerConfig { pub queue_capacity: Option, } /// Overall application configuration -#[derive(Debug, Deserialize, Clone)] +#[derive(Debug, Deserialize, Clone, Default)] pub struct AppConfig { pub observability: OtelConfig, pub sinks: SinkConfig, diff --git a/packages/obs/src/entry/audit.rs b/packages/obs/src/entry/audit.rs new file mode 100644 index 000000000..4e5099857 --- /dev/null +++ b/packages/obs/src/entry/audit.rs @@ -0,0 +1,88 @@ +use chrono::{DateTime, Utc}; +use serde::{Deserialize, Serialize}; +use serde_json::Value; +use std::collections::HashMap; + +/// ObjectVersion object version key/versionId +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct ObjectVersion { + #[serde(rename = "objectName")] + pub object_name: String, + #[serde(rename = "versionId", skip_serializing_if = "Option::is_none")] + pub version_id: Option, +} + +/// API details structure +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct ApiDetails { + #[serde(rename = "name", skip_serializing_if = "Option::is_none")] + pub name: Option, + #[serde(rename = "bucket", skip_serializing_if = "Option::is_none")] + pub bucket: Option, + #[serde(rename = "object", skip_serializing_if = "Option::is_none")] + pub object: Option, + #[serde(rename = "objects", skip_serializing_if = "Vec::is_empty", default)] + pub objects: Vec, + #[serde(rename = "status", skip_serializing_if = "Option::is_none")] + pub status: Option, + #[serde(rename = "statusCode", skip_serializing_if = "Option::is_none")] + pub status_code: Option, + #[serde(rename = "rx")] + pub input_bytes: i64, + #[serde(rename = "tx")] + pub output_bytes: i64, + #[serde(rename = "txHeaders", skip_serializing_if = "Option::is_none")] + pub header_bytes: Option, + #[serde(rename = "timeToFirstByte", skip_serializing_if = "Option::is_none")] + pub time_to_first_byte: Option, + #[serde(rename = "timeToFirstByteInNS", skip_serializing_if = "Option::is_none")] + pub time_to_first_byte_in_ns: Option, + #[serde(rename = "timeToResponse", skip_serializing_if = "Option::is_none")] + pub time_to_response: Option, + #[serde(rename = "timeToResponseInNS", skip_serializing_if = "Option::is_none")] + pub time_to_response_in_ns: Option, +} + +/// Entry - audit entry logs +#[derive(Debug, Serialize, Deserialize, Clone)] +pub struct AuditEntry { + pub version: String, + #[serde(rename = "deploymentid", skip_serializing_if = "Option::is_none")] + pub deployment_id: Option, + pub time: DateTime, + pub event: String, + + // Class of audit message - S3, admin ops, bucket management + #[serde(rename = "type", skip_serializing_if = "Option::is_none")] + pub entry_type: Option, + + // Deprecated, replaced by 'Event' + pub trigger: String, + pub api: ApiDetails, + #[serde(rename = "remotehost", skip_serializing_if = "Option::is_none")] + pub remote_host: Option, + #[serde(rename = "requestID", skip_serializing_if = "Option::is_none")] + pub request_id: Option, + #[serde(rename = "userAgent", skip_serializing_if = "Option::is_none")] + pub user_agent: Option, + #[serde(rename = "requestPath", skip_serializing_if = "Option::is_none")] + pub req_path: Option, + #[serde(rename = "requestHost", skip_serializing_if = "Option::is_none")] + pub req_host: Option, + #[serde(rename = "requestClaims", skip_serializing_if = "Option::is_none")] + pub req_claims: Option>, + #[serde(rename = "requestQuery", skip_serializing_if = "Option::is_none")] + pub req_query: Option>, + #[serde(rename = "requestHeader", skip_serializing_if = "Option::is_none")] + pub req_header: Option>, + #[serde(rename = "responseHeader", skip_serializing_if = "Option::is_none")] + pub resp_header: Option>, + #[serde(rename = "tags", skip_serializing_if = "Option::is_none")] + pub tags: Option>, + #[serde(rename = "accessKey", skip_serializing_if = "Option::is_none")] + pub access_key: Option, + #[serde(rename = "parentUser", skip_serializing_if = "Option::is_none")] + pub parent_user: Option, + #[serde(rename = "error", skip_serializing_if = "Option::is_none")] + pub error: Option, +} diff --git a/packages/obs/src/entry/base.rs b/packages/obs/src/entry/base.rs new file mode 100644 index 000000000..aa9d22ab7 --- /dev/null +++ b/packages/obs/src/entry/base.rs @@ -0,0 +1,114 @@ +use chrono::{DateTime, Utc}; +use serde::{Deserialize, Serialize}; +use serde_json::Value; +use std::collections::HashMap; + +/// ObjectVersion object version key/versionId +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ObjectVersion { + #[serde(rename = "objectName")] + pub object_name: String, + #[serde(rename = "versionId", skip_serializing_if = "Option::is_none")] + pub version_id: Option, +} + +/// Args - defines the arguments for the API +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct Args { + #[serde(rename = "bucket", skip_serializing_if = "Option::is_none")] + pub bucket: Option, + #[serde(rename = "object", skip_serializing_if = "Option::is_none")] + pub object: Option, + #[serde(rename = "versionId", skip_serializing_if = "Option::is_none")] + pub version_id: Option, + #[serde(rename = "objects", skip_serializing_if = "Option::is_none")] + pub objects: Option>, + #[serde(rename = "metadata", skip_serializing_if = "Option::is_none")] + pub metadata: Option>, +} + +/// Trace - defines the trace +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct Trace { + #[serde(rename = "message", skip_serializing_if = "Option::is_none")] + pub message: Option, + #[serde(rename = "source", skip_serializing_if = "Option::is_none")] + pub source: Option>, + #[serde(rename = "variables", skip_serializing_if = "Option::is_none")] + pub variables: Option>, +} + +/// API - defines the api type and its args +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct API { + #[serde(rename = "name", skip_serializing_if = "Option::is_none")] + pub name: Option, + #[serde(rename = "args", skip_serializing_if = "Option::is_none")] + pub args: Option, +} + +/// Log kind/level enum +#[derive(Debug, Clone, Serialize, Deserialize)] +pub enum LogKind { + #[serde(rename = "INFO")] + Info, + #[serde(rename = "WARNING")] + Warning, + #[serde(rename = "ERROR")] + Error, + #[serde(rename = "FATAL")] + Fatal, +} + +/// Entry - defines fields and values of each log entry +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct Entry { + #[serde(rename = "site", skip_serializing_if = "Option::is_none")] + pub site: Option, + + #[serde(rename = "deploymentid", skip_serializing_if = "Option::is_none")] + pub deployment_id: Option, + + pub level: LogKind, + + #[serde(rename = "errKind", skip_serializing_if = "Option::is_none")] + pub log_kind: Option, // Deprecated Jan 2024 + + pub time: DateTime, + + #[serde(rename = "api", skip_serializing_if = "Option::is_none")] + pub api: Option, + + #[serde(rename = "remotehost", skip_serializing_if = "Option::is_none")] + pub remote_host: Option, + + #[serde(rename = "host", skip_serializing_if = "Option::is_none")] + pub host: Option, + + #[serde(rename = "requestID", skip_serializing_if = "Option::is_none")] + pub request_id: Option, + + #[serde(rename = "userAgent", skip_serializing_if = "Option::is_none")] + pub user_agent: Option, + + #[serde(rename = "message", skip_serializing_if = "Option::is_none")] + pub message: Option, + + #[serde(rename = "error", skip_serializing_if = "Option::is_none")] + pub trace: Option, +} + +/// Info holds console log messages +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct Info { + #[serde(flatten)] + pub entry: Entry, + + pub console_msg: String, + + #[serde(rename = "node")] + pub node_name: String, + + #[serde(skip)] + pub err: Option, +} diff --git a/packages/obs/src/entry.rs b/packages/obs/src/entry/log.rs similarity index 95% rename from packages/obs/src/entry.rs rename to packages/obs/src/entry/log.rs index a8cef168a..6960fef21 100644 --- a/packages/obs/src/entry.rs +++ b/packages/obs/src/entry/log.rs @@ -52,6 +52,7 @@ pub struct LogEntry { pub level: SerializableLevel, // Log Level pub message: String, // Log messages pub source: String, // Log source (such as module name) + pub target: Option, // Log target pub request_id: Option, // Request ID (Common Server Fields) pub user_id: Option, // User ID (Common Server Fields) pub fields: Vec<(String, String)>, // Attached fields (key value pairs) @@ -63,6 +64,7 @@ impl LogEntry { level: Level, message: String, source: String, + target: Option, request_id: Option, user_id: Option, fields: Vec<(String, String)>, @@ -72,6 +74,7 @@ impl LogEntry { level: SerializableLevel::from(level), message, source, + target, request_id, user_id, fields, diff --git a/packages/obs/src/entry/mod.rs b/packages/obs/src/entry/mod.rs new file mode 100644 index 000000000..efc697e01 --- /dev/null +++ b/packages/obs/src/entry/mod.rs @@ -0,0 +1,3 @@ +pub(crate) mod audit; +mod base; +pub(crate) mod log; diff --git a/packages/obs/src/lib.rs b/packages/obs/src/lib.rs index 469547d61..079d9abfc 100644 --- a/packages/obs/src/lib.rs +++ b/packages/obs/src/lib.rs @@ -13,20 +13,53 @@ mod worker; pub use config::load_config; pub use config::{AppConfig, OtelConfig}; -pub use entry::{LogEntry, SerializableLevel}; +pub use entry::log::{LogEntry, SerializableLevel}; pub use logger::start_logger; +pub use logger::{ + ensure_logger_initialized, get_global_logger, init_global_logger, locked_logger, log_debug, log_error, log_info, log_trace, + log_warn, log_with_context, +}; pub use logger::{LogError, Logger}; pub use sink::Sink; +use std::sync::Arc; pub use telemetry::init_telemetry; +use tokio::sync::Mutex; pub use utils::{get_local_ip, get_local_ip_with_default}; pub use worker::start_worker; -/// Log module initialization function +/// Initialize the observability module /// -/// Return to Logger and Clean Guard -pub fn init_logging(config: AppConfig) -> (Logger, telemetry::OtelGuard) { +/// # Parameters +/// - `config`: Configuration information +/// +/// # Returns +/// A tuple containing the logger and the telemetry guard +/// +/// # Example +/// ``` +/// use rustfs_obs::{AppConfig, init_obs}; +/// +/// let config = AppConfig::default(); +/// let (logger, guard) = init_obs(config); +/// ``` +pub async fn init_obs(config: AppConfig) -> (Arc>, telemetry::OtelGuard) { let guard = init_telemetry(&config.observability); let sinks = sink::create_sinks(&config); - let logger = start_logger(&config, sinks); + let logger = init_global_logger(&config, sinks).await; (logger, guard) } + +/// Get the global logger instance +/// This function returns a reference to the global logger instance. +/// +/// # Returns +/// A reference to the global logger instance +/// +/// # Example +/// ``` +/// use rustfs_obs::get_logger; +/// let logger = get_logger(); +/// ``` +pub fn get_logger() -> &'static Arc> { + get_global_logger() +} diff --git a/packages/obs/src/logger.rs b/packages/obs/src/logger.rs index f2297c9ef..a3848cc3b 100644 --- a/packages/obs/src/logger.rs +++ b/packages/obs/src/logger.rs @@ -1,8 +1,14 @@ use crate::{AppConfig, LogEntry, SerializableLevel, Sink}; use std::sync::Arc; use tokio::sync::mpsc::{self, Receiver, Sender}; +use tokio::sync::{Mutex, OnceCell}; +use tracing_core::Level; + +// Add the global instance at the module level +static GLOBAL_LOGGER: OnceCell>> = OnceCell::const_new(); /// Server log processor +#[derive(Debug)] pub struct Logger { sender: Sender, // Log sending channel queue_capacity: usize, @@ -15,17 +21,21 @@ impl Logger { // Get queue capacity from configuration, or use default values 10000 let queue_capacity = config.logger.queue_capacity.unwrap_or(10000); let (sender, receiver) = mpsc::channel(queue_capacity); - ( - Logger { - sender, - queue_capacity, - }, - receiver, - ) + (Logger { sender, queue_capacity }, receiver) } - // Add a method to get queue capacity - pub fn queue_capacity(&self) -> usize { + /// get the queue capacity + /// This function returns the queue capacity. + /// # Returns + /// The queue capacity + /// # Example + /// ``` + /// use rustfs_obs::Logger; + /// async fn example(logger: &Logger) { + /// let _ = logger.get_queue_capacity(); + /// } + /// ``` + pub fn get_queue_capacity(&self) -> usize { self.queue_capacity } @@ -37,20 +47,20 @@ impl Logger { tracing::Span::current() .record("log_message", &entry.message) .record("source", &entry.source); + println!("target start is {:?}", &entry.target); + let target = if let Some(target) = &entry.target { + target.clone() + } else { + "server_logs".to_string() + }; - // Record queue utilization (if a certain threshold is exceeded) - let queue_len = self.sender.capacity(); - let utilization = queue_len as f64 / self.queue_capacity as f64; - if utilization > 0.8 { - tracing::warn!("Log queue utilization high: {:.1}%", utilization * 100.0); - } - + println!("target end is {:?}", target); // Generate independent Tracing Events with full LogEntry information // Generate corresponding events according to level match entry.level { - SerializableLevel(tracing::Level::ERROR) => { + SerializableLevel(Level::ERROR) => { tracing::error!( - target: "server_logs", + target = %target.clone(), timestamp = %entry.timestamp, message = %entry.message, source = %entry.source, @@ -59,9 +69,9 @@ impl Logger { fields = ?entry.fields ); } - SerializableLevel(tracing::Level::WARN) => { + SerializableLevel(Level::WARN) => { tracing::warn!( - target: "server_logs", + target = %target.clone(), timestamp = %entry.timestamp, message = %entry.message, source = %entry.source, @@ -70,9 +80,9 @@ impl Logger { fields = ?entry.fields ); } - SerializableLevel(tracing::Level::INFO) => { + SerializableLevel(Level::INFO) => { tracing::info!( - target: "server_logs", + target = %target.clone(), timestamp = %entry.timestamp, message = %entry.message, source = %entry.source, @@ -81,9 +91,9 @@ impl Logger { fields = ?entry.fields ); } - SerializableLevel(tracing::Level::DEBUG) => { + SerializableLevel(Level::DEBUG) => { tracing::debug!( - target: "server_logs", + target = %target.clone(), timestamp = %entry.timestamp, message = %entry.message, source = %entry.source, @@ -92,9 +102,9 @@ impl Logger { fields = ?entry.fields ); } - SerializableLevel(tracing::Level::TRACE) => { + SerializableLevel(Level::TRACE) => { tracing::trace!( - target: "server_logs", + target = %target.clone(), timestamp = %entry.timestamp, message = %entry.message, source = %entry.source, @@ -111,102 +121,52 @@ impl Logger { Err(mpsc::error::TrySendError::Full(entry)) => { // Processing strategy when queue is full tracing::warn!("Log queue full, applying backpressure"); - match tokio::time::timeout( - std::time::Duration::from_millis(500), - self.sender.send(entry), - ) - .await - { + match tokio::time::timeout(std::time::Duration::from_millis(500), self.sender.send(entry)).await { Ok(Ok(_)) => Ok(()), Ok(Err(_)) => Err(LogError::SendFailed("Channel closed")), Err(_) => Err(LogError::Timeout("Queue backpressure timeout")), } } - Err(mpsc::error::TrySendError::Closed(_)) => { - Err(LogError::SendFailed("Logger channel closed")) - } + Err(mpsc::error::TrySendError::Closed(_)) => Err(LogError::SendFailed("Logger channel closed")), } } - // Add convenient methods to simplify logging - // Fix the info() method, replacing None with an empty vector instead of the Option type - pub async fn info(&self, message: &str, source: &str) -> Result<(), LogError> { - self.log(LogEntry::new( - tracing::Level::INFO, - message.to_string(), - source.to_string(), - None, - None, - Vec::new(), // 使用空向量代替 None - )) - .await - } - - /// Add warn() method - pub async fn error(&self, message: &str, source: &str) -> Result<(), LogError> { - self.log(LogEntry::new( - tracing::Level::ERROR, - message.to_string(), - source.to_string(), - None, - None, - Vec::new(), - )) - .await - } - - /// Add warn() method - pub async fn warn(&self, message: &str, source: &str) -> Result<(), LogError> { - self.log(LogEntry::new( - tracing::Level::WARN, - message.to_string(), - source.to_string(), - None, - None, - Vec::new(), - )) - .await - } - - /// Add debug() method - pub async fn debug(&self, message: &str, source: &str) -> Result<(), LogError> { - self.log(LogEntry::new( - tracing::Level::DEBUG, - message.to_string(), - source.to_string(), - None, - None, - Vec::new(), - )) - .await - } - - /// Add trace() method - pub async fn trace(&self, message: &str, source: &str) -> Result<(), LogError> { - self.log(LogEntry::new( - tracing::Level::TRACE, - message.to_string(), - source.to_string(), - None, - None, - Vec::new(), - )) - .await - } - - // Add extension methods with context information for more flexibility - pub async fn info_with_context( + /// Write log with context information + /// This function writes log messages with context information. + /// + /// # Parameters + /// - `message`: Message to be logged + /// - `source`: Source of the log + /// - `request_id`: Request ID + /// - `user_id`: User ID + /// - `fields`: Additional fields + /// + /// # Returns + /// Result indicating whether the operation was successful + /// + /// # Example + /// ``` + /// use tracing_core::Level; + /// use rustfs_obs::Logger; + /// + /// async fn example(logger: &Logger) { + /// let _ = logger.write_with_context("This is an information message", "example",Level::INFO, Some("target".to_string()),Some("req-12345".to_string()), Some("user-6789".to_string()), vec![("endpoint".to_string(), "/api/v1/data".to_string())]).await; + /// } + pub async fn write_with_context( &self, message: &str, source: &str, + level: Level, + target: Option, request_id: Option, user_id: Option, fields: Vec<(String, String)>, ) -> Result<(), LogError> { self.log(LogEntry::new( - tracing::Level::INFO, + level, message.to_string(), source.to_string(), + target, request_id, user_id, fields, @@ -214,14 +174,69 @@ impl Logger { .await } - // Add elegant closing method + /// Write log + /// This function writes log messages. + /// # Parameters + /// - `message`: Message to be logged + /// - `source`: Source of the log + /// - `level`: Log level + /// + /// # Returns + /// Result indicating whether the operation was successful + /// + /// # Example + /// ``` + /// use rustfs_obs::Logger; + /// use tracing_core::Level; + /// + /// async fn example(logger: &Logger) { + /// let _ = logger.write("This is an information message", "example", Level::INFO).await; + /// } + /// ``` + pub async fn write(&self, message: &str, source: &str, level: Level) -> Result<(), LogError> { + self.log(LogEntry::new( + level, + message.to_string(), + source.to_string(), + None, + None, + None, + Vec::new(), + )) + .await + } + + /// Shutdown the logger + /// This function shuts down the logger. + /// + /// # Returns + /// Result indicating whether the operation was successful + /// + /// # Example + /// ``` + /// use rustfs_obs::Logger; + /// + /// async fn example(logger: Logger) { + /// let _ = logger.shutdown().await; + /// } + /// ``` pub async fn shutdown(self) -> Result<(), LogError> { drop(self.sender); //Close the sending end so that the receiver knows that there is no new message Ok(()) } } -// Define custom error type +/// Log error type +/// This enum defines the error types that can occur when logging. +/// It is used to provide more detailed error information. +/// # Example +/// ``` +/// use rustfs_obs::LogError; +/// use thiserror::Error; +/// +/// LogError::SendFailed("Failed to send log"); +/// LogError::Timeout("Operation timed out"); +/// ``` #[derive(Debug, thiserror::Error)] pub enum LogError { #[error("Failed to send log: {0}")] @@ -231,8 +246,245 @@ pub enum LogError { } /// Start the log module +/// This function starts the log module. +/// It initializes the logger and starts the worker to process logs. +/// # Parameters +/// - `config`: Configuration information +/// - `sinks`: A vector of Sink instances +/// # Returns +/// The global logger instance +/// # Example +/// ``` +/// use rustfs_obs::{AppConfig, start_logger}; +/// +/// let config = AppConfig::default(); +/// let sinks = vec![]; +/// let logger = start_logger(&config, sinks); +/// ``` pub fn start_logger(config: &AppConfig, sinks: Vec>) -> Logger { let (logger, receiver) = Logger::new(config); tokio::spawn(crate::worker::start_worker(receiver, sinks)); logger } + +/// Initialize the global logger instance +/// This function initializes the global logger instance and returns a reference to it. +/// If the logger has been initialized before, it will return the existing logger instance. +/// +/// # Parameters +/// - `config`: Configuration information +/// - `sinks`: A vector of Sink instances +/// +/// # Returns +/// A reference to the global logger instance +/// +/// # Example +/// ``` +/// use rustfs_obs::{AppConfig,init_global_logger}; +/// +/// let config = AppConfig::default(); +/// let sinks = vec![]; +/// let logger = init_global_logger(&config, sinks); +/// ``` +pub async fn init_global_logger(config: &AppConfig, sinks: Vec>) -> Arc> { + let logger = Arc::new(Mutex::new(start_logger(config, sinks))); + GLOBAL_LOGGER.set(logger.clone()).expect("Logger already initialized"); + logger +} + +/// Get the global logger instance +/// +/// This function returns a reference to the global logger instance. +/// +/// # Returns +/// A reference to the global logger instance +/// +/// # Example +/// ``` +/// use rustfs_obs::get_global_logger; +/// +/// let logger = get_global_logger(); +/// ``` +pub fn get_global_logger() -> &'static Arc> { + GLOBAL_LOGGER.get().expect("Logger not initialized") +} + +/// Get the global logger instance with a lock +/// This function returns a reference to the global logger instance with a lock. +/// It is used to ensure that the logger is thread-safe. +/// +/// # Returns +/// A reference to the global logger instance with a lock +/// +/// # Example +/// ``` +/// use rustfs_obs::locked_logger; +/// +/// async fn example() { +/// let logger = locked_logger().await; +/// } +/// ``` +pub async fn locked_logger() -> tokio::sync::MutexGuard<'static, Logger> { + get_global_logger().lock().await +} + +/// Initialize with default empty logger if needed (optional) +/// This function initializes the logger with a default empty logger if needed. +/// It is used to ensure that the logger is initialized before logging. +/// +/// # Returns +/// A reference to the global logger instance +/// +/// # Example +/// ``` +/// use rustfs_obs::ensure_logger_initialized; +/// +/// let logger = ensure_logger_initialized(); +/// ``` +pub fn ensure_logger_initialized() -> &'static Arc> { + if GLOBAL_LOGGER.get().is_none() { + let config = AppConfig::default(); + let sinks = vec![]; + let logger = Arc::new(Mutex::new(start_logger(&config, sinks))); + let _ = GLOBAL_LOGGER.set(logger); + } + GLOBAL_LOGGER.get().unwrap() +} + +/// Log information +/// This function logs information messages. +/// +/// # Parameters +/// - `message`: Message to be logged +/// - `source`: Source of the log +/// +/// # Returns +/// Result indicating whether the operation was successful +/// +/// # Example +/// ``` +/// use rustfs_obs::log_info; +/// +/// async fn example() { +/// let _ = log_info("This is an information message", "example").await; +/// } +/// ``` +pub async fn log_info(message: &str, source: &str) -> Result<(), LogError> { + get_global_logger().lock().await.write(message, source, Level::INFO).await +} + +/// Log error +/// This function logs error messages. +/// # Parameters +/// - `message`: Message to be logged +/// - `source`: Source of the log +/// # Returns +/// Result indicating whether the operation was successful +/// # Example +/// ``` +/// use rustfs_obs::log_error; +/// +/// async fn example() { +/// let _ = log_error("This is an error message", "example").await; +/// } +pub async fn log_error(message: &str, source: &str) -> Result<(), LogError> { + get_global_logger().lock().await.write(message, source, Level::ERROR).await +} + +/// Log warning +/// This function logs warning messages. +/// # Parameters +/// - `message`: Message to be logged +/// - `source`: Source of the log +/// # Returns +/// Result indicating whether the operation was successful +/// +/// # Example +/// ``` +/// use rustfs_obs::log_warn; +/// +/// async fn example() { +/// let _ = log_warn("This is a warning message", "example").await; +/// } +/// ``` +pub async fn log_warn(message: &str, source: &str) -> Result<(), LogError> { + get_global_logger().lock().await.write(message, source, Level::WARN).await +} + +/// Log debug +/// This function logs debug messages. +/// # Parameters +/// - `message`: Message to be logged +/// - `source`: Source of the log +/// # Returns +/// Result indicating whether the operation was successful +/// +/// # Example +/// ``` +/// use rustfs_obs::log_debug; +/// +/// async fn example() { +/// let _ = log_debug("This is a debug message", "example").await; +/// } +/// ``` +pub async fn log_debug(message: &str, source: &str) -> Result<(), LogError> { + get_global_logger().lock().await.write(message, source, Level::DEBUG).await +} + +/// Log trace +/// This function logs trace messages. +/// # Parameters +/// - `message`: Message to be logged +/// - `source`: Source of the log +/// +/// # Returns +/// Result indicating whether the operation was successful +/// +/// # Example +/// ``` +/// use rustfs_obs::log_trace; +/// +/// async fn example() { +/// let _ = log_trace("This is a trace message", "example").await; +/// } +/// ``` +pub async fn log_trace(message: &str, source: &str) -> Result<(), LogError> { + get_global_logger().lock().await.write(message, source, Level::TRACE).await +} + +/// Log with context information +/// This function logs messages with context information. +/// # Parameters +/// - `message`: Message to be logged +/// - `source`: Source of the log +/// - `level`: Log level +/// - `target`: Log target +/// - `request_id`: Request ID +/// - `user_id`: User ID +/// - `fields`: Additional fields +/// # Returns +/// Result indicating whether the operation was successful +/// # Example +/// ``` +/// use tracing_core::Level; +/// use rustfs_obs::log_with_context; +/// +/// async fn example() { +/// let _ = log_with_context("This is an information message", "example", Level::INFO, Some("target".to_string()), Some("req-12345".to_string()), Some("user-6789".to_string()), vec![("endpoint".to_string(), "/api/v1/data".to_string())]).await; +/// } +/// ``` +pub async fn log_with_context( + message: &str, + source: &str, + level: Level, + target: Option, + request_id: Option, + user_id: Option, + fields: Vec<(String, String)>, +) -> Result<(), LogError> { + get_global_logger() + .lock() + .await + .write_with_context(message, source, level, target, request_id, user_id, fields) + .await +} diff --git a/packages/obs/src/sink.rs b/packages/obs/src/sink.rs index 37785cfba..2363a69dd 100644 --- a/packages/obs/src/sink.rs +++ b/packages/obs/src/sink.rs @@ -258,10 +258,10 @@ impl FileSink { buffer_size: usize, flush_interval_ms: u64, flush_threshold: usize, - ) -> Result { + ) -> Result { let file = OpenOptions::new().append(true).create(true).open(&path).await?; - let writer = tokio::io::BufWriter::with_capacity(buffer_size, file); + let writer = io::BufWriter::with_capacity(buffer_size, file); let now = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) .unwrap() @@ -421,11 +421,10 @@ pub fn create_sinks(config: &AppConfig) -> Vec> { // Use synchronous file operations let file_result = std::fs::OpenOptions::new().append(true).create(true).open(&path); - match file_result { Ok(file) => { let buffer_size = config.sinks.file.buffer_size.unwrap_or(8192); - let writer = tokio::io::BufWriter::with_capacity(buffer_size, tokio::fs::File::from_std(file)); + let writer = io::BufWriter::with_capacity(buffer_size, tokio::fs::File::from_std(file)); let now = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) diff --git a/packages/obs/src/telemetry.rs b/packages/obs/src/telemetry.rs index fb1111303..cafe5badc 100644 --- a/packages/obs/src/telemetry.rs +++ b/packages/obs/src/telemetry.rs @@ -13,6 +13,8 @@ use opentelemetry_semantic_conventions::{ attribute::{DEPLOYMENT_ENVIRONMENT_NAME, NETWORK_LOCAL_ADDRESS, SERVICE_NAME, SERVICE_VERSION}, SCHEMA_URL, }; +use std::io::IsTerminal; +use tracing_error::ErrorLayer; use tracing_opentelemetry::{MetricsLayer, OpenTelemetryLayer}; use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt, EnvFilter, Layer}; @@ -155,63 +157,72 @@ pub fn init_telemetry(config: &OtelConfig) -> OtelGuard { let meter_provider = init_meter_provider(config); let tracer = tracer_provider.tracer(config.service_name.clone()); - let logger_provider = if config.endpoint.is_empty() { - SdkLoggerProvider::builder() - .with_resource(resource(config)) - .with_simple_exporter(opentelemetry_stdout::LogExporter::default()) - .build() - } else { - let exporter = opentelemetry_otlp::LogExporter::builder() - .with_tonic() - .with_endpoint(&config.endpoint) - .build() - .unwrap(); - SdkLoggerProvider::builder() - .with_resource(resource(config)) - .with_batch_exporter(exporter) - .with_batch_exporter(opentelemetry_stdout::LogExporter::default()) - .build() + // Initialize logger provider based on configuration + let logger_provider = { + let mut builder = SdkLoggerProvider::builder().with_resource(resource(config)); + + if config.endpoint.is_empty() { + // Use stdout exporter when no endpoint is configured + builder = builder.with_simple_exporter(opentelemetry_stdout::LogExporter::default()); + } else { + // Configure OTLP exporter when endpoint is provided + let exporter = opentelemetry_otlp::LogExporter::builder() + .with_tonic() + .with_endpoint(&config.endpoint) + .build() + .unwrap(); + + builder = builder.with_batch_exporter(exporter); + + // Add stdout exporter if requested + if config.use_stdout { + builder = builder.with_batch_exporter(opentelemetry_stdout::LogExporter::default()); + } + } + + builder.build() }; - let otel_layer = layer::OpenTelemetryTracingBridge::new(&logger_provider); - // For the OpenTelemetry layer, add a tracing filter to filter events from - // OpenTelemetry and its dependent crates (opentelemetry-otlp uses crates - // like reqwest/tonic etc.) from being sent back to OTel itself, thus - // preventing infinite telemetry generation. The filter levels are set as - // follows: - // - Allow `info` level and above by default. - // - Restrict `opentelemetry`, `hyper`, `tonic`, and `reqwest` completely. - // Note: This will also drop events from crates like `tonic` etc. even when - // they are used outside the OTLP Exporter. For more details, see: - // https://github.com/open-telemetry/opentelemetry-rust/issues/761 - let filter_otel = EnvFilter::new("info") - .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 = otel_layer.with_filter(filter_otel); + // Setup OpenTelemetryTracingBridge layer + let otel_layer = { + // Filter to prevent infinite telemetry loops + // This blocks events from OpenTelemetry and its dependent libraries (tonic, reqwest, etc.) + // from being sent back to OpenTelemetry itself + let filter_otel = EnvFilter::new("info") + .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()); + + layer::OpenTelemetryTracingBridge::new(&logger_provider).with_filter(filter_otel) + }; let registry = tracing_subscriber::registry() .with(tracing_subscriber::filter::LevelFilter::INFO) .with(OpenTelemetryLayer::new(tracer)) .with(MetricsLayer::new(meter_provider.clone())) .with(otel_layer); - if config.endpoint.is_empty() { - // Create a new tracing::Fmt layer to print the logs to stdout. It has a - // default filter of `info` level and above, and `debug` and above for logs - // from OpenTelemetry crates. The filter levels can be customized as needed. - 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 formatting layer + let enable_color = std::io::stdout().is_terminal(); + let fmt_layer = tracing_subscriber::fmt::layer() + .with_ansi(enable_color) + .with_thread_names(true) + .with_file(true) + .with_line_number(true); - registry - .with(tracing_subscriber::fmt::layer().with_ansi(true)) - .with(fmt_layer) - .init(); + // Creating a formatting layer with explicit type to avoid type mismatches + let fmt_layer = if config.endpoint.is_empty() { + // Local debug mode: add filter to show OpenTelemetry debug logs + let filter_fmt = EnvFilter::new("info").add_directive("opentelemetry=debug".parse().unwrap()); + fmt_layer.with_filter(filter_fmt) } else { - registry.with(tracing_subscriber::fmt::layer().with_ansi(false)).init(); - println!("Logs and meter,tracer enabled"); + // Production mode: use default filter settings + fmt_layer.with_filter(EnvFilter::new("info").add_directive("opentelemetry=off".parse().unwrap())) + }; + + registry.with(ErrorLayer::default()).with(fmt_layer).init(); + if !config.endpoint.is_empty() { + tracing::info!("OpenTelemetry telemetry initialized with OTLP endpoint: {}", config.endpoint); } OtelGuard { diff --git a/packages/obs/src/worker.rs b/packages/obs/src/worker.rs index a04327945..1c40f47f5 100644 --- a/packages/obs/src/worker.rs +++ b/packages/obs/src/worker.rs @@ -1,4 +1,4 @@ -use crate::{entry::LogEntry, sink::Sink}; +use crate::{entry::log::LogEntry, sink::Sink}; use std::sync::Arc; use tokio::sync::mpsc::Receiver; diff --git a/rustfs/src/main.rs b/rustfs/src/main.rs index a39e121ee..b87b16ca3 100644 --- a/rustfs/src/main.rs +++ b/rustfs/src/main.rs @@ -35,6 +35,7 @@ use hyper_util::{ }; use iam::init_iam_sys; use protos::proto_gen::node_service::node_service_server::NodeServiceServer; +use rustfs_obs::{init_obs, load_config}; use s3s::{host::MultiDomain, service::S3ServiceBuilder}; use service::hybrid; use std::{io::IsTerminal, net::SocketAddr}; @@ -45,6 +46,7 @@ use tracing::{debug, error, info, warn}; use tracing_error::ErrorLayer; use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt}; +#[allow(dead_code)] fn setup_tracing() { use tracing_subscriber::EnvFilter; @@ -89,8 +91,10 @@ fn main() -> Result<()> { let opt = config::Opt::parse(); //设置 trace - setup_tracing(); - + // setup_tracing(); + let config = load_config(Some("packages/obs/examples/config".to_string())); + // Initialize the logger asynchronously + let (_logger, _guard) = tokio::runtime::Runtime::new()?.block_on(async { init_obs(config).await }); //运行参数 run(opt) }