Skip to main content

solverforge_solver/manager/solver_manager/runtime/
pause.rs

1use 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}