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    // Every session a durable move intent names is owned by that intent,
109    // whether or not this startup resumes it, so reconciliation leaves it
110    // alone.
111    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            // A background download is the daemon's own work, not a session's,
130            // so these notices carry an empty session id and reach every
131            // workspace.
132            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    // Whatever is still in an in-flight lifecycle state now has no owner: the
177    // moves, the interrupted closes, and the checkpointing rows above are
178    // every operation that legitimately resumes. The list comes from the
179    // startup snapshot, so a session created after this cannot be caught by
180    // it, and the writes run off this path because they touch the database.
181    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    // Everything a client can use is initialized before this atomic
254    // publication. From here on every exit, including an error from the test
255    // hook or the event loop, flows through the same bounded epilogue.
256    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                            // A deferred copy found the agent working. That is
280                            // the normal state of a session in use, so it is
281                            // news, not a fault.
282                            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                    // Every session's view passes here whether or not anything is
360                    // attached, which is exactly what an automatic review needs to
361                    // see: the turn that just finished.
362                    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    // Idle exit and fallible loop exits do not arrive through the termination
374    // coordinator. Stop every daemon-owned task before closing the sole writer.
375    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
471/// Keep the event-loop failure as the primary result while still running and
472/// reporting every cleanup step. If the loop ended normally, the first
473/// cleanup failure becomes the daemon's result.
474pub(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
490/// Bounds the epilogue below.
491///
492/// The daemon leaves on its own long before this fires: a graceful exit
493/// returns from `run_daemon_process`, the process exits 0, and this task dies
494/// with the runtime. It exists so no unwinding step can hold the process open
495/// past the deadline its clients wait on, whatever the cause of the shutdown.
496pub(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        // The metadata file points clients at a process that is about to stop
504        // answering. Removing it is what the epilogue would have done.
505        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        // 128 + signal is reserved for exits that really were signalled.
511        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                    // Keep a controller loaded from the old config from being
528                    // installed after a concurrent id rename has committed.
529                    let _config_mutation = state.config_mutation.lock().await;
530                    match tokio::task::spawn_blocking(Controller::load).await {
531                        Ok(Ok(controller)) => {
532                            // Startup, force-stop, relocation, and the teardown
533                            // phase of close own the worker target. Graceful
534                            // close keeps polling only until it has released the
535                            // manager lease after sealing the relay.
536                            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                            // Prune the review host's retained transcripts to the
558                            // same live set, so a stopped or destroyed session's
559                            // MaterializedSession does not linger there forever.
560                            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                            // The one place divergence is classified. Every
573                            // read re-checks store compatibility, so an
574                            // incompatible migration reaches this branch
575                            // within one tick. A daemon that
576                            // cannot read its own store cannot serve anyone,
577                            // and its writer is already refusing work, so the
578                            // answer is the shutdown it already knows how to
579                            // perform.
580                            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}