brokk-mj-controller 2.23.2

Daemon-side controller, session manager, and web server for Mjolnir
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
use super::*;

/// PID of the process this daemon must not outlive, if one was requested.
///
/// Tests start daemons that no client ever attaches to, so idle exit cannot
/// retire them, and a test process that dies without unwinding never runs its
/// teardown. Naming an owner makes the daemon responsible for its own lifetime.
pub(super) fn owner_pid_to_watch() -> Result<Option<u32>> {
    let Some(value) = mj_core::config::env_override("DAEMON_OWNER_PID") else {
        return Ok(None);
    };
    let pid: u32 = value
        .trim()
        .parse()
        .map_err(|_| anyhow!("MJ_DAEMON_OWNER_PID must be a process id, but it is {value:?}"))?;
    ensure!(
        process_is_alive(pid),
        "MJ_DAEMON_OWNER_PID names process {pid}, which is not running"
    );
    Ok(Some(pid))
}

pub async fn run_daemon_process() -> Result<()> {
    // Checked before the store is locked so a bad value fails fast and leaves
    // no daemon state behind.
    let owner_pid = owner_pid_to_watch()?;
    let guard = ControllerStoreGuard::acquire()?;
    let database_writer = guard.start_database_writer()?;
    let epilogue_started = AtomicBool::new(false);
    let mut outcome = run_daemon_runtime(&epilogue_started, owner_pid).await;
    if !epilogue_started.load(Ordering::Acquire) {
        // Initialization failed before the runtime-owned epilogue existed.
        // The same process-level bound still applies to closing the writer.
        spawn_shutdown_watchdog();
    }
    let writer_shutdown = tokio::task::spawn_blocking(move || database_writer.shutdown())
        .await
        .context("database writer shutdown task panicked")
        .and_then(std::convert::identity);
    record_daemon_cleanup(&mut outcome, "shut down database writer", writer_shutdown);
    outcome
}

pub(super) async fn run_daemon_runtime(
    epilogue_started: &AtomicBool,
    owner_pid: Option<u32>,
) -> Result<()> {
    // Before this daemon can start a checkpoint, so every checkpoint file older
    // than this belongs to a checkpoint that no longer runs.
    let started = SystemTime::now();
    let startup_work = crate::upgrade::activity("daemon startup recovery")?;
    // Freeze worker sources before any session can be created or upgraded.
    // Copying binaries belongs on a blocking task, never the runtime event loop.
    tokio::task::spawn_blocking(crate::controller::pin_worker_binary_sources)
        .await
        .context("worker source snapshot task failed")??;
    // Lock files outlive the SSH masters they guarded (launch finding R3-11).
    tokio::task::spawn_blocking(crate::targets::SshSessions::remove_stale_master_locks)
        .await
        .context("SSH master lock cleanup task failed")?;
    Controller::recover_config_id_rename()?;
    let config = Config::load()?;
    crate::database::recover_interrupted_checkpointing_sessions(
        &chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
    )?;
    crate::controller::reconcile_managed_checkpoint_archives()?;

    let controller = tokio::task::spawn_blocking(|| {
        let mut controller = Controller::load()?;
        controller.prepare_persisted_sessions()?;
        Ok::<_, anyhow::Error>(controller)
    })
    .await
    .context("prepare persisted daemon sessions")??;
    // A local session an earlier release started from a profile home keeps
    // running from it until it is next staged. The link has to be in place
    // before any launch configuration is refreshed or any credential sync runs.
    {
        let state = controller.state.clone();
        tokio::task::spawn_blocking(move || {
            crate::controller::local_profile_homes::link_profile_homes_of_earlier_sessions(&state)
        })
        .await
        .context("earlier sessions' profile home link task failed")?;
    }
    // Earlier releases left each such session's project-memory replica in the
    // profile home, where nothing removed it when the session ended.
    {
        let config = config.clone();
        let state = controller.state.clone();
        tokio::task::spawn_blocking(move || {
            crate::controller::local_profile_homes::remove_replicas_left_in_profile_homes(
                &config,
                &state,
                &crate::targets::BoundedProcessExecutor::new(Duration::from_secs(15)),
            )
        })
        .await
        .context("leftover project-memory replica cleanup task failed")?;
    }
    let listener = TcpListener::bind(SocketAddr::new(IpAddr::V4(Ipv4Addr::LOCALHOST), 0))
        .await
        .context("bind Mjolnir daemon loopback endpoint")?;
    let metadata = DaemonMetadata {
        protocol_version: PROTOCOL_VERSION,
        pid: std::process::id(),
        address: listener.local_addr()?,
        token: random_hex::<32>()?,
        started_at: chrono::Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
        build_version: env!("CARGO_PKG_VERSION").to_owned(),
    };
    let workspaces = tokio::task::spawn_blocking(crate::database::list_workspaces)
        .await
        .context("daemon workspace load task panicked")??;
    let mut remote = if config.phone.enabled {
        Some(spawn_remote_session_manager()?)
    } else {
        None
    };

    // Start the primary manager after store preparation. Bootstrap failures
    // below have an awaited epilogue just like failures in the serving loop.
    let (delegation_tx, delegation_updates) = crate::session_manager::delegation_channel();
    let manager = crate::session_manager::spawn_session_manager_observed(Some(delegation_tx))?;
    let manager_targets = manager.targets;
    manager_targets.send_replace(dashboard_worker_targets(&controller));
    let manager_updates = manager.updates;
    let manager_control = manager.control.clone();
    let manager_shutdown = manager.shutdown;
    let mut recovery = crate::recovery::RecoveryCoordinator::spawn(manager_control.clone());
    let recovery_observer = recovery.observer();
    // Shares the recovery gate, so a recovery copy and a worker upgrade never
    // act on one session at the same time.
    let mut worker_upgrades = crate::worker_upgrade::WorkerUpgradeCoordinator::spawn(
        manager_control.clone(),
        &recovery_observer,
    );
    let state = Arc::new(RuntimeState::new(
        manager_control.clone(),
        Controller {
            config: controller.config.clone(),
            state: controller.state.clone(),
        },
        recovery_observer.clone(),
        worker_upgrades.observer(),
        workspaces,
    ));
    let cancellation = crate::termination::Coordinator::install().token();
    // Bootstrap already owns live managers. Capture its error so those owners
    // are shut down before the process-level writer can be closed.
    let bootstrap = async {
        // Restore review admission before prompts or external services start.
        state
            .review_host()
            .ready()
            .await
            .map_err(anyhow::Error::msg)?;
        let move_operations = blocking(crate::database::load_move_operations).await?;
        let move_sessions = move_operations
            .iter()
            .map(|operation| operation.selection.session_id.clone())
            .collect::<BTreeSet<_>>();
        let move_owned = state.recover_moves(move_operations)?;
        state.resume_retained_cleanups();
        state.restore_startup_deliveries(&cancellation).await?;
        Ok::<_, anyhow::Error>((move_sessions, move_owned))
    }
    .await;
    let (move_sessions, move_owned) = match bootstrap {
        Ok(ownership) => ownership,
        Err(error) => {
            epilogue_started.store(true, Ordering::Release);
            spawn_shutdown_watchdog();
            cancellation.cancel();
            let mut outcome = Err(error);
            record_daemon_cleanup(
                &mut outcome,
                "shut down bootstrap worker upgrade coordinator",
                worker_upgrades.shutdown().await,
            );
            record_daemon_cleanup(
                &mut outcome,
                "shut down bootstrap recovery coordinator",
                recovery.shutdown().await,
            );
            record_daemon_cleanup(
                &mut outcome,
                "shut down turn review host after bootstrap failure",
                state
                    .review_host()
                    .shutdown()
                    .await
                    .map_err(anyhow::Error::msg),
            );
            record_daemon_cleanup(
                &mut outcome,
                "cancel bootstrap lifecycle operations",
                state.cancel_and_wait_lifecycles().await,
            );
            record_daemon_cleanup(
                &mut outcome,
                "drain bootstrap startup deliveries",
                state.cancel_and_join_startup_prompts().await,
            );
            if let Some(remote) = remote.take() {
                record_daemon_cleanup(
                    &mut outcome,
                    "shut down bootstrap remote session manager",
                    remote.shutdown.shutdown().await,
                );
            }
            record_daemon_cleanup(
                &mut outcome,
                "shut down bootstrap session manager",
                manager_shutdown.shutdown().await,
            );
            return outcome;
        }
    };
    // Checkpoint files that a cancelled or interrupted checkpoint left in a
    // local worker root. Deleting them can take minutes after a long leak, so
    // this runs beside the daemon rather than before it serves. It stops at
    // shutdown, and the next start finishes it.
    let checkpoint_sweep = {
        let cancellation = cancellation.clone();
        tokio::task::spawn_blocking(move || {
            crate::controller::sweep_local_checkpoint_leftovers(
                &crate::targets::BoundedProcessExecutor::new(Duration::from_secs(15)),
                started,
                &|| cancellation.is_cancelled(),
            );
        })
    };
    let (mut manager_updates, continuation_task) =
        continuation::spawn(state.clone(), manager_updates, cancellation.clone());

    let (delegation_services, mut delegation_task) = super::delegation::spawn(
        state.clone(),
        manager_control.clone(),
        delegation_updates,
        cancellation.clone(),
    );

    let target_refresh = spawn_manager_target_refresher(
        manager_targets.clone(),
        cancellation.clone(),
        state.clone(),
    );
    let image_refresh = spawn_image_refresher(
        {
            let state = state.clone();
            move || state.with_config(crate::controller::image_refresh_plan)
        },
        {
            // A background download is the daemon's own work, not a session's,
            // so these notices carry an empty session id and reach every
            // workspace.
            let state = state.clone();
            move |report| {
                let text = match report {
                    crate::pollers::ImageRefreshReport::Started { host, image } => {
                        format!("Downloading image {image} for {host}\u{2026}")
                    }
                    crate::pollers::ImageRefreshReport::Pulled { host, image } => {
                        format!("Image {image} is ready on {host}.")
                    }
                    crate::pollers::ImageRefreshReport::Failed { host, image, error } => {
                        format!("Could not pull image {image} on {host}: {error}")
                    }
                };
                state.push_notice("", text);
            }
        },
        cancellation.clone(),
    );
    let exit_when_idle = mj_core::config::env_override_os("DAEMON_EXIT_WHEN_IDLE").is_some();
    let mut idle_tick = tokio::time::interval(Duration::from_millis(100));
    idle_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
    let mut owner_tick = tokio::time::interval(Duration::from_millis(500));
    owner_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
    let mut recovery_tick = tokio::time::interval(Duration::from_millis(250));
    recovery_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
    let mut background_policy_tick = tokio::time::interval(Duration::from_secs(1));
    background_policy_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
    // How often the harness readiness wait is looked at. The wait itself is
    // minutes long, so this only bounds how late a failure is noticed, and it
    // reads in-memory state rather than the store.
    let mut readiness_tick = tokio::time::interval(Duration::from_secs(5));
    readiness_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
    let (interrupted_close_tx, mut interrupted_close_rx) = tokio::sync::mpsc::unbounded_channel();
    let mut interrupted_close_tasks = Vec::new();
    for session_id in interrupted_suspend_session_ids(&controller) {
        if move_owned.contains(&session_id) {
            continue;
        }
        let recovery_state = state.clone();
        let recovery_shutdown = cancellation.clone();
        let updates = interrupted_close_tx.clone();
        let upgrade_work = startup_work.clone();
        let interrupted_close_task = tokio::spawn(async move {
            let _upgrade_work = upgrade_work;
            let result = tokio::select! {
                result = recovery_state.suspend_session(session_id.clone()) => result,
                () = recovery_shutdown.cancelled() => return,
            }
            .map(|()| crate::pollers::LifecycleSuccess::Closed)
            .map_err(|error| format!("{error:#}"));
            if updates
                .send(crate::pollers::LifecycleUpdate {
                    session_id,
                    result,
                    deferred_cleanup: false,
                })
                .is_err()
            {
                tracing::debug!("suspension recovery receiver stopped");
            }
        });
        interrupted_close_tasks.push(interrupted_close_task);
    }
    // Whatever is still in an in-flight lifecycle state now has no owner: the
    // moves, the interrupted closes, and the checkpointing rows above are
    // every operation that legitimately resumes. The list comes from the
    // startup snapshot, so a session created after this cannot be caught by
    // it, and the writes run off this path because they touch the database.
    let reconciliation = {
        let unowned = unowned_interrupted_lifecycles(
            &controller,
            &move_owned.union(&move_sessions).cloned().collect(),
        );
        (!unowned.is_empty()).then(|| {
            let state = state.clone();
            let upgrade_work = startup_work.clone();
            tokio::spawn(async move {
                let _upgrade_work = upgrade_work;
                let reconciled = tokio::task::spawn_blocking(move || {
                    let mut controller = Controller::load()?;
                    let mut reconciled = 0usize;
                    for (session_id, cause) in unowned {
                        match controller.fail_interrupted_lifecycle(&session_id, &cause) {
                            Ok(true) => {
                                tracing::warn!(%session_id, %cause, "session left in flight by a daemon restart marked failed");
                                reconciled += 1;
                            }
                            Ok(false) => {}
                            Err(error) => tracing::warn!(
                                %session_id,
                                error = format!("{error:#}"),
                                "could not reconcile an interrupted lifecycle state"
                            ),
                        }
                    }
                    anyhow::Ok(reconciled)
                })
                .await;
                match reconciled {
                    Ok(Ok(0)) => {}
                    Ok(Ok(_)) => refresh_runtime_controller(&state).await,
                    Ok(Err(error)) => tracing::warn!(
                        error = format!("{error:#}"),
                        "could not load the controller to reconcile interrupted lifecycles"
                    ),
                    Err(error) => {
                        tracing::warn!(%error, "interrupted lifecycle reconciliation task failed");
                    }
                }
            })
        })
    };
    // Rows left as tombstones by an older build, or by a discard that could
    // not finish before the daemon stopped. The list comes from the startup
    // snapshot, so a session that becomes lost after this is handled by the
    // view watcher instead.
    let tombstone_sweep = {
        let tombstones = tombstone_session_ids(&controller);
        (!tombstones.is_empty()).then(|| {
            let state = state.clone();
            let upgrade_work = startup_work.clone();
            tokio::spawn(async move {
                let _upgrade_work = upgrade_work;
                for session_id in tombstones {
                    state.discard_lost_session(session_id).await;
                }
            })
        })
    };
    let mut phone_publisher: Option<RemoteSessionPublisher> = None;
    let mut phone_targets = None;
    let mut prepared_targets = manager_targets.subscribe();
    let mut phone_task = None;
    let mut remote_request_bridge = None;
    if let Some(remote) = remote.take() {
        remote
            .targets
            .send_replace(prepared_targets.borrow_and_update().clone());
        phone_targets = Some(remote.targets.clone());
        phone_publisher = Some(remote.publisher.clone());
        remote_request_bridge = Some(spawn_remote_request_bridge(
            remote.requests,
            manager_control.clone(),
        ));
        phone_task = Some(spawn_phone_server(
            config.phone,
            cancellation.clone(),
            state.clone(),
            SessionManagerChannels {
                targets: remote.targets,
                control: remote.control,
                updates: remote.updates,
                shutdown: remote.shutdown,
            },
            delegation_services,
        ));
    } else {
        state.set_phone_status(WebViewerStatus::Disabled);
        state.web_viewer.publish(crate::server::WebViewerAccess::Unavailable("Web access is disabled. Enable [phone].enabled in your configuration, then restart the daemon.".into()));
    }
    let daemon_metadata_path = metadata_path();
    let mut client_tasks = tokio::task::JoinSet::new();
    let mut delegation_outcome = None;

    // Everything a client can use is initialized before this atomic
    // publication. From here on every exit, including an error from the test
    // hook or the event loop, flows through the same bounded epilogue.
    let mut outcome = async {
        write_metadata(&daemon_metadata_path, &metadata)?;
        reach_test_hook("daemon_metadata_before_listening").await?;
        drop(startup_work);
        loop {
            tokio::select! {
                _ = cancellation.cancelled() => break,
                changed = prepared_targets.changed(), if phone_targets.is_some() => {
                    changed.context("daemon worker target publication stopped")?;
                    if let Some(targets) = &phone_targets {
                        targets.send_replace(prepared_targets.borrow_and_update().clone());
                    }
                }
                result = &mut delegation_task => {
                    delegation_outcome = Some(result.map_err(anyhow::Error::from).and_then(|result| result));
                    anyhow::bail!("delegation coordinator stopped unexpectedly");
                }
                _ = idle_tick.tick(), if exit_when_idle && state.ever_attached.load(Ordering::Acquire) => {
                    state.prune_dead_clients();
                    if state.attachments().is_empty() {
                        break;
                    }
                }
                _ = owner_tick.tick(), if owner_pid.is_some() => {
                    if let Some(owner) = owner_pid
                        && !process_is_alive(owner)
                    {
                        tracing::info!(owner_pid = owner, "daemon owner process exited; shutting down");
                        break;
                    }
                }
                _ = recovery_tick.tick() => {
                    while let Some(result) = recovery.try_result() {
                        if let Err(error) = &result.outcome {
                            // A deferred copy found the agent working. That is
                            // the normal state of a session in use, so it is
                            // news, not a fault.
                            if result.deferred {
                                tracing::info!(session_id = %result.session_id, %error, "recovery copy deferred: agent is working");
                            } else {
                                tracing::warn!(session_id = %result.session_id, %error, "daemon recovery checkpoint failed");
                            }
                        }
                        refresh_runtime_controller(&state).await;
                    }
                    while let Some(result) = worker_upgrades.try_result() {
                        report_worker_upgrade(&state, &result);
                    }
                }
                _ = background_policy_tick.tick() => {
                    state.refresh_background_policies();
                }
                _ = readiness_tick.tick() => {
                    // A live session whose harness never advertised itself
                    // takes no prompt and reports no failure, so nothing else
                    // ever ends its wait (#1090). The store write runs off
                    // this loop.
                    for unready in state.sessions_without_a_usable_harness() {
                        let state = state.clone();
                        client_tasks.spawn(async move {
                            state.fail_unready_session(unready).await;
                        });
                    }
                }
                completed = interrupted_close_rx.recv() => {
                    if let Some(completed) = completed {
                        let recovered = completed.result.is_ok();
                        if let Err(error) = completed.result {
                            tracing::warn!(session_id = %completed.session_id, %error, "daemon could not resume interrupted close");
                        }
                        refresh_runtime_controller(&state).await;
                        if recovered && completed.deferred_cleanup
                            && let Err(error) = state.start_deferred_cleanup(completed.session_id.clone())
                        {
                            tracing::warn!(session_id = %completed.session_id, error = format!("{error:#}"), "could not continue cleanup after interrupted close");
                            state.push_notice(
                                &completed.session_id,
                                format!("Could not continue container storage cleanup: {error:#}"),
                            );
                        }
                    }
                }
                accepted = listener.accept() => {
                    let (stream, peer) = accepted.context("accept Mjolnir daemon client")?;
                    if !peer.ip().is_loopback() {
                        tracing::warn!(%peer, "rejected non-loopback daemon client");
                        continue;
                    }
                    let metadata = metadata.clone();
                    let state = state.clone();
                    let cancellation = cancellation.clone();
                    client_tasks.spawn(async move {
                        if let Err(error) = serve_client(stream, metadata, state, cancellation).await {
                            tracing::debug!(error = format!("{error:#}"), "daemon client disconnected");
                        }
                    });
                }
                completed = client_tasks.join_next(), if !client_tasks.is_empty() => {
                    if let Some(Err(error)) = completed {
                        tracing::warn!(%error, "daemon client task failed");
                    }
                }
                update = manager_updates.recv() => {
                    let Some(update) = update else {
                        bail!("controller daemon session manager stopped");
                    };
                    if let Some((detail, observed_updated_at)) =
                        state.missing_target_record(&update.session_id, &update.view)
                    {
                        let state = state.clone();
                        let session_id = update.session_id.clone();
                        client_tasks.spawn(async move {
                            if let Err(error) = state.persist_missing_target(
                                &session_id, detail, observed_updated_at,
                            ).await {
                                tracing::warn!(%session_id, %error, "could not persist missing worker target");
                                state.push_notice(&session_id, format!("Could not record missing session target: {error:#}"));
                            }
                        });
                    }
                    if let Some(publisher) = phone_publisher.as_ref()
                        && let Err(error) = publisher.try_publish(
                            update.session_id.clone(),
                            update.view.clone(),
                        )
                    {
                        tracing::warn!(%error, "phone session view bridge stopped");
                        phone_publisher = None;
                    }
                    // Every session's view passes here whether or not anything is
                    // attached, which is exactly what an automatic review needs to
                    // see: the turn that just finished.
                    // The continuation completion gate has already notified review.
                    state.publish_session(update.session_id, update.view).await?;
                }
            }
        }
        Ok(())
    }
    .await;

    epilogue_started.store(true, Ordering::Release);
    spawn_shutdown_watchdog();
    // Idle exit and fallible loop exits do not arrive through the termination
    // coordinator. Stop every daemon-owned task before closing the sole writer.
    cancellation.cancel();
    // Recovery copies and worker upgrades do not hold up a handoff, so some
    // may still be running. Stop them first: the next daemon starts them again.
    // The coordinators share one gate, and dropping either cancels both.
    record_daemon_cleanup(
        &mut outcome,
        "shut down worker upgrade coordinator",
        worker_upgrades.shutdown().await,
    );
    record_daemon_cleanup(
        &mut outcome,
        "shut down recovery coordinator",
        recovery.shutdown().await,
    );
    match continuation_task.await {
        Ok(Ok(())) => {}
        Ok(Err(error)) => tracing::error!(%error, "continuation service failed"),
        Err(error) => tracing::error!(%error, "continuation service task failed"),
    }
    record_daemon_cleanup(
        &mut outcome,
        "join delegation coordinator",
        match delegation_outcome {
            Some(result) => result,
            None => delegation_task
                .await
                .map_err(anyhow::Error::from)
                .and_then(|result| result),
        },
    );
    drop(interrupted_close_tx);
    record_daemon_cleanup(
        &mut outcome,
        "remove daemon metadata",
        remove_daemon_metadata(&daemon_metadata_path),
    );
    record_daemon_cleanup(
        &mut outcome,
        "shut down turn review host",
        state
            .review_host()
            .shutdown()
            .await
            .map_err(anyhow::Error::msg),
    );
    record_daemon_cleanup(
        &mut outcome,
        "join controller target refresher",
        target_refresh.await.map_err(anyhow::Error::new),
    );
    record_daemon_cleanup(
        &mut outcome,
        "join container image refresher",
        image_refresh.await.map_err(anyhow::Error::new),
    );
    if let Some(phone_task) = phone_task {
        record_daemon_cleanup(
            &mut outcome,
            "join phone server",
            phone_task.await.map_err(anyhow::Error::new),
        );
    }
    if let Some(remote_request_bridge) = remote_request_bridge {
        record_daemon_cleanup(
            &mut outcome,
            "join phone session request bridge",
            remote_request_bridge.await.map_err(anyhow::Error::new),
        );
    }
    client_tasks.abort_all();
    while let Some(result) = client_tasks.join_next().await {
        if let Err(error) = result
            && !error.is_cancelled()
        {
            record_daemon_cleanup(
                &mut outcome,
                "join daemon client task",
                Err(anyhow::Error::new(error)),
            );
        }
    }
    record_daemon_cleanup(
        &mut outcome,
        "cancel daemon lifecycle operations",
        state.cancel_and_wait_lifecycles().await,
    );
    record_daemon_cleanup(
        &mut outcome,
        "drain startup prompts",
        state.cancel_and_join_startup_prompts().await,
    );
    if let Some(reconciliation) = reconciliation {
        record_daemon_cleanup(
            &mut outcome,
            "join interrupted lifecycle reconciliation",
            reconciliation.await.map_err(anyhow::Error::new),
        );
    }
    if let Some(tombstone_sweep) = tombstone_sweep {
        record_daemon_cleanup(
            &mut outcome,
            "join lost-session discard sweep",
            tombstone_sweep.await.map_err(anyhow::Error::new),
        );
    }
    record_daemon_cleanup(
        &mut outcome,
        "join checkpoint leftover sweep",
        checkpoint_sweep.await.map_err(anyhow::Error::new),
    );
    for interrupted_close_task in interrupted_close_tasks {
        record_daemon_cleanup(
            &mut outcome,
            "join interrupted close recovery",
            interrupted_close_task.await.map_err(anyhow::Error::new),
        );
    }
    record_daemon_cleanup(
        &mut outcome,
        "shut down controller daemon session manager",
        manager_shutdown.shutdown().await,
    );
    outcome
}

/// Sessions whose record is nothing but a tombstone: they ended without a
/// target and without a checkpoint, so the record offers no action but its own
/// removal. `DestroyedWithDataLoss` only ever comes from an older build.
pub(super) fn tombstone_session_ids(controller: &Controller) -> Vec<String> {
    controller
        .state
        .sessions
        .values()
        .filter(|session| {
            matches!(
                session.state,
                SessionState::Lost | SessionState::DestroyedWithDataLoss
            )
        })
        .map(|session| session.id.clone())
        .collect()
}

pub(super) fn remove_daemon_metadata(path: &Path) -> Result<()> {
    match fs::remove_file(path) {
        Ok(()) => Ok(()),
        Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(()),
        Err(error) => Err(error).with_context(|| format!("remove {}", path.display())),
    }
}

/// Keep the event-loop failure as the primary result while still running and
/// reporting every cleanup step. If the loop ended normally, the first
/// cleanup failure becomes the daemon's result.
pub(super) fn record_daemon_cleanup(
    outcome: &mut Result<()>,
    operation: &'static str,
    cleanup: Result<()>,
) {
    let Err(error) = cleanup else {
        return;
    };
    let error = error.context(operation);
    if outcome.is_ok() {
        *outcome = Err(error);
    } else {
        tracing::warn!(error = format!("{error:#}"), "daemon cleanup step failed");
    }
}

/// Bounds the epilogue below.
///
/// The daemon leaves on its own long before this fires: a graceful exit
/// returns from `run_daemon_process`, the process exits 0, and this task dies
/// with the runtime. It exists so no unwinding step can hold the process open
/// past the deadline its clients wait on, whatever the cause of the shutdown.
pub(super) fn spawn_shutdown_watchdog() {
    tokio::spawn(async move {
        tokio::time::sleep(SHUTDOWN_FORCE_EXIT_TIMEOUT).await;
        tracing::error!(
            seconds = SHUTDOWN_FORCE_EXIT_TIMEOUT.as_secs(),
            "daemon shutdown did not finish in time; exiting"
        );
        // The metadata file points clients at a process that is about to stop
        // answering. Removing it is what the epilogue would have done.
        if let Err(error) = fs::remove_file(metadata_path())
            && error.kind() != std::io::ErrorKind::NotFound
        {
            tracing::warn!(%error, "could not remove daemon metadata before the forced exit");
        }
        // 128 + signal is reserved for exits that really were signalled.
        std::process::exit(1);
    });
}

pub(super) fn spawn_manager_target_refresher(
    targets: tokio::sync::watch::Sender<Vec<crate::session_manager::RelaySessionTarget>>,
    cancellation: CancellationToken,
    state: Arc<RuntimeState>,
) -> tokio::task::JoinHandle<()> {
    tokio::spawn(async move {
        let mut committed = state
            .committed
            .as_ref()
            .expect("daemon owns the writer")
            .clone();
        let mut revisions = state.revisions();
        let mut compatibility = tokio::time::interval(Duration::from_millis(500));
        compatibility.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
        let mut installed = None;
        loop {
            let inputs = {
                let owner = state.owner();
                if let Err(error) = owner.ensure_available() {
                    tracing::error!(%error, "daemon durable state is unavailable; shutting down");
                    cancellation.cancel();
                    return;
                }
                owner.pollable_worker_inputs()
            };
            if installed.as_ref() != Some(&inputs) {
                let preparation = inputs.clone();
                let refreshed =
                    match tokio::task::spawn_blocking(move || preparation.prepare()).await {
                        Ok(targets) => targets,
                        Err(error) => {
                            tracing::error!(%error, "worker target preparation failed");
                            cancellation.cancel();
                            return;
                        }
                    };
                let retained = refreshed
                    .iter()
                    .map(|target| target.session_id.clone())
                    .collect();
                {
                    let owner = state.owner();
                    // A lifecycle may have claimed a target while its commands
                    // were being prepared. Only the owner can authorize install.
                    if owner.pollable_worker_inputs() != inputs {
                        continue;
                    }
                    targets.send_if_modified(|current| {
                        if *current == refreshed {
                            false
                        } else {
                            *current = refreshed;
                            true
                        }
                    });
                }
                state.review_host().retain_sessions(retained);
                installed = Some(inputs);
            }
            tokio::select! {
                _ = cancellation.cancelled() => return,
                changed = committed.changed() => {
                    if changed.is_err() {
                        tracing::error!("database publication feed stopped");
                        cancellation.cancel();
                        return;
                    }
                    state.publish_revision();
                }
                changed = revisions.changed() => {
                    if changed.is_err() { return; }
                }
                _ = compatibility.tick() => {
                    let _config_mutation = state.config_mutation.lock().await;
                    let refreshed = tokio::task::spawn_blocking(|| {
                        crate::database::check_read_compatibility()?;
                        Config::load()
                    }).await;
                    match refreshed {
                        Ok(Ok(config)) => {
                            state.review_config.lock().unwrap_or_else(PoisonError::into_inner)
                                .clone_from(&config.review);
                            let changed = {
                                let mut owner = state.owner();
                                let changed = owner.controller().config != config;
                                owner.install_config(config);
                                changed
                            };
                            if changed { state.publish_revision(); }
                        }
                        Ok(Err(error)) => {
                            if error.chain().any(|cause| cause.downcast_ref::<StoreSchemaMismatch>().is_some()) {
                                tracing::error!(%error, "daemon store schema diverged underneath the daemon; shutting down");
                                cancellation.cancel();
                                return;
                            }
                            tracing::warn!(error = format!("{error:#}"), "could not refresh daemon configuration");
                        }
                        Err(error) => {
                            tracing::error!(%error, "daemon configuration reader failed");
                            cancellation.cancel();
                            return;
                        }
                    }
                }
            }
        }
    })
}

pub(super) async fn refresh_runtime_controller(state: &RuntimeState) {
    if let Err(error) = state.reload_controller().await {
        tracing::warn!(
            error = format!("{error:#}"),
            "could not refresh daemon controller state"
        );
    }
}