diff --git a/src/core/actions/analyze_project_dependencies.rs b/src/core/actions/analyze_project_dependencies.rs index 604a751..83de17f 100644 --- a/src/core/actions/analyze_project_dependencies.rs +++ b/src/core/actions/analyze_project_dependencies.rs @@ -1,5 +1,6 @@ use anyhow::Result; use futures::future; +use tracing::instrument::WithSubscriber; use crate::core::{ application::Application, @@ -30,60 +31,64 @@ pub async fn run(app: &Application, project_id: String) -> Result<()> { let scanners = app.ecosystems(); let handle = app.handle(); - tokio::spawn(async move { - let futures = scanners.iter().map(|scanner| async { - let mut discovered_dependencies = scanner - .discover_project_dependencies(snapshot.as_ref()) - .await?; + tokio::spawn( + async move { + let futures = scanners.iter().map(|scanner| async { + let mut discovered_dependencies = scanner + .discover_project_dependencies(snapshot.as_ref()) + .await?; - // We do not require that 'discover_project_dependencies' returns unique dependencies. - // It is very likely that there are duplicates. That happens quite often in monorepos - // where multiple packages have the same dependency. To not waste time querying - // dependency update options, make the deps unique. - discovered_dependencies.sort(); - discovered_dependencies.dedup(); + // We do not require that 'discover_project_dependencies' returns unique dependencies. + // It is very likely that there are duplicates. That happens quite often in monorepos + // where multiple packages have the same dependency. To not waste time querying + // dependency update options, make the deps unique. + discovered_dependencies.sort(); + discovered_dependencies.dedup(); - // For each of the discovered dependencies, see if it can be updated. This requires - // contacting the ecosystem package registry, and lots of custom code to determine - // what the version candidates are. - let results = - future::try_join_all(discovered_dependencies.into_iter().map(|dep| async { - let dependency_update_options = - scanner.query_dependency_update_options(&dep).await?; + // For each of the discovered dependencies, see if it can be updated. This requires + // contacting the ecosystem package registry, and lots of custom code to determine + // what the version candidates are. + let results = + future::try_join_all(discovered_dependencies.into_iter().map(|dep| async { + let dependency_update_options = + scanner.query_dependency_update_options(&dep).await?; - Ok::<_, anyhow::Error>(AnalyzedProjectDependency { - discovered_dependency: dep, - dependency_update_options, - }) - })) - .await?; + Ok::<_, anyhow::Error>(AnalyzedProjectDependency { + discovered_dependency: dep, + dependency_update_options, + }) + })) + .await?; - Ok::<_, anyhow::Error>(results) - }); + Ok::<_, anyhow::Error>(results) + }); - match future::try_join_all(futures).await { - Ok(results) => { - let analyzed_project_dependencies = results.into_iter().flatten().collect(); - let payload = crate::core::message::Payload::PersistAnalyzedProjectDependencies { - project_id: project_id.clone(), - scan_result: AnalyzedProjectDependencies { - analyzed_project_dependencies, - }, - }; + match future::try_join_all(futures).await { + Ok(results) => { + let analyzed_project_dependencies = results.into_iter().flatten().collect(); + let payload = + crate::core::message::Payload::PersistAnalyzedProjectDependencies { + project_id: project_id.clone(), + scan_result: AnalyzedProjectDependencies { + analyzed_project_dependencies, + }, + }; - // We send it to the application mailbox using an ad-hoc message_id - if let Err(e) = handle.send(pk(), payload).await { - tracing::error!( - "Failed to send PersistAnalyzedProjectDependencies message: {:#}", - e - ); + // We send it to the application mailbox using an ad-hoc message_id + if let Err(e) = handle.send(pk(), payload).await { + tracing::error!( + "Failed to send PersistAnalyzedProjectDependencies message: {:#}", + e + ); + } + } + Err(e) => { + tracing::error!("Failed to analyze project {}: {:#}", project_id, e); } - } - Err(e) => { - tracing::error!("Failed to analyze project {}: {:#}", project_id, e); } } - }); + .with_current_subscriber(), + ); Ok(()) } diff --git a/src/core/application.rs b/src/core/application.rs index 95fc8b5..2e73db8 100644 --- a/src/core/application.rs +++ b/src/core/application.rs @@ -52,13 +52,29 @@ impl Application { pub fn start(self) -> Handle { let handle = self.handle.clone(); - tokio::spawn(async move { - tracing::info!("Event loop active"); + let filter = tracing_subscriber::EnvFilter::try_from_default_env() + .unwrap_or_else(|_| tracing_subscriber::EnvFilter::new("info")); - if let Err(e) = self.run().await { - tracing::error!("Application loop exited with error: {:#}", e); + let log_layer = super::logger::CoreLogLayer::new(handle.events.clone()); + + use tracing_subscriber::layer::SubscriberExt; + let subscriber = tracing_subscriber::registry() + .with(filter) + .with(log_layer) + .with(tracing_subscriber::fmt::layer()); + let dispatch = tracing::Dispatch::new(subscriber); + + use tracing::instrument::WithSubscriber; + tokio::spawn( + async move { + tracing::info!("Event loop active"); + + if let Err(e) = self.run().await { + tracing::error!("Application loop exited with error: {:#}", e); + } } - }); + .with_subscriber(dispatch), + ); handle } diff --git a/src/core/logger.rs b/src/core/logger.rs new file mode 100644 index 0000000..53fcfd3 --- /dev/null +++ b/src/core/logger.rs @@ -0,0 +1,54 @@ +use tokio::sync::broadcast; +use tracing::Subscriber; +use tracing_subscriber::Layer; + +use crate::core::event::Event; + +pub struct CoreLogLayer { + tx: broadcast::Sender, +} + +impl CoreLogLayer { + pub fn new(tx: broadcast::Sender) -> Self { + Self { tx } + } +} + +impl Layer for CoreLogLayer { + fn on_event( + &self, + event: &tracing::Event<'_>, + _ctx: tracing_subscriber::layer::Context<'_, S>, + ) { + let mut visitor = StringVisitor::new(); + event.record(&mut visitor); + + let trace_event = Event::Trace { + level: *event.metadata().level(), + message: visitor.message, + }; + + // We ignore send errors because there might not be any subscribers yet + let _ = self.tx.send(trace_event); + } +} + +struct StringVisitor { + message: String, +} + +impl StringVisitor { + fn new() -> Self { + Self { + message: String::new(), + } + } +} + +impl tracing::field::Visit for StringVisitor { + fn record_debug(&mut self, field: &tracing::field::Field, value: &dyn std::fmt::Debug) { + if field.name() == "message" { + self.message = format!("{:?}", value); + } + } +} diff --git a/src/core/mod.rs b/src/core/mod.rs index f9d6a54..d8329cd 100644 --- a/src/core/mod.rs +++ b/src/core/mod.rs @@ -5,5 +5,6 @@ pub mod database; pub mod engine; pub mod event; pub mod http_agent; +pub mod logger; pub mod message; pub mod version_resolver;