@@ -17,6 +17,7 @@ use crate::external_sources::{
1717} ;
1818use crate :: infrastructure:: get_path_manager_arc;
1919use crate :: util:: errors:: { OpenBitFunError , OpenBitFunResult } ;
20+ use futures:: { stream, StreamExt } ;
2021use log:: { debug, error, warn} ;
2122use openbitfun_agent_runtime:: skills:: {
2223 annotate_shadowed_skills, build_mode_skill_infos, filter_candidates_for_mode,
@@ -57,6 +58,9 @@ const MAX_OPENCODE_CONFIGURED_POLICY_BYTES: usize = 64 * 1024;
5758const OPENCODE_CONFIGURED_PRIORITY_BAND : usize =
5859 MAX_OPENCODE_CONFIGURED_SKILL_ROOTS * MAX_OPENCODE_CONFIGURED_SKILLS_PER_ROOT ;
5960
61+ // Bound remote IO across the whole scan, including workspaces with many roots.
62+ const REMOTE_SKILL_SCAN_CONCURRENCY : usize = 4 ;
63+
6064const DEEP_RESEARCH_AGENT_ID : & str = "DeepResearch" ;
6165const DEEP_RESEARCH_SKILL_NAME : & str = "deep-research" ;
6266
@@ -1244,84 +1248,102 @@ impl SkillRegistry {
12441248 fs : & dyn WorkspaceFileSystem ,
12451249 remote_root : & str ,
12461250 ) -> Vec < SkillCandidate > {
1247- let mut roots = Vec :: new ( ) ;
12481251 let root = remote_root. trim_end_matches ( '/' ) ;
1249- for ( priority, spec) in PROJECT_SKILL_ROOTS . iter ( ) . enumerate ( ) {
1250- let path = format ! ( "{}/{}/{}" , root, spec. parent, spec. subdir) ;
1251- if fs. is_dir ( & path) . await . unwrap_or ( false ) {
1252- roots. push ( RemoteSkillRootEntry {
1252+ // Finish directory discovery before scanning files, so the two bounded
1253+ // stages cannot multiply the number of simultaneous SFTP operations.
1254+ // `buffered` preserves source precedence and sorted directory order even
1255+ // when responses complete in a different order.
1256+ let root_scans = PROJECT_SKILL_ROOTS
1257+ . iter ( )
1258+ . enumerate ( )
1259+ . map ( |( priority, spec) | async move {
1260+ let path = format ! ( "{}/{}/{}" , root, spec. parent, spec. subdir) ;
1261+ let entry = RemoteSkillRootEntry {
12531262 path,
12541263 slot : spec. slot ,
12551264 source_id : spec. source_id ,
12561265 source_label : spec. source_label ,
12571266 priority,
1258- } ) ;
1259- }
1260- }
1261-
1262- let mut skills = Vec :: new ( ) ;
1263- for entry in roots {
1264- let mut entries = match fs. read_dir ( & entry. path ) . await {
1265- Ok ( value) => value,
1266- Err ( _) => continue ,
1267- } ;
1268- sort_remote_dir_entries ( & mut entries) ;
1269-
1270- for item in entries {
1271- if !item. is_dir || item. is_symlink {
1272- continue ;
1273- }
1274-
1275- let Some ( dir_name) = normalize_remote_skill_dir_name ( & item. path ) else {
1276- continue ;
12771267 } ;
1268+ let mut items = if fs. is_dir ( & entry. path ) . await . unwrap_or ( false ) {
1269+ fs. read_dir ( & entry. path ) . await . unwrap_or_default ( )
1270+ } else {
1271+ Vec :: new ( )
1272+ } ;
1273+ sort_remote_dir_entries ( & mut items) ;
1274+ ( entry, items)
1275+ } )
1276+ . collect :: < Vec < _ > > ( ) ;
1277+ let roots = stream:: iter ( root_scans)
1278+ . buffered ( REMOTE_SKILL_SCAN_CONCURRENCY )
1279+ . collect :: < Vec < _ > > ( )
1280+ . await ;
1281+
1282+ let directories = roots. iter ( ) . flat_map ( |( entry, items) | {
1283+ items
1284+ . iter ( )
1285+ . filter ( |item| item. is_dir && !item. is_symlink )
1286+ . map ( move |item| ( entry, item) )
1287+ } ) ;
1288+ let skill_scans = directories
1289+ . map ( |( entry, item) | async move {
1290+ let dir_name = normalize_remote_skill_dir_name ( & item. path ) ?;
12781291 let skill_md_path = format ! ( "{}/SKILL.md" , item. path. trim_end_matches( '/' ) ) ;
12791292 if !fs. is_file ( & skill_md_path) . await . unwrap_or ( false ) {
1280- continue ;
1293+ return None ;
12811294 }
1282-
1283- match fs. read_file_text ( & skill_md_path) . await {
1284- Ok ( content) => match Self :: parse_skill_markdown (
1285- item. path . clone ( ) ,
1286- & content,
1287- SkillLocation :: Project ,
1288- false ,
1289- entry. slot ,
1290- ) {
1291- Ok ( mut skill_data) => {
1292- Self :: apply_remote_openai_policy ( & mut skill_data, fs, & item. path ) . await ;
1293- skill_data. dir_name = dir_name;
1294- skills. push ( SkillCandidate :: from_data (
1295- skill_data,
1296- entry. slot ,
1297- entry. source_id ,
1298- entry. source_label ,
1299- PROJECT_SKILL_KEY_PREFIX ,
1300- entry. priority ,
1301- false ,
1302- ) ) ;
1303- }
1304- Err ( error) => {
1305- error ! ( "Failed to parse SKILL.md in {}: {}" , item. path, error) ;
1306- }
1307- } ,
1295+ let content = match fs. read_file_text ( & skill_md_path) . await {
1296+ Ok ( content) => content,
13081297 Err ( error) => {
13091298 debug ! ( "Failed to read {}: {}" , skill_md_path, error) ;
1299+ return None ;
13101300 }
1311- }
1312- }
1313- }
1314-
1315- skills
1301+ } ;
1302+ let mut skill_data = match Self :: parse_skill_markdown (
1303+ item. path . clone ( ) ,
1304+ & content,
1305+ SkillLocation :: Project ,
1306+ false ,
1307+ entry. slot ,
1308+ ) {
1309+ Ok ( data) => data,
1310+ Err ( error) => {
1311+ error ! ( "Failed to parse SKILL.md in {}: {}" , item. path, error) ;
1312+ return None ;
1313+ }
1314+ } ;
1315+ Self :: apply_remote_openai_policy ( & mut skill_data, fs, & item. path ) . await ;
1316+ skill_data. dir_name = dir_name;
1317+ Some ( SkillCandidate :: from_data (
1318+ skill_data,
1319+ entry. slot ,
1320+ entry. source_id ,
1321+ entry. source_label ,
1322+ PROJECT_SKILL_KEY_PREFIX ,
1323+ entry. priority ,
1324+ false ,
1325+ ) )
1326+ } )
1327+ . collect :: < Vec < _ > > ( ) ;
1328+ stream:: iter ( skill_scans)
1329+ . buffered ( REMOTE_SKILL_SCAN_CONCURRENCY )
1330+ . collect :: < Vec < _ > > ( )
1331+ . await
1332+ . into_iter ( )
1333+ . flatten ( )
1334+ . collect ( )
13161335 }
13171336
13181337 async fn scan_skill_candidates_for_remote_workspace (
13191338 & self ,
13201339 fs : & dyn WorkspaceFileSystem ,
13211340 remote_root : & str ,
13221341 ) -> Vec < SkillCandidate > {
1323- let mut skills = self . scan_skill_candidates_for_workspace ( None ) . await ;
1324- skills. extend ( Self :: scan_remote_project_skills ( fs, remote_root) . await ) ;
1342+ let ( mut skills, project_skills) = tokio:: join!(
1343+ self . scan_skill_candidates_for_workspace( None ) ,
1344+ Self :: scan_remote_project_skills( fs, remote_root) ,
1345+ ) ;
1346+ skills. extend ( project_skills) ;
13251347 skills
13261348 }
13271349
@@ -2212,3 +2234,158 @@ mod opencode_configured_skill_tests {
22122234 std:: os:: windows:: fs:: symlink_dir ( target, link) . is_ok ( )
22132235 }
22142236}
2237+
2238+ #[ cfg( test) ]
2239+ mod remote_scan_tests {
2240+ use super :: SkillRegistry ;
2241+ use crate :: agentic:: workspace:: { WorkspaceDirEntry , WorkspaceFileSystem } ;
2242+ use std:: sync:: atomic:: { AtomicUsize , Ordering } ;
2243+ use std:: time:: { Duration , Instant } ;
2244+
2245+ #[ derive( Default ) ]
2246+ struct DelayedFs {
2247+ active : AtomicUsize ,
2248+ peak : AtomicUsize ,
2249+ calls : AtomicUsize ,
2250+ }
2251+
2252+ impl DelayedFs {
2253+ async fn round_trip ( & self ) {
2254+ let active = self . active . fetch_add ( 1 , Ordering :: SeqCst ) + 1 ;
2255+ self . peak . fetch_max ( active, Ordering :: SeqCst ) ;
2256+ self . calls . fetch_add ( 1 , Ordering :: SeqCst ) ;
2257+ tokio:: time:: sleep ( Duration :: from_millis ( 10 ) ) . await ;
2258+ self . active . fetch_sub ( 1 , Ordering :: SeqCst ) ;
2259+ }
2260+ }
2261+
2262+ #[ async_trait:: async_trait]
2263+ impl WorkspaceFileSystem for DelayedFs {
2264+ async fn read_file ( & self , path : & str ) -> anyhow:: Result < Vec < u8 > > {
2265+ Ok ( self . read_file_text ( path) . await ?. into_bytes ( ) )
2266+ }
2267+ async fn read_file_text ( & self , path : & str ) -> anyhow:: Result < String > {
2268+ self . round_trip ( ) . await ;
2269+ if path. ends_with ( "openai.yaml" ) {
2270+ return Ok ( "policy:\n allow_implicit_invocation: false\n " . into ( ) ) ;
2271+ }
2272+ let name = path. rsplit ( '/' ) . nth ( 1 ) . unwrap ( ) ;
2273+ Ok ( format ! (
2274+ "---\n name: {name}\n description: {path}\n ---\n Body\n "
2275+ ) )
2276+ }
2277+ async fn write_file ( & self , _: & str , _: & [ u8 ] ) -> anyhow:: Result < ( ) > {
2278+ anyhow:: bail!( "read-only fixture" )
2279+ }
2280+ async fn exists ( & self , path : & str ) -> anyhow:: Result < bool > {
2281+ self . is_file ( path) . await
2282+ }
2283+ async fn is_file ( & self , path : & str ) -> anyhow:: Result < bool > {
2284+ self . round_trip ( ) . await ;
2285+ Ok ( path. ends_with ( "SKILL.md" ) || path. ends_with ( "skill-00/agents/openai.yaml" ) )
2286+ }
2287+ async fn is_dir ( & self , path : & str ) -> anyhow:: Result < bool > {
2288+ self . round_trip ( ) . await ;
2289+ Ok ( path. contains ( "/.openbitfun/" ) || path. contains ( "/.codex/" ) )
2290+ }
2291+ async fn read_dir ( & self , path : & str ) -> anyhow:: Result < Vec < WorkspaceDirEntry > > {
2292+ self . round_trip ( ) . await ;
2293+ Ok ( ( 0 ..13 )
2294+ . rev ( )
2295+ . map ( |index| WorkspaceDirEntry {
2296+ name : format ! ( "skill-{index:02}" ) ,
2297+ path : format ! ( "{path}/skill-{index:02}" ) ,
2298+ is_dir : true ,
2299+ is_symlink : index == 12 ,
2300+ modified : None ,
2301+ } )
2302+ . collect ( ) )
2303+ }
2304+ }
2305+
2306+ #[ tokio:: test]
2307+ async fn remote_scan_preserves_order_and_policy_with_bounded_io ( ) {
2308+ let fs = DelayedFs :: default ( ) ;
2309+ let start = Instant :: now ( ) ;
2310+ let skills = SkillRegistry :: scan_remote_project_skills ( & fs, "/remote/project/" ) . await ;
2311+ eprintln ! (
2312+ "remote scan: {:?}, {} logical IO calls, peak {}" ,
2313+ start. elapsed( ) ,
2314+ fs. calls. load( Ordering :: SeqCst ) ,
2315+ fs. peak. load( Ordering :: SeqCst )
2316+ ) ;
2317+ assert_eq ! ( skills. len( ) , 24 ) ;
2318+ for group in skills. chunks ( 12 ) {
2319+ assert_eq ! (
2320+ group
2321+ . iter( )
2322+ . map( |skill| skill. info. name. clone( ) )
2323+ . collect:: <Vec <_>>( ) ,
2324+ ( 0 ..12 )
2325+ . map( |index| format!( "skill-{index:02}" ) )
2326+ . collect:: <Vec <_>>( )
2327+ ) ;
2328+ assert ! ( !group[ 0 ] . info. allow_implicit_invocation) ;
2329+ assert ! ( group[ 1 ] . info. allow_implicit_invocation) ;
2330+ }
2331+ assert ! ( skills[ 0 ] . priority < skills[ 12 ] . priority) ;
2332+ assert_eq ! ( fs. calls. load( Ordering :: SeqCst ) , 82 ) ;
2333+ assert_eq ! ( fs. active. load( Ordering :: SeqCst ) , 0 ) ;
2334+ assert ! ( fs. peak. load( Ordering :: SeqCst ) > 1 ) ;
2335+ assert ! ( fs. peak. load( Ordering :: SeqCst ) <= super :: REMOTE_SKILL_SCAN_CONCURRENCY ) ;
2336+
2337+ // The same project catalog on disk must retain the remote scan's source
2338+ // precedence and invocation policy, without involving user-global skills.
2339+ let local_root = tempfile:: tempdir ( ) . unwrap ( ) ;
2340+ for parent in [ ".openbitfun" , ".codex" ] {
2341+ for index in 0 ..12 {
2342+ let name = format ! ( "skill-{index:02}" ) ;
2343+ let dir = local_root. path ( ) . join ( parent) . join ( "skills" ) . join ( & name) ;
2344+ std:: fs:: create_dir_all ( & dir) . unwrap ( ) ;
2345+ std:: fs:: write (
2346+ dir. join ( "SKILL.md" ) ,
2347+ format ! ( "---\n name: {name}\n description: fixture\n ---\n Body\n " ) ,
2348+ )
2349+ . unwrap ( ) ;
2350+ if index == 0 {
2351+ std:: fs:: create_dir_all ( dir. join ( "agents" ) ) . unwrap ( ) ;
2352+ std:: fs:: write (
2353+ dir. join ( "agents/openai.yaml" ) ,
2354+ "policy:\n allow_implicit_invocation: false\n " ,
2355+ )
2356+ . unwrap ( ) ;
2357+ }
2358+ }
2359+ }
2360+ let start = Instant :: now ( ) ;
2361+ let mut local = Vec :: new ( ) ;
2362+ for entry in SkillRegistry :: get_project_skill_roots ( local_root. path ( ) ) {
2363+ local. extend ( SkillRegistry :: scan_skills_in_dir ( & entry) . await ) ;
2364+ }
2365+ eprintln ! (
2366+ "local project scan: {:?}, {} skills" ,
2367+ start. elapsed( ) ,
2368+ local. len( )
2369+ ) ;
2370+ let catalog = |candidates : Vec < openbitfun_agent_runtime:: skills:: SkillCandidate > | {
2371+ candidates
2372+ . into_iter ( )
2373+ . map ( |candidate| {
2374+ (
2375+ candidate. info . key ,
2376+ candidate. info . name ,
2377+ candidate. info . allow_implicit_invocation ,
2378+ candidate. priority ,
2379+ )
2380+ } )
2381+ . collect :: < Vec < _ > > ( )
2382+ } ;
2383+ // Local directory enumeration is sorted by the shared resolver later.
2384+ local. sort_by ( |a, b| {
2385+ a. priority
2386+ . cmp ( & b. priority )
2387+ . then ( a. info . name . cmp ( & b. info . name ) )
2388+ } ) ;
2389+ assert_eq ! ( catalog( local) , catalog( skills) ) ;
2390+ }
2391+ }
0 commit comments