mj_controller/pollers/
lifecycle.rs1use super::*;
2
3pub enum LifecycleSuccess {
4 Created,
5 Resumed {
6 profile_id: String,
7 target_id: String,
8 },
9 Moved(mj_core::state::MoveOutcome),
10 Closed,
11 ForceStopped,
12 DestroyedStopped,
13 ForceDestroyed,
14}
15
16pub struct LifecycleUpdate {
17 pub session_id: String,
18 pub result: std::result::Result<LifecycleSuccess, String>,
19 pub deferred_cleanup: bool,
20}
21
22pub fn is_interrupted_close(session: &SessionRecord) -> bool {
26 matches!(
27 session.state,
28 SessionState::Closing | SessionState::Destroying
29 ) && session.target.is_some()
30}
31
32pub fn interrupted_suspend_session_ids(controller: &Controller) -> Vec<String> {
33 controller
34 .state
35 .sessions
36 .values()
37 .filter(|session| is_interrupted_close(session))
38 .map(|session| session.id.clone())
39 .collect()
40}
41
42pub fn interrupted_destroy_session_ids(controller: &Controller) -> Vec<String> {
47 controller
48 .state
49 .sessions
50 .values()
51 .filter(|session| {
52 session.state == SessionState::Error
53 && session.target.is_some()
54 && session.last_error.as_deref().is_some_and(|error| {
55 error.starts_with(mj_core::state::DESTRUCTION_FAILURE_PREFIX)
56 })
57 })
58 .map(|session| session.id.clone())
59 .collect()
60}
61
62pub fn interrupted_lifecycle_cause(session: &SessionRecord) -> Option<String> {
76 match session.state {
77 SessionState::Provisioning => Some(
82 "the daemon stopped while this session was provisioning; anything it created \
83 is offered by `mj recover scan`"
84 .to_owned(),
85 ),
86 SessionState::Closing => Some(
91 "the daemon stopped while this session was closing, and it has no target left \
92 to close"
93 .to_owned(),
94 ),
95 SessionState::Destroying => Some(
96 "the daemon stopped while this session was being torn down, and it has no \
97 target left to remove"
98 .to_owned(),
99 ),
100 SessionState::StartupCleanup
103 | SessionState::Checkpointing
104 | SessionState::Running
105 | SessionState::Disconnected
106 | SessionState::Parked
107 | SessionState::Stopped
108 | SessionState::Lost
109 | SessionState::Error
110 | SessionState::DestroyedWithDataLoss => None,
111 }
112}
113
114pub fn unowned_interrupted_lifecycles(
118 controller: &Controller,
119 owned: &std::collections::BTreeSet<String>,
120) -> Vec<(String, String)> {
121 controller
122 .state
123 .sessions
124 .values()
125 .filter(|session| !owned.contains(&session.id) && !is_interrupted_close(session))
126 .filter_map(|session| {
127 interrupted_lifecycle_cause(session).map(|cause| (session.id.clone(), cause))
128 })
129 .collect()
130}
131
132pub fn reserve_recovery_or_cancel(
133 observer: &crate::recovery_gate::RecoveryObserver,
134 session_id: &str,
135 cancelled: &AtomicBool,
136) -> Result<crate::recovery_gate::RecoveryReservation> {
137 let reservation = observer.reserve(session_id);
138 observer.cancel_busy(session_id);
142 while observer.is_busy(session_id) {
143 if cancelled.load(Ordering::Acquire) {
144 bail!("operation cancelled while waiting for recovery copy");
145 }
146 std::thread::sleep(Duration::from_millis(25));
147 }
148 Ok(reservation)
149}
150
151pub fn project_worker_title(
152 controller: &mut Controller,
153 update: &WorkerPollUpdate,
154) -> Option<Option<String>> {
155 let snapshot = update.view.snapshot.as_ref()?;
156 let session = controller.state.sessions.get_mut(&update.session_id)?;
157 let title = snapshot.resolved_title();
158 if session.acp_session_title == title {
159 return None;
160 }
161 session.acp_session_title = title.clone();
162 Some(title)
163}
164
165pub fn apply_worker_record_update(controller: &mut Controller, update: &WorkerPollUpdate) {
166 let Some(title) = project_worker_title(controller, update) else {
167 return;
168 };
169 let session_id = update.session_id.clone();
170 tokio::spawn(async move {
171 let result = tokio::task::spawn_blocking(move || {
172 crate::database::set_session_acp_title(&session_id, title.as_deref())
173 })
174 .await;
175 match result {
176 Ok(Ok(())) => {}
177 Ok(Err(error)) => tracing::warn!(%error, "could not persist relay title"),
178 Err(error) => tracing::warn!(%error, "relay title persistence task failed"),
179 }
180 });
181}