@@ -19,6 +19,7 @@ const RECOVERY_OWNER_ID: &str = "runner-exec-recovery";
1919const RECOVERY_CHILD_KIND : & str = "runner_exec_recovery_child" ;
2020const RECOVERY_OWNER_LEASE : Duration = Duration :: from_secs ( 30 ) ;
2121const RECOVERY_LEASE_HEARTBEAT : Duration = Duration :: from_secs ( 1 ) ;
22+ const MAX_SOURCE_RECOVERY_DEFERRALS : u64 = 3 ;
2223#[ derive( Clone ) ]
2324struct RecoveryWorker {
2425 id : String ,
@@ -223,7 +224,9 @@ fn reconcile_terminal_runner_exec_runs_with_owner(
223224 if error. details . get ( "http_status" ) . and_then ( Value :: as_u64) != Some ( 404 ) {
224225 unavailable_endpoints. insert ( endpoint) ;
225226 }
226- deferred += 1 ;
227+ if defer_or_fail_recovery_source ( & store, run, & worker. token , job_id, & error) ? {
228+ deferred += 1 ;
229+ }
227230 continue ;
228231 }
229232 } ;
@@ -337,6 +340,54 @@ fn reconcile_terminal_runner_exec_runs_with_owner(
337340 Ok ( ( reconciled, deferred) )
338341}
339342
343+ /// A runner that remains unavailable cannot be retried by every unrelated
344+ /// command forever. Keep a small, inspectable retry budget, then terminalize
345+ /// the source with the exact blocker retained in its durable metadata.
346+ fn defer_or_fail_recovery_source (
347+ store : & ObservationStore ,
348+ run : & homeboy_core:: observation:: RunRecord ,
349+ child_token : & str ,
350+ runner_job_id : & str ,
351+ error : & Error ,
352+ ) -> Result < bool > {
353+ let attempts = run
354+ . metadata_json
355+ . pointer ( "/runner_exec_recovery/deferral_count" )
356+ . and_then ( Value :: as_u64)
357+ . unwrap_or ( 0 )
358+ + 1 ;
359+ let mut metadata = run. metadata_json . clone ( ) ;
360+ metadata. as_object_mut ( ) . map ( |metadata| {
361+ metadata. remove ( "runner_exec_source_lease" ) ;
362+ } ) ;
363+ metadata[ "runner_exec_recovery" ] = json ! ( {
364+ "schema" : "homeboy/runner-exec-recovery/v1" ,
365+ "phase" : if attempts >= MAX_SOURCE_RECOVERY_DEFERRALS { "blocked" } else { "deferred" } ,
366+ "deferral_count" : attempts,
367+ "max_deferrals" : MAX_SOURCE_RECOVERY_DEFERRALS ,
368+ "reason" : error. message,
369+ "details" : error. details,
370+ "inspection_action" : format!( "homeboy runs show {}" , run. id) ,
371+ } ) ;
372+ if attempts >= MAX_SOURCE_RECOVERY_DEFERRALS {
373+ metadata[ "runner_terminal_projection" ] = json ! ( {
374+ "state" : "recovery_blocked" ,
375+ "classification" : "runner_unavailable_after_bounded_recovery" ,
376+ "runner_id" : run. metadata_json[ "runner_id" ] ,
377+ "runner_job_id" : runner_job_id,
378+ } ) ;
379+ store. fail_running_runner_exec_recovery_source (
380+ & run. id ,
381+ child_token,
382+ runner_job_id,
383+ metadata,
384+ ) ?;
385+ return Ok ( false ) ;
386+ }
387+ store. defer_running_runner_exec_recovery_source ( & run. id , child_token, metadata) ?;
388+ Ok ( true )
389+ }
390+
340391fn endpoint_identity ( session : & crate :: RunnerSession ) -> String {
341392 session
342393 . remote_daemon_address
@@ -1012,6 +1063,102 @@ mod tests {
10121063 } ) ;
10131064 }
10141065
1066+ #[ test]
1067+ fn repeated_deferred_recovery_converges_without_stale_children ( ) {
1068+ with_isolated_home ( |_| {
1069+ let run_id = "unavailable-source" ;
1070+ let runner_job_id = "unavailable-job" ;
1071+ homeboy_agents:: agent_task_lifecycle:: record_runner_exec_job_identity (
1072+ run_id,
1073+ "unavailable-runner" ,
1074+ runner_job_id,
1075+ "/workspace" ,
1076+ & [ ] ,
1077+ )
1078+ . expect ( "source" ) ;
1079+ let store = ObservationStore :: open_initialized ( ) . expect ( "store" ) ;
1080+ let error = Error :: validation_invalid_argument (
1081+ "runner" ,
1082+ "runner has no persisted daemon session for recovery" ,
1083+ Some ( "unavailable-runner" . to_string ( ) ) ,
1084+ None ,
1085+ ) ;
1086+
1087+ let mut child_ids = BTreeSet :: new ( ) ;
1088+ for attempt in 1 ..=MAX_SOURCE_RECOVERY_DEFERRALS {
1089+ let owner = schedule_terminal_runner_exec_recovery ( )
1090+ . expect ( "schedule" )
1091+ . expect ( "owner" ) ;
1092+ let work = run_scheduled_terminal_runner_exec_recovery (
1093+ & owner. owner_id ,
1094+ & owner. owner_token ,
1095+ )
1096+ . expect ( "schedule child" )
1097+ . expect ( "owner work" ) ;
1098+ assert_eq ! ( work. children. len( ) , 1 ) ;
1099+ let child = & work. children [ 0 ] ;
1100+ child_ids. insert ( child. child_id . clone ( ) ) ;
1101+ let source = store. get_run ( run_id) . expect ( "read" ) . expect ( "source" ) ;
1102+ assert_eq ! (
1103+ defer_or_fail_recovery_source(
1104+ & store,
1105+ & source,
1106+ & child. child_token,
1107+ runner_job_id,
1108+ & error,
1109+ )
1110+ . expect( "record bounded deferral" ) ,
1111+ attempt < MAX_SOURCE_RECOVERY_DEFERRALS ,
1112+ ) ;
1113+ let child_record = store
1114+ . get_run ( & child. child_id )
1115+ . expect ( "read" )
1116+ . expect ( "child" ) ;
1117+ store
1118+ . finish_running_run_with_owner_token (
1119+ & child. child_id ,
1120+ & child. child_token ,
1121+ RunStatus :: Pass ,
1122+ child_record. metadata_json ,
1123+ )
1124+ . expect ( "finish child" ) ;
1125+ finish_scheduled_terminal_runner_exec_recovery (
1126+ & owner. owner_id ,
1127+ & owner. owner_token ,
1128+ 1 ,
1129+ 0 ,
1130+ usize:: from ( attempt < MAX_SOURCE_RECOVERY_DEFERRALS ) ,
1131+ )
1132+ . expect ( "finish owner" ) ;
1133+ }
1134+
1135+ let source = store. get_run ( run_id) . expect ( "read" ) . expect ( "source" ) ;
1136+ assert_eq ! ( source. status, RunStatus :: Fail . as_str( ) ) ;
1137+ assert_eq ! (
1138+ source. metadata_json[ "runner_exec_recovery" ] [ "deferral_count" ] ,
1139+ MAX_SOURCE_RECOVERY_DEFERRALS
1140+ ) ;
1141+ assert_eq ! (
1142+ source. metadata_json[ "runner_terminal_projection" ] [ "classification" ] ,
1143+ "runner_unavailable_after_bounded_recovery"
1144+ ) ;
1145+ let reader = ObservationStore :: open_scheduler_reader ( ) . expect ( "reader" ) ;
1146+ assert ! ( recovery_candidates( & reader) . expect( "candidates" ) . is_empty( ) ) ;
1147+ let children = store
1148+ . list_runs ( RunListFilter {
1149+ kind : Some ( RECOVERY_CHILD_KIND . to_string ( ) ) ,
1150+ ..RunListFilter :: default ( )
1151+ } )
1152+ . expect ( "list children" ) ;
1153+ assert_eq ! ( child_ids. len( ) , 1 , "retries reuse one child identity" ) ;
1154+ assert_eq ! ( children. len( ) , 1 , "retries do not accumulate children" ) ;
1155+ assert_ne ! ( children[ 0 ] . status, RunStatus :: Running . as_str( ) ) ;
1156+ assert ! ( schedule_terminal_runner_exec_recovery( )
1157+ . expect( "final schedule" )
1158+ . is_none( ) ) ;
1159+ } ) ;
1160+ }
1161+
10151162 #[ test]
10161163 fn owner_schedules_one_durable_child_per_source_within_its_budget ( ) {
10171164 with_isolated_home ( |_| {
0 commit comments