From b6fc12b2a12e0717da5f6f4daf56e7aefcc11496 Mon Sep 17 00:00:00 2001 From: jif-oai Date: Wed, 24 Jun 2026 16:21:13 +0100 Subject: [PATCH 1/5] Pipeline lexical workspace metadata probes --- codex-rs/core/src/agents_md.rs | 114 +++++++++++++-------------- codex-rs/core/src/agents_md_tests.rs | 71 ++++++++++++++++- codex-rs/core/src/git_info_tests.rs | 16 ++++ codex-rs/git-utils/src/info.rs | 48 ++++++----- 4 files changed, 162 insertions(+), 87 deletions(-) diff --git a/codex-rs/core/src/agents_md.rs b/codex-rs/core/src/agents_md.rs index 73dfbeb548e3..cf7cd500882a 100644 --- a/codex-rs/core/src/agents_md.rs +++ b/codex-rs/core/src/agents_md.rs @@ -28,6 +28,7 @@ use codex_exec_server::ExecutorFileSystem; use codex_extension_api::UserInstructions; use codex_utils_absolute_path::AbsolutePathBuf; use codex_utils_path_uri::PathUri; +use futures::future::join_all; use std::io; use toml::Value as TomlValue; use tracing::error; @@ -104,13 +105,6 @@ async fn read_agents_md( break; } - match fs.get_metadata(&p, /*sandbox*/ None).await { - Ok(metadata) if !metadata.is_file => continue, - Ok(_) => {} - Err(err) if err.kind() == io::ErrorKind::NotFound => continue, - Err(err) => return Err(err), - } - let mut data = match fs.read_file(&p, /*sandbox*/ None).await { Ok(data) => data, Err(err) if err.kind() == io::ErrorKind::NotFound => continue, @@ -177,68 +171,74 @@ async fn agents_md_paths( default_project_root_markers() } }; - let mut project_root = None; - if !project_root_markers.is_empty() { - for current in dir.ancestors() { - for marker in &project_root_markers { - let marker_path = current + let ancestors = dir.ancestors().collect::>(); + let candidate_filenames = candidate_filenames(config); + let marker_probe_count = ancestors.len() * project_root_markers.len(); + let mut probes = + Vec::with_capacity(marker_probe_count + ancestors.len() * candidate_filenames.len()); + for ancestor in &ancestors { + for marker in &project_root_markers { + probes.push( + ancestor .join(marker) - .map_err(|err| io::Error::new(io::ErrorKind::InvalidInput, err))?; - let marker_exists = match fs.get_metadata(&marker_path, /*sandbox*/ None).await { - Ok(_) => true, - Err(err) if err.kind() == io::ErrorKind::NotFound => false, - Err(err) => return Err(err), - }; - if marker_exists { - project_root = Some(current.clone()); - break; - } - } - if project_root.is_some() { - break; - } + .map_err(|err| io::Error::new(io::ErrorKind::InvalidInput, err))?, + ); } } - - let search_dirs: Vec = if let Some(root) = project_root { - let mut dirs = Vec::new(); - let mut cursor = dir.clone(); - loop { - dirs.push(cursor.clone()); - if cursor == root { - break; - } - let Some(parent) = cursor.parent() else { - break; - }; - cursor = parent; + let mut candidate_paths = Vec::with_capacity(ancestors.len() * candidate_filenames.len()); + for ancestor in &ancestors { + for name in &candidate_filenames { + candidate_paths.push( + ancestor + .join(name) + .map_err(|err| io::Error::new(io::ErrorKind::InvalidInput, err))?, + ); } - dirs.reverse(); - dirs - } else { - vec![dir] - }; + } + probes.extend(candidate_paths.iter().cloned()); - let mut found: Vec = Vec::new(); - let candidate_filenames = candidate_filenames(config); - for d in search_dirs { - for name in &candidate_filenames { - let candidate = d - .join(name) - .map_err(|err| io::Error::new(io::ErrorKind::InvalidInput, err))?; - match fs.get_metadata(&candidate, /*sandbox*/ None).await { - Ok(md) if md.is_file => { - found.push(candidate); + let mut metadata = join_all( + probes + .iter() + .map(|path| fs.get_metadata(path, /*sandbox*/ None)), + ) + .await; + let candidate_metadata = metadata.split_off(marker_probe_count); + + let mut project_root_index = None; + if !project_root_markers.is_empty() { + for (probe_index, metadata) in metadata.into_iter().enumerate() { + match metadata { + Ok(_) => { + project_root_index = Some(probe_index / project_root_markers.len()); break; } - Ok(_) => {} - Err(err) if err.kind() == io::ErrorKind::NotFound => continue, + Err(err) if err.kind() == io::ErrorKind::NotFound => {} Err(err) => return Err(err), } } } - Ok(found) + let search_through_index = project_root_index.unwrap_or(0); + let mut selected = vec![None; search_through_index + 1]; + for (probe_index, (candidate, metadata)) in candidate_paths + .into_iter() + .zip(candidate_metadata) + .enumerate() + { + let ancestor_index = probe_index / candidate_filenames.len(); + if ancestor_index > search_through_index || selected[ancestor_index].is_some() { + continue; + } + match metadata { + Ok(metadata) if metadata.is_file => selected[ancestor_index] = Some(candidate), + Ok(_) => {} + Err(err) if err.kind() == io::ErrorKind::NotFound => {} + Err(err) => return Err(err), + } + } + + Ok(selected.into_iter().rev().flatten().collect()) } fn candidate_filenames(config: &Config) -> Vec<&str> { diff --git a/codex-rs/core/src/agents_md_tests.rs b/codex-rs/core/src/agents_md_tests.rs index 0122910c3b94..636e4c30a920 100644 --- a/codex-rs/core/src/agents_md_tests.rs +++ b/codex-rs/core/src/agents_md_tests.rs @@ -30,6 +30,8 @@ use std::ops::Deref; use std::ops::DerefMut; use std::path::PathBuf; use std::sync::Arc; +use std::sync::atomic::AtomicUsize; +use std::sync::atomic::Ordering; use tempfile::TempDir; #[derive(Clone, Copy)] @@ -41,6 +43,14 @@ enum InjectedFailure { struct FailingFileSystem { path: AbsolutePathBuf, failure: InjectedFailure, + metadata_calls: Arc, +} + +#[derive(Default)] +struct MetadataCallCounts { + active_calls: AtomicUsize, + max_active_calls: AtomicUsize, + scalar_calls: AtomicUsize, } impl FailingFileSystem { @@ -88,12 +98,30 @@ impl FailingFileSystem { path: &PathUri, sandbox: Option<&FileSystemSandboxContext>, ) -> io::Result { - if path.to_abs_path()? == self.path + let path_abs = path.to_abs_path()?; + self.metadata_calls + .scalar_calls + .fetch_add(1, Ordering::Relaxed); + let active_calls = self + .metadata_calls + .active_calls + .fetch_add(1, Ordering::Relaxed) + + 1; + self.metadata_calls + .max_active_calls + .fetch_max(active_calls, Ordering::Relaxed); + tokio::task::yield_now().await; + let result = if path_abs == self.path && let InjectedFailure::Metadata(kind) = self.failure { - return Err(io::Error::new(kind, "injected metadata failure")); - } - LOCAL_FS.get_metadata(path, sandbox).await + Err(io::Error::new(kind, "injected metadata failure")) + } else { + LOCAL_FS.get_metadata(path, sandbox).await + }; + self.metadata_calls + .active_calls + .fetch_sub(1, Ordering::Relaxed); + result } async fn read_directory( @@ -600,6 +628,7 @@ async fn read_agents_md_propagates_metadata_errors() { let fs = FailingFileSystem { path: marker_path, failure: InjectedFailure::Metadata(io::ErrorKind::PermissionDenied), + metadata_calls: Arc::default(), }; let cwd = config.cwd.clone(); @@ -618,6 +647,7 @@ async fn read_agents_md_propagates_read_errors() { let fs = FailingFileSystem { path: config.cwd.join("AGENTS.md"), failure: InjectedFailure::Read(io::ErrorKind::PermissionDenied), + metadata_calls: Arc::default(), }; let cwd = config.cwd.clone(); @@ -636,6 +666,7 @@ async fn read_agents_md_ignores_files_removed_after_discovery() { let fs = FailingFileSystem { path: config.cwd.join("AGENTS.md"), failure: InjectedFailure::Read(io::ErrorKind::NotFound), + metadata_calls: Arc::default(), }; let cwd = config.cwd.clone(); @@ -646,6 +677,38 @@ async fn read_agents_md_ignores_files_removed_after_discovery() { assert_eq!(loaded, None); } +#[tokio::test] +async fn read_agents_md_pipelines_all_lexical_metadata_probes() { + let tmp = tempfile::tempdir().expect("tempdir"); + fs::write(tmp.path().join(".git"), "").unwrap(); + fs::write(tmp.path().join("AGENTS.md"), "project doc").unwrap(); + let nested = tmp.path().join("nested"); + fs::create_dir(&nested).unwrap(); + + let mut config = make_config(&tmp, /*limit*/ 4096, /*instructions*/ None).await; + config.cwd = nested.abs(); + let metadata_calls = Arc::new(MetadataCallCounts::default()); + let fs = FailingFileSystem { + path: config.cwd.join("unused"), + failure: InjectedFailure::Read(io::ErrorKind::PermissionDenied), + metadata_calls: Arc::clone(&metadata_calls), + }; + let cwd = PathUri::from_abs_path(&config.cwd); + + let loaded = read_agents_md(&config.config, &fs, "local", &cwd) + .await + .expect("project instructions") + .expect("project instructions"); + + let expected_probe_count = cwd.ancestors().count() * 3; + assert_eq!(loaded.text(), "project doc"); + assert_eq!( + metadata_calls.scalar_calls.load(Ordering::Relaxed), + expected_probe_count + ); + assert!(metadata_calls.max_active_calls.load(Ordering::Relaxed) > 1); +} + /// When `cwd` is nested inside a repo, the search should locate AGENTS.md /// placed at the repository root (identified by `.git`). #[tokio::test] diff --git a/codex-rs/core/src/git_info_tests.rs b/codex-rs/core/src/git_info_tests.rs index e172cd5f20fe..04f28c2c9a3f 100644 --- a/codex-rs/core/src/git_info_tests.rs +++ b/codex-rs/core/src/git_info_tests.rs @@ -501,6 +501,22 @@ async fn get_git_repo_root_with_fs_detects_gitdir_pointer() { ); } +#[tokio::test] +async fn get_git_repo_root_with_fs_starts_at_parent_for_file() { + let tmp = TempDir::new().expect("tempdir"); + let proj = tmp.path().join("proj"); + let nested = proj.join("nested"); + std::fs::create_dir_all(proj.join(".git")).unwrap(); + std::fs::create_dir_all(&nested).unwrap(); + let file = nested.join("file.txt"); + std::fs::write(&file, "contents").unwrap(); + + assert_eq!( + get_git_repo_root_with_fs(LOCAL_FS.as_ref(), &file.abs()).await, + Some(proj.abs()) + ); +} + #[tokio::test] async fn resolve_root_git_project_for_trust_regular_repo_returns_repo_root() { let temp_dir = TempDir::new().expect("Failed to create temp dir"); diff --git a/codex-rs/git-utils/src/info.rs b/codex-rs/git-utils/src/info.rs index 48d229df8db5..b87ad87317ab 100644 --- a/codex-rs/git-utils/src/info.rs +++ b/codex-rs/git-utils/src/info.rs @@ -47,14 +47,28 @@ pub async fn get_git_repo_root_with_fs( fs: &dyn ExecutorFileSystem, cwd: &AbsolutePathBuf, ) -> Option { - let cwd_uri = PathUri::from_abs_path(cwd); - let base = match fs.get_metadata(&cwd_uri, /*sandbox*/ None).await { - Ok(metadata) if metadata.is_directory => cwd.clone(), - _ => cwd.parent()?, - }; - find_ancestor_git_entry_with_fs(fs, &base) - .await - .map(|(repo_root, _)| repo_root) + let ancestors = cwd.ancestors().collect::>(); + let mut probes = Vec::with_capacity(ancestors.len() + 1); + probes.push(PathUri::from_abs_path(cwd)); + probes.extend( + ancestors + .iter() + .map(|ancestor| PathUri::from_abs_path(&ancestor.join(".git"))), + ); + + let mut metadata = join_all( + probes + .iter() + .map(|path| fs.get_metadata(path, /*sandbox*/ None)), + ) + .await + .into_iter(); + let cwd_is_directory = matches!(metadata.next()?, Ok(metadata) if metadata.is_directory); + ancestors + .into_iter() + .zip(metadata) + .skip(usize::from(!cwd_is_directory)) + .find_map(|(ancestor, metadata)| metadata.is_ok().then_some(ancestor)) } /// Timeout for git commands to prevent freezing on large repositories @@ -853,24 +867,6 @@ fn find_ancestor_git_entry(base_dir: &Path) -> Option<(PathBuf, PathBuf)> { None } -async fn find_ancestor_git_entry_with_fs( - fs: &dyn ExecutorFileSystem, - base_dir: &AbsolutePathBuf, -) -> Option<(AbsolutePathBuf, AbsolutePathBuf)> { - for dir in base_dir.ancestors() { - let dot_git = dir.join(".git"); - let dot_git_uri = PathUri::from_abs_path(&dot_git); - if fs - .get_metadata(&dot_git_uri, /*sandbox*/ None) - .await - .is_ok() - { - return Some((dir, dot_git)); - } - } - None -} - /// Returns a list of local git branches. /// Includes the default branch at the beginning of the list, if it exists. pub async fn local_git_branches(cwd: &Path) -> Vec { From a993329848a2d5f066b95a18171ae82fea083b9b Mon Sep 17 00:00:00 2001 From: jif-oai Date: Wed, 24 Jun 2026 17:43:24 +0100 Subject: [PATCH 2/5] Bound lexical root discovery --- codex-rs/core/src/agents_md.rs | 99 ++++++++------------ codex-rs/core/src/agents_md_tests.rs | 130 ++++++++++++++++++++++++--- codex-rs/file-system/src/find_up.rs | 69 ++++++++++++++ codex-rs/file-system/src/lib.rs | 5 ++ codex-rs/git-utils/src/info.rs | 37 ++++---- 5 files changed, 246 insertions(+), 94 deletions(-) create mode 100644 codex-rs/file-system/src/find_up.rs diff --git a/codex-rs/core/src/agents_md.rs b/codex-rs/core/src/agents_md.rs index cf7cd500882a..b2b3cfea20e7 100644 --- a/codex-rs/core/src/agents_md.rs +++ b/codex-rs/core/src/agents_md.rs @@ -26,9 +26,10 @@ use codex_config::merge_toml_values; use codex_config::project_root_markers_from_config; use codex_exec_server::ExecutorFileSystem; use codex_extension_api::UserInstructions; +use codex_file_system::FindUpErrorPolicy; +use codex_file_system::find_nearest_ancestor_with_markers; use codex_utils_absolute_path::AbsolutePathBuf; use codex_utils_path_uri::PathUri; -use futures::future::join_all; use std::io; use toml::Value as TomlValue; use tracing::error; @@ -171,74 +172,52 @@ async fn agents_md_paths( default_project_root_markers() } }; - let ancestors = dir.ancestors().collect::>(); - let candidate_filenames = candidate_filenames(config); - let marker_probe_count = ancestors.len() * project_root_markers.len(); - let mut probes = - Vec::with_capacity(marker_probe_count + ancestors.len() * candidate_filenames.len()); - for ancestor in &ancestors { - for marker in &project_root_markers { - probes.push( - ancestor - .join(marker) - .map_err(|err| io::Error::new(io::ErrorKind::InvalidInput, err))?, - ); - } - } - let mut candidate_paths = Vec::with_capacity(ancestors.len() * candidate_filenames.len()); - for ancestor in &ancestors { - for name in &candidate_filenames { - candidate_paths.push( - ancestor - .join(name) - .map_err(|err| io::Error::new(io::ErrorKind::InvalidInput, err))?, - ); + let project_root = find_nearest_ancestor_with_markers( + fs, + &dir, + project_root_markers, + FindUpErrorPolicy::Propagate, + /*sandbox*/ None, + ) + .await?; + let search_dirs = if let Some(root) = project_root { + let mut dirs = Vec::new(); + let mut cursor = dir.clone(); + loop { + dirs.push(cursor.clone()); + if cursor == root { + break; + } + let Some(parent) = cursor.parent() else { + break; + }; + cursor = parent; } - } - probes.extend(candidate_paths.iter().cloned()); + dirs.reverse(); + dirs + } else { + vec![dir] + }; - let mut metadata = join_all( - probes - .iter() - .map(|path| fs.get_metadata(path, /*sandbox*/ None)), - ) - .await; - let candidate_metadata = metadata.split_off(marker_probe_count); - - let mut project_root_index = None; - if !project_root_markers.is_empty() { - for (probe_index, metadata) in metadata.into_iter().enumerate() { - match metadata { - Ok(_) => { - project_root_index = Some(probe_index / project_root_markers.len()); + let mut found = Vec::new(); + let candidate_filenames = candidate_filenames(config); + for directory in search_dirs { + for name in &candidate_filenames { + let candidate = directory + .join(name) + .map_err(|err| io::Error::new(io::ErrorKind::InvalidInput, err))?; + match fs.get_metadata(&candidate, /*sandbox*/ None).await { + Ok(metadata) if metadata.is_file => { + found.push(candidate); break; } + Ok(_) => {} Err(err) if err.kind() == io::ErrorKind::NotFound => {} Err(err) => return Err(err), } } } - - let search_through_index = project_root_index.unwrap_or(0); - let mut selected = vec![None; search_through_index + 1]; - for (probe_index, (candidate, metadata)) in candidate_paths - .into_iter() - .zip(candidate_metadata) - .enumerate() - { - let ancestor_index = probe_index / candidate_filenames.len(); - if ancestor_index > search_through_index || selected[ancestor_index].is_some() { - continue; - } - match metadata { - Ok(metadata) if metadata.is_file => selected[ancestor_index] = Some(candidate), - Ok(_) => {} - Err(err) if err.kind() == io::ErrorKind::NotFound => {} - Err(err) => return Err(err), - } - } - - Ok(selected.into_iter().rev().flatten().collect()) + Ok(found) } fn candidate_filenames(config: &Config) -> Vec<&str> { diff --git a/codex-rs/core/src/agents_md_tests.rs b/codex-rs/core/src/agents_md_tests.rs index 636e4c30a920..ca5e289a3416 100644 --- a/codex-rs/core/src/agents_md_tests.rs +++ b/codex-rs/core/src/agents_md_tests.rs @@ -30,6 +30,7 @@ use std::ops::Deref; use std::ops::DerefMut; use std::path::PathBuf; use std::sync::Arc; +use std::sync::Mutex; use std::sync::atomic::AtomicUsize; use std::sync::atomic::Ordering; use tempfile::TempDir; @@ -37,6 +38,7 @@ use tempfile::TempDir; #[derive(Clone, Copy)] enum InjectedFailure { Metadata(io::ErrorKind), + MetadataPending, Read(io::ErrorKind), } @@ -50,7 +52,7 @@ struct FailingFileSystem { struct MetadataCallCounts { active_calls: AtomicUsize, max_active_calls: AtomicUsize, - scalar_calls: AtomicUsize, + paths: Mutex>, } impl FailingFileSystem { @@ -100,8 +102,10 @@ impl FailingFileSystem { ) -> io::Result { let path_abs = path.to_abs_path()?; self.metadata_calls - .scalar_calls - .fetch_add(1, Ordering::Relaxed); + .paths + .lock() + .expect("metadata paths lock") + .push(path.clone()); let active_calls = self .metadata_calls .active_calls @@ -111,12 +115,16 @@ impl FailingFileSystem { .max_active_calls .fetch_max(active_calls, Ordering::Relaxed); tokio::task::yield_now().await; - let result = if path_abs == self.path - && let InjectedFailure::Metadata(kind) = self.failure - { - Err(io::Error::new(kind, "injected metadata failure")) - } else { - LOCAL_FS.get_metadata(path, sandbox).await + let result = match self.failure { + InjectedFailure::Metadata(kind) if path_abs == self.path => { + Err(io::Error::new(kind, "injected metadata failure")) + } + InjectedFailure::MetadataPending if path_abs == self.path => { + std::future::pending().await + } + InjectedFailure::Metadata(_) + | InjectedFailure::MetadataPending + | InjectedFailure::Read(_) => LOCAL_FS.get_metadata(path, sandbox).await, }; self.metadata_calls .active_calls @@ -678,7 +686,7 @@ async fn read_agents_md_ignores_files_removed_after_discovery() { } #[tokio::test] -async fn read_agents_md_pipelines_all_lexical_metadata_probes() { +async fn read_agents_md_pipelines_root_markers_before_candidate_search() { let tmp = tempfile::tempdir().expect("tempdir"); fs::write(tmp.path().join(".git"), "").unwrap(); fs::write(tmp.path().join("AGENTS.md"), "project doc").unwrap(); @@ -700,13 +708,107 @@ async fn read_agents_md_pipelines_all_lexical_metadata_probes() { .expect("project instructions") .expect("project instructions"); - let expected_probe_count = cwd.ancestors().count() * 3; assert_eq!(loaded.text(), "project doc"); + let max_active_calls = metadata_calls.max_active_calls.load(Ordering::Relaxed); + assert!(max_active_calls > 1); + assert!(max_active_calls <= 8); + let candidate_paths = metadata_calls + .paths + .lock() + .expect("metadata paths lock") + .iter() + .filter(|path| path.basename().as_deref() != Some(".git")) + .cloned() + .collect::>(); + assert_eq!( + candidate_paths, + vec![ + PathUri::from_abs_path(&tmp.path().join(LOCAL_AGENTS_MD_FILENAME).abs()), + PathUri::from_abs_path(&tmp.path().join(DEFAULT_AGENTS_MD_FILENAME).abs()), + cwd.join(LOCAL_AGENTS_MD_FILENAME).expect("override path"), + cwd.join(DEFAULT_AGENTS_MD_FILENAME).expect("agents path"), + ] + ); +} + +#[tokio::test] +async fn marker_search_does_not_wait_for_a_higher_ancestor() { + let tmp = tempfile::tempdir().expect("tempdir"); + fs::write(tmp.path().join(".git"), "").unwrap(); + fs::write(tmp.path().join("AGENTS.md"), "project doc").unwrap(); + let nested = tmp.path().join("nested"); + fs::create_dir(&nested).unwrap(); + + let mut config = make_config(&tmp, /*limit*/ 4096, /*instructions*/ None).await; + config.cwd = nested.abs(); + let pending_marker = tmp + .path() + .parent() + .expect("tempdir parent") + .join(".git") + .abs(); + let fs = FailingFileSystem { + path: pending_marker, + failure: InjectedFailure::MetadataPending, + metadata_calls: Arc::default(), + }; + let cwd = PathUri::from_abs_path(&config.cwd); + + let paths = tokio::time::timeout( + std::time::Duration::from_secs(1), + super::agents_md_paths(&config.config, &cwd, &fs), + ) + .await + .expect("nearest marker should complete") + .expect("AGENTS.md discovery"); + + assert_eq!( + paths, + vec![PathUri::from_abs_path( + &tmp.path().join(DEFAULT_AGENTS_MD_FILENAME).abs() + )] + ); +} + +#[tokio::test] +async fn empty_project_root_markers_only_probe_cwd_candidates() { + let tmp = tempfile::tempdir().expect("tempdir"); + fs::write(tmp.path().join("AGENTS.md"), "parent doc").unwrap(); + let nested = tmp.path().join("nested"); + fs::create_dir(&nested).unwrap(); + fs::write(nested.join("AGENTS.md"), "cwd doc").unwrap(); + + let mut config = make_config_with_project_root_markers( + &tmp, + /*limit*/ 4096, + /*instructions*/ None, + &[], + ) + .await; + config.cwd = nested.abs(); + let metadata_calls = Arc::new(MetadataCallCounts::default()); + let fs = FailingFileSystem { + path: config.cwd.join("unused"), + failure: InjectedFailure::Read(io::ErrorKind::PermissionDenied), + metadata_calls: Arc::clone(&metadata_calls), + }; + let cwd = PathUri::from_abs_path(&config.cwd); + + let paths = super::agents_md_paths(&config.config, &cwd, &fs) + .await + .expect("AGENTS.md discovery"); + + let override_path = cwd.join(LOCAL_AGENTS_MD_FILENAME).expect("override path"); + let agents_path = cwd.join(DEFAULT_AGENTS_MD_FILENAME).expect("agents path"); + assert_eq!(paths, vec![agents_path.clone()]); assert_eq!( - metadata_calls.scalar_calls.load(Ordering::Relaxed), - expected_probe_count + metadata_calls + .paths + .lock() + .expect("metadata paths lock") + .clone(), + vec![override_path, agents_path] ); - assert!(metadata_calls.max_active_calls.load(Ordering::Relaxed) > 1); } /// When `cwd` is nested inside a repo, the search should locate AGENTS.md diff --git a/codex-rs/file-system/src/find_up.rs b/codex-rs/file-system/src/find_up.rs new file mode 100644 index 000000000000..b3e26c5570fd --- /dev/null +++ b/codex-rs/file-system/src/find_up.rs @@ -0,0 +1,69 @@ +use crate::ExecutorFileSystem; +use crate::FileSystemResult; +use crate::FileSystemSandboxContext; +use codex_utils_path_uri::PathUri; +use futures::StreamExt; +use std::io; + +const MAX_CONCURRENT_PROBES: usize = 8; + +/// Controls how an upward marker search handles metadata errors other than `NotFound`. +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum FindUpErrorPolicy { + /// Return the first error in lexical search order. + Propagate, + /// Treat errors as missing markers and continue searching. + Ignore, +} + +/// Finds the nearest ancestor containing one of the provided marker names. +/// +/// Marker paths are probed in lexical order from `start` toward the filesystem root. A bounded +/// number of ordinary metadata calls are kept in flight so remote filesystems can pipeline them +/// without requiring a batch protocol operation. +pub async fn find_nearest_ancestor_with_markers( + file_system: &dyn ExecutorFileSystem, + start: &PathUri, + markers: Vec, + error_policy: FindUpErrorPolicy, + sandbox: Option<&FileSystemSandboxContext>, +) -> FileSystemResult> { + let mut ancestors = start.ancestors(); + let mut ancestor = ancestors.next(); + let mut marker_index = 0; + let probes = std::iter::from_fn(move || { + let current_ancestor = ancestor.clone()?; + let marker = markers.get(marker_index)?; + let marker_path = current_ancestor + .join(marker) + .map_err(|err| io::Error::new(io::ErrorKind::InvalidInput, err)); + + marker_index += 1; + if marker_index == markers.len() { + marker_index = 0; + ancestor = ancestors.next(); + } + + Some((current_ancestor, marker_path)) + }); + let mut results = futures::stream::iter(probes) + .map(|(ancestor, marker_path)| async move { + let marker_path = marker_path?; + match file_system.get_metadata(&marker_path, sandbox).await { + Ok(_) => Ok(Some(ancestor)), + Err(err) if err.kind() == io::ErrorKind::NotFound => Ok(None), + Err(err) => match error_policy { + FindUpErrorPolicy::Propagate => Err(err), + FindUpErrorPolicy::Ignore => Ok(None), + }, + } + }) + .buffered(MAX_CONCURRENT_PROBES); + + while let Some(result) = results.next().await { + if let Some(ancestor) = result? { + return Ok(Some(ancestor)); + } + } + Ok(None) +} diff --git a/codex-rs/file-system/src/lib.rs b/codex-rs/file-system/src/lib.rs index 3ae4b56753a7..5b809b8dab5a 100644 --- a/codex-rs/file-system/src/lib.rs +++ b/codex-rs/file-system/src/lib.rs @@ -1,3 +1,8 @@ +mod find_up; + +pub use find_up::FindUpErrorPolicy; +pub use find_up::find_nearest_ancestor_with_markers; + use bytes::Bytes; use codex_protocol::config_types::WindowsSandboxLevel; use codex_protocol::models::ManagedFileSystemPermissions; diff --git a/codex-rs/git-utils/src/info.rs b/codex-rs/git-utils/src/info.rs index b87ad87317ab..32d1b6f0a90b 100644 --- a/codex-rs/git-utils/src/info.rs +++ b/codex-rs/git-utils/src/info.rs @@ -5,6 +5,8 @@ use std::path::Path; use std::path::PathBuf; use codex_file_system::ExecutorFileSystem; +use codex_file_system::FindUpErrorPolicy; +use codex_file_system::find_nearest_ancestor_with_markers; use codex_utils_absolute_path::AbsolutePathBuf; use codex_utils_path_uri::PathUri; use futures::future::join_all; @@ -47,28 +49,23 @@ pub async fn get_git_repo_root_with_fs( fs: &dyn ExecutorFileSystem, cwd: &AbsolutePathBuf, ) -> Option { - let ancestors = cwd.ancestors().collect::>(); - let mut probes = Vec::with_capacity(ancestors.len() + 1); - probes.push(PathUri::from_abs_path(cwd)); - probes.extend( - ancestors - .iter() - .map(|ancestor| PathUri::from_abs_path(&ancestor.join(".git"))), - ); - - let mut metadata = join_all( - probes - .iter() - .map(|path| fs.get_metadata(path, /*sandbox*/ None)), + let cwd_uri = PathUri::from_abs_path(cwd); + let base = match fs.get_metadata(&cwd_uri, /*sandbox*/ None).await { + Ok(metadata) if metadata.is_directory => cwd.clone(), + _ => cwd.parent()?, + }; + let base_uri = PathUri::from_abs_path(&base); + find_nearest_ancestor_with_markers( + fs, + &base_uri, + vec![".git".to_string()], + FindUpErrorPolicy::Ignore, + /*sandbox*/ None, ) .await - .into_iter(); - let cwd_is_directory = matches!(metadata.next()?, Ok(metadata) if metadata.is_directory); - ancestors - .into_iter() - .zip(metadata) - .skip(usize::from(!cwd_is_directory)) - .find_map(|(ancestor, metadata)| metadata.is_ok().then_some(ancestor)) + .ok()?? + .to_abs_path() + .ok() } /// Timeout for git commands to prevent freezing on large repositories From 8c514430ff17b0c2d7d3609a668f9b38e2550373 Mon Sep 17 00:00:00 2001 From: jif-oai Date: Wed, 24 Jun 2026 18:34:44 +0100 Subject: [PATCH 3/5] Preserve native paths in Git root discovery --- codex-rs/core/src/git_info_tests.rs | 18 +++++++++ codex-rs/file-system/src/find_up.rs | 60 +++++++++++++++++++++++++++-- codex-rs/file-system/src/lib.rs | 1 + codex-rs/git-utils/src/info.rs | 11 ++---- 4 files changed, 79 insertions(+), 11 deletions(-) diff --git a/codex-rs/core/src/git_info_tests.rs b/codex-rs/core/src/git_info_tests.rs index 04f28c2c9a3f..3da69c0e3785 100644 --- a/codex-rs/core/src/git_info_tests.rs +++ b/codex-rs/core/src/git_info_tests.rs @@ -11,6 +11,7 @@ use codex_utils_path::normalize_for_path_comparison; use core_test_support::PathBufExt; use core_test_support::PathExt; use core_test_support::skip_if_sandbox; +use pretty_assertions::assert_eq; use std::fs; #[cfg(unix)] use std::os::unix::fs::PermissionsExt; @@ -517,6 +518,23 @@ async fn get_git_repo_root_with_fs_starts_at_parent_for_file() { ); } +#[cfg(windows)] +#[tokio::test] +async fn get_git_repo_root_with_fs_supports_windows_namespace_paths() { + let tmp = TempDir::new().expect("tempdir"); + let repo = tmp.path().join("repo"); + std::fs::create_dir_all(repo.join(".git")).unwrap(); + std::fs::create_dir_all(repo.join("nested")).unwrap(); + + let namespace_repo = PathBuf::from(format!(r"\\?\{}", repo.display())); + let namespace_nested = namespace_repo.join("nested"); + + assert_eq!( + get_git_repo_root_with_fs(LOCAL_FS.as_ref(), &namespace_nested.abs()).await, + Some(namespace_repo.abs()) + ); +} + #[tokio::test] async fn resolve_root_git_project_for_trust_regular_repo_returns_repo_root() { let temp_dir = TempDir::new().expect("Failed to create temp dir"); diff --git a/codex-rs/file-system/src/find_up.rs b/codex-rs/file-system/src/find_up.rs index b3e26c5570fd..c0278569f7a1 100644 --- a/codex-rs/file-system/src/find_up.rs +++ b/codex-rs/file-system/src/find_up.rs @@ -1,6 +1,7 @@ use crate::ExecutorFileSystem; use crate::FileSystemResult; use crate::FileSystemSandboxContext; +use codex_utils_absolute_path::AbsolutePathBuf; use codex_utils_path_uri::PathUri; use futures::StreamExt; use std::io; @@ -28,15 +29,66 @@ pub async fn find_nearest_ancestor_with_markers( error_policy: FindUpErrorPolicy, sandbox: Option<&FileSystemSandboxContext>, ) -> FileSystemResult> { - let mut ancestors = start.ancestors(); + find_nearest_ancestor( + file_system, + start.clone(), + markers, + PathUri::parent, + |ancestor, marker| { + ancestor + .join(marker) + .map_err(|err| io::Error::new(io::ErrorKind::InvalidInput, err)) + }, + error_policy, + sandbox, + ) + .await +} + +/// Finds the nearest native ancestor containing one of the provided marker names. +/// +/// Ancestors and marker paths remain native until each complete probe is converted to a URI. This +/// preserves paths that require an opaque [`PathUri`] fallback. +pub async fn find_nearest_native_ancestor_with_markers( + file_system: &dyn ExecutorFileSystem, + start: &AbsolutePathBuf, + markers: Vec, + error_policy: FindUpErrorPolicy, + sandbox: Option<&FileSystemSandboxContext>, +) -> FileSystemResult> { + find_nearest_ancestor( + file_system, + start.clone(), + markers, + AbsolutePathBuf::parent, + |ancestor, marker| Ok(PathUri::from_abs_path(&ancestor.join(marker))), + error_policy, + sandbox, + ) + .await +} + +async fn find_nearest_ancestor( + file_system: &dyn ExecutorFileSystem, + start: P, + markers: Vec, + parent: Parent, + mut marker_path: MarkerPath, + error_policy: FindUpErrorPolicy, + sandbox: Option<&FileSystemSandboxContext>, +) -> FileSystemResult> +where + P: Clone + Send, + Parent: FnMut(&P) -> Option

+ Send, + MarkerPath: FnMut(&P, &str) -> FileSystemResult + Send, +{ + let mut ancestors = std::iter::successors(Some(start), parent); let mut ancestor = ancestors.next(); let mut marker_index = 0; let probes = std::iter::from_fn(move || { let current_ancestor = ancestor.clone()?; let marker = markers.get(marker_index)?; - let marker_path = current_ancestor - .join(marker) - .map_err(|err| io::Error::new(io::ErrorKind::InvalidInput, err)); + let marker_path = marker_path(¤t_ancestor, marker); marker_index += 1; if marker_index == markers.len() { diff --git a/codex-rs/file-system/src/lib.rs b/codex-rs/file-system/src/lib.rs index 5b809b8dab5a..e4aaf7d39d59 100644 --- a/codex-rs/file-system/src/lib.rs +++ b/codex-rs/file-system/src/lib.rs @@ -2,6 +2,7 @@ mod find_up; pub use find_up::FindUpErrorPolicy; pub use find_up::find_nearest_ancestor_with_markers; +pub use find_up::find_nearest_native_ancestor_with_markers; use bytes::Bytes; use codex_protocol::config_types::WindowsSandboxLevel; diff --git a/codex-rs/git-utils/src/info.rs b/codex-rs/git-utils/src/info.rs index 32d1b6f0a90b..d14cd1e0b1fd 100644 --- a/codex-rs/git-utils/src/info.rs +++ b/codex-rs/git-utils/src/info.rs @@ -6,7 +6,7 @@ use std::path::PathBuf; use codex_file_system::ExecutorFileSystem; use codex_file_system::FindUpErrorPolicy; -use codex_file_system::find_nearest_ancestor_with_markers; +use codex_file_system::find_nearest_native_ancestor_with_markers; use codex_utils_absolute_path::AbsolutePathBuf; use codex_utils_path_uri::PathUri; use futures::future::join_all; @@ -54,18 +54,15 @@ pub async fn get_git_repo_root_with_fs( Ok(metadata) if metadata.is_directory => cwd.clone(), _ => cwd.parent()?, }; - let base_uri = PathUri::from_abs_path(&base); - find_nearest_ancestor_with_markers( + find_nearest_native_ancestor_with_markers( fs, - &base_uri, + &base, vec![".git".to_string()], FindUpErrorPolicy::Ignore, /*sandbox*/ None, ) .await - .ok()?? - .to_abs_path() - .ok() + .ok()? } /// Timeout for git commands to prevent freezing on large repositories From 56176751a2f55dcc51a1bb171ee2e6d0814f8acf Mon Sep 17 00:00:00 2001 From: jif-oai Date: Wed, 24 Jun 2026 21:16:07 +0100 Subject: [PATCH 4/5] Add root discovery regression coverage --- codex-rs/core/src/agents_md_tests.rs | 96 +++++++------------- codex-rs/core/src/git_info_tests.rs | 130 +++++++++++++++++++++++++++ 2 files changed, 161 insertions(+), 65 deletions(-) diff --git a/codex-rs/core/src/agents_md_tests.rs b/codex-rs/core/src/agents_md_tests.rs index ca5e289a3416..69cd43496195 100644 --- a/codex-rs/core/src/agents_md_tests.rs +++ b/codex-rs/core/src/agents_md_tests.rs @@ -31,8 +31,6 @@ use std::ops::DerefMut; use std::path::PathBuf; use std::sync::Arc; use std::sync::Mutex; -use std::sync::atomic::AtomicUsize; -use std::sync::atomic::Ordering; use tempfile::TempDir; #[derive(Clone, Copy)] @@ -50,8 +48,6 @@ struct FailingFileSystem { #[derive(Default)] struct MetadataCallCounts { - active_calls: AtomicUsize, - max_active_calls: AtomicUsize, paths: Mutex>, } @@ -106,16 +102,7 @@ impl FailingFileSystem { .lock() .expect("metadata paths lock") .push(path.clone()); - let active_calls = self - .metadata_calls - .active_calls - .fetch_add(1, Ordering::Relaxed) - + 1; - self.metadata_calls - .max_active_calls - .fetch_max(active_calls, Ordering::Relaxed); - tokio::task::yield_now().await; - let result = match self.failure { + match self.failure { InjectedFailure::Metadata(kind) if path_abs == self.path => { Err(io::Error::new(kind, "injected metadata failure")) } @@ -125,11 +112,7 @@ impl FailingFileSystem { InjectedFailure::Metadata(_) | InjectedFailure::MetadataPending | InjectedFailure::Read(_) => LOCAL_FS.get_metadata(path, sandbox).await, - }; - self.metadata_calls - .active_calls - .fetch_sub(1, Ordering::Relaxed); - result + } } async fn read_directory( @@ -685,52 +668,6 @@ async fn read_agents_md_ignores_files_removed_after_discovery() { assert_eq!(loaded, None); } -#[tokio::test] -async fn read_agents_md_pipelines_root_markers_before_candidate_search() { - let tmp = tempfile::tempdir().expect("tempdir"); - fs::write(tmp.path().join(".git"), "").unwrap(); - fs::write(tmp.path().join("AGENTS.md"), "project doc").unwrap(); - let nested = tmp.path().join("nested"); - fs::create_dir(&nested).unwrap(); - - let mut config = make_config(&tmp, /*limit*/ 4096, /*instructions*/ None).await; - config.cwd = nested.abs(); - let metadata_calls = Arc::new(MetadataCallCounts::default()); - let fs = FailingFileSystem { - path: config.cwd.join("unused"), - failure: InjectedFailure::Read(io::ErrorKind::PermissionDenied), - metadata_calls: Arc::clone(&metadata_calls), - }; - let cwd = PathUri::from_abs_path(&config.cwd); - - let loaded = read_agents_md(&config.config, &fs, "local", &cwd) - .await - .expect("project instructions") - .expect("project instructions"); - - assert_eq!(loaded.text(), "project doc"); - let max_active_calls = metadata_calls.max_active_calls.load(Ordering::Relaxed); - assert!(max_active_calls > 1); - assert!(max_active_calls <= 8); - let candidate_paths = metadata_calls - .paths - .lock() - .expect("metadata paths lock") - .iter() - .filter(|path| path.basename().as_deref() != Some(".git")) - .cloned() - .collect::>(); - assert_eq!( - candidate_paths, - vec![ - PathUri::from_abs_path(&tmp.path().join(LOCAL_AGENTS_MD_FILENAME).abs()), - PathUri::from_abs_path(&tmp.path().join(DEFAULT_AGENTS_MD_FILENAME).abs()), - cwd.join(LOCAL_AGENTS_MD_FILENAME).expect("override path"), - cwd.join(DEFAULT_AGENTS_MD_FILENAME).expect("agents path"), - ] - ); -} - #[tokio::test] async fn marker_search_does_not_wait_for_a_higher_ancestor() { let tmp = tempfile::tempdir().expect("tempdir"); @@ -770,6 +707,35 @@ async fn marker_search_does_not_wait_for_a_higher_ancestor() { ); } +#[tokio::test] +async fn project_root_marker_search_continues_beyond_concurrency_window() { + const NESTING_DEPTH: usize = 9; + + let tmp = tempfile::tempdir().expect("tempdir"); + fs::write(tmp.path().join(".git"), "").unwrap(); + fs::write(tmp.path().join("AGENTS.md"), "project doc").unwrap(); + let mut nested = tmp.path().to_path_buf(); + for depth in 0..NESTING_DEPTH { + nested.push(format!("nested-{depth}")); + } + fs::create_dir_all(&nested).unwrap(); + + let mut config = make_config(&tmp, /*limit*/ 4096, /*instructions*/ None).await; + config.cwd = nested.abs(); + let cwd = PathUri::from_abs_path(&config.cwd); + + let paths = super::agents_md_paths(&config.config, &cwd, LOCAL_FS.as_ref()) + .await + .expect("AGENTS.md discovery"); + + assert_eq!( + paths, + vec![PathUri::from_abs_path( + &tmp.path().join(DEFAULT_AGENTS_MD_FILENAME).abs() + )] + ); +} + #[tokio::test] async fn empty_project_root_markers_only_probe_cwd_candidates() { let tmp = tempfile::tempdir().expect("tempdir"); diff --git a/codex-rs/core/src/git_info_tests.rs b/codex-rs/core/src/git_info_tests.rs index 3da69c0e3785..d056a6e62607 100644 --- a/codex-rs/core/src/git_info_tests.rs +++ b/codex-rs/core/src/git_info_tests.rs @@ -1,4 +1,14 @@ +use codex_exec_server::CopyOptions; +use codex_exec_server::CreateDirectoryOptions; +use codex_exec_server::ExecutorFileSystem; +use codex_exec_server::ExecutorFileSystemFuture; +use codex_exec_server::FileMetadata; +use codex_exec_server::FileSystemReadStream; +use codex_exec_server::FileSystemResult; +use codex_exec_server::FileSystemSandboxContext; use codex_exec_server::LOCAL_FS; +use codex_exec_server::ReadDirectoryEntry; +use codex_exec_server::RemoveOptions; use codex_git_utils::GitInfo; use codex_git_utils::GitSha; use codex_git_utils::collect_git_info; @@ -8,17 +18,120 @@ use codex_git_utils::git_diff_to_remote; use codex_git_utils::recent_commits; use codex_git_utils::resolve_root_git_project_for_trust; use codex_utils_path::normalize_for_path_comparison; +use codex_utils_path_uri::PathUri; use core_test_support::PathBufExt; use core_test_support::PathExt; use core_test_support::skip_if_sandbox; use pretty_assertions::assert_eq; use std::fs; +use std::io; #[cfg(unix)] use std::os::unix::fs::PermissionsExt; use std::path::PathBuf; use tempfile::TempDir; use tokio::process::Command; +struct FailingMetadataFileSystem { + path: PathUri, +} + +impl FailingMetadataFileSystem { + fn unsupported() -> FileSystemResult { + Err(io::Error::new( + io::ErrorKind::Unsupported, + "operation is not used by Git root discovery", + )) + } +} + +impl ExecutorFileSystem for FailingMetadataFileSystem { + fn canonicalize<'a>( + &'a self, + _path: &'a PathUri, + _sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, PathUri> { + Box::pin(async { Self::unsupported() }) + } + + fn read_file<'a>( + &'a self, + _path: &'a PathUri, + _sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, Vec> { + Box::pin(async { Self::unsupported() }) + } + + fn read_file_stream<'a>( + &'a self, + _path: &'a PathUri, + _sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, FileSystemReadStream> { + Box::pin(async { Self::unsupported() }) + } + + fn write_file<'a>( + &'a self, + _path: &'a PathUri, + _contents: Vec, + _sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, ()> { + Box::pin(async { Self::unsupported() }) + } + + fn create_directory<'a>( + &'a self, + _path: &'a PathUri, + _options: CreateDirectoryOptions, + _sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, ()> { + Box::pin(async { Self::unsupported() }) + } + + fn get_metadata<'a>( + &'a self, + path: &'a PathUri, + sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, FileMetadata> { + Box::pin(async move { + if path == &self.path { + Err(io::Error::new( + io::ErrorKind::PermissionDenied, + "injected metadata failure", + )) + } else { + LOCAL_FS.get_metadata(path, sandbox).await + } + }) + } + + fn read_directory<'a>( + &'a self, + _path: &'a PathUri, + _sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, Vec> { + Box::pin(async { Self::unsupported() }) + } + + fn remove<'a>( + &'a self, + _path: &'a PathUri, + _options: RemoveOptions, + _sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, ()> { + Box::pin(async { Self::unsupported() }) + } + + fn copy<'a>( + &'a self, + _source_path: &'a PathUri, + _destination_path: &'a PathUri, + _options: CopyOptions, + _sandbox: Option<&'a FileSystemSandboxContext>, + ) -> ExecutorFileSystemFuture<'a, ()> { + Box::pin(async { Self::unsupported() }) + } +} + // Helper function to create a test git repository async fn create_test_git_repo(temp_dir: &TempDir) -> PathBuf { let repo_path = temp_dir.path().join("repo"); @@ -518,6 +631,23 @@ async fn get_git_repo_root_with_fs_starts_at_parent_for_file() { ); } +#[tokio::test] +async fn get_git_repo_root_with_fs_ignores_metadata_errors() { + let tmp = TempDir::new().expect("tempdir"); + let proj = tmp.path().join("proj"); + let nested = proj.join("nested"); + std::fs::create_dir_all(proj.join(".git")).unwrap(); + std::fs::create_dir_all(&nested).unwrap(); + let fs = FailingMetadataFileSystem { + path: PathUri::from_abs_path(&nested.join(".git").abs()), + }; + + assert_eq!( + get_git_repo_root_with_fs(&fs, &nested.abs()).await, + Some(proj.abs()) + ); +} + #[cfg(windows)] #[tokio::test] async fn get_git_repo_root_with_fs_supports_windows_namespace_paths() { From 6f076303fc542a07f8e4bbb7ce079fd5704a9731 Mon Sep 17 00:00:00 2001 From: jif-oai Date: Wed, 24 Jun 2026 21:18:30 +0100 Subject: [PATCH 5/5] Make root probe concurrency test deterministic --- codex-rs/core/src/agents_md_tests.rs | 52 ++++++++++++++++++++++++++-- 1 file changed, 50 insertions(+), 2 deletions(-) diff --git a/codex-rs/core/src/agents_md_tests.rs b/codex-rs/core/src/agents_md_tests.rs index 69cd43496195..6551347aaf3e 100644 --- a/codex-rs/core/src/agents_md_tests.rs +++ b/codex-rs/core/src/agents_md_tests.rs @@ -32,10 +32,12 @@ use std::path::PathBuf; use std::sync::Arc; use std::sync::Mutex; use tempfile::TempDir; +use tokio::sync::Notify; #[derive(Clone, Copy)] enum InjectedFailure { Metadata(io::ErrorKind), + MetadataBlocked, MetadataPending, Read(io::ErrorKind), } @@ -49,6 +51,8 @@ struct FailingFileSystem { #[derive(Default)] struct MetadataCallCounts { paths: Mutex>, + started: Notify, + release: Notify, } impl FailingFileSystem { @@ -102,14 +106,20 @@ impl FailingFileSystem { .lock() .expect("metadata paths lock") .push(path.clone()); + self.metadata_calls.started.notify_one(); match self.failure { InjectedFailure::Metadata(kind) if path_abs == self.path => { Err(io::Error::new(kind, "injected metadata failure")) } + InjectedFailure::MetadataBlocked if path_abs == self.path => { + self.metadata_calls.release.notified().await; + LOCAL_FS.get_metadata(path, sandbox).await + } InjectedFailure::MetadataPending if path_abs == self.path => { std::future::pending().await } InjectedFailure::Metadata(_) + | InjectedFailure::MetadataBlocked | InjectedFailure::MetadataPending | InjectedFailure::Read(_) => LOCAL_FS.get_metadata(path, sandbox).await, } @@ -708,8 +718,9 @@ async fn marker_search_does_not_wait_for_a_higher_ancestor() { } #[tokio::test] -async fn project_root_marker_search_continues_beyond_concurrency_window() { +async fn project_root_marker_search_pipelines_bounded_window_and_continues() { const NESTING_DEPTH: usize = 9; + const CONCURRENCY_LIMIT: usize = 8; let tmp = tempfile::tempdir().expect("tempdir"); fs::write(tmp.path().join(".git"), "").unwrap(); @@ -723,9 +734,46 @@ async fn project_root_marker_search_continues_beyond_concurrency_window() { let mut config = make_config(&tmp, /*limit*/ 4096, /*instructions*/ None).await; config.cwd = nested.abs(); let cwd = PathUri::from_abs_path(&config.cwd); + let metadata_calls = Arc::new(MetadataCallCounts::default()); + let fs = FailingFileSystem { + path: config.cwd.join(".git"), + failure: InjectedFailure::MetadataBlocked, + metadata_calls: Arc::clone(&metadata_calls), + }; + + let search = + tokio::spawn(async move { super::agents_md_paths(&config.config, &cwd, &fs).await }); + tokio::time::timeout(std::time::Duration::from_secs(5), async { + loop { + let started = metadata_calls.started.notified(); + if metadata_calls + .paths + .lock() + .expect("metadata paths lock") + .len() + >= CONCURRENCY_LIMIT + { + break; + } + started.await; + } + }) + .await + .expect("initial marker window should start"); + assert_eq!( + metadata_calls + .paths + .lock() + .expect("metadata paths lock") + .len(), + CONCURRENCY_LIMIT + ); - let paths = super::agents_md_paths(&config.config, &cwd, LOCAL_FS.as_ref()) + metadata_calls.release.notify_one(); + let paths = tokio::time::timeout(std::time::Duration::from_secs(5), search) .await + .expect("marker search should complete") + .expect("marker search task") .expect("AGENTS.md discovery"); assert_eq!(