1use super::*;
2
3pub(super) fn owner_pid_to_watch() -> Result<Option<u32>> {
9 let Some(value) = mj_core::config::env_override("DAEMON_OWNER_PID") else {
10 return Ok(None);
11 };
12 let pid: u32 = value
13 .trim()
14 .parse()
15 .map_err(|_| anyhow!("MJ_DAEMON_OWNER_PID must be a process id, but it is {value:?}"))?;
16 ensure!(
17 process_is_alive(pid),
18 "MJ_DAEMON_OWNER_PID names process {pid}, which is not running"
19 );
20 Ok(Some(pid))
21}
22
23pub async fn run_daemon_process() -> Result<()> {
24 let owner_pid = owner_pid_to_watch()?;
27 let guard = ControllerStoreGuard::acquire()?;
28 let database_writer = guard.start_database_writer()?;
29 let epilogue_started = AtomicBool::new(false);
30 let mut outcome = run_daemon_runtime(&epilogue_started, owner_pid).await;
31 if !epilogue_started.load(Ordering::Acquire) {
32 spawn_shutdown_watchdog();
35 }
36 let writer_shutdown = tokio::task::spawn_blocking(move || database_writer.shutdown())
37 .await
38 .context("database writer shutdown task panicked")
39 .and_then(std::convert::identity);
40 record_daemon_cleanup(&mut outcome, "shut down database writer", writer_shutdown);
41 outcome
42}
43
44pub(super) async fn run_daemon_runtime(
45 epilogue_started: &AtomicBool,
46 owner_pid: Option<u32>,
47) -> Result<()> {
48 tokio::task::spawn_blocking(crate::controller::pin_worker_binary_sources)
51 .await
52 .context("worker source snapshot task failed")??;
53 Controller::recover_config_id_rename()?;
54 let config = Config::load()?;
55 crate::database::recover_interrupted_checkpointing_sessions(
56 &chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
57 )?;
58 crate::controller::reconcile_managed_checkpoint_archives()?;
59
60 let controller = Controller::load()?;
61 let listener = TcpListener::bind(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0))
62 .await
63 .context("bind Mjolnir daemon loopback endpoint")?;
64 let metadata = DaemonMetadata {
65 protocol_version: PROTOCOL_VERSION,
66 pid: std::process::id(),
67 address: listener.local_addr()?,
68 token: random_hex::<32>()?,
69 started_at: chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
70 build_version: env!("CARGO_PKG_VERSION").to_owned(),
71 };
72 let workspaces = tokio::task::spawn_blocking(crate::database::list_workspaces)
73 .await
74 .context("daemon workspace load task panicked")??;
75 let mut remote = if config.phone.enabled {
76 Some(spawn_remote_session_manager()?)
77 } else {
78 None
79 };
80
81 let manager = spawn_session_manager()?;
84 let manager_targets = manager.targets;
85 manager_targets.send_replace(dashboard_worker_targets(&controller));
86 let mut manager_updates = manager.updates;
87 let manager_control = manager.control.clone();
88 let manager_shutdown = manager.shutdown;
89 let mut recovery = crate::recovery::RecoveryCoordinator::spawn(manager_control.clone());
90 let recovery_observer = recovery.observer();
91 let mut worker_upgrades = crate::worker_upgrade::WorkerUpgradeCoordinator::spawn(
94 manager_control.clone(),
95 &recovery_observer,
96 );
97 let state = Arc::new(RuntimeState::new(
98 manager_control.clone(),
99 Controller {
100 config: controller.config.clone(),
101 state: controller.state.clone(),
102 },
103 recovery_observer.clone(),
104 worker_upgrades.observer(),
105 workspaces,
106 ));
107 let move_operations = blocking(crate::database::load_move_operations).await?;
108 let move_owned = state.recover_moves(move_operations)?;
109 state.resume_retained_cleanups();
110 let cancellation = crate::termination::Coordinator::install().token();
111 let target_refresh = spawn_manager_target_refresher(
112 manager_targets.clone(),
113 cancellation.clone(),
114 state.clone(),
115 );
116 let image_refresh = spawn_image_refresher(
117 {
118 let state = state.clone();
119 move || state.with_config(crate::controller::image_refresh_plan)
120 },
121 cancellation.clone(),
122 );
123 let exit_when_idle = mj_core::config::env_override_os("DAEMON_EXIT_WHEN_IDLE").is_some();
124 let mut idle_tick = tokio::time::interval(Duration::from_millis(100));
125 idle_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
126 let mut owner_tick = tokio::time::interval(Duration::from_millis(500));
127 owner_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
128 let mut recovery_tick = tokio::time::interval(Duration::from_millis(250));
129 recovery_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
130 let (interrupted_close_tx, mut interrupted_close_rx) = tokio::sync::mpsc::unbounded_channel();
131 let mut interrupted_close_cancellations = Vec::new();
132 let mut interrupted_close_tasks = Vec::new();
133 for session_id in interrupted_close_session_ids(&controller) {
134 if move_owned.contains(&session_id) {
135 continue;
136 }
137 let interrupted_cancellation = Arc::new(AtomicBool::new(false));
138 let interrupted_close_task = spawn_interrupted_close_recovery(
139 session_id,
140 manager_control.clone(),
141 recovery_observer.clone(),
142 interrupted_cancellation.clone(),
143 interrupted_close_tx.clone(),
144 None,
145 );
146 interrupted_close_cancellations.push(interrupted_cancellation);
147 interrupted_close_tasks.push(interrupted_close_task);
148 }
149 let mut phone_publisher: Option<RemoteSessionPublisher> = None;
150 let mut phone_task = None;
151 let mut remote_request_bridge = None;
152 if let Some(remote) = remote.take() {
153 remote
154 .targets
155 .send_replace(dashboard_worker_targets(&controller));
156 phone_publisher = Some(remote.publisher.clone());
157 remote_request_bridge = Some(spawn_remote_request_bridge(
158 remote.requests,
159 manager_control.clone(),
160 ));
161 phone_task = Some(spawn_phone_server(
162 config.phone,
163 cancellation.clone(),
164 state.clone(),
165 SessionManagerChannels {
166 targets: remote.targets,
167 control: remote.control,
168 updates: remote.updates,
169 shutdown: remote.shutdown,
170 },
171 ));
172 } else {
173 state.set_phone_status(WebViewerStatus::Disabled);
174 state.web_viewer.publish(crate::server::WebViewerAccess::Unavailable("Web access is disabled. Enable [phone].enabled in your configuration, then restart the daemon.".into()));
175 }
176 let daemon_metadata_path = metadata_path();
177 let mut client_tasks = tokio::task::JoinSet::new();
178
179 let mut outcome = async {
183 write_metadata(&daemon_metadata_path, &metadata)?;
184 reach_test_hook("daemon_metadata_before_listening").await?;
185 loop {
186 tokio::select! {
187 _ = cancellation.cancelled() => break,
188 _ = idle_tick.tick(), if exit_when_idle && state.ever_attached.load(Ordering::Acquire) => {
189 state.prune_dead_clients();
190 if state.attachments().is_empty() {
191 break;
192 }
193 }
194 _ = owner_tick.tick(), if owner_pid.is_some() => {
195 if let Some(owner) = owner_pid
196 && !process_is_alive(owner)
197 {
198 tracing::info!(owner_pid = owner, "daemon owner process exited; shutting down");
199 break;
200 }
201 }
202 _ = recovery_tick.tick() => {
203 while let Some(result) = recovery.try_result() {
204 if let Err(error) = &result.outcome {
205 if result.deferred {
209 tracing::info!(session_id = %result.session_id, %error, "recovery copy deferred: agent is working");
210 } else {
211 tracing::warn!(session_id = %result.session_id, %error, "daemon recovery checkpoint failed");
212 }
213 }
214 refresh_runtime_controller(&state).await;
215 }
216 while let Some(result) = worker_upgrades.try_result() {
217 report_worker_upgrade(&state, &result);
218 }
219 }
220 completed = interrupted_close_rx.recv() => {
221 if let Some(completed) = completed {
222 let recovered = completed.result.is_ok();
223 if let Err(error) = completed.result {
224 tracing::warn!(session_id = %completed.session_id, %error, "daemon could not resume interrupted close");
225 }
226 refresh_runtime_controller(&state).await;
227 if recovered && completed.deferred_cleanup
228 && let Err(error) = state.start_deferred_cleanup(completed.session_id.clone())
229 {
230 tracing::warn!(session_id = %completed.session_id, error = format!("{error:#}"), "could not continue cleanup after interrupted close");
231 state.push_notice(
232 &completed.session_id,
233 format!("Could not continue container storage cleanup: {error:#}"),
234 );
235 }
236 }
237 }
238 accepted = listener.accept() => {
239 let (stream, peer) = accepted.context("accept Mjolnir daemon client")?;
240 if !peer.ip().is_loopback() {
241 tracing::warn!(%peer, "rejected non-loopback daemon client");
242 continue;
243 }
244 let metadata = metadata.clone();
245 let state = state.clone();
246 let cancellation = cancellation.clone();
247 client_tasks.spawn(async move {
248 if let Err(error) = serve_client(stream, metadata, state, cancellation).await {
249 tracing::debug!(error = format!("{error:#}"), "daemon client disconnected");
250 }
251 });
252 }
253 completed = client_tasks.join_next(), if !client_tasks.is_empty() => {
254 if let Some(Err(error)) = completed {
255 tracing::warn!(%error, "daemon client task failed");
256 }
257 }
258 update = manager_updates.recv() => {
259 let Some(update) = update else {
260 bail!("controller daemon session manager stopped");
261 };
262 if let Some((detail, observed_updated_at)) =
263 state.missing_target_record(&update.session_id, &update.view)
264 {
265 let state = state.clone();
266 let session_id = update.session_id.clone();
267 client_tasks.spawn(async move {
268 if let Err(error) = state.persist_missing_target(
269 &session_id, detail, observed_updated_at,
270 ).await {
271 tracing::warn!(%session_id, %error, "could not persist missing worker target");
272 state.push_notice(&session_id, format!("Could not record missing session target: {error:#}"));
273 }
274 });
275 }
276 if let Some(publisher) = phone_publisher.as_ref()
277 && let Err(error) = publisher.try_publish(
278 update.session_id.clone(),
279 update.view.clone(),
280 )
281 {
282 tracing::warn!(%error, "phone session view bridge stopped");
283 phone_publisher = None;
284 }
285 state.review_host().observe(&update.session_id, &update.view);
289 state.publish_session(update.session_id, update.view).await?;
290 }
291 }
292 }
293 Ok(())
294 }
295 .await;
296
297 epilogue_started.store(true, Ordering::Release);
298 spawn_shutdown_watchdog();
299 cancellation.cancel();
302 for interrupted_cancellation in interrupted_close_cancellations {
303 interrupted_cancellation.store(true, Ordering::Release);
304 }
305 drop(interrupted_close_tx);
306 record_daemon_cleanup(
307 &mut outcome,
308 "remove daemon metadata",
309 remove_daemon_metadata(&daemon_metadata_path),
310 );
311 record_daemon_cleanup(
312 &mut outcome,
313 "shut down turn review host",
314 state
315 .review_host()
316 .shutdown()
317 .await
318 .map_err(anyhow::Error::msg),
319 );
320 record_daemon_cleanup(
321 &mut outcome,
322 "join controller target refresher",
323 target_refresh.await.map_err(anyhow::Error::new),
324 );
325 record_daemon_cleanup(
326 &mut outcome,
327 "join container image refresher",
328 image_refresh.await.map_err(anyhow::Error::new),
329 );
330 if let Some(phone_task) = phone_task {
331 record_daemon_cleanup(
332 &mut outcome,
333 "join phone server",
334 phone_task.await.map_err(anyhow::Error::new),
335 );
336 }
337 if let Some(remote_request_bridge) = remote_request_bridge {
338 record_daemon_cleanup(
339 &mut outcome,
340 "join phone session request bridge",
341 remote_request_bridge.await.map_err(anyhow::Error::new),
342 );
343 }
344 client_tasks.abort_all();
345 while let Some(result) = client_tasks.join_next().await {
346 if let Err(error) = result
347 && !error.is_cancelled()
348 {
349 record_daemon_cleanup(
350 &mut outcome,
351 "join daemon client task",
352 Err(anyhow::Error::new(error)),
353 );
354 }
355 }
356 record_daemon_cleanup(
357 &mut outcome,
358 "cancel daemon lifecycle operations",
359 state.cancel_and_wait_lifecycles().await,
360 );
361 for interrupted_close_task in interrupted_close_tasks {
362 record_daemon_cleanup(
363 &mut outcome,
364 "join interrupted close recovery",
365 interrupted_close_task.await.map_err(anyhow::Error::new),
366 );
367 }
368 drop(recovery);
369 record_daemon_cleanup(
370 &mut outcome,
371 "shut down controller daemon session manager",
372 manager_shutdown.shutdown().await,
373 );
374 outcome
375}
376
377pub(super) fn remove_daemon_metadata(path: &Path) -> Result<()> {
378 match fs::remove_file(path) {
379 Ok(()) => Ok(()),
380 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
381 Err(error) => Err(error).with_context(|| format!("remove {}", path.display())),
382 }
383}
384
385pub(super) fn record_daemon_cleanup(
389 outcome: &mut Result<()>,
390 operation: &'static str,
391 cleanup: Result<()>,
392) {
393 let Err(error) = cleanup else {
394 return;
395 };
396 let error = error.context(operation);
397 if outcome.is_ok() {
398 *outcome = Err(error);
399 } else {
400 tracing::warn!(error = format!("{error:#}"), "daemon cleanup step failed");
401 }
402}
403
404pub(super) fn spawn_shutdown_watchdog() {
411 tokio::spawn(async move {
412 tokio::time::sleep(SHUTDOWN_FORCE_EXIT_TIMEOUT).await;
413 tracing::error!(
414 seconds = SHUTDOWN_FORCE_EXIT_TIMEOUT.as_secs(),
415 "daemon shutdown did not finish in time; exiting"
416 );
417 if let Err(error) = fs::remove_file(metadata_path())
420 && error.kind() != std::io::ErrorKind::NotFound
421 {
422 tracing::warn!(%error, "could not remove daemon metadata before the forced exit");
423 }
424 std::process::exit(1);
426 });
427}
428
429pub(super) fn spawn_manager_target_refresher(
430 targets: tokio::sync::watch::Sender<Vec<crate::session_manager::RelaySessionTarget>>,
431 cancellation: CancellationToken,
432 state: Arc<RuntimeState>,
433) -> tokio::task::JoinHandle<()> {
434 tokio::spawn(async move {
435 let mut interval = tokio::time::interval(Duration::from_millis(500));
436 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
437 loop {
438 tokio::select! {
439 _ = cancellation.cancelled() => return,
440 _ = interval.tick() => {
441 let _config_mutation = state.config_mutation.lock().await;
444 match tokio::task::spawn_blocking(Controller::load).await {
445 Ok(Ok(controller)) => {
446 let lifecycle_sessions =
451 state.worker_poll_exclusion_session_ids(&controller);
452 let refreshed = dashboard_worker_targets_excluding(
453 &controller,
454 &lifecycle_sessions,
455 );
456 let changed = {
457 let mut review = state
458 .review_config
459 .lock()
460 .unwrap_or_else(PoisonError::into_inner);
461 review.clone_from(&controller.config.review);
462 drop(review);
463 let mut current = state
464 .controller
465 .lock()
466 .unwrap_or_else(PoisonError::into_inner);
467 let changed = current.config != controller.config;
468 *current = controller;
469 changed
470 };
471 state.review_host().retain_sessions(
475 refreshed
476 .iter()
477 .map(|target| target.session_id.clone())
478 .collect(),
479 );
480 targets.send_replace(refreshed);
481 if changed {
482 state.publish_revision();
483 }
484 }
485 Ok(Err(error)) => {
486 if let Some(mismatch) = error
495 .chain()
496 .find_map(|cause| cause.downcast_ref::<StoreSchemaMismatch>())
497 {
498 tracing::error!(
499 found = mismatch.found,
500 supported = mismatch.supported,
501 error = %mismatch,
502 "daemon store schema diverged underneath the daemon; shutting down"
503 );
504 cancellation.cancel();
505 return;
506 }
507 tracing::warn!(error = format!("{error:#}"), "could not refresh daemon session targets");
508 }
509 Err(error) => {
510 tracing::error!(%error, "daemon target refresh task failed");
511 return;
512 }
513 }
514 }
515 }
516 }
517 })
518}
519
520pub(super) async fn refresh_runtime_controller(state: &RuntimeState) {
521 if let Err(error) = state.reload_controller().await {
522 tracing::warn!(
523 error = format!("{error:#}"),
524 "could not refresh daemon controller state"
525 );
526 }
527}