Skip to main content

mj_controller/daemon/
process.rs

1use super::*;
2
3/// PID of the process this daemon must not outlive, if one was requested.
4///
5/// Tests start daemons that no client ever attaches to, so idle exit cannot
6/// retire them, and a test process that dies without unwinding never runs its
7/// teardown. Naming an owner makes the daemon responsible for its own lifetime.
8pub(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    // Checked before the store is locked so a bad value fails fast and leaves
25    // no daemon state behind.
26    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        // Initialization failed before the runtime-owned epilogue existed.
33        // The same process-level bound still applies to closing the writer.
34        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    // Freeze worker sources before any session can be created or upgraded.
49    // Copying binaries belongs on a blocking task, never the runtime event loop.
50    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    // Start the primary manager last: every remaining fallible operation is
82    // inside `outcome`, so its owner always reaches the awaited epilogue.
83    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    // Shares the recovery gate, so a recovery copy and a worker upgrade never
92    // act on one session at the same time.
93    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    // Everything a client can use is initialized before this atomic
180    // publication. From here on every exit, including an error from the test
181    // hook or the event loop, flows through the same bounded epilogue.
182    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                            // A deferred copy found the agent working. That is
206                            // the normal state of a session in use, so it is
207                            // news, not a fault.
208                            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                    // Every session's view passes here whether or not anything is
286                    // attached, which is exactly what an automatic review needs to
287                    // see: the turn that just finished.
288                    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    // Idle exit and fallible loop exits do not arrive through the termination
300    // coordinator. Stop every daemon-owned task before closing the sole writer.
301    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
385/// Keep the event-loop failure as the primary result while still running and
386/// reporting every cleanup step. If the loop ended normally, the first
387/// cleanup failure becomes the daemon's result.
388pub(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
404/// Bounds the epilogue below.
405///
406/// The daemon leaves on its own long before this fires: a graceful exit
407/// returns from `run_daemon_process`, the process exits 0, and this task dies
408/// with the runtime. It exists so no unwinding step can hold the process open
409/// past the deadline its clients wait on, whatever the cause of the shutdown.
410pub(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        // The metadata file points clients at a process that is about to stop
418        // answering. Removing it is what the epilogue would have done.
419        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        // 128 + signal is reserved for exits that really were signalled.
425        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                    // Keep a controller loaded from the old config from being
442                    // installed after a concurrent id rename has committed.
443                    let _config_mutation = state.config_mutation.lock().await;
444                    match tokio::task::spawn_blocking(Controller::load).await {
445                        Ok(Ok(controller)) => {
446                            // Startup, force-stop, relocation, and the teardown
447                            // phase of close own the worker target. Graceful
448                            // close keeps polling only until it has released the
449                            // manager lease after sealing the relay.
450                            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                            // Prune the review host's retained transcripts to the
472                            // same live set, so a stopped or destroyed session's
473                            // MaterializedSession does not linger there forever.
474                            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                            // The one place divergence is classified. Every
487                            // read re-checks store compatibility, so an
488                            // incompatible migration reaches this branch
489                            // within one tick. A daemon that
490                            // cannot read its own store cannot serve anyone,
491                            // and its writer is already refusing work, so the
492                            // answer is the shutdown it already knows how to
493                            // perform.
494                            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}