Parallelize skills list cwd loading

This commit is contained in:
xli-oai
2026-05-06 15:35:44 -07:00
parent bdf075769d
commit b1ec596a53
2 changed files with 135 additions and 71 deletions

View File

@@ -1,4 +1,5 @@
use super::*;
use futures::StreamExt;
#[derive(Clone)]
pub(crate) struct CatalogRequestProcessor {
@@ -9,6 +10,8 @@ pub(crate) struct CatalogRequestProcessor {
pub(super) workspace_settings_cache: Arc<workspace_settings::WorkspaceSettingsCache>,
}
const SKILLS_LIST_CWD_CONCURRENCY: usize = 8;
fn skills_to_info(
skills: &[codex_core::skills::SkillMetadata],
disabled_paths: &HashSet<AbsolutePathBuf>,
@@ -441,82 +444,107 @@ impl CatalogRequestProcessor {
.environment_manager()
.default_environment()
.map(|environment| environment.get_filesystem());
let mut data = Vec::new();
for cwd in cwds {
let cwd_started_at = Instant::now();
let resolve_cwd_config_started_at = Instant::now();
let (cwd_abs, config_layer_stack) = match self.resolve_cwd_config(&cwd).await {
Ok(resolved) => resolved,
Err(message) => {
let mut data = futures::stream::iter(cwds.into_iter().enumerate())
.map(|(index, cwd)| {
let config = &config;
let extra_roots_by_cwd = &extra_roots_by_cwd;
let fs = fs.clone();
let plugins_manager = &plugins_manager;
let skills_manager = &skills_manager;
async move {
let cwd_started_at = Instant::now();
let resolve_cwd_config_started_at = Instant::now();
let (cwd_abs, config_layer_stack) =
match self.resolve_cwd_config(&cwd).await {
Ok(resolved) => resolved,
Err(message) => {
warn!(
cwd = %cwd.display(),
total_ms = cwd_started_at.elapsed().as_millis(),
resolve_cwd_config_ms = resolve_cwd_config_started_at.elapsed().as_millis(),
"skills/list cwd timing failed to resolve cwd config"
);
let error_path = cwd.clone();
return (
index,
codex_app_server_protocol::SkillsListEntry {
cwd,
skills: Vec::new(),
errors: vec![
codex_app_server_protocol::SkillErrorInfo {
path: error_path,
message,
},
],
},
);
}
};
let resolve_cwd_config_ms =
resolve_cwd_config_started_at.elapsed().as_millis();
let extra_roots = extra_roots_by_cwd
.get(&cwd)
.map_or(&[][..], std::vec::Vec::as_slice);
let effective_skill_roots_started_at = Instant::now();
let effective_skill_roots = if workspace_codex_plugins_enabled {
let plugins_input = config.plugins_config_input();
plugins_manager
.effective_skill_roots_for_layer_stack(
&config_layer_stack,
&plugins_input,
)
.await
} else {
Vec::new()
};
let effective_skill_roots_ms =
effective_skill_roots_started_at.elapsed().as_millis();
let effective_skill_root_count = effective_skill_roots.len();
let skills_input = codex_core::skills::SkillsLoadInput::new(
cwd_abs.clone(),
effective_skill_roots,
config_layer_stack,
config.bundled_skills_enabled(),
);
let load_skills_started_at = Instant::now();
let outcome = skills_manager
.skills_for_cwd_with_extra_user_roots(
&skills_input,
force_reload,
extra_roots,
fs,
)
.await;
let load_skills_ms = load_skills_started_at.elapsed().as_millis();
let errors = errors_to_info(&outcome.errors);
let skills = skills_to_info(&outcome.skills, &outcome.disabled_paths);
warn!(
cwd = %cwd.display(),
total_ms = cwd_started_at.elapsed().as_millis(),
resolve_cwd_config_ms = resolve_cwd_config_started_at.elapsed().as_millis(),
"skills/list cwd timing failed to resolve cwd config"
resolve_cwd_config_ms,
effective_skill_roots_ms,
load_skills_ms,
extra_root_count = extra_roots.len(),
effective_skill_root_count,
skill_count = skills.len(),
error_count = errors.len(),
"skills/list cwd timing"
);
let error_path = cwd.clone();
data.push(codex_app_server_protocol::SkillsListEntry {
cwd,
skills: Vec::new(),
errors: vec![codex_app_server_protocol::SkillErrorInfo {
path: error_path,
message,
}],
});
continue;
(
index,
codex_app_server_protocol::SkillsListEntry {
cwd,
skills,
errors,
},
)
}
};
let resolve_cwd_config_ms = resolve_cwd_config_started_at.elapsed().as_millis();
let extra_roots = extra_roots_by_cwd
.get(&cwd)
.map_or(&[][..], std::vec::Vec::as_slice);
let effective_skill_roots_started_at = Instant::now();
let effective_skill_roots = if workspace_codex_plugins_enabled {
let plugins_input = config.plugins_config_input();
plugins_manager
.effective_skill_roots_for_layer_stack(&config_layer_stack, &plugins_input)
.await
} else {
Vec::new()
};
let effective_skill_roots_ms = effective_skill_roots_started_at.elapsed().as_millis();
let effective_skill_root_count = effective_skill_roots.len();
let skills_input = codex_core::skills::SkillsLoadInput::new(
cwd_abs.clone(),
effective_skill_roots,
config_layer_stack,
config.bundled_skills_enabled(),
);
let load_skills_started_at = Instant::now();
let outcome = skills_manager
.skills_for_cwd_with_extra_user_roots(
&skills_input,
force_reload,
extra_roots,
fs.clone(),
)
.await;
let load_skills_ms = load_skills_started_at.elapsed().as_millis();
let errors = errors_to_info(&outcome.errors);
let skills = skills_to_info(&outcome.skills, &outcome.disabled_paths);
warn!(
cwd = %cwd.display(),
total_ms = cwd_started_at.elapsed().as_millis(),
resolve_cwd_config_ms,
effective_skill_roots_ms,
load_skills_ms,
extra_root_count = extra_roots.len(),
effective_skill_root_count,
skill_count = skills.len(),
error_count = errors.len(),
"skills/list cwd timing"
);
data.push(codex_app_server_protocol::SkillsListEntry {
cwd,
skills,
errors,
});
}
})
.buffer_unordered(SKILLS_LIST_CWD_CONCURRENCY)
.collect::<Vec<_>>()
.await;
data.sort_unstable_by_key(|(index, _)| *index);
let data = data.into_iter().map(|(_, entry)| entry).collect::<Vec<_>>();
let skill_count = data.iter().map(|entry| entry.skills.len()).sum::<usize>();
let error_count = data.iter().map(|entry| entry.errors.len()).sum::<usize>();
warn!(

View File

@@ -530,6 +530,42 @@ async fn skills_list_accepts_relative_cwds() -> Result<()> {
Ok(())
}
#[tokio::test]
async fn skills_list_preserves_requested_cwd_order() -> Result<()> {
let codex_home = TempDir::new()?;
let first_cwd = TempDir::new()?;
let second_cwd = TempDir::new()?;
let mut mcp = McpProcess::new(codex_home.path()).await?;
timeout(DEFAULT_TIMEOUT, mcp.initialize()).await??;
let request_id = mcp
.send_skills_list_request(SkillsListParams {
cwds: vec![
first_cwd.path().to_path_buf(),
second_cwd.path().to_path_buf(),
],
force_reload: true,
per_cwd_extra_user_roots: None,
})
.await?;
let response: JSONRPCResponse = timeout(
DEFAULT_TIMEOUT,
mcp.read_stream_until_response_message(RequestId::Integer(request_id)),
)
.await??;
let SkillsListResponse { data } = to_response(response)?;
assert_eq!(
data.iter().map(|entry| &entry.cwd).collect::<Vec<_>>(),
vec![
&first_cwd.path().to_path_buf(),
&second_cwd.path().to_path_buf(),
],
);
Ok(())
}
#[tokio::test]
async fn skills_list_ignores_per_cwd_extra_roots_for_unknown_cwd() -> Result<()> {
let codex_home = TempDir::new()?;