6363//! those packages, and abandoning them would make it re-issue requests
6464//! the plan has already paid for.
6565//! * It never requests PAST a position whose own fetch was an
66- //! availability failure until the loop has consumed that position. So
66+ //! availability failure until the loop has consumed that position, and
67+ //! once it has, only the loop's own position until the service answers
68+ //! well again (the loop's breaker is then one failure from opening). So
6769//! an outage part-way down a list the window has already widened over
6870//! costs what the serial loop paid, as long as the task is running
6971//! ahead of the loop (the usual case: the loop stops to write between
@@ -188,7 +190,8 @@ struct Lookahead {
188190 at : AtomicUsize ,
189191 /// How far past `at` the task may run: one until the service has
190192 /// answered once, then [`SLOW_START`], growing by one per good answer
191- /// up to the whole window and falling back on availability failures.
193+ /// up to the whole window, falling back on availability failures and
194+ /// to one when the loop passes a failure.
192195 reach : AtomicUsize ,
193196 /// Lowest position whose own fetch was an availability failure and
194197 /// that the loop has not consumed yet; nothing past it is started
@@ -233,10 +236,16 @@ impl Lookahead {
233236 fn arrive ( & self , position : usize ) {
234237 // Past the failure the barrier stands at: the loop consumed that
235238 // package and went on, so its breaker did not end the run and the
236- // task may speculate again. (A failure landing concurrently just
237- // re-sets the line; all of this is advisory, and every outcome is
238- // still decided at the loop's own call.)
239+ // task may speculate again — but only from one position. The loop's
240+ // breaker is now a failure from opening, so a reach still sized by
241+ // the answers from BEFORE the failure would aim retry ladders at
242+ // packages the serial loop never asks for; only the position the
243+ // loop is at may start until the service answers well again. (A
244+ // failure landing concurrently just re-sets the line; all of this
245+ // is advisory, and every outcome is still decided at the loop's own
246+ // call.)
239247 if position > self . barrier . load ( Ordering :: Relaxed ) {
248+ self . reach . store ( 1 , Ordering :: Relaxed ) ;
240249 self . barrier . store ( usize:: MAX , Ordering :: Relaxed ) ;
241250 }
242251 self . at . store ( position, Ordering :: Relaxed ) ;
@@ -839,6 +848,50 @@ mod tests {
839848 ( out, c. vendor_outage_count ( ) )
840849 }
841850
851+ /// [`run`] with `plan` attached and the loop pausing before each call
852+ /// after the first — as the real loop stops to write between packages
853+ /// — until `ready` holds for the plan's lookahead. That pins where the
854+ /// task stands relative to the loop, instead of leaving it to how the
855+ /// runner schedules the mock server's delays.
856+ async fn run_paced (
857+ server : & MockServer ,
858+ plan : & [ usize ] ,
859+ calls : & [ usize ] ,
860+ ready : impl Fn ( & Lookahead ) -> bool ,
861+ ) -> ( Vec < String > , u32 ) {
862+ let c = client ( & server. uri ( ) ) ;
863+ let guard = c. prefetch_vendor_packages (
864+ plan. iter ( ) . map ( |& i| uuid ( i) ) . collect ( ) ,
865+ false ,
866+ None ,
867+ None ,
868+ 4 ,
869+ ) ;
870+ let look = Arc :: clone ( & guard. plan . look ) ;
871+ let mut out = Vec :: new ( ) ;
872+ for ( n, & i) in calls. iter ( ) . enumerate ( ) {
873+ if n > 0 {
874+ tokio:: time:: timeout ( Duration :: from_secs ( 30 ) , async {
875+ loop {
876+ let notified = look. moved . notified ( ) ;
877+ tokio:: pin!( notified) ;
878+ notified. as_mut ( ) . enable ( ) ;
879+ if ready ( & look) {
880+ return ;
881+ }
882+ notified. await ;
883+ }
884+ } )
885+ . await
886+ . expect ( "the prefetch never reached the state the loop waits for" ) ;
887+ }
888+ out. push ( summary (
889+ & c. fetch_vendor_package ( & uuid ( i) , false , None , None ) . await ,
890+ ) ) ;
891+ }
892+ ( out, c. vendor_outage_count ( ) )
893+ }
894+
842895 async fn assert_matches_serial ( scripts : & [ Script ] , plan : & [ usize ] , calls : & [ usize ] ) {
843896 let server = serve ( scripts) . await ;
844897 let serial = run ( & server, None , calls) . await ;
@@ -861,10 +914,12 @@ mod tests {
861914 /// whose download the prefetch had already started.
862915 ///
863916 /// And it costs the service the same REQUESTS, not only the same
864- /// outcomes. The task here is running ahead of a loop still on the
865- /// granted packages, so nothing past the failure is ever started, and
866- /// what was in flight when the task stopped is delivered instead of
867- /// being dropped and re-issued live by the loop.
917+ /// outcomes. The task here runs ahead of the loop: before each call the
918+ /// loop waits until the task has seen the failure it is about to meet.
919+ /// So nothing past the failure is started while the loop is still on
920+ /// the granted packages, and once the loop has consumed it and moved
921+ /// on, only the package it is at is requested — the serial loop's next
922+ /// ladder, which opens both breakers.
868923 #[ tokio:: test]
869924 async fn an_outage_mid_list_opens_the_breaker_at_the_same_package ( ) {
870925 use Script :: * ;
@@ -890,7 +945,11 @@ mod tests {
890945 assert_eq ! ( count, 2 ) ;
891946 let serial_requests = request_log ( & server) . await . len ( ) ;
892947
893- assert_eq ! ( run( & server, Some ( & all) , & all) . await , ( serial, count) ) ;
948+ let failure_landed = |look : & Lookahead | look. barrier . load ( Ordering :: Relaxed ) != usize:: MAX ;
949+ assert_eq ! (
950+ run_paced( & server, & all, & all, failure_landed) . await ,
951+ ( serial, count)
952+ ) ;
894953 assert_eq ! (
895954 request_log( & server) . await . len( ) - serial_requests,
896955 serial_requests,
0 commit comments