diff --git a/Cargo.lock b/Cargo.lock index ef3176efb..f833d5ceb 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5617,7 +5617,6 @@ dependencies = [ "flatbuffers", "futures", "futures-util", - "h2", "http", "http-body", "hyper", @@ -5627,7 +5626,6 @@ dependencies = [ "lazy_static", "local-ip-address", "lock", - "log", "madmin", "matchit 0.8.6", "mime", @@ -5658,6 +5656,7 @@ dependencies = [ "tower 0.5.2", "tower-http", "tracing", + "tracing-core", "tracing-error", "tracing-subscriber", "transform-stream", diff --git a/ecstore/src/endpoints.rs b/ecstore/src/endpoints.rs index 8dbee9928..3bca94804 100644 --- a/ecstore/src/endpoints.rs +++ b/ecstore/src/endpoints.rs @@ -1,4 +1,4 @@ -use tracing::warn; +use tracing::{instrument, warn}; use crate::{ disk::endpoint::{Endpoint, EndpointType}, @@ -407,7 +407,7 @@ pub struct PoolEndpoints { pub platform: String, } -/// list of list of endpoints +/// list of endpoints #[derive(Debug, Clone, Default)] pub struct EndpointServerPools(pub Vec); @@ -532,6 +532,8 @@ impl EndpointServerPools { nodes } + + #[instrument] pub fn hosts_sorted(&self) -> Vec> { let (mut peers, local) = self.peers(); @@ -604,7 +606,6 @@ impl EndpointServerPools { #[cfg(test)] mod test { - use super::*; use std::path::Path; diff --git a/ecstore/src/global.rs b/ecstore/src/global.rs index 8735249c9..f6fc4bdb8 100644 --- a/ecstore/src/global.rs +++ b/ecstore/src/global.rs @@ -61,7 +61,7 @@ pub fn get_global_deployment_id() -> Option { pub fn set_global_endpoints(eps: Vec) { GLOBAL_Endpoints .set(EndpointServerPools::from(eps)) - .expect("GLOBAL_Endpoints set faild") + .expect("GLOBAL_Endpoints set failed") } pub fn get_global_endpoints() -> EndpointServerPools { diff --git a/iam/src/lib.rs b/iam/src/lib.rs index 65484724c..364627a29 100644 --- a/iam/src/lib.rs +++ b/iam/src/lib.rs @@ -2,11 +2,11 @@ use auth::Credentials; use ecstore::error::{Error, Result}; use ecstore::store::ECStore; use error::Error as IamError; -use log::debug; use manager::IamCache; use std::sync::{Arc, OnceLock}; use store::object::ObjectStore; use sys::IamSys; +use tracing::{debug, instrument}; pub mod cache; mod format; @@ -58,6 +58,7 @@ pub fn get_global_action_cred() -> Option { GLOBAL_ACTIVE_CRED.get().cloned() } +#[instrument] pub async fn init_iam_sys(ecstore: Arc) -> Result<()> { debug!("init iam system"); let s = IamCache::new(ObjectStore::new(ecstore)).await; diff --git a/packages/obs/Cargo.toml b/packages/obs/Cargo.toml index 86a9fb6c4..af3e3480d 100644 --- a/packages/obs/Cargo.toml +++ b/packages/obs/Cargo.toml @@ -23,14 +23,14 @@ opentelemetry = { workspace = true } opentelemetry-appender-tracing = { workspace = true, features = ["experimental_use_tracing_span_context", "experimental_metadata_attributes"] } opentelemetry_sdk = { workspace = true, features = ["rt-tokio"] } opentelemetry-stdout = { workspace = true } -opentelemetry-otlp = { workspace = true, features = ["grpc-tonic", "metrics"] } +opentelemetry-otlp = { workspace = true, features = ["grpc-tonic", "gzip-tonic"] } opentelemetry-semantic-conventions = { workspace = true, features = ["semconv_experimental"] } serde = { workspace = true } -tracing = { workspace = true } +tracing = { workspace = true, features = ["std", "attributes"] } 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"] } +tracing-subscriber = { workspace = true, features = ["registry", "std", "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 } @@ -41,10 +41,10 @@ local-ip-address = { workspace = true } [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 = { workspace = true } +opentelemetry_sdk = { workspace = true, features = ["rt-tokio"] } +opentelemetry-stdout = { workspace = true } +opentelemetry-otlp = { workspace = true, features = ["grpc-tonic"] } opentelemetry-semantic-conventions = { workspace = true, features = ["semconv_experimental"] } tokio = { workspace = true, features = ["full"] } tracing = { workspace = true, features = ["std", "attributes"] } diff --git a/packages/obs/examples/config.toml b/packages/obs/examples/config.toml index 2ceac4397..135dd39ac 100644 --- a/packages/obs/examples/config.toml +++ b/packages/obs/examples/config.toml @@ -1,9 +1,9 @@ [observability] endpoint = "http://localhost:4317" use_stdout = true -sample_ratio = 0.5 +sample_ratio = 1 meter_interval = 30 -service_name = "rustfs_obs_service" +service_name = "rustfs_obs" service_version = "0.1.0" deployment_environment = "develop" diff --git a/packages/obs/src/telemetry.rs b/packages/obs/src/telemetry.rs index cafe5badc..914071903 100644 --- a/packages/obs/src/telemetry.rs +++ b/packages/obs/src/telemetry.rs @@ -129,7 +129,7 @@ fn init_tracer_provider(config: &OtelConfig) -> SdkTracerProvider { let tracer_provider = if config.endpoint.is_empty() { builder - .with_simple_exporter(opentelemetry_stdout::SpanExporter::default()) + .with_batch_exporter(opentelemetry_stdout::SpanExporter::default()) .build() } else { let exporter = opentelemetry_otlp::SpanExporter::builder() diff --git a/rustfs/Cargo.toml b/rustfs/Cargo.toml index 530cc57ae..640baea76 100644 --- a/rustfs/Cargo.toml +++ b/rustfs/Cargo.toml @@ -16,7 +16,7 @@ workspace = true [dependencies] madmin.workspace = true -log.workspace = true +#log.workspace = true async-trait.workspace = true bytes.workspace = true clap.workspace = true @@ -25,7 +25,7 @@ ecstore.workspace = true flatbuffers.workspace = true futures.workspace = true futures-util.workspace = true -h2 = "0.4.7" +#h2 = "0.4.7" hyper.workspace = true hyper-util.workspace = true http.workspace = true @@ -57,6 +57,7 @@ tokio-stream.workspace = true tonic = { version = "0.12.3", features = ["gzip"] } tonic-reflection.workspace = true tower.workspace = true +tracing-core = { workspace = true } tracing-error.workspace = true tracing-subscriber.workspace = true transform-stream.workspace = true diff --git a/rustfs/src/main.rs b/rustfs/src/main.rs index 6f4fee1ee..787ee8c99 100644 --- a/rustfs/src/main.rs +++ b/rustfs/src/main.rs @@ -36,14 +36,16 @@ use hyper_util::{ use iam::init_iam_sys; use once_cell::sync::OnceCell; use protos::proto_gen::node_service::node_service_server::NodeServiceServer; -use rustfs_obs::{init_obs, load_config}; +use rustfs_obs::{init_obs, load_config, BaseLogEntry, ServerLogEntry}; use s3s::{host::MultiDomain, service::S3ServiceBuilder}; use service::hybrid; +use std::time::SystemTime; use std::{io::IsTerminal, net::SocketAddr}; use tokio::net::TcpListener; use tonic::{metadata::MetadataValue, Request, Status}; use tower_http::cors::CorsLayer; -use tracing::{debug, error, info, warn}; +use tracing::{debug, error, info, info_span, warn}; +use tracing_core::Level; use tracing_error::ErrorLayer; use tracing_subscriber::{layer::SubscriberExt, util::SubscriberInitExt}; @@ -101,28 +103,43 @@ fn print_server_info() { info!("Docs: {}", cfg.doc()); } -fn main() -> Result<()> { +#[tokio::main] +async fn main() -> Result<()> { // Parse the obtained parameters let opt = config::Opt::parse(); - + println!("config: {:?}", &opt); // 设置 trace // setup_tracing(); let config = load_config(Some(opt.clone().obs_config)); // Initialize Observability - let (_logger, guard) = tokio::runtime::Runtime::new()?.block_on(async { init_obs(config).await }); - + let (logger, guard) = init_obs(config).await; + // let (logger, guard) = tokio::runtime::Runtime::new()?.block_on(async { init_obs(config).await }); + // let _rr = tokio::runtime::Runtime::new()?.block_on(async { + let start_time = SystemTime::now(); + let base_entry = BaseLogEntry::new() + .timestamp(chrono::DateTime::from(start_time)) + .message(Some("main init obs end".to_string())) + .request_id(Some("main".to_string())); + let server_entry = ServerLogEntry::new(Level::INFO, "main_server_entry".to_string()) + .with_base(base_entry) + .user_id(Some("user_id".to_string())); + let _r = logger.lock().await.log_server_entry(server_entry).await; + // }); // Pack and store the guard GLOBAL_GUARD.set(TracingGuard(Box::new(guard))).unwrap_or_else(|_| { error!("Unable to set global tracing guard"); }); // Run parameters - run(opt) + run(opt).await } - -#[tokio::main] +// +// #[tokio::main] async fn run(opt: config::Opt) -> Result<()> { + let span = info_span!("trace-main-run"); + let _enter = span.enter(); + debug!("opt: {:?}", &opt); let mut server_addr = net::check_local_server_addr(opt.address.as_str()).unwrap(); @@ -259,7 +276,7 @@ async fn run(opt: config::Opt) -> Result<()> { match res { Ok(conn) => conn, Err(err) => { - tracing::error!("error accepting connection: {err}"); + error!("error accepting connection: {err}"); continue; } } @@ -278,10 +295,10 @@ async fn run(opt: config::Opt) -> Result<()> { tokio::select! { () = graceful.shutdown() => { - tracing::debug!("Gracefully shutdown!"); + debug!("Gracefully shutdown!"); }, () = tokio::time::sleep(std::time::Duration::from_secs(10)) => { - tracing::debug!("Waited 10 seconds for graceful shutdown, aborting..."); + debug!("Waited 10 seconds for graceful shutdown, aborting..."); } } }); diff --git a/rustfs/src/storage/ecfs.rs b/rustfs/src/storage/ecfs.rs index edb936dc3..7dfcb46c1 100644 --- a/rustfs/src/storage/ecfs.rs +++ b/rustfs/src/storage/ecfs.rs @@ -42,7 +42,6 @@ use http::HeaderMap; use iam::policy::action::Action; use iam::policy::action::S3Action; use lazy_static::lazy_static; -use log::warn; use s3s::dto::*; use s3s::s3_error; use s3s::S3Error; @@ -57,6 +56,7 @@ use tokio_util::io::StreamReader; use tracing::debug; use tracing::error; use tracing::info; +use tracing::warn; use transform_stream::AsyncTryStream; use uuid::Uuid; @@ -217,7 +217,7 @@ impl S3 for FS { #[tracing::instrument(level = "debug", skip(self, req))] async fn delete_bucket(&self, req: S3Request) -> S3Result> { let input = req.input; - // TODO: DeleteBucketInput 没有force参数? + // TODO: DeleteBucketInput 没有 force 参数? let Some(store) = new_object_layer_fn() else { return Err(S3Error::with_message(S3ErrorCode::InternalError, "Not init".to_string())); };