11use std:: {
22 borrow:: Cow ,
33 collections:: VecDeque ,
4- sync:: {
5- Arc ,
6- atomic:: { AtomicUsize , Ordering } ,
7- mpsc:: { self , Receiver , Sender } ,
8- } ,
4+ sync:: mpsc:: { self , Receiver , Sender } ,
95 thread,
106 time:: Duration ,
117} ;
@@ -33,7 +29,6 @@ use super::{
3329} ;
3430
3531const MAX_CACHED_FRAMES : usize = 48 ;
36- const MAX_IN_FLIGHT_PREFETCHES : usize = 3 ;
3732const PREFETCH_FORWARD_STEPS : usize = 3 ;
3833const PREFETCH_GRACE_PERIOD : Duration = Duration :: from_millis ( 12 ) ;
3934const PAN_ATLAS_MIN_POINTS : usize = 2_000 ;
@@ -50,9 +45,8 @@ pub(super) struct PlotFrameCacheKey {
5045pub ( super ) struct PlotFrameCache {
5146 entries : VecDeque < CachedPlotFrame > ,
5247 pub ( super ) last : Option < CachedPlotFrame > ,
53- prefetch_tx : Sender < PrefetchResult > ,
54- prefetch_rx : Receiver < PrefetchResult > ,
55- in_flight : Arc < AtomicUsize > ,
48+ prefetch_tx : Sender < PrefetchJob > ,
49+ prefetch_rx : Receiver < PrefetchEvent > ,
5650 queued : Vec < PlotFrameCacheKey > ,
5751 next_image_id : u32 ,
5852 next_transmit_batch : u64 ,
@@ -98,15 +92,41 @@ struct PrefetchResult {
9892 frame : Result < CachedPlotFrame , String > ,
9993}
10094
95+ #[ derive( Debug , Clone , Copy , PartialEq , Eq ) ]
96+ enum PrefetchMode {
97+ Batch ,
98+ Pan ,
99+ }
100+
101+ #[ derive( Debug ) ]
102+ struct PrefetchJob {
103+ scene : PlotScene ,
104+ requests : Vec < PrefetchRequest > ,
105+ mode : PrefetchMode ,
106+ }
107+
108+ impl PrefetchJob {
109+ fn keys ( & self ) -> Vec < PlotFrameCacheKey > {
110+ self . requests . iter ( ) . map ( |request| request. key ) . collect ( )
111+ }
112+ }
113+
114+ #[ derive( Debug ) ]
115+ enum PrefetchEvent {
116+ Cancelled ( Vec < PlotFrameCacheKey > ) ,
117+ Result ( PrefetchResult ) ,
118+ }
119+
101120impl Default for PlotFrameCache {
102121 fn default ( ) -> Self {
103- let ( prefetch_tx, prefetch_rx) = mpsc:: channel ( ) ;
122+ let ( prefetch_tx, prefetch_job_rx) = mpsc:: channel ( ) ;
123+ let ( prefetch_event_tx, prefetch_rx) = mpsc:: channel ( ) ;
124+ spawn_prefetch_worker ( prefetch_job_rx, prefetch_event_tx) ;
104125 Self {
105126 entries : VecDeque :: new ( ) ,
106127 last : None ,
107128 prefetch_tx,
108129 prefetch_rx,
109- in_flight : Arc :: new ( AtomicUsize :: new ( 0 ) ) ,
110130 queued : Vec :: new ( ) ,
111131 next_image_id : 1 ,
112132 next_transmit_batch : 0 ,
@@ -161,9 +181,6 @@ impl PlotFrameCache {
161181 if protocol == Protocol :: Blocks {
162182 return ;
163183 }
164- if self . in_flight . load ( Ordering :: Relaxed ) >= MAX_IN_FLIGHT_PREFETCHES {
165- return ;
166- }
167184
168185 let Some ( action) = recent_action else {
169186 return ;
@@ -208,13 +225,13 @@ impl PlotFrameCache {
208225 ) {
209226 let keys = self . keys_with_image_ids ( keys) ;
210227 if scene. total_points ( ) >= PAN_ATLAS_MIN_POINTS {
211- self . spawn_pan_prefetch ( scene. clone ( ) , keys) ;
228+ self . schedule_prefetch ( scene. clone ( ) , keys, PrefetchMode :: Pan ) ;
212229 } else {
213- self . spawn_batch_prefetch ( scene. clone ( ) , keys) ;
230+ self . schedule_prefetch ( scene. clone ( ) , keys, PrefetchMode :: Batch ) ;
214231 }
215232 } else {
216233 let keys = self . keys_with_image_ids ( keys) ;
217- self . spawn_batch_prefetch ( scene. clone ( ) , keys) ;
234+ self . schedule_prefetch ( scene. clone ( ) , keys, PrefetchMode :: Batch ) ;
218235 }
219236 }
220237
@@ -223,6 +240,11 @@ impl PlotFrameCache {
223240 || self . entries . iter ( ) . any ( |cached| cached. key == key)
224241 }
225242
243+ #[ cfg( test) ]
244+ pub ( super ) fn has_cached_or_queued_key ( & self , key : PlotFrameCacheKey ) -> bool {
245+ self . contains_key ( key) || self . queued . contains ( & key)
246+ }
247+
226248 pub ( super ) fn drain_transmit_payloads ( & mut self , max_count : usize ) -> Vec < String > {
227249 self . collect_prefetches ( ) ;
228250 let mut payloads = Vec :: new ( ) ;
@@ -291,11 +313,20 @@ impl PlotFrameCache {
291313 }
292314
293315 fn collect_prefetches ( & mut self ) {
294- while let Ok ( result) = self . prefetch_rx . try_recv ( ) {
295- self . queued . retain ( |key| * key != result. key ) ;
296- if let Ok ( mut frame) = result. frame {
297- discard_stale_transmit_payload ( & mut frame, self . min_transmit_priority ) ;
298- self . insert_prefetched ( frame) ;
316+ while let Ok ( event) = self . prefetch_rx . try_recv ( ) {
317+ match event {
318+ PrefetchEvent :: Cancelled ( keys) => {
319+ for key in keys {
320+ self . queued . retain ( |queued| * queued != key) ;
321+ }
322+ }
323+ PrefetchEvent :: Result ( result) => {
324+ self . queued . retain ( |key| * key != result. key ) ;
325+ if let Ok ( mut frame) = result. frame {
326+ discard_stale_transmit_payload ( & mut frame, self . min_transmit_priority ) ;
327+ self . insert_prefetched ( frame) ;
328+ }
329+ }
299330 }
300331 }
301332 }
@@ -320,42 +351,32 @@ impl PlotFrameCache {
320351 image_id
321352 }
322353
323- fn spawn_batch_prefetch ( & mut self , scene : PlotScene , keys : Vec < PrefetchRequest > ) {
324- self . discard_stale_transmit_payloads ( & keys) ;
325- self . queued . extend ( keys. iter ( ) . map ( |request| request. key ) ) ;
326- self . in_flight . fetch_add ( 1 , Ordering :: Relaxed ) ;
327- let tx = self . prefetch_tx . clone ( ) ;
328- let in_flight = Arc :: clone ( & self . in_flight ) ;
329- thread:: spawn ( move || {
330- thread:: sleep ( PREFETCH_GRACE_PERIOD ) ;
331- for request in keys {
332- let frame = render_plot_frame_for_key (
333- & scene,
334- request. key ,
335- Some ( request. image_id ) ,
336- request. transmit_priority ,
337- )
338- . map_err ( |error| error. to_string ( ) ) ;
339- let key = request. key ;
340- let _ = tx. send ( PrefetchResult { key, frame } ) ;
341- }
342- in_flight. fetch_sub ( 1 , Ordering :: Relaxed ) ;
343- } ) ;
344- }
345-
346- fn spawn_pan_prefetch ( & mut self , scene : PlotScene , keys : Vec < PrefetchRequest > ) {
347- self . discard_stale_transmit_payloads ( & keys) ;
348- self . queued . extend ( keys. iter ( ) . map ( |request| request. key ) ) ;
349- self . in_flight . fetch_add ( 1 , Ordering :: Relaxed ) ;
350- let tx = self . prefetch_tx . clone ( ) ;
351- let in_flight = Arc :: clone ( & self . in_flight ) ;
352- thread:: spawn ( move || {
353- thread:: sleep ( PREFETCH_GRACE_PERIOD ) ;
354- for ( key, frame) in render_pan_prefetch_frames ( & scene, & keys) {
355- let _ = tx. send ( PrefetchResult { key, frame } ) ;
354+ fn schedule_prefetch (
355+ & mut self ,
356+ scene : PlotScene ,
357+ requests : Vec < PrefetchRequest > ,
358+ mode : PrefetchMode ,
359+ ) {
360+ self . discard_stale_transmit_payloads ( & requests) ;
361+ self . queued
362+ . extend ( requests. iter ( ) . map ( |request| request. key ) ) ;
363+ let queued_keys = requests
364+ . iter ( )
365+ . map ( |request| request. key )
366+ . collect :: < Vec < _ > > ( ) ;
367+ if self
368+ . prefetch_tx
369+ . send ( PrefetchJob {
370+ scene,
371+ requests,
372+ mode,
373+ } )
374+ . is_err ( )
375+ {
376+ for key in queued_keys {
377+ self . queued . retain ( |queued| * queued != key) ;
356378 }
357- in_flight. fetch_sub ( 1 , Ordering :: Relaxed ) ;
358- } ) ;
379+ }
359380 }
360381
361382 fn render_plot_frame (
@@ -405,6 +426,55 @@ fn discard_stale_transmit_payload(frame: &mut CachedPlotFrame, min_priority: u64
405426 }
406427}
407428
429+ fn spawn_prefetch_worker ( job_rx : Receiver < PrefetchJob > , event_tx : Sender < PrefetchEvent > ) {
430+ thread:: spawn ( move || {
431+ while let Ok ( job) = job_rx. recv ( ) {
432+ let job = newest_prefetch_job_after_grace ( job, & job_rx, & event_tx) ;
433+ render_prefetch_job ( job, & event_tx) ;
434+ }
435+ } ) ;
436+ }
437+
438+ fn newest_prefetch_job_after_grace (
439+ mut job : PrefetchJob ,
440+ job_rx : & Receiver < PrefetchJob > ,
441+ event_tx : & Sender < PrefetchEvent > ,
442+ ) -> PrefetchJob {
443+ thread:: sleep ( PREFETCH_GRACE_PERIOD ) ;
444+ let mut cancelled = Vec :: new ( ) ;
445+ while let Ok ( next) = job_rx. try_recv ( ) {
446+ cancelled. extend ( job. keys ( ) ) ;
447+ job = next;
448+ }
449+ if !cancelled. is_empty ( ) {
450+ let _ = event_tx. send ( PrefetchEvent :: Cancelled ( cancelled) ) ;
451+ }
452+ job
453+ }
454+
455+ fn render_prefetch_job ( job : PrefetchJob , event_tx : & Sender < PrefetchEvent > ) {
456+ match job. mode {
457+ PrefetchMode :: Batch => {
458+ for request in job. requests {
459+ let frame = render_plot_frame_for_key (
460+ & job. scene ,
461+ request. key ,
462+ Some ( request. image_id ) ,
463+ request. transmit_priority ,
464+ )
465+ . map_err ( |error| error. to_string ( ) ) ;
466+ let key = request. key ;
467+ let _ = event_tx. send ( PrefetchEvent :: Result ( PrefetchResult { key, frame } ) ) ;
468+ }
469+ }
470+ PrefetchMode :: Pan => {
471+ for ( key, frame) in render_pan_prefetch_frames ( & job. scene , & job. requests ) {
472+ let _ = event_tx. send ( PrefetchEvent :: Result ( PrefetchResult { key, frame } ) ) ;
473+ }
474+ }
475+ }
476+ }
477+
408478pub ( super ) fn render_plot_frame (
409479 scene : & PlotScene ,
410480 kind : PlotKind ,
0 commit comments