1use super::diagnostics::ServingPhase;
2use super::*;
3use std::time::Instant;
4
5pub(super) fn owner_pid_to_watch() -> Result<Option<u32>> {
11 let Some(value) = mj_core::config::env_override("DAEMON_OWNER_PID") else {
12 return Ok(None);
13 };
14 let pid: u32 = value
15 .trim()
16 .parse()
17 .map_err(|_| anyhow!("MJ_DAEMON_OWNER_PID must be a process id, but it is {value:?}"))?;
18 ensure!(
19 process_is_alive(pid),
20 "MJ_DAEMON_OWNER_PID names process {pid}, which is not running"
21 );
22 Ok(Some(pid))
23}
24
25pub async fn run_daemon_process() -> Result<()> {
26 let owner_pid = owner_pid_to_watch()?;
29 let guard = ControllerStoreGuard::acquire()?;
30 let database_writer = guard.start_database_writer()?;
31 let epilogue_started = AtomicBool::new(false);
32 let mut outcome = run_daemon_runtime(&epilogue_started, owner_pid).await;
33 if !epilogue_started.load(Ordering::Acquire) {
34 spawn_shutdown_watchdog();
37 }
38 let writer_shutdown = tokio::task::spawn_blocking(move || database_writer.shutdown())
39 .await
40 .context("database writer shutdown task panicked")
41 .and_then(std::convert::identity);
42 record_daemon_cleanup(&mut outcome, "shut down database writer", writer_shutdown);
43 outcome
44}
45
46pub(super) async fn run_daemon_runtime(
47 epilogue_started: &AtomicBool,
48 owner_pid: Option<u32>,
49) -> Result<()> {
50 let progress = super::diagnostics::DaemonProgressMonitor::start()?;
51 let started = SystemTime::now();
54 let startup_work = crate::upgrade::activity("daemon startup recovery")?;
55 tokio::task::spawn_blocking(crate::controller::pin_worker_binary_sources)
58 .await
59 .context("worker source snapshot task failed")??;
60 tokio::task::spawn_blocking(crate::targets::SshSessions::remove_stale_master_locks)
62 .await
63 .context("SSH master lock cleanup task failed")?;
64 Controller::recover_config_id_rename()?;
65 let config = Config::load()?;
66 crate::database::recover_interrupted_checkpointing_sessions(
67 &chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
68 )?;
69 crate::controller::reconcile_managed_checkpoint_archives()?;
70
71 let controller = tokio::task::spawn_blocking(|| {
72 let mut controller = Controller::load()?;
73 controller.prepare_persisted_sessions()?;
74 Ok::<_, anyhow::Error>(controller)
75 })
76 .await
77 .context("prepare persisted daemon sessions")??;
78 {
82 let state = controller.state.clone();
83 tokio::task::spawn_blocking(move || {
84 crate::controller::local_profile_homes::link_profile_homes_of_earlier_sessions(&state)
85 })
86 .await
87 .context("earlier sessions' profile home link task failed")?;
88 }
89 {
92 let config = config.clone();
93 let state = controller.state.clone();
94 tokio::task::spawn_blocking(move || {
95 crate::controller::local_profile_homes::remove_replicas_left_in_profile_homes(
96 &config,
97 &state,
98 &crate::targets::BoundedProcessExecutor::new(Duration::from_secs(15)),
99 )
100 })
101 .await
102 .context("leftover project-memory replica cleanup task failed")?;
103 }
104 let listener = TcpListener::bind(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0))
105 .await
106 .context("bind Mjolnir daemon loopback endpoint")?;
107 let metadata = DaemonMetadata {
108 protocol_version: PROTOCOL_VERSION,
109 pid: std::process::id(),
110 address: listener.local_addr()?,
111 token: random_hex::<32>()?,
112 started_at: chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
113 build_version: mj_client::build_identity::this_build().published(),
116 };
117 let workspaces = tokio::task::spawn_blocking(crate::database::list_workspaces)
118 .await
119 .context("daemon workspace load task panicked")??;
120 let mut remote = if config.phone.enabled {
121 Some(spawn_remote_session_manager()?)
122 } else {
123 None
124 };
125
126 let (delegation_tx, delegation_updates) = crate::session_manager::delegation_channel();
129 let manager = crate::session_manager::spawn_session_manager_observed(Some(delegation_tx))?;
130 let manager_targets = manager.targets;
131 manager_targets.send_replace(dashboard_worker_targets(&controller));
132 let manager_updates = manager.updates;
133 let manager_control = manager.control.clone();
134 let manager_shutdown = manager.shutdown;
135 let mut recovery = crate::recovery::RecoveryCoordinator::spawn(manager_control.clone());
136 let recovery_observer = recovery.observer();
137 let mut worker_upgrades = crate::worker_upgrade::WorkerUpgradeCoordinator::spawn(
140 manager_control.clone(),
141 &recovery_observer,
142 );
143 let state = Arc::new(RuntimeState::new(
144 manager_control.clone(),
145 Controller {
146 config: controller.config.clone(),
147 state: controller.state.clone(),
148 },
149 recovery_observer.clone(),
150 worker_upgrades.observer(),
151 workspaces,
152 ));
153 state.owner().install_relay_sessions(
156 manager_targets
157 .borrow()
158 .iter()
159 .map(|target| target.session_id.clone())
160 .collect(),
161 );
162 let cancellation = crate::termination::Coordinator::install().token();
163 let bootstrap = async {
166 let move_operations = blocking(crate::database::load_move_operations).await?;
167 let move_sessions = move_operations
168 .iter()
169 .map(|operation| operation.selection.session_id.clone())
170 .collect::<BTreeSet<_>>();
171 let move_owned = state.recover_moves(move_operations)?;
172 state.resume_retained_cleanups();
173 state.resume_startup_cleanups(true);
174 state.restore_startup_deliveries(&cancellation).await?;
175 Ok::<_, anyhow::Error>((move_sessions, move_owned))
176 }
177 .await;
178 let (move_sessions, move_owned) = match bootstrap {
179 Ok(ownership) => ownership,
180 Err(error) => {
181 epilogue_started.store(true, Ordering::Release);
182 spawn_shutdown_watchdog();
183 cancellation.cancel();
184 let mut outcome = Err(error);
185 record_daemon_cleanup(
186 &mut outcome,
187 "shut down bootstrap worker upgrade coordinator",
188 worker_upgrades.shutdown().await,
189 );
190 record_daemon_cleanup(
191 &mut outcome,
192 "shut down bootstrap recovery coordinator",
193 recovery.shutdown().await,
194 );
195 record_daemon_cleanup(
196 &mut outcome,
197 "shut down turn review host after bootstrap failure",
198 state
199 .review_host()
200 .shutdown()
201 .await
202 .map_err(anyhow::Error::msg),
203 );
204 record_daemon_cleanup(
205 &mut outcome,
206 "cancel bootstrap lifecycle operations",
207 state.cancel_and_wait_lifecycles().await,
208 );
209 record_daemon_cleanup(
210 &mut outcome,
211 "drain bootstrap startup deliveries",
212 state.cancel_and_join_startup_prompts().await,
213 );
214 if let Some(remote) = remote.take() {
215 record_daemon_cleanup(
216 &mut outcome,
217 "shut down bootstrap remote session manager",
218 remote.shutdown.shutdown().await,
219 );
220 }
221 record_daemon_cleanup(
222 &mut outcome,
223 "shut down bootstrap session manager",
224 manager_shutdown.shutdown().await,
225 );
226 return outcome;
227 }
228 };
229 let checkpoint_sweep = {
234 let cancellation = cancellation.clone();
235 tokio::task::spawn_blocking(move || {
236 crate::controller::sweep_local_checkpoint_leftovers(
237 &crate::targets::BoundedProcessExecutor::new(Duration::from_secs(15)),
238 started,
239 &|| cancellation.is_cancelled(),
240 );
241 })
242 };
243 let mut update_feed = continuation::spawn(state.clone(), manager_updates, cancellation.clone());
244
245 let (delegation_services, mut delegation_task) = super::delegation::spawn(
246 state.clone(),
247 manager_control.clone(),
248 delegation_updates,
249 cancellation.clone(),
250 );
251
252 let cpu_publication = {
255 let mut receiver = manager.session_cpu;
256 let state = state.clone();
257 let cancellation = cancellation.clone();
258 tokio::spawn(async move {
259 let mut next_publication = tokio::time::Instant::now();
260 loop {
261 tokio::select! {
262 _ = cancellation.cancelled() => return,
263 changed = receiver.changed() => {
264 if changed.is_err() {
265 tracing::error!("session CPU publication channel closed");
266 return;
267 }
268 }
269 }
270 tokio::select! {
273 _ = cancellation.cancelled() => return,
274 _ = tokio::time::sleep_until(next_publication) => {}
275 }
276 receiver.borrow_and_update();
277 state.publish_revision();
278 next_publication = tokio::time::Instant::now() + Duration::from_secs(10);
279 }
280 })
281 };
282 let project_catalog = tokio::spawn(state.projects().run(cancellation.child_token()));
283
284 let target_refresh = spawn_manager_target_refresher(
285 manager_targets.clone(),
286 cancellation.clone(),
287 state.clone(),
288 );
289 let image_refresh = spawn_image_refresher(
290 {
291 let state = state.clone();
292 move || state.with_config(crate::controller::image_refresh_plan)
293 },
294 {
295 let state = state.clone();
299 move |report| {
300 let text = match report {
301 crate::pollers::ImageRefreshReport::Started { host, image } => {
302 format!("Downloading image {image} for {host}\u{2026}")
303 }
304 crate::pollers::ImageRefreshReport::Pulled { host, image } => {
305 format!("Image {image} is ready on {host}.")
306 }
307 crate::pollers::ImageRefreshReport::Failed { host, image, error } => {
308 format!("Could not pull image {image} on {host}: {error}")
309 }
310 };
311 state.push_notice("", text);
312 }
313 },
314 cancellation.clone(),
315 );
316 let exit_when_idle = mj_core::config::env_override_os("DAEMON_EXIT_WHEN_IDLE").is_some();
317 let mut idle_tick = tokio::time::interval(Duration::from_millis(100));
318 idle_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
319 let mut owner_tick = tokio::time::interval(Duration::from_millis(500));
320 owner_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
321 let mut recovery_tick = tokio::time::interval(Duration::from_millis(250));
322 recovery_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
323 let mut background_policy_tick = tokio::time::interval(Duration::from_secs(1));
324 background_policy_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
325 let mut readiness_tick = tokio::time::interval(Duration::from_secs(5));
329 readiness_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
330 let mut progress_tick = tokio::time::interval(Duration::from_secs(1));
331 progress_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
332 let (interrupted_close_tx, mut interrupted_close_rx) = tokio::sync::mpsc::unbounded_channel();
333 let mut interrupted_close_tasks = Vec::new();
334 for session_id in interrupted_suspend_session_ids(&controller) {
335 if move_owned.contains(&session_id) {
336 continue;
337 }
338 let recovery_state = state.clone();
339 let recovery_shutdown = cancellation.clone();
340 let updates = interrupted_close_tx.clone();
341 let upgrade_work = startup_work.clone();
342 let interrupted_close_task = tokio::spawn(async move {
343 let _upgrade_work = upgrade_work;
344 let result = tokio::select! {
345 result = recovery_state.suspend_session(session_id.clone()) => result,
346 () = recovery_shutdown.cancelled() => return,
347 }
348 .map(|()| crate::pollers::LifecycleSuccess::Closed)
349 .map_err(|error| format!("{error:#}"));
350 if updates
351 .send(crate::pollers::LifecycleUpdate {
352 session_id,
353 result,
354 deferred_cleanup: false,
355 })
356 .is_err()
357 {
358 tracing::debug!("suspension recovery receiver stopped");
359 }
360 });
361 interrupted_close_tasks.push(interrupted_close_task);
362 }
363 for session_id in interrupted_destroy_session_ids(&controller) {
368 let recovery_state = state.clone();
369 let recovery_shutdown = cancellation.clone();
370 let updates = interrupted_close_tx.clone();
371 let upgrade_work = startup_work.clone();
372 interrupted_close_tasks.push(tokio::spawn(async move {
373 let _upgrade_work = upgrade_work;
374 let result = tokio::select! {
375 result = recovery_state.force_destroy_session(
376 session_id.clone(),
377 crate::daemon::BranchDisposition::Keep,
378 ) => result,
379 () = recovery_shutdown.cancelled() => return,
380 }
381 .map(|()| crate::pollers::LifecycleSuccess::ForceDestroyed)
382 .map_err(|error| format!("{error:#}"));
383 let _ = updates.send(crate::pollers::LifecycleUpdate {
384 session_id,
385 result,
386 deferred_cleanup: false,
387 });
388 }));
389 }
390 let reconciliation = {
396 let unowned = unowned_interrupted_lifecycles(
397 &controller,
398 &move_owned.union(&move_sessions).cloned().collect(),
399 );
400 (!unowned.is_empty()).then(|| {
401 let state = state.clone();
402 let upgrade_work = startup_work.clone();
403 tokio::spawn(async move {
404 let _upgrade_work = upgrade_work;
405 let reconciled = tokio::task::spawn_blocking(move || {
406 let mut controller = Controller::load()?;
407 let mut reconciled = 0usize;
408 for (session_id, cause) in unowned {
409 match controller.fail_interrupted_lifecycle(&session_id, &cause) {
410 Ok(true) => {
411 tracing::warn!(%session_id, %cause, "session left in flight by a daemon restart marked failed");
412 reconciled += 1;
413 }
414 Ok(false) => {}
415 Err(error) => tracing::warn!(
416 %session_id,
417 error = format!("{error:#}"),
418 "could not reconcile an interrupted lifecycle state"
419 ),
420 }
421 }
422 anyhow::Ok(reconciled)
423 })
424 .await;
425 match reconciled {
426 Ok(Ok(0)) => {}
427 Ok(Ok(_)) => refresh_runtime_controller(&state).await,
428 Ok(Err(error)) => tracing::warn!(
429 error = format!("{error:#}"),
430 "could not load the controller to reconcile interrupted lifecycles"
431 ),
432 Err(error) => {
433 tracing::warn!(%error, "interrupted lifecycle reconciliation task failed");
434 }
435 }
436 })
437 })
438 };
439 let tombstone_sweep = {
444 let tombstones = tombstone_session_ids(&controller);
445 (!tombstones.is_empty()).then(|| {
446 let state = state.clone();
447 let upgrade_work = startup_work.clone();
448 tokio::spawn(async move {
449 let _upgrade_work = upgrade_work;
450 for session_id in tombstones {
451 state.discard_lost_session(session_id).await;
452 }
453 })
454 })
455 };
456 let mut phone_publisher: Option<RemoteSessionPublisher> = None;
457 let mut phone_targets = None;
458 let mut prepared_targets = manager_targets.subscribe();
459 let mut phone_task = None;
460 let mut remote_request_bridge = None;
461 if let Some(remote) = remote.take() {
462 remote
463 .targets
464 .send_replace(prepared_targets.borrow_and_update().clone());
465 phone_targets = Some(remote.targets.clone());
466 phone_publisher = Some(remote.publisher.clone());
467 remote_request_bridge = Some(spawn_remote_request_bridge(
468 remote.requests,
469 manager_control.clone(),
470 ));
471 phone_task = Some(spawn_phone_server(
472 config.phone,
473 cancellation.clone(),
474 state.clone(),
475 SessionManagerChannels {
476 session_cpu: remote.control.session_cpu.clone(),
477 targets: remote.targets,
478 control: remote.control,
479 updates: remote.updates,
480 shutdown: remote.shutdown,
481 },
482 delegation_services,
483 ));
484 } else {
485 state.set_phone_status(WebViewerStatus::Disabled);
486 state.web_viewer.publish(crate::server::WebViewerAccess::Unavailable("Web access is disabled. Enable [phone].enabled in your configuration, then restart the daemon.".into()));
487 }
488 let daemon_metadata_path = metadata_path();
489 let mut client_tasks = tokio::task::JoinSet::new();
490 let mut delegation_outcome = None;
491 let mut cache_task = tokio::spawn(crate::controller::mbx::service::run(
492 state.clone(),
493 cancellation.clone(),
494 ));
495
496 let mut cache_outcome = None;
497
498 let mut outcome = async {
502 write_metadata(&daemon_metadata_path, &metadata)?;
503 reach_test_hook("daemon_metadata_before_listening").await?;
504 drop(startup_work);
505 loop {
506 progress.serving_tick();
507 tokio::select! {
508 _ = cancellation.cancelled() => break,
509 _ = progress_tick.tick() => {},
510 changed = prepared_targets.changed(), if phone_targets.is_some() => {
511 progress.phase(ServingPhase::Targets);
512 changed.context("daemon worker target publication stopped")?;
513 if let Some(targets) = &phone_targets {
514 targets.send_replace(prepared_targets.borrow_and_update().clone());
515 }
516 }
517 result = &mut cache_task => {
518 cache_outcome = Some(result.context("machine cache service failed").and_then(|result| result));
519 anyhow::bail!("machine build cache service stopped unexpectedly");
520 }
521 result = &mut delegation_task => {
522 delegation_outcome = Some(result.map_err(anyhow::Error::from).and_then(|result| result));
523 anyhow::bail!("delegation coordinator stopped unexpectedly");
524 }
525 _ = idle_tick.tick(), if exit_when_idle && state.ever_attached.load(Ordering::Acquire) => {
526 progress.phase(ServingPhase::IdleCheck);
527 state.prune_dead_clients();
528 if state.attachments().is_empty() {
529 break;
530 }
531 }
532 _ = owner_tick.tick(), if owner_pid.is_some() => {
533 progress.phase(ServingPhase::OwnerCheck);
534 if let Some(owner) = owner_pid
535 && !process_is_alive(owner)
536 {
537 tracing::info!(owner_pid = owner, "daemon owner process exited; shutting down");
538 break;
539 }
540 }
541 _ = recovery_tick.tick() => {
542 progress.phase(ServingPhase::Recovery);
543 while let Some(result) = recovery.try_result() {
544 if let Err(error) = &result.outcome {
545 if result.deferred {
549 tracing::info!(session_id = %result.session_id, %error, "recovery copy deferred: agent is working");
550 } else if result.cancelled {
551 tracing::info!(session_id = %result.session_id, %error, reason = "superseded by suspend", "daemon recovery checkpoint cancelled");
555 } else {
556 tracing::warn!(session_id = %result.session_id, %error, "daemon recovery checkpoint failed");
557 }
558 }
559 refresh_runtime_controller(&state).await;
560 }
561 while let Some(result) = worker_upgrades.try_result() {
562 report_worker_upgrade(&state, &result);
563 }
564 }
565 _ = background_policy_tick.tick() => {
566 progress.phase(ServingPhase::BackgroundPolicy);
567 state.refresh_background_policies();
568 state.resume_startup_cleanups(false);
569 }
570 _ = readiness_tick.tick() => {
571 progress.phase(ServingPhase::Readiness);
572 for unready in state.sessions_without_a_usable_harness() {
577 let state = state.clone();
578 client_tasks.spawn(async move {
579 state.fail_unready_session(unready).await;
580 });
581 }
582 }
583 completed = interrupted_close_rx.recv() => {
584 progress.phase(ServingPhase::SuspensionRecovery);
585 if let Some(completed) = completed {
586 let recovered = completed.result.is_ok();
587 if let Err(error) = completed.result {
588 tracing::warn!(session_id = %completed.session_id, %error, "daemon could not resume interrupted close");
589 }
590 refresh_runtime_controller(&state).await;
591 if recovered && completed.deferred_cleanup
592 && let Err(error) = state.start_deferred_cleanup(completed.session_id.clone())
593 {
594 tracing::warn!(session_id = %completed.session_id, error = format!("{error:#}"), "could not continue cleanup after interrupted close");
595 state.push_notice(
596 &completed.session_id,
597 format!("Could not continue container storage cleanup: {error:#}"),
598 );
599 }
600 }
601 }
602 accepted = listener.accept() => {
603 progress.phase(ServingPhase::Client);
604 let (stream, peer) = accepted.context("accept Mjolnir daemon client")?;
605 if !peer.ip().is_loopback() {
606 tracing::warn!(%peer, "rejected non-loopback daemon client");
607 continue;
608 }
609 let metadata = metadata.clone();
610 let state = state.clone();
611 let cancellation = cancellation.clone();
612 client_tasks.spawn(async move {
613 if let Err(error) = serve_client(stream, metadata, state, cancellation).await {
614 tracing::debug!(error = format!("{error:#}"), "daemon client disconnected");
615 }
616 });
617 }
618 completed = client_tasks.join_next(), if !client_tasks.is_empty() => {
619 if let Some(Err(error)) = completed {
620 tracing::warn!(%error, "daemon client task failed");
621 }
622 }
623 update = update_feed.next() => {
624 progress.phase(ServingPhase::SessionUpdate);
625 let Some(update) = update? else { break };
629 if let Some((detail, observed_updated_at)) =
630 state.missing_target_record(&update.session_id, &update.view)
631 {
632 let state = state.clone();
633 let session_id = update.session_id.clone();
634 client_tasks.spawn(async move {
635 if let Err(error) = state.persist_missing_target(
636 &session_id, detail, observed_updated_at,
637 ).await {
638 tracing::warn!(%session_id, %error, "could not persist missing worker target");
639 state.push_notice(&session_id, format!("Could not record missing session target: {error:#}"));
640 }
641 });
642 }
643 if let Some(publisher) = phone_publisher.as_ref()
644 && let Err(error) = publisher.try_publish(
645 update.session_id.clone(),
646 update.view.clone(),
647 )
648 {
649 tracing::warn!(%error, "phone session view bridge stopped");
650 phone_publisher = None;
651 }
652 state.publish_session(update.session_id, update.view).await?;
657 }
658 }
659 }
660 Ok(())
661 }
662 .await;
663
664 progress.phase(ServingPhase::Shutdown);
665 let mut epilogue = EpilogueClock::start();
666 epilogue_started.store(true, Ordering::Release);
667 spawn_shutdown_watchdog();
668 cancellation.cancel();
671 epilogue.record(
672 &mut outcome,
673 "join machine build cache applications",
674 match cache_outcome {
675 Some(result) => result,
676 None => cache_task
677 .await
678 .context("machine cache service failed")
679 .and_then(|result| result),
680 },
681 );
682 epilogue.record(
686 &mut outcome,
687 "shut down worker upgrade coordinator",
688 worker_upgrades.shutdown().await,
689 );
690 epilogue.record(
691 &mut outcome,
692 "shut down recovery coordinator",
693 recovery.shutdown().await,
694 );
695 epilogue.record(
696 &mut outcome,
697 "join continuation service",
698 update_feed.join().await,
699 );
700 epilogue.record(
701 &mut outcome,
702 "join delegation coordinator",
703 match delegation_outcome {
704 Some(result) => result,
705 None => delegation_task
706 .await
707 .map_err(anyhow::Error::from)
708 .and_then(|result| result),
709 },
710 );
711 drop(interrupted_close_tx);
712 epilogue.record(
713 &mut outcome,
714 "remove daemon metadata",
715 remove_daemon_metadata(&daemon_metadata_path),
716 );
717 epilogue.record(
718 &mut outcome,
719 "shut down turn review host",
720 state
721 .review_host()
722 .shutdown()
723 .await
724 .map_err(anyhow::Error::msg),
725 );
726 epilogue.record(
727 &mut outcome,
728 "join CPU publication",
729 cpu_publication.await.map_err(anyhow::Error::from),
730 );
731 record_daemon_cleanup(
732 &mut outcome,
733 "join project discovery",
734 project_catalog.await.map_err(anyhow::Error::from),
735 );
736 epilogue.record(
737 &mut outcome,
738 "join controller target refresher",
739 target_refresh.await.map_err(anyhow::Error::new),
740 );
741 epilogue.record(
742 &mut outcome,
743 "join container image refresher",
744 image_refresh.await.map_err(anyhow::Error::new),
745 );
746 if let Some(phone_task) = phone_task {
747 epilogue.record(
748 &mut outcome,
749 "join phone server",
750 phone_task.await.map_err(anyhow::Error::new),
751 );
752 }
753 if let Some(remote_request_bridge) = remote_request_bridge {
754 epilogue.record(
755 &mut outcome,
756 "join phone session request bridge",
757 remote_request_bridge.await.map_err(anyhow::Error::new),
758 );
759 }
760 client_tasks.abort_all();
761 while let Some(result) = client_tasks.join_next().await {
762 if let Err(error) = result
763 && !error.is_cancelled()
764 {
765 record_daemon_cleanup(
766 &mut outcome,
767 "join daemon client task",
768 Err(anyhow::Error::new(error)),
769 );
770 }
771 }
772 epilogue.record(&mut outcome, "join daemon client tasks", Ok(()));
773 epilogue.record(
774 &mut outcome,
775 "cancel daemon lifecycle operations",
776 state.cancel_and_wait_lifecycles().await,
777 );
778 epilogue.record(
779 &mut outcome,
780 "drain startup prompts",
781 state.cancel_and_join_startup_prompts().await,
782 );
783 if let Some(reconciliation) = reconciliation {
784 epilogue.record(
785 &mut outcome,
786 "join interrupted lifecycle reconciliation",
787 reconciliation.await.map_err(anyhow::Error::new),
788 );
789 }
790 if let Some(tombstone_sweep) = tombstone_sweep {
791 epilogue.record(
792 &mut outcome,
793 "join lost-session discard sweep",
794 tombstone_sweep.await.map_err(anyhow::Error::new),
795 );
796 }
797 epilogue.record(
798 &mut outcome,
799 "join checkpoint leftover sweep",
800 checkpoint_sweep.await.map_err(anyhow::Error::new),
801 );
802 for interrupted_close_task in interrupted_close_tasks {
803 epilogue.record(
804 &mut outcome,
805 "join interrupted close recovery",
806 interrupted_close_task.await.map_err(anyhow::Error::new),
807 );
808 }
809 epilogue.record(
810 &mut outcome,
811 "shut down controller daemon session manager",
812 manager_shutdown.shutdown().await,
813 );
814 epilogue.finish();
815 outcome
816}
817
818pub(super) fn tombstone_session_ids(controller: &Controller) -> Vec<String> {
822 controller
823 .state
824 .sessions
825 .values()
826 .filter(|session| {
827 matches!(
828 session.state,
829 SessionState::Lost | SessionState::DestroyedWithDataLoss
830 )
831 })
832 .map(|session| session.id.clone())
833 .collect()
834}
835
836pub(super) fn remove_daemon_metadata(path: &Path) -> Result<()> {
837 match fs::remove_file(path) {
838 Ok(()) => Ok(()),
839 Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
840 Err(error) => Err(error).with_context(|| format!("remove {}", path.display())),
841 }
842}
843
844struct EpilogueClock {
850 began: Instant,
851 step: Instant,
852}
853
854impl EpilogueClock {
855 const LOGGED_STEP: Duration = Duration::from_millis(250);
857
858 fn start() -> Self {
859 let now = Instant::now();
860 tracing::info!("daemon shutdown began");
861 Self {
862 began: now,
863 step: now,
864 }
865 }
866
867 fn record(&mut self, outcome: &mut Result<()>, operation: &'static str, cleanup: Result<()>) {
870 let took = self.step.elapsed();
871 if took >= Self::LOGGED_STEP {
872 tracing::info!(
873 step = operation,
874 duration_ms = took.as_millis(),
875 "daemon shutdown step finished"
876 );
877 }
878 record_daemon_cleanup(outcome, operation, cleanup);
879 self.step = Instant::now();
880 }
881
882 fn finish(self) {
883 tracing::info!(
884 duration_ms = self.began.elapsed().as_millis(),
885 "daemon shutdown finished"
886 );
887 }
888}
889
890pub(super) fn record_daemon_cleanup(
891 outcome: &mut Result<()>,
892 operation: &'static str,
893 cleanup: Result<()>,
894) {
895 let Err(error) = cleanup else {
896 return;
897 };
898 let error = error.context(operation);
899 if outcome.is_ok() {
900 *outcome = Err(error);
901 } else {
902 tracing::warn!(error = format!("{error:#}"), "daemon cleanup step failed");
903 }
904}
905
906pub(super) fn spawn_shutdown_watchdog() {
913 tokio::spawn(async move {
914 tokio::time::sleep(SHUTDOWN_FORCE_EXIT_TIMEOUT).await;
915 tracing::error!(
916 seconds = SHUTDOWN_FORCE_EXIT_TIMEOUT.as_secs(),
917 "daemon shutdown did not finish in time; exiting"
918 );
919 if let Err(error) = fs::remove_file(metadata_path())
922 && error.kind() != std::io::ErrorKind::NotFound
923 {
924 tracing::warn!(%error, "could not remove daemon metadata before the forced exit");
925 }
926 std::process::exit(1);
928 });
929}
930
931pub(super) fn spawn_manager_target_refresher(
932 targets: tokio::sync::watch::Sender<Vec<crate::session_manager::RelaySessionTarget>>,
933 cancellation: CancellationToken,
934 state: Arc<RuntimeState>,
935) -> tokio::task::JoinHandle<()> {
936 tokio::spawn(async move {
937 let mut committed = state
938 .committed
939 .as_ref()
940 .expect("daemon owns the writer")
941 .clone();
942 let mut revisions = state.revisions();
943 let mut compatibility = tokio::time::interval(Duration::from_millis(500));
944 compatibility.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
945 let mut installed = None;
946 loop {
947 let inputs = {
948 let owner = state.owner();
949 if let Err(error) = owner.ensure_available() {
950 tracing::error!(%error, "daemon durable state is unavailable; shutting down");
951 cancellation.cancel();
952 return;
953 }
954 owner.pollable_worker_inputs()
955 };
956 if installed.as_ref() != Some(&inputs) {
957 let preparation = inputs.clone();
958 let refreshed =
959 match tokio::task::spawn_blocking(move || preparation.prepare()).await {
960 Ok(targets) => targets,
961 Err(error) => {
962 tracing::error!(%error, "worker target preparation failed");
963 cancellation.cancel();
964 return;
965 }
966 };
967 let retained = refreshed
968 .iter()
969 .map(|target| target.session_id.clone())
970 .collect::<BTreeSet<_>>();
971 let views_dropped = {
972 let mut owner = state.owner();
973 if owner.pollable_worker_inputs() != inputs {
976 continue;
977 }
978 targets.send_if_modified(|current| {
979 if *current == refreshed {
980 false
981 } else {
982 *current = refreshed;
983 true
984 }
985 });
986 owner.install_relay_sessions(retained.clone())
989 };
990 if views_dropped {
991 state.publish_revision();
992 }
993 state.review_host().retain_sessions(retained);
994 installed = Some(inputs);
995 }
996 tokio::select! {
997 _ = cancellation.cancelled() => return,
998 changed = committed.changed() => {
999 if changed.is_err() {
1000 tracing::error!("database publication feed stopped");
1001 cancellation.cancel();
1002 return;
1003 }
1004 state.publish_revision();
1005 }
1006 changed = revisions.changed() => {
1007 if changed.is_err() { return; }
1008 }
1009 _ = compatibility.tick() => {
1010 let _config_mutation = state.config_mutation.lock().await;
1011 let refreshed = tokio::task::spawn_blocking(|| {
1012 crate::database::check_read_compatibility()?;
1013 Config::load()
1014 }).await;
1015 match refreshed {
1016 Ok(Ok(config)) => {
1017 state.review_config.lock().unwrap_or_else(PoisonError::into_inner)
1018 .clone_from(&config.review);
1019 let changed = {
1020 let mut owner = state.owner();
1021 let changed = owner.controller().config != config;
1022 owner.install_config(config);
1023 changed
1024 };
1025 if changed { state.publish_revision(); }
1026 }
1027 Ok(Err(error)) => {
1028 if error.chain().any(|cause| cause.downcast_ref::<StoreSchemaMismatch>().is_some()) {
1029 tracing::error!(%error, "daemon store schema diverged underneath the daemon; shutting down");
1030 cancellation.cancel();
1031 return;
1032 }
1033 tracing::warn!(error = format!("{error:#}"), "could not refresh daemon configuration");
1034 }
1035 Err(error) => {
1036 tracing::error!(%error, "daemon configuration reader failed");
1037 cancellation.cancel();
1038 return;
1039 }
1040 }
1041 }
1042 }
1043 }
1044 })
1045}
1046
1047pub(super) async fn refresh_runtime_controller(state: &RuntimeState) {
1048 if let Err(error) = state.reload_controller().await {
1049 tracing::warn!(
1050 error = format!("{error:#}"),
1051 "could not refresh daemon controller state"
1052 );
1053 }
1054}