use super::diagnostics::ServingPhase;
use super::*;
use std::time::Instant;
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<()> {
let owner_pid = owner_pid_to_watch()?;
let opening = Instant::now();
let guard = ControllerStoreGuard::acquire()?;
let database_writer = guard.start_database_writer()?;
log_startup_phase("open the store", opening);
let epilogue_started = AtomicBool::new(false);
let mut outcome = run_daemon_runtime(&epilogue_started, owner_pid).await;
if !epilogue_started.load(Ordering::Acquire) {
spawn_shutdown_watchdog();
}
let writer_started = Instant::now();
let writer_shutdown = tokio::task::spawn_blocking(move || database_writer.shutdown())
.await
.context("database writer shutdown task panicked")
.and_then(std::convert::identity);
let writer_took = writer_started.elapsed();
if writer_took >= Duration::from_millis(250) {
tracing::info!(
duration_ms = writer_took.as_millis(),
"database writer shut down"
);
}
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<()> {
let progress = super::diagnostics::DaemonProgressMonitor::start()?;
let started = SystemTime::now();
let startup_work = crate::upgrade::activity("daemon startup recovery")?;
tokio::task::spawn_blocking(crate::controller::pin_worker_binary_sources)
.await
.context("worker source snapshot task failed")??;
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 loading = Instant::now();
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")??;
log_startup_phase("load persisted sessions", loading);
{
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")?;
}
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: mj_client::build_identity::this_build().published(),
};
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
};
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();
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,
));
state.owner().install_relay_sessions(
manager_targets
.borrow()
.iter()
.map(|target| target.session_id.clone())
.collect(),
);
let cancellation = crate::termination::Coordinator::install().token();
let bootstrapping = Instant::now();
let bootstrap = async {
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 mut lifecycle_owned = state.recover_moves(move_operations)?;
let restart_intents = blocking(crate::database::load_session_restarts).await?;
lifecycle_owned.extend(
restart_intents
.iter()
.map(|intent| intent.session_id.clone()),
);
lifecycle_owned.extend(state.recover_restarts(restart_intents)?);
lifecycle_owned.extend(move_sessions.iter().cloned());
state.resume_retained_cleanups();
state.resume_startup_cleanups(true);
state.restore_startup_deliveries(&cancellation).await?;
Ok::<_, anyhow::Error>((move_sessions, lifecycle_owned))
}
.await;
log_startup_phase("recover moves and startup deliveries", bootstrapping);
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;
}
};
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 update_feed = 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 cpu_publication = {
let mut receiver = manager.session_cpu;
let state = state.clone();
let cancellation = cancellation.clone();
tokio::spawn(async move {
let mut next_publication = tokio::time::Instant::now();
loop {
tokio::select! {
_ = cancellation.cancelled() => return,
changed = receiver.changed() => {
if changed.is_err() {
tracing::error!("session CPU publication channel closed");
return;
}
}
}
tokio::select! {
_ = cancellation.cancelled() => return,
_ = tokio::time::sleep_until(next_publication) => {}
}
receiver.borrow_and_update();
state.publish_revision();
next_publication = tokio::time::Instant::now() + Duration::from_secs(10);
}
})
};
let project_catalog = tokio::spawn(state.projects().run(cancellation.child_token()));
let capacity_service =
super::capacity::spawn_capacity_service(state.clone(), 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)
},
{
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);
let mut readiness_tick = tokio::time::interval(Duration::from_secs(5));
readiness_tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
let mut progress_tick = tokio::time::interval(Duration::from_secs(1));
progress_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);
}
for session_id in interrupted_destroy_session_ids(&controller) {
let recovery_state = state.clone();
let recovery_shutdown = cancellation.clone();
let updates = interrupted_close_tx.clone();
let upgrade_work = startup_work.clone();
interrupted_close_tasks.push(tokio::spawn(async move {
let _upgrade_work = upgrade_work;
let result = tokio::select! {
result = recovery_state.force_destroy_session(
session_id.clone(),
crate::daemon::BranchDisposition::Keep,
) => result,
() = recovery_shutdown.cancelled() => return,
}
.map(|()| crate::pollers::LifecycleSuccess::ForceDestroyed)
.map_err(|error| format!("{error:#}"));
let _ = updates.send(crate::pollers::LifecycleUpdate {
session_id,
result,
deferred_cleanup: false,
});
}));
}
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");
}
}
})
})
};
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 {
session_cpu: remote.control.session_cpu.clone(),
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;
let mut cache_task = tokio::spawn(crate::controller::mbx::service::run(
state.clone(),
cancellation.clone(),
));
let mut cache_outcome = None;
let mut outcome = async {
write_metadata(&daemon_metadata_path, &metadata)?;
reach_test_hook("daemon_metadata_before_listening").await?;
tracing::info!("the daemon is serving");
drop(startup_work);
loop {
progress.serving_tick();
tokio::select! {
_ = cancellation.cancelled() => break,
_ = progress_tick.tick() => {},
changed = prepared_targets.changed(), if phone_targets.is_some() => {
progress.phase(ServingPhase::Targets);
changed.context("daemon worker target publication stopped")?;
if let Some(targets) = &phone_targets {
targets.send_replace(prepared_targets.borrow_and_update().clone());
}
}
result = &mut cache_task => {
cache_outcome = Some(result.context("machine cache service failed").and_then(|result| result));
anyhow::bail!("machine build cache service stopped unexpectedly");
}
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) => {
progress.phase(ServingPhase::IdleCheck);
state.prune_dead_clients();
if state.attachments().is_empty() {
break;
}
}
_ = owner_tick.tick(), if owner_pid.is_some() => {
progress.phase(ServingPhase::OwnerCheck);
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() => {
progress.phase(ServingPhase::Recovery);
while let Some(result) = recovery.try_result() {
if let Err(error) = &result.outcome {
if result.deferred {
tracing::info!(session_id = %result.session_id, %error, "recovery copy deferred: agent is working");
} else if result.cancelled {
tracing::info!(session_id = %result.session_id, %error, reason = "superseded by suspend", "daemon recovery checkpoint cancelled");
} 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).await;
}
}
_ = background_policy_tick.tick() => {
progress.phase(ServingPhase::BackgroundPolicy);
state.refresh_background_policies();
state.resume_startup_cleanups(false);
}
_ = readiness_tick.tick() => {
progress.phase(ServingPhase::Readiness);
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() => {
progress.phase(ServingPhase::SuspensionRecovery);
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() => {
progress.phase(ServingPhase::Client);
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 = update_feed.next() => {
progress.phase(ServingPhase::SessionUpdate);
let Some(update) = update? else { break };
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;
}
state.publish_session(update.session_id, update.view).await?;
}
}
}
Ok(())
}
.await;
progress.phase(ServingPhase::Shutdown);
let mut epilogue = EpilogueClock::start();
epilogue_started.store(true, Ordering::Release);
spawn_shutdown_watchdog();
cancellation.cancel();
epilogue.record(
&mut outcome,
"join machine build cache applications",
match cache_outcome {
Some(result) => result,
None => cache_task
.await
.context("machine cache service failed")
.and_then(|result| result),
},
);
epilogue.record(
&mut outcome,
"shut down worker upgrade coordinator",
worker_upgrades.shutdown().await,
);
epilogue.record(
&mut outcome,
"shut down recovery coordinator",
recovery.shutdown().await,
);
epilogue.record(
&mut outcome,
"join continuation service",
update_feed.join().await,
);
epilogue.record(
&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);
epilogue.record(
&mut outcome,
"remove daemon metadata",
remove_daemon_metadata(&daemon_metadata_path),
);
epilogue.record(
&mut outcome,
"shut down turn review host",
state
.review_host()
.shutdown()
.await
.map_err(anyhow::Error::msg),
);
epilogue.record(
&mut outcome,
"join CPU publication",
cpu_publication.await.map_err(anyhow::Error::from),
);
record_daemon_cleanup(
&mut outcome,
"join project discovery",
project_catalog.await.map_err(anyhow::Error::from),
);
epilogue.record(
&mut outcome,
"join controller target refresher",
target_refresh.await.map_err(anyhow::Error::new),
);
epilogue.record(
&mut outcome,
"join container image refresher",
image_refresh.await.map_err(anyhow::Error::new),
);
epilogue.record(
&mut outcome,
"join capacity service",
capacity_service.await.map_err(anyhow::Error::new),
);
if let Some(phone_task) = phone_task {
epilogue.record(
&mut outcome,
"join phone server",
phone_task.await.map_err(anyhow::Error::new),
);
}
if let Some(remote_request_bridge) = remote_request_bridge {
epilogue.record(
&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)),
);
}
}
epilogue.record(&mut outcome, "join daemon client tasks", Ok(()));
epilogue.record(
&mut outcome,
"cancel daemon lifecycle operations",
state.cancel_and_wait_lifecycles().await,
);
epilogue.record(
&mut outcome,
"drain startup prompts",
state.cancel_and_join_startup_prompts().await,
);
if let Some(reconciliation) = reconciliation {
epilogue.record(
&mut outcome,
"join interrupted lifecycle reconciliation",
reconciliation.await.map_err(anyhow::Error::new),
);
}
if let Some(tombstone_sweep) = tombstone_sweep {
epilogue.record(
&mut outcome,
"join lost-session discard sweep",
tombstone_sweep.await.map_err(anyhow::Error::new),
);
}
epilogue.record(
&mut outcome,
"join checkpoint leftover sweep",
checkpoint_sweep.await.map_err(anyhow::Error::new),
);
for interrupted_close_task in interrupted_close_tasks {
epilogue.record(
&mut outcome,
"join interrupted close recovery",
interrupted_close_task.await.map_err(anyhow::Error::new),
);
}
epilogue.record(
&mut outcome,
"shut down controller daemon session manager",
manager_shutdown.shutdown().await,
);
epilogue.finish();
outcome
}
fn log_startup_phase(phase: &'static str, started: Instant) {
let took = started.elapsed();
if took >= EpilogueClock::LOGGED_STEP {
tracing::info!(
phase,
duration_ms = took.as_millis(),
"daemon startup phase finished"
);
}
}
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())),
}
}
struct EpilogueClock {
began: Instant,
step: Instant,
}
impl EpilogueClock {
const LOGGED_STEP: Duration = Duration::from_millis(250);
fn start() -> Self {
let now = Instant::now();
tracing::info!("daemon shutdown began");
Self {
began: now,
step: now,
}
}
fn record(&mut self, outcome: &mut Result<()>, operation: &'static str, cleanup: Result<()>) {
let took = self.step.elapsed();
if took >= Self::LOGGED_STEP {
tracing::info!(
step = operation,
duration_ms = took.as_millis(),
"daemon shutdown step finished"
);
}
record_daemon_cleanup(outcome, operation, cleanup);
self.step = Instant::now();
}
fn finish(self) {
tracing::info!(
duration_ms = self.began.elapsed().as_millis(),
"daemon shutdown finished"
);
}
}
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");
}
}
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"
);
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");
}
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;
let mut credential_installed = None;
let mut auth = crate::pollers::ProviderAuthCache::default();
let mut auth_profiles = None;
let mut staged_sources = 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()
};
let config = inputs.controller().config;
if auth_profiles.as_ref() != Some(&config.profiles) {
let mut cache = auth.clone();
let profiles = config.profiles.clone();
match tokio::task::spawn_blocking(move || {
cache.refresh(&config);
cache
})
.await
{
Ok(cache) => {
auth = cache;
auth_profiles = Some(profiles);
}
Err(error) => {
tracing::error!(%error, "provider refresh task stopped");
cancellation.cancel();
return;
}
}
}
if !installed
.as_ref()
.is_some_and(|current| Arc::ptr_eq(current, &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::<BTreeSet<_>>();
let views_dropped = {
let mut owner = state.owner();
if !Arc::ptr_eq(&owner.pollable_worker_inputs(), &inputs) {
continue;
}
targets.send_if_modified(|current| {
if *current == refreshed {
false
} else {
*current = refreshed;
true
}
});
owner.install_relay_sessions(retained.clone())
};
if views_dropped {
state.publish_revision();
}
state.review_host().retain_sessions(retained);
installed = Some(inputs.clone());
}
let closing = state.owner().close_requested.clone();
if !staged_sources
.as_ref()
.is_some_and(|(current, _)| Arc::ptr_eq(current, &inputs))
{
let staging = inputs.clone();
match tokio::task::spawn_blocking(move || staging.staged_credentials()).await {
Ok(staged) => staged_sources = Some((inputs.clone(), staged)),
Err(error) => {
tracing::error!(%error, "credential staging observation stopped");
cancellation.cancel();
return;
}
}
}
let staged = &staged_sources
.as_ref()
.expect("credential sources observed")
.1;
let credential_key = (
inputs.clone(),
auth.schemes.clone(),
closing.clone(),
staged.clone(),
);
if !credential_installed.as_ref().is_some_and(
|(current, schemes, closed, previous_staged)| {
Arc::ptr_eq(current, &inputs)
&& *schemes == auth.schemes
&& *closed == closing
&& previous_staged == staged
},
) {
let preparation = inputs.clone();
let schemes = auth.schemes.clone();
let staged = staged.clone();
let prepared = match tokio::task::spawn_blocking(move || {
preparation.prepare_credentials(&schemes, &staged)
})
.await
{
Ok(targets) => targets,
Err(error) => {
tracing::error!(%error, "credential target preparation stopped");
cancellation.cancel();
return;
}
};
if !state.owner().install_credentials(
&inputs,
&closing,
prepared,
&state.credential_targets,
) {
continue;
}
credential_installed = Some(credential_key);
}
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 mut cache = auth.clone();
let staging = inputs.clone();
let refreshed = tokio::task::spawn_blocking(move || {
crate::database::check_read_compatibility()?;
let config = Config::load()?;
cache.refresh(&config);
let staged = staging.staged_credentials();
Ok::<_, anyhow::Error>((config, cache, staged))
}).await;
match refreshed {
Ok(Ok((config, cache, staged))) => {
staged_sources = Some((inputs.clone(), staged));
auth = cache;
auth_profiles = Some(config.profiles.clone());
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"
);
}
}