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_sessions = move_operations
112 .iter()
113 .map(|operation| operation.selection.session_id.clone())
114 .collect::<BTreeSet<_>>();
115 let move_owned = state.recover_moves(move_operations)?;
116 state.resume_retained_cleanups();
117 let cancellation = crate::termination::Coordinator::install().token();
118 let target_refresh = spawn_manager_target_refresher(
119 manager_targets.clone(),
120 cancellation.clone(),
121 state.clone(),
122 );
123 let image_refresh = spawn_image_refresher(
124 {
125 let state = state.clone();
126 move || state.with_config(crate::controller::image_refresh_plan)
127 },
128 {
129 let state = state.clone();
133 move |report| {
134 let text = match report {
135 crate::pollers::ImageRefreshReport::Started { host, image } => {
136 format!("Downloading image {image} for {host}\u{2026}")
137 }
138 crate::pollers::ImageRefreshReport::Pulled { host, image } => {
139 format!("Image {image} is ready on {host}.")
140 }
141 crate::pollers::ImageRefreshReport::Failed { host, image, error } => {
142 format!("Could not pull image {image} on {host}: {error}")
143 }
144 };
145 state.push_notice("", text);
146 }
147 },
148 cancellation.clone(),
149 );
150 let exit_when_idle = mj_core::config::env_override_os("DAEMON_EXIT_WHEN_IDLE").is_some();
151 let mut idle_tick = tokio::time::interval(Duration::from_millis(100));
152 idle_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
153 let mut owner_tick = tokio::time::interval(Duration::from_millis(500));
154 owner_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
155 let mut recovery_tick = tokio::time::interval(Duration::from_millis(250));
156 recovery_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
157 let (interrupted_close_tx, mut interrupted_close_rx) = tokio::sync::mpsc::unbounded_channel();
158 let mut interrupted_close_cancellations = Vec::new();
159 let mut interrupted_close_tasks = Vec::new();
160 for session_id in interrupted_close_session_ids(&controller) {
161 if move_owned.contains(&session_id) {
162 continue;
163 }
164 let interrupted_cancellation = Arc::new(AtomicBool::new(false));
165 let interrupted_close_task = spawn_interrupted_close_recovery(
166 session_id,
167 manager_control.clone(),
168 recovery_observer.clone(),
169 interrupted_cancellation.clone(),
170 interrupted_close_tx.clone(),
171 None,
172 );
173 interrupted_close_cancellations.push(interrupted_cancellation);
174 interrupted_close_tasks.push(interrupted_close_task);
175 }
176 let reconciliation = {
182 let unowned = unowned_interrupted_lifecycles(
183 &controller,
184 &move_owned.union(&move_sessions).cloned().collect(),
185 );
186 (!unowned.is_empty()).then(|| {
187 let state = state.clone();
188 tokio::spawn(async move {
189 let reconciled = tokio::task::spawn_blocking(move || {
190 let mut controller = Controller::load()?;
191 let mut reconciled = 0usize;
192 for (session_id, cause) in unowned {
193 match controller.fail_interrupted_lifecycle(&session_id, &cause) {
194 Ok(true) => {
195 tracing::warn!(%session_id, %cause, "session left in flight by a daemon restart marked failed");
196 reconciled += 1;
197 }
198 Ok(false) => {}
199 Err(error) => tracing::warn!(
200 %session_id,
201 error = format!("{error:#}"),
202 "could not reconcile an interrupted lifecycle state"
203 ),
204 }
205 }
206 anyhow::Ok(reconciled)
207 })
208 .await;
209 match reconciled {
210 Ok(Ok(0)) => {}
211 Ok(Ok(_)) => refresh_runtime_controller(&state).await,
212 Ok(Err(error)) => tracing::warn!(
213 error = format!("{error:#}"),
214 "could not load the controller to reconcile interrupted lifecycles"
215 ),
216 Err(error) => {
217 tracing::warn!(%error, "interrupted lifecycle reconciliation task failed");
218 }
219 }
220 })
221 })
222 };
223 let mut phone_publisher: Option<RemoteSessionPublisher> = None;
224 let mut phone_task = None;
225 let mut remote_request_bridge = None;
226 if let Some(remote) = remote.take() {
227 remote
228 .targets
229 .send_replace(dashboard_worker_targets(&controller));
230 phone_publisher = Some(remote.publisher.clone());
231 remote_request_bridge = Some(spawn_remote_request_bridge(
232 remote.requests,
233 manager_control.clone(),
234 ));
235 phone_task = Some(spawn_phone_server(
236 config.phone,
237 cancellation.clone(),
238 state.clone(),
239 SessionManagerChannels {
240 targets: remote.targets,
241 control: remote.control,
242 updates: remote.updates,
243 shutdown: remote.shutdown,
244 },
245 ));
246 } else {
247 state.set_phone_status(WebViewerStatus::Disabled);
248 state.web_viewer.publish(crate::server::WebViewerAccess::Unavailable("Web access is disabled. Enable [phone].enabled in your configuration, then restart the daemon.".into()));
249 }
250 let daemon_metadata_path = metadata_path();
251 let mut client_tasks = tokio::task::JoinSet::new();
252
253 let mut outcome = async {
257 write_metadata(&daemon_metadata_path, &metadata)?;
258 reach_test_hook("daemon_metadata_before_listening").await?;
259 loop {
260 tokio::select! {
261 _ = cancellation.cancelled() => break,
262 _ = idle_tick.tick(), if exit_when_idle && state.ever_attached.load(Ordering::Acquire) => {
263 state.prune_dead_clients();
264 if state.attachments().is_empty() {
265 break;
266 }
267 }
268 _ = owner_tick.tick(), if owner_pid.is_some() => {
269 if let Some(owner) = owner_pid
270 && !process_is_alive(owner)
271 {
272 tracing::info!(owner_pid = owner, "daemon owner process exited; shutting down");
273 break;
274 }
275 }
276 _ = recovery_tick.tick() => {
277 while let Some(result) = recovery.try_result() {
278 if let Err(error) = &result.outcome {
279 if result.deferred {
283 tracing::info!(session_id = %result.session_id, %error, "recovery copy deferred: agent is working");
284 } else {
285 tracing::warn!(session_id = %result.session_id, %error, "daemon recovery checkpoint failed");
286 }
287 }
288 refresh_runtime_controller(&state).await;
289 }
290 while let Some(result) = worker_upgrades.try_result() {
291 report_worker_upgrade(&state, &result);
292 }
293 }
294 completed = interrupted_close_rx.recv() => {
295 if let Some(completed) = completed {
296 let recovered = completed.result.is_ok();
297 if let Err(error) = completed.result {
298 tracing::warn!(session_id = %completed.session_id, %error, "daemon could not resume interrupted close");
299 }
300 refresh_runtime_controller(&state).await;
301 if recovered && completed.deferred_cleanup
302 && let Err(error) = state.start_deferred_cleanup(completed.session_id.clone())
303 {
304 tracing::warn!(session_id = %completed.session_id, error = format!("{error:#}"), "could not continue cleanup after interrupted close");
305 state.push_notice(
306 &completed.session_id,
307 format!("Could not continue container storage cleanup: {error:#}"),
308 );
309 }
310 }
311 }
312 accepted = listener.accept() => {
313 let (stream, peer) = accepted.context("accept Mjolnir daemon client")?;
314 if !peer.ip().is_loopback() {
315 tracing::warn!(%peer, "rejected non-loopback daemon client");
316 continue;
317 }
318 let metadata = metadata.clone();
319 let state = state.clone();
320 let cancellation = cancellation.clone();
321 client_tasks.spawn(async move {
322 if let Err(error) = serve_client(stream, metadata, state, cancellation).await {
323 tracing::debug!(error = format!("{error:#}"), "daemon client disconnected");
324 }
325 });
326 }
327 completed = client_tasks.join_next(), if !client_tasks.is_empty() => {
328 if let Some(Err(error)) = completed {
329 tracing::warn!(%error, "daemon client task failed");
330 }
331 }
332 update = manager_updates.recv() => {
333 let Some(update) = update else {
334 bail!("controller daemon session manager stopped");
335 };
336 if let Some((detail, observed_updated_at)) =
337 state.missing_target_record(&update.session_id, &update.view)
338 {
339 let state = state.clone();
340 let session_id = update.session_id.clone();
341 client_tasks.spawn(async move {
342 if let Err(error) = state.persist_missing_target(
343 &session_id, detail, observed_updated_at,
344 ).await {
345 tracing::warn!(%session_id, %error, "could not persist missing worker target");
346 state.push_notice(&session_id, format!("Could not record missing session target: {error:#}"));
347 }
348 });
349 }
350 if let Some(publisher) = phone_publisher.as_ref()
351 && let Err(error) = publisher.try_publish(
352 update.session_id.clone(),
353 update.view.clone(),
354 )
355 {
356 tracing::warn!(%error, "phone session view bridge stopped");
357 phone_publisher = None;
358 }
359 state.review_host().observe(&update.session_id, &update.view);
363 state.publish_session(update.session_id, update.view).await?;
364 }
365 }
366 }
367 Ok(())
368 }
369 .await;
370
371 epilogue_started.store(true, Ordering::Release);
372 spawn_shutdown_watchdog();
373 cancellation.cancel();
376 for interrupted_cancellation in interrupted_close_cancellations {
377 interrupted_cancellation.store(true, Ordering::Release);
378 }
379 drop(interrupted_close_tx);
380 record_daemon_cleanup(
381 &mut outcome,
382 "remove daemon metadata",
383 remove_daemon_metadata(&daemon_metadata_path),
384 );
385 record_daemon_cleanup(
386 &mut outcome,
387 "shut down turn review host",
388 state
389 .review_host()
390 .shutdown()
391 .await
392 .map_err(anyhow::Error::msg),
393 );
394 record_daemon_cleanup(
395 &mut outcome,
396 "join controller target refresher",
397 target_refresh.await.map_err(anyhow::Error::new),
398 );
399 record_daemon_cleanup(
400 &mut outcome,
401 "join container image refresher",
402 image_refresh.await.map_err(anyhow::Error::new),
403 );
404 if let Some(phone_task) = phone_task {
405 record_daemon_cleanup(
406 &mut outcome,
407 "join phone server",
408 phone_task.await.map_err(anyhow::Error::new),
409 );
410 }
411 if let Some(remote_request_bridge) = remote_request_bridge {
412 record_daemon_cleanup(
413 &mut outcome,
414 "join phone session request bridge",
415 remote_request_bridge.await.map_err(anyhow::Error::new),
416 );
417 }
418 client_tasks.abort_all();
419 while let Some(result) = client_tasks.join_next().await {
420 if let Err(error) = result
421 && !error.is_cancelled()
422 {
423 record_daemon_cleanup(
424 &mut outcome,
425 "join daemon client task",
426 Err(anyhow::Error::new(error)),
427 );
428 }
429 }
430 record_daemon_cleanup(
431 &mut outcome,
432 "cancel daemon lifecycle operations",
433 state.cancel_and_wait_lifecycles().await,
434 );
435 record_daemon_cleanup(
436 &mut outcome,
437 "drain startup prompts",
438 state.cancel_and_join_startup_prompts().await,
439 );
440 if let Some(reconciliation) = reconciliation {
441 record_daemon_cleanup(
442 &mut outcome,
443 "join interrupted lifecycle reconciliation",
444 reconciliation.await.map_err(anyhow::Error::new),
445 );
446 }
447 for interrupted_close_task in interrupted_close_tasks {
448 record_daemon_cleanup(
449 &mut outcome,
450 "join interrupted close recovery",
451 interrupted_close_task.await.map_err(anyhow::Error::new),
452 );
453 }
454 drop(recovery);
455 record_daemon_cleanup(
456 &mut outcome,
457 "shut down controller daemon session manager",
458 manager_shutdown.shutdown().await,
459 );
460 outcome
461}
462
463pub(super) fn remove_daemon_metadata(path: &Path) -> Result<()> {
464 match fs::remove_file(path) {
465 Ok(()) => Ok(()),
466 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
467 Err(error) => Err(error).with_context(|| format!("remove {}", path.display())),
468 }
469}
470
471pub(super) fn record_daemon_cleanup(
475 outcome: &mut Result<()>,
476 operation: &'static str,
477 cleanup: Result<()>,
478) {
479 let Err(error) = cleanup else {
480 return;
481 };
482 let error = error.context(operation);
483 if outcome.is_ok() {
484 *outcome = Err(error);
485 } else {
486 tracing::warn!(error = format!("{error:#}"), "daemon cleanup step failed");
487 }
488}
489
490pub(super) fn spawn_shutdown_watchdog() {
497 tokio::spawn(async move {
498 tokio::time::sleep(SHUTDOWN_FORCE_EXIT_TIMEOUT).await;
499 tracing::error!(
500 seconds = SHUTDOWN_FORCE_EXIT_TIMEOUT.as_secs(),
501 "daemon shutdown did not finish in time; exiting"
502 );
503 if let Err(error) = fs::remove_file(metadata_path())
506 && error.kind() != std::io::ErrorKind::NotFound
507 {
508 tracing::warn!(%error, "could not remove daemon metadata before the forced exit");
509 }
510 std::process::exit(1);
512 });
513}
514
515pub(super) fn spawn_manager_target_refresher(
516 targets: tokio::sync::watch::Sender<Vec<crate::session_manager::RelaySessionTarget>>,
517 cancellation: CancellationToken,
518 state: Arc<RuntimeState>,
519) -> tokio::task::JoinHandle<()> {
520 tokio::spawn(async move {
521 let mut interval = tokio::time::interval(Duration::from_millis(500));
522 interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
523 loop {
524 tokio::select! {
525 _ = cancellation.cancelled() => return,
526 _ = interval.tick() => {
527 let _config_mutation = state.config_mutation.lock().await;
530 match tokio::task::spawn_blocking(Controller::load).await {
531 Ok(Ok(controller)) => {
532 let lifecycle_sessions =
537 state.worker_poll_exclusion_session_ids(&controller);
538 let refreshed = dashboard_worker_targets_excluding(
539 &controller,
540 &lifecycle_sessions,
541 );
542 let changed = {
543 let mut review = state
544 .review_config
545 .lock()
546 .unwrap_or_else(PoisonError::into_inner);
547 review.clone_from(&controller.config.review);
548 drop(review);
549 let mut current = state
550 .controller
551 .lock()
552 .unwrap_or_else(PoisonError::into_inner);
553 let changed = current.config != controller.config;
554 *current = controller;
555 changed
556 };
557 state.review_host().retain_sessions(
561 refreshed
562 .iter()
563 .map(|target| target.session_id.clone())
564 .collect(),
565 );
566 targets.send_replace(refreshed);
567 if changed {
568 state.publish_revision();
569 }
570 }
571 Ok(Err(error)) => {
572 if let Some(mismatch) = error
581 .chain()
582 .find_map(|cause| cause.downcast_ref::<StoreSchemaMismatch>())
583 {
584 tracing::error!(
585 found = mismatch.found,
586 supported = mismatch.supported,
587 error = %mismatch,
588 "daemon store schema diverged underneath the daemon; shutting down"
589 );
590 cancellation.cancel();
591 return;
592 }
593 tracing::warn!(error = format!("{error:#}"), "could not refresh daemon session targets");
594 }
595 Err(error) => {
596 tracing::error!(%error, "daemon target refresh task failed");
597 return;
598 }
599 }
600 }
601 }
602 }
603 })
604}
605
606pub(super) async fn refresh_runtime_controller(state: &RuntimeState) {
607 if let Err(error) = state.reload_controller().await {
608 tracing::warn!(
609 error = format!("{error:#}"),
610 "could not refresh daemon controller state"
611 );
612 }
613}