4040 "end_ts" : 105 , "steady_start_ts" : 100 }
4141
4242
43- def ranges (packing , node = (), pod = (), active = ()):
44- """A query_prometheus_range fake answering packing, PSI and active actors.
43+ def ranges (packing , node = (), pod = (), active = (), ws = (), actor_ws = () ):
44+ """A query_prometheus_range fake: packing, PSI, active actors, working set .
4545
4646 Pod PSI only answers the pod-slice selector, so a query that loses the id
4747 regex (and would also sum the pause container) gets nothing back. Node PSI
@@ -56,6 +56,11 @@ def fake(_url, query, *_args, **_kwargs):
5656 return packing
5757 if "ate_actor_stats_sampled_actors" in query :
5858 return list (active )
59+ if "ate_actor_stats_memory_working_set_bytes" in query :
60+ return list (actor_ws )
61+ # Before the PSI branch: the node working set query has container="node".
62+ if "container_memory_working_set_bytes" in query :
63+ return list (ws )
5964 if 'container="node"' in query :
6065 return list (node ) if worker_nodes in query and pod_slice in query else []
6166 return list (pod ) if pod_slice in query else []
@@ -244,7 +249,11 @@ def test_packing_and_checkpoint_math(self, mock_instant, mock_range):
244249 "values" : [[100 , "0.5" ], [105 , "0.5" ]]} for p in "ab" ]
245250 mock_range .side_effect = ranges (
246251 [partial , full , idle ], node = quiet , pod = pods ,
247- active = [{"metric" : {"instance" : "atelet-a" }, "values" : [[135 , "3" ]]}])
252+ active = [{"metric" : {"instance" : "atelet-a" }, "values" : [[135 , "3" ]]}],
253+ ws = [{"metric" : {"instance" : "node-a" },
254+ "values" : [[100 , str (2 ** 30 )], [105 , str (2 ** 30 )]]}],
255+ actor_ws = [{"metric" : {"instance" : "atelet-a" },
256+ "values" : [[135 , str (6 * 2 ** 30 )]]}])
248257 # The default 70s lag reads the snapshot window at 135 and 140.
249258 mock_instant .side_effect = snapshot_prom (
250259 "11.5" ,
@@ -278,14 +287,14 @@ def test_packing_and_checkpoint_math(self, mock_instant, mock_range):
278287
279288 sleep .assert_called_once_with (65 ) # until end 105 + 70, from 110
280289 self .assertEqual (summary ["metadata" ]["atelet_lag_s" ], 70 )
281- # Packing and PSI are not shifted; active actors are, like snapshots.
282- windows = {("sampled_actors " in c .args [1 ], c .args [2 :4 ])
290+ # Packing, PSI and cAdvisor are not shifted; atelet stats are, like snapshots.
291+ windows = {("ate_actor_stats " in c .args [1 ], c .args [2 :4 ])
283292 for c in mock_range .call_args_list }
284293 self .assertEqual (windows , {(False , (100 , 105 )), (True , (135 , 140 ))})
285294 self .assertEqual (
286295 set (summary ),
287296 {"metadata" , "cluster_packing" , "node_psi" , "pod_psi" , "snapshots" ,
288- "active_actors" },
297+ "active_actors" , "working_set" },
289298 )
290299 packing = summary ["cluster_packing" ]
291300 self .assertEqual (packing ["summary" ]["p50" ], 0.8 ) # 3 + 1 busy / 5 workers
@@ -305,7 +314,7 @@ def test_packing_and_checkpoint_math(self, mock_instant, mock_range):
305314 self .assertEqual (row ["metric" ], "server_summary" )
306315 m = row ["measurements" ]
307316 # Every source answered, so a key wired to a wrong name would read None.
308- self .assertEqual (len (m ), 49 )
317+ self .assertEqual (len (m ), 74 )
309318 self .assertEqual ([k for k , v in m .items () if v is None ], [])
310319 # Every value a string, so one row's types match every other row's.
311320 self .assertTrue (all (isinstance (v , str ) for v in m .values ()))
@@ -317,6 +326,11 @@ def test_packing_and_checkpoint_math(self, mock_instant, mock_range):
317326 self .assertEqual (m ["restore_mean_s" ], "0.5" )
318327 self .assertEqual (m ["checkpoint_mean_s" ], "2.0" )
319328 self .assertEqual (m ["active_actors_p99" ], "3.0" )
329+ # Read at 135, moved back to 100; 6 GiB over the 3 actors on atelet-a.
330+ actor = summary ["working_set" ]["actor" ]
331+ self .assertEqual (actor ["timeseries" ][0 ]["timestamp" ], 100 )
332+ self .assertEqual (m ["working_set_actor_p50_gb" ], "6.0" )
333+ self .assertEqual (m ["working_set_per_actor_p50_gb" ], "2.0" )
320334
321335 @mock .patch ("server_telemetry.query_prometheus_range" )
322336 @mock .patch ("server_telemetry.query_prometheus_instant" )
@@ -466,6 +480,43 @@ def test_steady_window_bounds_packing_and_psi(self, mock_instant, mock_range):
466480 self .assertEqual (out ["node_psi" ]["cpu_stall_pct" ]["max" ], 0.5 )
467481 self .assertEqual (out ["pod_psi" ]["cpu_stall_pct" ]["max" ], 0.5 )
468482
483+ @mock .patch ("server_telemetry.query_prometheus_range" )
484+ def test_working_set_gib_and_steady_window (self , mock_range ):
485+ gib = 2 ** 30
486+ mock_range .side_effect = ranges ([], ws = [
487+ {"metric" : {"instance" : "node-a" },
488+ "values" : [[90 , str (8 * gib )], [100 , str (gib )], [110 , str (3 * gib )]]},
489+ {"metric" : {"instance" : "node-b" }, "values" : [[110 , str (gib )]]},
490+ ])
491+ out = server_telemetry ._measure_working_set ("http://p" , 90 , 110 , 100 , 110 )
492+ node = out ["node" ]
493+ # The 8 GiB sample at 90 is ramp-up: in the timeseries, not the summary.
494+ self .assertEqual ((node ["summary" ]["p50" ], node ["summary" ]["max" ]), (1.0 , 3.0 ))
495+ self .assertEqual (node ["count" ], 2 )
496+ self .assertEqual (node ["timeseries" ], [
497+ {"timestamp" : 90 , "total_gb" : 8.0 , "max_gb" : 8.0 },
498+ {"timestamp" : 100 , "total_gb" : 1.0 , "max_gb" : 1.0 },
499+ {"timestamp" : 110 , "total_gb" : 4.0 , "max_gb" : 3.0 },
500+ ])
501+ self .assertEqual (set (out ), {"node" , "ateom" , "atelet" , "actor" , "per_actor" })
502+
503+ @mock .patch ("server_telemetry.query_prometheus_range" )
504+ def test_working_set_null_not_zero (self , mock_range ):
505+ mock_range .side_effect = ranges ([])
506+ out = server_telemetry ._measure_working_set ("http://p" , 90 , 110 , 100 , 110 )
507+ for block in out .values ():
508+ self .assertEqual (block , {"summary" : NO_PERCENTILES , "count" : None ,
509+ "timeseries" : []})
510+
511+ def test_working_set_queries_keep_live_workers (self ):
512+ q = server_telemetry .WORKING_SET_QUERIES
513+ live = "rate(container_cpu_usage_seconds_total"
514+ for block in q .values ():
515+ self .assertTrue (block .startswith ("max by (" )) # a restart has 2 series
516+ self .assertIn (live , block )
517+ self .assertIn ('container="node"' , q ["node" ])
518+ self .assertIn ("and on (pod)" , q ["ateom" ])
519+
469520 @mock .patch ("server_telemetry.query_prometheus_instant" )
470521 def test_failed_delta_read_is_unknown (self , mock_instant ):
471522 # A failed delta query is not taken as a series born in the window.
0 commit comments