3131from engine .api .server import create_app , set_engine , ws_manager
3232from engine .baseline .baseline_model import BaselineModel
3333from engine .buffer .time_series_buffer import TimeSeriesBuffer
34- from engine .collector .system_collector import SystemCollector
34+ from engine .collector .system_collector import ProcessInfo , SystemCollector
3535from engine .config import COLLECTION_INTERVAL_SEC , DATASTORE_DIR
3636from engine .runtime_info import (
3737 allocate_listen_port ,
@@ -207,9 +207,16 @@ def set_user_preferences(self, body: dict) -> dict:
207207 return {"ok" : True , "preferences" : prefs .to_dict ()}
208208
209209 def _apply_safeguard (
210- self , prefs : UserPreferences , top_procs : list [ProcessImpact ]
210+ self ,
211+ prefs : UserPreferences ,
212+ snapshot_procs : tuple [ProcessInfo , ...] | list [ProcessInfo ],
213+ tracked : list [ProcessImpact ],
211214 ) -> None :
212- """Terminate configured processes when an alert has just fired."""
215+ """Terminate configured processes when an alert has just fired.
216+
217+ Matches against the current snapshot top processes (instant load) plus
218+ sustained tracker list, not only the small set used for alert RCA.
219+ """
213220 if not prefs .safeguard_enabled :
214221 return
215222 targets = {
@@ -219,20 +226,40 @@ def _apply_safeguard(
219226 }
220227 if not targets :
221228 return
229+
230+ # Snapshot first (ranked by current CPU+mem), then tracker extras — dedupe by PID.
231+ ordered : list [tuple [int , str ]] = []
232+ seen : set [int ] = set ()
233+ for p in snapshot_procs :
234+ if p .pid not in seen :
235+ seen .add (p .pid )
236+ ordered .append ((p .pid , p .name ))
237+ for proc in tracked :
238+ if proc .pid not in seen :
239+ seen .add (proc .pid )
240+ ordered .append ((proc .pid , proc .name ))
241+
222242 terminated = 0
223- for proc in top_procs [: 15 ] :
224- key = canonical_process_name (proc . name )
243+ for pid , name in ordered :
244+ key = canonical_process_name (name )
225245 if key not in targets :
226246 continue
227- res = self .process_action (proc .pid , "terminate" )
228- terminated += 1
229- logger .warning (
230- "Safeguard: terminate %s (pid %s) -> %s" ,
231- proc .name ,
232- proc .pid ,
233- res ,
234- )
235- if terminated >= 5 :
247+ res = self .process_action (pid , "terminate" )
248+ if res .get ("ok" ):
249+ terminated += 1
250+ logger .warning (
251+ "Safeguard: terminated %s (pid %s)" ,
252+ name ,
253+ pid ,
254+ )
255+ else :
256+ logger .warning (
257+ "Safeguard: could not terminate %s (pid %s): %s" ,
258+ name ,
259+ pid ,
260+ res .get ("error" , res ),
261+ )
262+ if terminated >= 8 :
236263 break
237264
238265 async def _broadcast_state (self ) -> None :
@@ -339,7 +366,10 @@ async def run(self) -> None:
339366 stress , top_procs , recent_events , prefs = prefs
340367 )
341368 if alert is not None :
342- self ._apply_safeguard (prefs , top_procs )
369+ tracked_for_safeguard = self .process_tracker .get_top_consumers (50 )
370+ self ._apply_safeguard (
371+ prefs , snapshot .processes , tracked_for_safeguard
372+ )
343373
344374 # 11. Broadcast via WebSocket
345375 await self ._broadcast_state ()
0 commit comments