@@ -334,6 +334,7 @@ pub fn submit_request(
334334 let receipt_path = receipt_path ( & receipt_dir, & request. request_id ) ?;
335335 let request_path = request_path ( & request_dir, & request. request_id ) ?;
336336 let _lock = lock_request_scope ( & request_dir) ?;
337+ prune_request_temp_files ( & request_dir) ?;
337338 if !request_path. exists ( ) && durable_request_capacity_is_full ( & request_dir) ? {
338339 return Err ( ObserveAdmissionBackpressure {
339340 limit : MAX_PENDING_OBSERVE_REQUESTS ,
@@ -351,7 +352,10 @@ pub fn submit_request(
351352pub ( crate ) fn prepare_scope ( scope : & SupervisorScope ) -> anyhow:: Result < ( ) > {
352353 ensure_private_directory ( scope. root ( ) ) ?;
353354 ensure_private_directory ( & scope. observe_request_dir ( ) ) ?;
354- ensure_private_directory ( & scope. observe_receipt_dir ( ) )
355+ ensure_private_directory ( & scope. observe_receipt_dir ( ) ) ?;
356+ let request_dir = scope. observe_request_dir ( ) ;
357+ let _lock = lock_request_scope ( & request_dir) ?;
358+ prune_request_temp_files ( & request_dir)
355359}
356360
357361#[ derive( Debug ) ]
@@ -364,6 +368,25 @@ pub(crate) struct PendingRequestRecord {
364368pub ( crate ) fn scan_requests ( dir : & Path ) -> ( Vec < PendingRequestRecord > , Vec < String > ) {
365369 let mut records = Vec :: new ( ) ;
366370 let mut errors = Vec :: new ( ) ;
371+ {
372+ let _lock = match lock_request_scope ( dir) {
373+ Ok ( lock) => lock,
374+ Err ( error) => {
375+ errors. push ( format ! (
376+ "locking observe requests {}: {error:#}" ,
377+ dir. display( )
378+ ) ) ;
379+ return ( records, errors) ;
380+ }
381+ } ;
382+ if let Err ( error) = prune_request_temp_files ( dir) {
383+ errors. push ( format ! (
384+ "pruning observe request temps {}: {error:#}" ,
385+ dir. display( )
386+ ) ) ;
387+ return ( records, errors) ;
388+ }
389+ }
367390 let entries = match fs:: read_dir ( dir) {
368391 Ok ( entries) => entries,
369392 Err ( error) if error. kind ( ) == std:: io:: ErrorKind :: NotFound => return ( records, errors) ,
@@ -375,12 +398,17 @@ pub(crate) fn scan_requests(dir: &Path) -> (Vec<PendingRequestRecord>, Vec<Strin
375398 return ( records, errors) ;
376399 }
377400 } ;
378- for entry in entries. take ( MAX_PENDING_OBSERVE_REQUESTS ) . flatten ( ) {
401+ for ( entry, request_id) in entries
402+ . flatten ( )
403+ . filter_map ( |entry| {
404+ let name = entry. file_name ( ) ;
405+ let request_id = final_request_id ( name. to_str ( ) ?) ?. to_owned ( ) ;
406+ Some ( ( entry, request_id) )
407+ } )
408+ . take ( MAX_PENDING_OBSERVE_REQUESTS )
409+ {
379410 let name = entry. file_name ( ) ;
380- let Some ( name) = name. to_str ( ) else { continue } ;
381- let Some ( request_id) = final_request_id ( name) else {
382- continue ;
383- } ;
411+ let name = name. to_string_lossy ( ) ;
384412 let path = entry. path ( ) ;
385413 let parsed = ( || -> anyhow:: Result < PendingRequestRecord > {
386414 let metadata = entry. metadata ( ) ?;
@@ -552,7 +580,31 @@ fn normalize_diagnostic(diagnostic: Option<String>) -> Option<String> {
552580
553581fn durable_request_capacity_is_full ( dir : & Path ) -> anyhow:: Result < bool > {
554582 let entries = fs:: read_dir ( dir) . with_context ( || format ! ( "list {}" , dir. display( ) ) ) ?;
555- Ok ( entries. take ( MAX_PENDING_OBSERVE_REQUESTS ) . count ( ) >= MAX_PENDING_OBSERVE_REQUESTS )
583+ Ok ( entries
584+ . flatten ( )
585+ . filter ( |entry| {
586+ entry
587+ . file_name ( )
588+ . to_str ( )
589+ . and_then ( final_request_id)
590+ . is_some ( )
591+ } )
592+ . take ( MAX_PENDING_OBSERVE_REQUESTS )
593+ . count ( )
594+ >= MAX_PENDING_OBSERVE_REQUESTS )
595+ }
596+
597+ fn prune_request_temp_files ( dir : & Path ) -> anyhow:: Result < ( ) > {
598+ let entries = fs:: read_dir ( dir) . with_context ( || format ! ( "list {}" , dir. display( ) ) ) ?;
599+ for entry in entries. flatten ( ) {
600+ let name = entry. file_name ( ) ;
601+ let Some ( name) = name. to_str ( ) else { continue } ;
602+ if name. starts_with ( ".observe-" ) {
603+ fs:: remove_file ( entry. path ( ) )
604+ . with_context ( || format ! ( "remove stale observe request temp {name:?}" ) ) ?;
605+ }
606+ }
607+ Ok ( ( ) )
556608}
557609
558610fn lock_request_scope ( request_dir : & Path ) -> anyhow:: Result < fs:: File > {
@@ -696,12 +748,32 @@ mod tests {
696748 }
697749
698750 #[ test]
699- fn final_basename_filter_ignores_temps_and_rejects_traversal ( ) {
751+ fn final_basename_filter_prunes_crash_temps_without_starving_requests ( ) {
700752 let temp = tempfile:: tempdir ( ) . unwrap ( ) ;
701- fs:: write ( temp. path ( ) . join ( ".observe-temp" ) , b"not json" ) . unwrap ( ) ;
753+ let request = ObserveRequest :: new ( "h.worker" . into ( ) , "queue" . into ( ) , None , None ) . unwrap ( ) ;
754+ let path = request_path ( temp. path ( ) , & request. request_id ) . unwrap ( ) ;
755+ write_json_atomically_no_fsync ( & path, & request) . unwrap ( ) ;
756+ for index in 0 ..MAX_PENDING_OBSERVE_REQUESTS {
757+ fs:: write (
758+ temp. path ( ) . join ( format ! ( ".observe-crash-{index}" ) ) ,
759+ b"incomplete" ,
760+ )
761+ . unwrap ( ) ;
762+ }
702763 fs:: write ( temp. path ( ) . join ( "sibling" ) , b"not json" ) . unwrap ( ) ;
764+
703765 let ( records, errors) = scan_requests ( temp. path ( ) ) ;
704- assert ! ( records. is_empty( ) && errors. is_empty( ) ) ;
766+
767+ assert ! ( errors. is_empty( ) ) ;
768+ assert_eq ! ( records. len( ) , 1 ) ;
769+ assert_eq ! ( records[ 0 ] . request, request) ;
770+ assert ! (
771+ fs:: read_dir( temp. path( ) )
772+ . unwrap( )
773+ . flatten( )
774+ . all( |entry| !entry. file_name( ) . to_string_lossy( ) . starts_with( ".observe-" ) )
775+ ) ;
776+ assert ! ( !durable_request_capacity_is_full( temp. path( ) ) . unwrap( ) ) ;
705777 for invalid in [ "../escape" , "a/b" , ".hidden" , "x.json" , "with space" ] {
706778 assert ! (
707779 request_path( temp. path( ) , invalid) . is_err( ) ,
0 commit comments