solverforge_solver/manager/solver_manager/runtime/
pause.rs1use std::sync::atomic::Ordering;
2
3use solverforge_core::domain::PlanningSolution;
4
5use super::{EventKind, SolverRuntime};
6use crate::manager::solver_manager::slot::SLOT_PAUSED;
7use crate::manager::{SolverEvent, SolverLifecycleState};
8use crate::scope::{ProgressCallback, SolverScope};
9use crate::stats::SolverTelemetry;
10
11impl<S: PlanningSolution> SolverRuntime<S> {
12 pub(crate) fn pause_if_requested<D, ProgressCb>(
13 &self,
14 solver_scope: &mut SolverScope<'_, S, D, ProgressCb>,
15 ) where
16 D: solverforge_scoring::Director<S>,
17 ProgressCb: ProgressCallback<S>,
18 {
19 if !self.slot.pause_requested.load(Ordering::Acquire) || self.is_cancel_requested() {
20 return;
21 }
22
23 solver_scope.pause_timers();
24 let current_score = solver_scope.current_score().copied();
25 let telemetry = solver_scope.stats().snapshot();
26 if solver_scope.best_solution_publication_enabled() {
27 let solution = solver_scope.score_director().clone_working_solution();
28 let best_score = solver_scope.best_score().copied();
29 let _ = self.pause_with_snapshot(solution, current_score, best_score, telemetry);
30 } else {
31 let _ = self.pause_without_snapshot(current_score, telemetry);
32 }
33 solver_scope.resume_timers();
34 }
35
36 pub fn pause_with_snapshot(
37 &self,
38 solution: S,
39 current_score: Option<S::Score>,
40 best_score: Option<S::Score>,
41 telemetry: SolverTelemetry,
42 ) -> bool {
43 if !self.slot.pause_requested.load(Ordering::Acquire) || self.is_cancel_requested() {
44 return false;
45 }
46
47 let (telemetry, candidate_trace) = telemetry.split_candidate_trace();
48 let published = self.slot.with_publication(|sender, record| {
49 if !self.slot.pause_requested.load(Ordering::Acquire)
50 || self.is_cancel_requested()
51 || self.current_state().is_terminal()
52 {
53 return false;
54 }
55 self.slot.state.store(SLOT_PAUSED, Ordering::SeqCst);
56 let terminal_reason = record.terminal_reason;
57 record.checkpoint_available = true;
58 record.current_score = current_score;
59 record.best_score = best_score;
60 record.publish_telemetry_with_candidate_trace(telemetry.clone(), candidate_trace);
61
62 let snapshot_revision = record.push_snapshot(crate::manager::SolverSnapshot {
63 job_id: self.job_id,
64 snapshot_revision: 0,
65 lifecycle_state: SolverLifecycleState::Paused,
66 terminal_reason,
67 current_score,
68 best_score,
69 telemetry: record.telemetry.clone(),
70 solution,
71 });
72 let metadata = record.next_metadata(
73 self.job_id,
74 SolverLifecycleState::Paused,
75 Some(snapshot_revision),
76 );
77 if let Some(sender) = sender {
78 let _ = sender.send(SolverEvent::Paused { metadata });
79 }
80 true
81 });
82 if !published {
83 return false;
84 }
85 self.wait_for_resume(current_score, best_score, telemetry)
86 }
87
88 fn pause_without_snapshot(
89 &self,
90 current_score: Option<S::Score>,
91 telemetry: SolverTelemetry,
92 ) -> bool {
93 let (telemetry, candidate_trace) = telemetry.split_candidate_trace();
94 let published = self.slot.with_publication(|sender, record| {
95 if !self.slot.pause_requested.load(Ordering::Acquire)
96 || self.is_cancel_requested()
97 || self.current_state().is_terminal()
98 {
99 return false;
100 }
101 self.slot.state.store(SLOT_PAUSED, Ordering::SeqCst);
102 record.checkpoint_available = false;
103 record.current_score = current_score;
104 record.best_score = None;
105 record.publish_telemetry_with_candidate_trace(telemetry.clone(), candidate_trace);
106 let metadata = record.next_metadata(self.job_id, SolverLifecycleState::Paused, None);
107 if let Some(sender) = sender {
108 let _ = sender.send(SolverEvent::Paused { metadata });
109 }
110 true
111 });
112 if !published {
113 return false;
114 }
115 self.wait_for_resume(current_score, None, telemetry)
116 }
117
118 fn wait_for_resume(
119 &self,
120 current_score: Option<S::Score>,
121 best_score: Option<S::Score>,
122 telemetry: SolverTelemetry,
123 ) -> bool {
124 let mut guard = self.slot.pause_gate.lock().unwrap();
125 while self.slot.pause_requested.load(Ordering::Acquire) && !self.is_cancel_requested() {
126 guard = self.slot.pause_condvar.wait(guard).unwrap();
127 }
128 drop(guard);
129 if self.is_cancel_requested() {
130 return false;
131 }
132 self.emit_non_snapshot_event(current_score, best_score, telemetry, EventKind::Resumed);
133 true
134 }
135}
136
137#[cfg(test)]
138mod tests {
139 use std::sync::atomic::Ordering;
140 use std::thread;
141 use std::time::{Duration, Instant};
142
143 use solverforge_core::score::SoftScore;
144
145 use super::*;
146 use crate::manager::solver_manager::slot::{JobSlot, SLOT_SOLVING};
147
148 #[derive(Clone, Debug)]
149 struct PauseTestSolution {
150 score: Option<SoftScore>,
151 }
152
153 impl PlanningSolution for PauseTestSolution {
154 type Score = SoftScore;
155
156 fn score(&self) -> Option<Self::Score> {
157 self.score
158 }
159
160 fn set_score(&mut self, score: Option<Self::Score>) {
161 self.score = score;
162 }
163 }
164
165 #[test]
166 fn incomplete_pause_exposes_no_checkpoint_or_solution_snapshot() {
167 let slot = Box::leak(Box::new(JobSlot::<PauseTestSolution>::new()));
168 slot.state.store(SLOT_SOLVING, Ordering::Release);
169 slot.pause_requested.store(true, Ordering::Release);
170 let runtime = SolverRuntime::new(7, slot);
171 let worker = thread::spawn(move || {
172 runtime.pause_without_snapshot(Some(SoftScore::of(0)), SolverTelemetry::default())
173 });
174
175 let deadline = Instant::now() + Duration::from_secs(2);
176 while slot.state.load(Ordering::Acquire) != SLOT_PAUSED {
177 assert!(Instant::now() < deadline, "pause did not settle");
178 thread::yield_now();
179 }
180 {
181 let record = slot.record.lock().unwrap();
182 assert!(!record.checkpoint_available);
183 assert!(record.latest_snapshot_revision.is_none());
184 assert!(record.snapshots.is_empty());
185 assert!(record.best_score.is_none());
186 }
187
188 slot.pause_requested.store(false, Ordering::Release);
189 slot.pause_condvar.notify_all();
190 assert!(worker.join().expect("pause worker must resume"));
191 }
192}