@@ -10,6 +10,14 @@ use crate::utils::fs::is_dir;
1010#[ cfg( test) ]
1111mod oracle;
1212
13+ /// How many `.pom` paths the parallel parse takes at a time. Every phase
14+ /// stays in walk order whatever the chunk, so this only bounds peak
15+ /// memory: a real `~/.m2` holds 10-50k artifacts (corporate caches many
16+ /// times that, and a GAV dir commonly holds more than one `.pom`), and
17+ /// buffering every path and every parse result before the dedup would put
18+ /// tens of megabytes on a scan that runs eight other crawlers beside it.
19+ const POM_PARSE_CHUNK : usize = 1024 ;
20+
1321// ---------------------------------------------------------------------------
1422// POM XML minimal parser
1523// ---------------------------------------------------------------------------
@@ -575,59 +583,82 @@ impl MavenCrawler {
575583 /// Uses `walkdir` to recursively find `.pom` files, then extracts
576584 /// coordinates from the POM content or falls back to directory path parsing.
577585 ///
578- /// Three phases: the (serial) walk collects the `.pom` paths in walk
579- /// order, the reads and parses run through [`par_map`] — in parallel on
580- /// the walk pool's threads, and on the calling thread when no walk
581- /// thread could be spawned (each POM's coordinates depend on that file
582- /// alone) — and the PURL dedup runs serially in walk order, so the
583- /// first-seen version dir wins and packages come out exactly as the
584- /// one-at-a-time scan's did.
586+ /// Three phases per chunk of the walk : the (serial) walk collects
587+ /// `.pom` paths in walk order, the reads and parses run through
588+ /// [`par_map`] — in parallel on the walk pool's threads, and on the
589+ /// calling thread when no walk thread could be spawned (each POM's
590+ /// coordinates depend on that file alone) — and the PURL dedup runs
591+ /// serially in walk order, so the first-seen version dir wins and
592+ /// packages come out exactly as the one-at-a-time scan's did.
585593 fn scan_maven_repo ( & self , repo_path : & Path , seen : & mut HashSet < String > ) -> Vec < CrawledPackage > {
594+ self . scan_maven_repo_chunked ( repo_path, seen, POM_PARSE_CHUNK )
595+ }
596+
597+ /// [`Self::scan_maven_repo`] over an explicit chunk size, so tests can
598+ /// cross the chunk boundary on a small fixture.
599+ fn scan_maven_repo_chunked (
600+ & self ,
601+ repo_path : & Path ,
602+ seen : & mut HashSet < String > ,
603+ chunk : usize ,
604+ ) -> Vec < CrawledPackage > {
605+ let chunk = chunk. max ( 1 ) ;
606+ let mut results = Vec :: new ( ) ;
586607 let mut poms: Vec < PathBuf > = Vec :: new ( ) ;
587- for entry in walkdir:: WalkDir :: new ( repo_path)
608+ let mut walk = walkdir:: WalkDir :: new ( repo_path)
588609 . follow_links ( false )
589610 . into_iter ( )
590- . filter_map ( |e| e. ok ( ) )
591- {
592- if !entry. file_type ( ) . is_file ( ) {
593- continue ;
594- }
595- let path = entry. path ( ) ;
596- if path. extension ( ) . is_none_or ( |ext| ext != "pom" ) {
597- continue ;
611+ . filter_map ( |e| e. ok ( ) ) ;
612+
613+ loop {
614+ for entry in walk. by_ref ( ) {
615+ if !entry. file_type ( ) . is_file ( ) {
616+ continue ;
617+ }
618+ let path = entry. path ( ) ;
619+ if path. extension ( ) . is_none_or ( |ext| ext != "pom" ) {
620+ continue ;
621+ }
622+ if path. parent ( ) . is_none ( ) {
623+ continue ;
624+ }
625+ poms. push ( entry. into_path ( ) ) ;
626+ if poms. len ( ) >= chunk {
627+ break ;
628+ }
598629 }
599- if path . parent ( ) . is_none ( ) {
600- continue ;
630+ if poms . is_empty ( ) {
631+ break ;
601632 }
602- poms. push ( entry. into_path ( ) ) ;
603- }
604633
605- let parsed: Vec < Option < ( String , String , String ) > > = par_map ( & poms, |path| {
606- let version_dir = path. parent ( ) ?;
607- // Try POM parsing first, fall back to directory path parsing
608- std:: fs:: read_to_string ( path)
609- . ok ( )
610- . and_then ( |content| parse_pom_group_artifact_version ( & content) )
611- . or_else ( || parse_path_coordinates ( version_dir, repo_path) )
612- } ) ;
613-
614- let mut results = Vec :: new ( ) ;
615- for ( path, coords) in poms. iter ( ) . zip ( parsed) {
616- let Some ( version_dir) = path. parent ( ) else {
617- continue ;
618- } ;
619- if let Some ( ( group_id, artifact_id, version) ) = coords {
620- let purl = crate :: utils:: purl:: build_maven_purl ( & group_id, & artifact_id, & version) ;
621- if seen. insert ( purl. clone ( ) ) {
622- results. push ( CrawledPackage {
623- name : artifact_id,
624- version,
625- namespace : Some ( group_id) ,
626- purl,
627- path : version_dir. to_path_buf ( ) ,
628- } ) ;
634+ let parsed: Vec < Option < ( String , String , String ) > > = par_map ( & poms, |path| {
635+ let version_dir = path. parent ( ) ?;
636+ // Try POM parsing first, fall back to directory path parsing
637+ std:: fs:: read_to_string ( path)
638+ . ok ( )
639+ . and_then ( |content| parse_pom_group_artifact_version ( & content) )
640+ . or_else ( || parse_path_coordinates ( version_dir, repo_path) )
641+ } ) ;
642+
643+ for ( path, coords) in poms. iter ( ) . zip ( parsed) {
644+ let Some ( version_dir) = path. parent ( ) else {
645+ continue ;
646+ } ;
647+ if let Some ( ( group_id, artifact_id, version) ) = coords {
648+ let purl =
649+ crate :: utils:: purl:: build_maven_purl ( & group_id, & artifact_id, & version) ;
650+ if seen. insert ( purl. clone ( ) ) {
651+ results. push ( CrawledPackage {
652+ name : artifact_id,
653+ version,
654+ namespace : Some ( group_id) ,
655+ purl,
656+ path : version_dir. to_path_buf ( ) ,
657+ } ) ;
658+ }
629659 }
630660 }
661+ poms. clear ( ) ;
631662 }
632663
633664 results
@@ -1755,6 +1786,37 @@ mod tests {
17551786 }
17561787 }
17571788
1789+ /// Chunking the parallel parse changes nothing: the walk, the
1790+ /// parse and the dedup each stay in walk order, so every chunk
1791+ /// size — one POM at a time included — produces the serial
1792+ /// oracle's rows, with the first-seen version dir still winning
1793+ /// across a chunk boundary.
1794+ #[ tokio:: test]
1795+ async fn every_parse_chunk_size_matches_the_serial_oracle ( ) {
1796+ let mut total = 0 ;
1797+ for seed in 0 ..16u64 {
1798+ let tmp = tempfile:: tempdir ( ) . unwrap ( ) ;
1799+ let mut perms = PermGuard :: default ( ) ;
1800+ let mut rng = Rng :: new ( seed) ;
1801+ let root = tmp. path ( ) . join ( "repository" ) ;
1802+ repo ( & mut rng, & root, & tmp. path ( ) . join ( "outside" ) , & mut perms) ;
1803+ perms. apply ( ) ;
1804+ let options = CrawlerOptions {
1805+ cwd : tmp. path ( ) . to_path_buf ( ) ,
1806+ global : false ,
1807+ global_prefix : Some ( root. clone ( ) ) ,
1808+ } ;
1809+ let old = LegacyMavenCrawler :: crawl_all ( & options) . await ;
1810+ for chunk in [ 0usize , 1 , 2 , 3 , 7 , 4096 ] {
1811+ let mut seen = HashSet :: new ( ) ;
1812+ let found = MavenCrawler . scan_maven_repo_chunked ( & root, & mut seen, chunk) ;
1813+ assert_eq ! ( rows( & found) , rows( & old) , "seed {seed}, chunk {chunk}" ) ;
1814+ }
1815+ total += old. len ( ) ;
1816+ }
1817+ assert ! ( total > 50 , "vacuous fixtures: {total}" ) ;
1818+ }
1819+
17581820 /// The no-walk-pool fallback (the OS refused even one walk thread,
17591821 /// so `run_walk` runs the walk on the calling blocking-pool thread)
17601822 /// scans exactly as the serial oracle does. A bare rayon iterator
0 commit comments