diff --git a/crates/turborepo-lib/src/run/mod.rs b/crates/turborepo-lib/src/run/mod.rs index 7c51674fb8d98..2b83a14747133 100644 --- a/crates/turborepo-lib/src/run/mod.rs +++ b/crates/turborepo-lib/src/run/mod.rs @@ -38,8 +38,9 @@ use turborepo_run_summary::{ObservabilityHandle, RunTracker}; use turborepo_scm::{RepoGitIndex, SCM}; use turborepo_signals::{listeners::get_signal, ShutdownReason, SignalHandler}; use turborepo_task_hash::{ - collect_global_file_hash_inputs, get_external_deps_hash, get_internal_deps_hash, - global_hash::GLOBAL_CACHE_KEY, GlobalHashableInputs, PackageInputsHashes, + collect_global_file_hash_inputs, compute_external_deps_hashes, get_external_deps_hash, + get_internal_deps_hash, global_hash::GLOBAL_CACHE_KEY, GlobalHashableInputs, + PackageInputsHashes, }; use turborepo_telemetry::events::generic::GenericEventBuilder; use turborepo_types::{EnvMode, UIMode}; @@ -970,10 +971,12 @@ impl Run { let is_monorepo = !self.opts.run_opts.single_package; - // Run three expensive I/O-bound operations concurrently using rayon::scope: + // Run four expensive operations concurrently using rayon::scope: // 1. Package file hashing - walks every package's files and computes hashes // 2. Internal deps hashing - walks root internal dependency packages // 3. Global file hash inputs - globwalks global deps and hashes them + // 4. External deps hashing - hashes every package's transitive lockfile + // dependencies (consumed later by the task hasher) // // These are completely independent and dominate the pre-execution phase. // Running them in parallel can significantly reduce wall-clock time. @@ -987,6 +990,7 @@ impl Run { let mut file_hash_result = None; let mut internal_deps_result = None; let mut global_file_result = None; + let mut external_deps_hashes = None; let _hash_scope_span = tracing::info_span!("hash_scope").entered(); crate::rayon_compat::block_in_place(|| { @@ -1036,6 +1040,14 @@ impl Run { &self.scm, )); }); + if is_monorepo { + s.spawn(|_| { + let _span = + tracing::info_span!("compute_external_deps_hashes_task").entered(); + external_deps_hashes = + Some(compute_external_deps_hashes(self.pkg_dep_graph.packages())); + }); + } }); }); @@ -1116,6 +1128,7 @@ impl Run { ui_sender, is_watch, self.micro_frontend_configs.as_ref(), + external_deps_hashes, ) .await; diff --git a/crates/turborepo-lib/src/task_graph/visitor/mod.rs b/crates/turborepo-lib/src/task_graph/visitor/mod.rs index f51755b6378cf..2d1e81b043368 100644 --- a/crates/turborepo-lib/src/task_graph/visitor/mod.rs +++ b/crates/turborepo-lib/src/task_graph/visitor/mod.rs @@ -220,6 +220,7 @@ impl<'a> Visitor<'a> { ui_sender: Option, is_watch: bool, micro_frontends_configs: Option<&'a MicrofrontendsConfigs>, + external_deps_hashes: Option>, ) -> Self { let (task_hasher, color_cache, grouping_layer) = { let _span = tracing::info_span!("visitor_new").entered(); @@ -232,9 +233,15 @@ impl<'a> Visitor<'a> { global_env_patterns, ); - crate::rayon_compat::block_in_place(|| { - task_hasher.precompute_external_deps_hashes(package_graph.packages()); - }); + // The caller may have computed the external dependency hashes + // concurrently with other startup work; fall back to computing + // them here if not. + match external_deps_hashes { + Some(cache) => task_hasher.set_external_deps_hash_cache(cache), + None => crate::rayon_compat::block_in_place(|| { + task_hasher.precompute_external_deps_hashes(package_graph.packages()); + }), + } let color_cache = ColorSelector::default(); diff --git a/crates/turborepo-task-hash/src/lib.rs b/crates/turborepo-task-hash/src/lib.rs index ada707318ffba..91e7b24ed85ae 100644 --- a/crates/turborepo-task-hash/src/lib.rs +++ b/crates/turborepo-task-hash/src/lib.rs @@ -131,6 +131,7 @@ impl PackageInputsHashes { inputs: &'b TaskInputs, } + let collect_span = tracing::info_span!("collect_task_hash_keys").entered(); let mut task_infos = Vec::new(); for task in all_tasks { let TaskNode::Task(task_id) = task else { @@ -185,27 +186,43 @@ impl PackageInputsHashes { unique_hash_keys = unique_keys.len(), "file hash deduplication" ); + drop(collect_span); - // Phase 2: Compute file hashes in parallel across unique keys. + // Phase 2: Compute file hashes in parallel across unique keys. The + // summary hash of each `FileHashes` is computed here too, once per + // unique key, so distribution below never re-hashes for the many + // tasks that share a key. // EMFILE (too many open files) errors are handled via retry-with-backoff // in the globwalk and hash_objects layers, so we can safely parallelize // all keys on rayon without worrying about fd exhaustion. - let file_hash_results: Vec, Error>> = unique_keys + let hash_span = tracing::info_span!("hash_unique_inputs").entered(); + let file_hash_results: Vec, String), Error>> = unique_keys .into_par_iter() .map(|(package_path, globs, default, eager)| { - if !eager { - return Ok(Arc::new(FileHashes(Vec::new()))); - } - - file_hashes_for_inputs(scm, repo_root, &package_path, &globs, default, repo_index) + let file_hashes = if !eager { + Arc::new(FileHashes(Vec::new())) + } else { + file_hashes_for_inputs( + scm, + repo_root, + &package_path, + &globs, + default, + repo_index, + )? + }; + let hash = file_hashes.as_ref().hash(); + Ok((file_hashes, hash)) }) .collect(); - let file_hash_results: Vec> = file_hash_results + let file_hash_results: Vec<(Arc, String)> = file_hash_results .into_iter() .collect::, _>>()?; + drop(hash_span); // Phase 3: Distribute shared results to individual tasks. + let _span = tracing::info_span!("distribute_task_file_hashes").entered(); let mut hashes = HashMap::with_capacity(task_infos.len()); let mut expanded_hashes = if needs_expanded_hashes { HashMap::with_capacity(task_infos.len()) @@ -215,11 +232,9 @@ impl PackageInputsHashes { for (i, info) in task_infos.into_iter().enumerate() { let key_idx = task_key_map[i]; - let file_hashes = &file_hash_results[key_idx]; + let (file_hashes, hash) = &file_hash_results[key_idx]; - let hash = file_hashes.as_ref().hash(); - - hashes.insert(info.task_id.clone(), hash); + hashes.insert(info.task_id.clone(), hash.clone()); if needs_expanded_hashes || info.inputs.has_deferred_inputs() { expanded_hashes.insert(info.task_id, Arc::clone(file_hashes)); } @@ -232,6 +247,22 @@ impl PackageInputsHashes { } } +/// Compute the external dependency hash for every workspace in parallel. +/// Many tasks share the same package, so hashing once per package avoids +/// re-sorting transitive dependencies for every task. +#[tracing::instrument(skip_all)] +pub fn compute_external_deps_hashes<'b>( + workspaces: impl Iterator, +) -> HashMap { + let ws: Vec<_> = workspaces.collect(); + ws.par_iter() + .map(|(name, info)| { + let hash = get_external_deps_hash(&info.transitive_dependencies); + (name.as_str().to_owned(), hash) + }) + .collect() +} + #[derive(Default, Debug, Clone)] pub struct TaskHashTracker { state: Arc>, @@ -315,14 +346,15 @@ impl<'a, R: RunOptsHashInfo> TaskHasher<'a, R> { if self.run_opts.single_package() { return; } - let ws: Vec<_> = workspaces.collect(); - self.external_deps_hash_cache = ws - .par_iter() - .map(|(name, info)| { - let hash = get_external_deps_hash(&info.transitive_dependencies); - (name.as_str().to_owned(), hash) - }) - .collect(); + self.external_deps_hash_cache = compute_external_deps_hashes(workspaces); + } + + /// Install an externally computed dependency-hash cache (see + /// [`compute_external_deps_hashes`]). Lets callers compute the cache + /// concurrently with other startup work instead of serially during + /// hasher construction. + pub fn set_external_deps_hash_cache(&mut self, cache: HashMap) { + self.external_deps_hash_cache = cache; } #[tracing::instrument(skip(self, task_definition, task_env_mode, workspace, dependency_set))]