Skip to main content

mj_controller/daemon/
process.rs

1use super::diagnostics::ServingPhase;
2use super::*;
3use std::time::Instant;
4
5/// PID of the process this daemon must not outlive, if one was requested.
6///
7/// Tests start daemons that no client ever attaches to, so idle exit cannot
8/// retire them, and a test process that dies without unwinding never runs its
9/// teardown. Naming an owner makes the daemon responsible for its own lifetime.
10pub(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    // Checked before the store is locked so a bad value fails fast and leaves
27    // no daemon state behind.
28    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        // Initialization failed before the runtime-owned epilogue existed.
35        // The same process-level bound still applies to closing the writer.
36        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    // Before this daemon can start a checkpoint, so every checkpoint file older
52    // than this belongs to a checkpoint that no longer runs.
53    let started = SystemTime::now();
54    let startup_work = crate::upgrade::activity("daemon startup recovery")?;
55    // Freeze worker sources before any session can be created or upgraded.
56    // Copying binaries belongs on a blocking task, never the runtime event loop.
57    tokio::task::spawn_blocking(crate::controller::pin_worker_binary_sources)
58        .await
59        .context("worker source snapshot task failed")??;
60    // Lock files outlive the SSH masters they guarded (launch finding R3-11).
61    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    // A local session an earlier release started from a profile home keeps
79    // running from it until it is next staged. The link has to be in place
80    // before any launch configuration is refreshed or any credential sync runs.
81    {
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    // Earlier releases left each such session's project-memory replica in the
90    // profile home, where nothing removed it when the session ended.
91    {
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        // The version, with the revision and build times as semver build
114        // metadata that older clients ignore; see `mj_client::build_identity`.
115        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    // Start the primary manager after store preparation. Bootstrap failures
127    // below have an awaited epilogue just like failures in the serving loop.
128    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    // Shares the recovery gate, so a recovery copy and a worker upgrade never
138    // act on one session at the same time.
139    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    // Views are published only from the serving loop below, so the first
154    // targets are installed before any can arrive.
155    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    // Bootstrap already owns live managers. Capture its error so those owners
164    // are shut down before the process-level writer can be closed.
165    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    // Checkpoint files that a cancelled or interrupted checkpoint left in a
230    // local worker root. Deleting them can take minutes after a long leak, so
231    // this runs beside the daemon rather than before it serves. It stops at
232    // shutdown, and the next start finishes it.
233    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    // CPU publication is bounded and reconstructed from workers. It holds no
253    // admission during handoff and exits with the daemon's cancellation token.
254    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                // Publish the first update promptly, then coalesce other actors'
271                // updates. A separate polling phase adds a needless ten seconds.
272                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            // A background download is the daemon's own work, not a session's,
296            // so these notices carry an empty session id and reach every
297            // workspace.
298            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    // How often the harness readiness wait is looked at. The wait itself is
326    // minutes long, so this only bounds how late a failure is noticed, and it
327    // reads in-memory state rather than the store.
328    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    // A destroy that failed left its record in `Error` with the target still
364    // present; it is finished the way `mj destroy` finishes it. The branch is
365    // kept, because whether to delete it was that command's choice and is not
366    // recorded.
367    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    // Whatever is still in an in-flight lifecycle state now has no owner: the
391    // moves, the interrupted closes, and the checkpointing rows above are
392    // every operation that legitimately resumes. The list comes from the
393    // startup snapshot, so a session created after this cannot be caught by
394    // it, and the writes run off this path because they touch the database.
395    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    // Rows left as tombstones by an older build, or by a discard that could
440    // not finish before the daemon stopped. The list comes from the startup
441    // snapshot, so a session that becomes lost after this is handled by the
442    // view watcher instead.
443    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    // Everything a client can use is initialized before this atomic
499    // publication. From here on every exit, including an error from the test
500    // hook or the event loop, flows through the same bounded epilogue.
501    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                            // A deferred copy found the agent working. That is
546                            // the normal state of a session in use, so it is
547                            // news, not a fault.
548                            if result.deferred {
549                                tracing::info!(session_id = %result.session_id, %error, "recovery copy deferred: agent is working");
550                            } else if result.cancelled {
551                                // The daemon asked for this cancel (a suspend
552                                // or close superseded the copy), so it is not
553                                // a failure of the checkpoint.
554                                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                    // A live session whose harness never advertised itself
573                    // takes no prompt and reports no failure, so nothing else
574                    // ever ends its wait (#1090). The store write runs off
575                    // this loop.
576                    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                    // The feed says why it ended. A shutdown also stops the
626                    // continuation service and closes this channel, so a
627                    // closed channel is not evidence of a failure.
628                    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                    // Every session's view passes here whether or not anything is
653                    // attached, which is exactly what an automatic review needs to
654                    // see: the turn that just finished.
655                    // The continuation completion gate has already notified review.
656                    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    // Idle exit and fallible loop exits do not arrive through the termination
669    // coordinator. Stop every daemon-owned task before closing the sole writer.
670    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    // Recovery copies and worker upgrades do not hold up a handoff, so some
683    // may still be running. Stop them first: the next daemon starts them again.
684    // The coordinators share one gate, and dropping either cancels both.
685    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
818/// Sessions whose record is nothing but a tombstone: they ended without a
819/// target and without a checkpoint, so the record offers no action but its own
820/// removal. `DestroyedWithDataLoss` only ever comes from an older build.
821pub(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
844/// Keep the event-loop failure as the primary result while still running and
845/// reporting every cleanup step. If the loop ended normally, the first
846/// cleanup failure becomes the daemon's result.
847/// Times the shutdown epilogue, so a slow exit, such as the end of an upgrade
848/// handoff, records which step took the time.
849struct EpilogueClock {
850    began: Instant,
851    step: Instant,
852}
853
854impl EpilogueClock {
855    /// Steps shorter than this are not logged.
856    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    /// Record `cleanup`, the result of the step that just finished. Steps run
868    /// one after another, so the time since the previous one is this step's.
869    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
906/// Bounds the epilogue below.
907///
908/// The daemon leaves on its own long before this fires: a graceful exit
909/// returns from `run_daemon_process`, the process exits 0, and this task dies
910/// with the runtime. It exists so no unwinding step can hold the process open
911/// past the deadline its clients wait on, whatever the cause of the shutdown.
912pub(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        // The metadata file points clients at a process that is about to stop
920        // answering. Removing it is what the epilogue would have done.
921        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        // 128 + signal is reserved for exits that really were signalled.
927        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                    // A lifecycle may have claimed a target while its commands
974                    // were being prepared. Only the owner can authorize install.
975                    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                    // Same lock as the install, so a view is published only
987                    // while its actor is still wanted.
988                    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}