Skip to main content

scv_server/
restart.rs

1//! Planned restarts into a newly installed release, and the notices around
2//! them.
3//!
4//! `scv restart --when-idle` (run by the feature-flow deploy script after
5//! `cargo install`) asks the daemon to restart into the binary now at its own
6//! path. The daemon checks that binary runs, then waits until the delegation
7//! that asked has finished (for a live child that serves a whole
8//! conversation, until its turn has ended and, for a nested SCV, its own
9//! background jobs have been reported to it) and its report is stored, and no
10//! owner message is being answered, or until the request's deadline. It then
11//! records a plan, keeps a copy of its own binary for rollback, and starts a
12//! watchdog outside its own cgroup (`systemd-run`). The watchdog restarts the
13//! unit, checks that the new release comes up with the channels that were
14//! connected before, and otherwise puts the previous binary back when both
15//! releases share a config layout. The daemon that starts next announces the
16//! outcome in the chat that asked, or through the notify list.
17//!
18//! The same notifier tells the owner about restarts after a crash and about
19//! accounts that stay disconnected.
20
21use std::{
22    collections::HashMap,
23    path::{Path, PathBuf},
24    sync::{Arc, LazyLock, Mutex as SyncMutex, PoisonError, Weak, atomic::AtomicBool},
25    time::Duration,
26};
27
28use anyhow::{Context, Result, anyhow, bail};
29use scv_channels::hub::{Hub, Origin, Restart};
30use scv_client::Layout;
31use scv_protocol::{ComponentState, DaemonCommand, RestartInfo};
32use scv_tools::{background::BackgroundJobs, delegation::DelegationRegistry};
33use serde::{Deserialize, Serialize};
34use tokio::sync::Mutex;
35use tokio_util::sync::CancellationToken;
36
37use crate::components::Components;
38use crate::config::Instance;
39
40/// Where configuration and state files live and how they are shaped. Bump it
41/// when a release reads or writes them in a way the previous release cannot:
42/// a rollback between releases with different layouts is refused.
43pub(crate) const CONFIG_LAYOUT: u32 = 1;
44
45const DEFAULT_MAX_WAIT: u64 = 10 * 60;
46const MAX_WAIT_LIMIT: u64 = 60 * 60;
47/// How long the watchdog gives a new release to report its version and
48/// reconnect the channels that were connected before.
49const VERIFY_SECONDS: u64 = 180;
50/// How long the watchdog waits for a rolled-back release to come back.
51const ROLLBACK_SECONDS: u64 = 90;
52/// Checks a restart must pass in a row before it goes ahead, a second apart,
53/// so a job that just finished has time to start its report.
54const CLEAR_CHECKS: u32 = 2;
55/// An account disconnected this long gets a notice through another account.
56const DOWN_NOTICE_AFTER: Duration = Duration::from_secs(10 * 60);
57const MONITOR_INTERVAL: Duration = Duration::from_secs(30);
58/// A plan restarted this long ago no longer explains interrupted work.
59const RESTART_CONTEXT_MAX_AGE: u64 = 60 * 60;
60/// How long a restart that goes ahead waits for mail actions under way.
61const MAIL_DRAIN: Duration = Duration::from_secs(60);
62
63/// What a binary reports about itself for a planned restart.
64#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
65pub struct BuildInfo {
66    pub(crate) version: String,
67    pub(crate) config_layout: u32,
68}
69
70/// This binary's build information, printed by `scv build-info`.
71pub fn build_info() -> BuildInfo {
72    BuildInfo {
73        version: env!("CARGO_PKG_VERSION").into(),
74        config_layout: CONFIG_LAYOUT,
75    }
76}
77
78#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
79#[serde(rename_all = "snake_case")]
80pub(crate) enum PlanState {
81    /// Waiting for the requesting work to end.
82    Waiting,
83    /// The watchdog is restarting the unit and checking the new release.
84    Restarting,
85    /// The new release came up with its channels.
86    Verified,
87    /// The new release failed and the previous binary was put back.
88    RolledBack,
89    /// The new release failed and was not rolled back, or the restart could
90    /// not start.
91    Failed,
92}
93
94/// The delegation that asked for a restart.
95#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
96pub(crate) struct Requester {
97    pub(crate) handle: String,
98    pub(crate) session: String,
99}
100
101/// A planned restart, saved in `<home>/state/update.json` (mode 0600) and
102/// shared by the daemon that plans it, the watchdog, and the next daemon.
103#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
104pub(crate) struct Plan {
105    pub(crate) id: String,
106    pub(crate) state: PlanState,
107    pub(crate) from_version: String,
108    pub(crate) to_version: String,
109    #[serde(default, skip_serializing_if = "Option::is_none")]
110    pub(crate) commit: Option<String>,
111    pub(crate) from_layout: u32,
112    pub(crate) to_layout: u32,
113    pub(crate) unit: String,
114    /// The daemon's executable, where the new release was installed.
115    pub(crate) binary: PathBuf,
116    /// A copy of the release the daemon ran, for rollback.
117    #[serde(default, skip_serializing_if = "Option::is_none")]
118    pub(crate) previous: Option<PathBuf>,
119    #[serde(default, skip_serializing_if = "Option::is_none")]
120    pub(crate) requester: Option<Requester>,
121    /// The chat that asked, which hears the outcome.
122    #[serde(default, skip_serializing_if = "Option::is_none")]
123    pub(crate) origin: Option<Origin>,
124    /// Accounts connected when the restart went ahead; the new release
125    /// must reconnect them.
126    #[serde(default, skip_serializing_if = "Vec::is_empty")]
127    pub(crate) expected: Vec<String>,
128    pub(crate) requested_unix: u64,
129    pub(crate) deadline_unix: u64,
130    #[serde(default, skip_serializing_if = "Option::is_none")]
131    pub(crate) restart_unix: Option<u64>,
132    /// The restart went ahead at the deadline while work still ran.
133    #[serde(default, skip_serializing_if = "std::ops::Not::not")]
134    pub(crate) waited_out: bool,
135    /// Why the new release failed, for the announcement.
136    #[serde(default, skip_serializing_if = "Option::is_none")]
137    pub(crate) detail: Option<String>,
138    /// How long the watchdog gives the new release.
139    #[serde(default = "default_verify_seconds")]
140    pub(crate) verify_seconds: u64,
141}
142
143fn default_verify_seconds() -> u64 {
144    VERIFY_SECONDS
145}
146
147impl Plan {
148    fn info(&self, waiting_for: Option<String>) -> RestartInfo {
149        RestartInfo {
150            to_version: self.to_version.clone(),
151            waiting_for,
152            requester: self.requester.as_ref().map(|r| r.handle.clone()),
153            origin: self.origin.as_ref().map(|origin| origin.component.clone()),
154            deadline_unix_seconds: self.deadline_unix,
155        }
156    }
157
158    fn label(&self) -> String {
159        match &self.commit {
160            Some(commit) => format!("v{} ({commit})", self.to_version),
161            None => format!("v{}", self.to_version),
162        }
163    }
164}
165
166pub(crate) fn load_plan(path: &Path) -> Result<Option<Plan>> {
167    match std::fs::read(path) {
168        Ok(bytes) => {
169            Ok(Some(serde_json::from_slice(&bytes).with_context(|| {
170                format!("parse restart plan {}", path.display())
171            })?))
172        }
173        Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(None),
174        Err(error) => Err(error).with_context(|| format!("read {}", path.display())),
175    }
176}
177
178pub(crate) fn save_plan(path: &Path, plan: &Plan) -> Result<()> {
179    write_private(path, &serde_json::to_vec_pretty(plan)?)
180}
181
182fn write_private(path: &Path, bytes: &[u8]) -> Result<()> {
183    let parent = path
184        .parent()
185        .ok_or_else(|| anyhow!("{} has no parent", path.display()))?;
186    std::fs::create_dir_all(parent).with_context(|| format!("create {}", parent.display()))?;
187    scv_client::fs::replace_private(path, bytes)
188        .with_context(|| format!("write {}", path.display()))
189}
190
191fn unix_now() -> u64 {
192    std::time::SystemTime::now()
193        .duration_since(std::time::UNIX_EPOCH)
194        .map_or(0, |elapsed| elapsed.as_secs())
195}
196
197/// The daemon's executable path. Linux names an executable that was
198/// replaced on disk `<path> (deleted)`; the path is where the new one is.
199fn own_executable() -> Result<PathBuf> {
200    std::env::current_exe()
201        .map(strip_deleted)
202        .context("locate the daemon's executable")
203}
204
205fn strip_deleted(path: PathBuf) -> PathBuf {
206    match path
207        .to_str()
208        .and_then(|text| text.strip_suffix(" (deleted)"))
209    {
210        Some(stripped) => PathBuf::from(stripped),
211        None => path,
212    }
213}
214
215/// Whether this process runs in `unit`'s cgroup.
216fn runs_as_unit(unit: &str) -> bool {
217    let suffix = format!("/{unit}");
218    std::fs::read_to_string("/proc/self/cgroup")
219        .is_ok_and(|text| text.lines().any(|line| line.ends_with(&suffix)))
220}
221
222/// Run `binary build-info` and parse what it reports.
223async fn probe(binary: &Path) -> Result<BuildInfo> {
224    let output = tokio::time::timeout(
225        Duration::from_secs(10),
226        tokio::process::Command::new(binary)
227            .arg("build-info")
228            .stdin(std::process::Stdio::null())
229            .kill_on_drop(true)
230            .output(),
231    )
232    .await
233    .map_err(|_| anyhow!("it did not answer within 10 seconds"))??;
234    if !output.status.success() {
235        bail!("it exited with {}", output.status);
236    }
237    serde_json::from_slice(&output.stdout).context("it printed no build information")
238}
239
240// ---------------------------------------------------------------------------
241// Sessions' own activity, which a restart waits out for the session that
242// asked (its report turn, or a turn the TUI started).
243
244struct SessionActivity {
245    busy: AtomicBool,
246    background: Option<Weak<BackgroundJobs>>,
247}
248
249static SESSIONS: LazyLock<SyncMutex<HashMap<String, Arc<SessionActivity>>>> =
250    LazyLock::new(Default::default);
251
252/// A daemon session's entry in the activity table while it lives.
253pub(crate) struct SessionTracker {
254    id: String,
255    activity: Arc<SessionActivity>,
256}
257
258impl SessionTracker {
259    pub(crate) fn new(id: &str, background: Option<&Arc<BackgroundJobs>>) -> Self {
260        let activity = Arc::new(SessionActivity {
261            busy: AtomicBool::new(false),
262            background: background.map(Arc::downgrade),
263        });
264        SESSIONS
265            .lock()
266            .unwrap_or_else(PoisonError::into_inner)
267            .insert(id.to_owned(), Arc::clone(&activity));
268        Self {
269            id: id.to_owned(),
270            activity,
271        }
272    }
273
274    /// A turn runs, or a finished background job waits for its report turn.
275    pub(crate) fn set_busy(&self, busy: bool) {
276        self.activity
277            .busy
278            .store(busy, std::sync::atomic::Ordering::Release);
279    }
280}
281
282impl Drop for SessionTracker {
283    fn drop(&mut self) {
284        SESSIONS
285            .lock()
286            .unwrap_or_else(PoisonError::into_inner)
287            .remove(&self.id);
288    }
289}
290
291/// Whether session `id` has work a planned restart waits out: a turn, or a
292/// background job still running or not yet reported (one whose report turn
293/// failed and waits to be tried again included).
294fn session_busy(id: &str) -> bool {
295    let activity = SESSIONS
296        .lock()
297        .unwrap_or_else(PoisonError::into_inner)
298        .get(id)
299        .cloned();
300    activity.is_some_and(|activity| {
301        activity.busy.load(std::sync::atomic::Ordering::Acquire)
302            || activity
303                .background
304                .as_ref()
305                .and_then(Weak::upgrade)
306                .is_some_and(|jobs| jobs.pending() > 0)
307    })
308}
309
310// ---------------------------------------------------------------------------
311// Notices: where a message nobody asked for goes.
312
313/// A place a notice may go: an account, and the chat partner there, or the
314/// account's owner.
315#[derive(Debug, Clone, PartialEq, Eq)]
316struct Candidate {
317    component: String,
318    peer: Option<String>,
319}
320
321#[derive(Debug, Clone, PartialEq, Eq)]
322enum Pick {
323    Send {
324        component: String,
325        peer: String,
326    },
327    /// An earlier candidate may still connect.
328    Wait,
329    Nothing,
330}
331
332/// The first candidate that is connected and whose peer is known. Until the
333/// grace period is over, a candidate that may still connect keeps its place
334/// ahead of later ones.
335fn pick(
336    candidates: &[Candidate],
337    states: &HashMap<String, ComponentState>,
338    owner: &dyn Fn(&str) -> Option<Option<String>>,
339    exclude: Option<&str>,
340    grace_over: bool,
341) -> Pick {
342    for candidate in candidates {
343        if exclude == Some(candidate.component.as_str()) {
344            continue;
345        }
346        let registered = owner(&candidate.component);
347        let peer = candidate
348            .peer
349            .clone()
350            .or_else(|| registered.clone().flatten());
351        match (states.get(&candidate.component), &registered, peer) {
352            (Some(ComponentState::Connected), Some(_), Some(peer)) => {
353                return Pick::Send {
354                    component: candidate.component.clone(),
355                    peer,
356                };
357            }
358            (
359                Some(
360                    ComponentState::Starting
361                    | ComponentState::Connected
362                    | ComponentState::Disconnected
363                    | ComponentState::Backoff,
364                ),
365                _,
366                _,
367            ) if !grace_over => return Pick::Wait,
368            _ => {}
369        }
370    }
371    Pick::Nothing
372}
373
374/// The human name of a component's channel.
375fn channel_title(component: &str) -> &str {
376    match component.split(':').next() {
377        Some("wechat") => "WeChat",
378        Some("feishu") => "Feishu",
379        Some("email") => "Email",
380        Some(other) => other,
381        None => component,
382    }
383}
384
385/// Where component states come from.
386#[derive(Clone)]
387enum States {
388    Components(Weak<Mutex<Components>>),
389    #[cfg(test)]
390    Fixed(Arc<SyncMutex<HashMap<String, ComponentState>>>),
391}
392
393impl States {
394    async fn get(&self) -> HashMap<String, ComponentState> {
395        match self {
396            Self::Components(components) => match components.upgrade() {
397                Some(components) => components
398                    .lock()
399                    .await
400                    .status()
401                    .components
402                    .into_iter()
403                    .filter(|health| health.enabled)
404                    .map(|health| (health.id, health.state))
405                    .collect(),
406                None => HashMap::new(),
407            },
408            #[cfg(test)]
409            Self::Fixed(states) => states.lock().unwrap().clone(),
410        }
411    }
412}
413
414/// Sends notices to the owner through the hub.
415#[derive(Clone)]
416pub(crate) struct Notifier {
417    hub: Arc<Hub>,
418    states: States,
419    /// How long an account ahead in line may take to connect.
420    grace: Duration,
421    /// When an undeliverable notice is dropped.
422    give_up: Duration,
423    poll: Duration,
424    /// Where the notify list is configured.
425    instance: Instance,
426    /// The notify list; `None` reads it from the user configuration.
427    #[cfg(test)]
428    list: Option<Vec<String>>,
429}
430
431impl Notifier {
432    pub(crate) fn new(
433        instance: Instance,
434        hub: Arc<Hub>,
435        components: Weak<Mutex<Components>>,
436    ) -> Self {
437        Self {
438            hub,
439            instance,
440            states: States::Components(components),
441            grace: Duration::from_secs(120),
442            give_up: Duration::from_secs(15 * 60),
443            poll: Duration::from_secs(2),
444            #[cfg(test)]
445            list: None,
446        }
447    }
448
449    /// A notifier that sees `states` and the notify `list`, for other
450    /// modules' tests.
451    #[cfg(test)]
452    pub(crate) fn fixed(
453        hub: &Arc<Hub>,
454        list: Vec<String>,
455        states: HashMap<String, ComponentState>,
456    ) -> Self {
457        Self {
458            hub: Arc::clone(hub),
459            states: States::Fixed(Arc::new(SyncMutex::new(states))),
460            grace: Duration::from_millis(200),
461            give_up: Duration::from_secs(5),
462            poll: Duration::from_millis(20),
463            instance: crate::test_support::test_instance("/unused"),
464            list: Some(list),
465        }
466    }
467
468    fn notify_list(&self) -> Vec<String> {
469        #[cfg(test)]
470        if let Some(list) = &self.list {
471            return list.clone();
472        }
473        self.instance.load_user().map_or_else(
474            |error| {
475                tracing::warn!(
476                    "Notices use the owner's last chat; configuration failed: {error:#}"
477                );
478                Vec::new()
479            },
480            |config| config.notify.owner,
481        )
482    }
483
484    /// The owner chat a notice would go to right now, without waiting for
485    /// an account to connect: the first connected notify target, or else
486    /// the chat the owner last wrote from.
487    pub(crate) async fn owner_chat(&self) -> Option<Origin> {
488        let states = self.states.get().await;
489        let owner = |component: &str| self.hub.owner(component);
490        match pick(&self.candidates(), &states, &owner, None, true) {
491            Pick::Send { component, peer } => Some(Origin { component, peer }),
492            Pick::Wait | Pick::Nothing => None,
493        }
494    }
495
496    /// The notify list, or else the chat the owner last wrote from; never a
497    /// mail chat, which carries only mail, and never a mailbox.
498    fn candidates(&self) -> Vec<Candidate> {
499        let list = self.notify_list();
500        let candidates: Vec<Candidate> = if list.is_empty() {
501            self.hub
502                .last_owner()
503                .map(|last| Candidate {
504                    component: last.component,
505                    peer: Some(last.peer),
506                })
507                .into_iter()
508                .collect()
509        } else {
510            list.into_iter()
511                .map(|component| Candidate {
512                    component,
513                    peer: None,
514                })
515                .collect()
516        };
517        candidates
518            .into_iter()
519            .filter(|candidate| {
520                !self.hub.is_mail_chat(&candidate.component)
521                    && !candidate.component.starts_with("email:")
522            })
523            .collect()
524    }
525
526    /// Store `text` for the `origin` chat, or, when it is not given or does
527    /// not connect in time, for the first reachable notify target other than
528    /// `exclude`. Returns where it went.
529    pub(crate) async fn deliver(
530        &self,
531        origin: Option<&Origin>,
532        text: &str,
533        exclude: Option<&str>,
534        cancel: &CancellationToken,
535    ) -> Option<String> {
536        let started = tokio::time::Instant::now();
537        let fallback = self.candidates();
538        // The asking chat alone first; the notify targets once its grace is
539        // over, saying why the answer comes there.
540        let mut phase = match origin {
541            Some(origin) => (
542                vec![Candidate {
543                    component: origin.component.clone(),
544                    peer: Some(origin.peer.clone()),
545                }],
546                None,
547                text.to_owned(),
548            ),
549            None => (fallback.clone(), exclude, text.to_owned()),
550        };
551        let mut phase_started = started;
552        loop {
553            let states = self.states.get().await;
554            let grace_over = phase_started.elapsed() >= self.grace;
555            if let Some(origin) = origin
556                && grace_over
557                && phase.1.is_none()
558            {
559                phase = (
560                    fallback.clone(),
561                    Some(origin.component.as_str()),
562                    format!(
563                        "(You asked on {}, which is not connected, so this comes here.) {text}",
564                        channel_title(&origin.component)
565                    ),
566                );
567                phase_started = tokio::time::Instant::now();
568                continue;
569            }
570            let (candidates, exclude, text) = &phase;
571            let owner = |component: &str| self.hub.owner(component);
572            match pick(candidates, &states, &owner, *exclude, grace_over) {
573                Pick::Send { component, peer } => {
574                    match self.hub.notify(&component, &peer, text).await {
575                        Ok(()) => return Some(component),
576                        Err(error) => tracing::warn!("Notice to {component} not stored: {error}"),
577                    }
578                }
579                Pick::Nothing if grace_over => {
580                    tracing::warn!("No connected account can take this notice: {text}");
581                    return None;
582                }
583                Pick::Wait | Pick::Nothing => {}
584            }
585            if started.elapsed() >= self.give_up {
586                tracing::warn!("Gave up delivering a notice: {text}");
587                return None;
588            }
589            tokio::select! {
590                () = cancel.cancelled() => return None,
591                () = tokio::time::sleep(self.poll) => {}
592            }
593        }
594    }
595}
596
597// ---------------------------------------------------------------------------
598// The daemon side: requests, waiting, and handing over to the watchdog.
599
600/// How the restart is carried out once it may go ahead.
601enum Launcher {
602    /// A watchdog unit started with `systemd-run`.
603    Systemd,
604    /// Tests record the plan instead.
605    #[cfg(test)]
606    Record(Arc<SyncMutex<Vec<Plan>>>),
607}
608
609/// Plans restarts for the daemon.
610pub(crate) struct Restarter {
611    launcher: Launcher,
612    instance: Instance,
613    hub: Arc<Hub>,
614    registry: Arc<DelegationRegistry>,
615    notifier: Notifier,
616    components: Weak<Mutex<Components>>,
617    cancel: CancellationToken,
618    /// The plan being waited on or carried out, and what it waits for.
619    current: SyncMutex<Option<(Plan, Option<String>)>>,
620}
621
622impl Restarter {
623    pub(crate) fn new(
624        instance: Instance,
625        hub: Arc<Hub>,
626        registry: Arc<DelegationRegistry>,
627        components: &Arc<Mutex<Components>>,
628        cancel: CancellationToken,
629    ) -> Arc<Self> {
630        Arc::new(Self {
631            launcher: Launcher::Systemd,
632            notifier: Notifier::new(
633                instance.clone(),
634                Arc::clone(&hub),
635                Arc::downgrade(components),
636            ),
637            instance,
638            hub,
639            registry,
640            components: Arc::downgrade(components),
641            cancel,
642            current: SyncMutex::new(None),
643        })
644    }
645
646    pub(crate) fn notifier(&self) -> &Notifier {
647        &self.notifier
648    }
649
650    /// The scheduled restart, for status replies.
651    pub(crate) fn info(&self) -> Option<RestartInfo> {
652        self.current
653            .lock()
654            .unwrap_or_else(PoisonError::into_inner)
655            .as_ref()
656            .map(|(plan, waiting)| plan.info(waiting.clone()))
657    }
658
659    /// Handle `restart_when_idle`. The error is shown to the caller.
660    pub(crate) async fn request(
661        self: &Arc<Self>,
662        command: DaemonCommand,
663    ) -> std::result::Result<RestartInfo, String> {
664        let DaemonCommand::RestartWhenIdle {
665            version,
666            commit,
667            parent,
668            max_wait_seconds,
669        } = command
670        else {
671            return Err("not a restart request".into());
672        };
673        if let Some(info) = self.info() {
674            return if version.as_deref().is_none_or(|v| v == info.to_version) {
675                Ok(info)
676            } else {
677                Err(format!(
678                    "a restart into v{} is already scheduled",
679                    info.to_version
680                ))
681            };
682        }
683        let unit = self.instance.layout.service_name();
684        if !runs_as_unit(&unit) {
685            return Err(format!(
686                "this daemon does not run as {unit}, so it cannot restart itself; \
687                 restart it yourself"
688            ));
689        }
690        let binary = own_executable().map_err(|error| format!("{error:#}"))?;
691        let installed = probe(&binary).await.map_err(|error| {
692            format!(
693                "the binary at {} does not run ({error:#}); not restarting",
694                binary.display()
695            )
696        })?;
697        if let Some(version) = &version
698            && version != &installed.version
699        {
700            return Err(format!(
701                "{} reports v{}, not v{version}; not restarting",
702                binary.display(),
703                installed.version
704            ));
705        }
706        let requester = parent.as_deref().and_then(|chain| self.requester(chain));
707        let now = unix_now();
708        let wait = max_wait_seconds
709            .unwrap_or(DEFAULT_MAX_WAIT)
710            .min(MAX_WAIT_LIMIT);
711        let plan = Plan {
712            id: uuid::Uuid::new_v4().simple().to_string()[..8].to_owned(),
713            state: PlanState::Waiting,
714            from_version: env!("CARGO_PKG_VERSION").into(),
715            to_version: installed.version,
716            commit: commit.filter(|commit| !commit.trim().is_empty()),
717            from_layout: CONFIG_LAYOUT,
718            to_layout: installed.config_layout,
719            unit,
720            previous: None,
721            binary,
722            requester,
723            origin: None,
724            expected: Vec::new(),
725            requested_unix: now,
726            deadline_unix: now + wait,
727            restart_unix: None,
728            waited_out: false,
729            detail: None,
730            verify_seconds: VERIFY_SECONDS,
731        };
732        // An owner confirmation step would go here, before the plan is armed.
733        self.arm(plan)
734    }
735
736    /// Save `plan` and wait for it in the background.
737    fn arm(self: &Arc<Self>, mut plan: Plan) -> std::result::Result<RestartInfo, String> {
738        plan.origin = plan
739            .requester
740            .as_ref()
741            .and_then(|requester| self.hub.origin(&requester.session));
742        save_plan(&self.instance.layout.update_plan(), &plan)
743            .map_err(|error| format!("{error:#}"))?;
744        let waiting = self.waiting_for(&plan);
745        let info = plan.info(waiting.clone());
746        *self.current.lock().unwrap_or_else(PoisonError::into_inner) =
747            Some((plan.clone(), waiting));
748        tracing::info!(
749            "Restart into v{} scheduled; waiting at most {} seconds",
750            plan.to_version,
751            plan.deadline_unix.saturating_sub(plan.requested_unix)
752        );
753        let restarter = Arc::clone(self);
754        tokio::spawn(async move { restarter.wait_and_restart(plan).await });
755        Ok(info)
756    }
757
758    /// The delegation of this daemon named in a `SCV_PARENT` chain.
759    fn requester(&self, chain: &str) -> Option<Requester> {
760        self.registry.own_run(chain).map(|run| Requester {
761            handle: run.handle,
762            session: run.session,
763        })
764    }
765
766    /// What the restart still waits for, or `None` when it may go ahead.
767    fn waiting_for(&self, plan: &Plan) -> Option<String> {
768        if let Some(requester) = &plan.requester {
769            // A per-turn run works while its processes live; a live child
770            // (nested SCV, ACP agent) keeps its process for the whole
771            // conversation, so it works only while a turn runs on it or,
772            // for a nested SCV, background jobs of its own still run or wait
773            // to be reported to it.
774            let working = self
775                .registry
776                .list(true)
777                .into_iter()
778                .find(|entry| entry.record.handle == requester.handle && entry.working());
779            if let Some(entry) = working {
780                return Some(if entry.record.idle_since_unix.is_some() {
781                    format!("{}'s background jobs", requester.handle)
782                } else {
783                    format!("{} to finish", requester.handle)
784                });
785            }
786            if session_busy(&requester.session) || self.hub.session_work(&requester.session) > 0 {
787                return Some(format!("{}'s report", requester.handle));
788            }
789        }
790        if self.hub.owner_claims() > 0 {
791            return Some("an owner message to be answered".into());
792        }
793        if self.hub.mail_executing() > 0 {
794            return Some("a mail action to finish".into());
795        }
796        None
797    }
798
799    async fn wait_and_restart(self: Arc<Self>, mut plan: Plan) {
800        let mut clear = 0;
801        loop {
802            tokio::select! {
803                // The daemon is stopping: the next one finds the plan waiting.
804                () = self.cancel.cancelled() => return,
805                () = tokio::time::sleep(Duration::from_secs(1)) => {}
806            }
807            let waiting = self.waiting_for(&plan);
808            clear = if waiting.is_none() { clear + 1 } else { 0 };
809            if let Some((_, current)) = self
810                .current
811                .lock()
812                .unwrap_or_else(PoisonError::into_inner)
813                .as_mut()
814            {
815                current.clone_from(&waiting);
816            }
817            if clear >= CLEAR_CHECKS {
818                break;
819            }
820            if unix_now() >= plan.deadline_unix {
821                tracing::warn!(
822                    "Restarting into v{} at its deadline while waiting for {}",
823                    plan.to_version,
824                    waiting.as_deref().unwrap_or("work")
825                );
826                plan.waited_out = true;
827                break;
828            }
829        }
830        // No mail action starts from here on; one under way finishes first,
831        // for at most a minute.
832        self.hub.set_mail_drain(true);
833        let drained = tokio::time::Instant::now() + MAIL_DRAIN;
834        while self.hub.mail_executing() > 0 && tokio::time::Instant::now() < drained {
835            tokio::time::sleep(Duration::from_millis(200)).await;
836        }
837        if let Err(error) = self.hand_over(&mut plan).await {
838            self.hub.set_mail_drain(false);
839            tracing::error!("Restart into v{} did not start: {error:#}", plan.to_version);
840            plan.state = PlanState::Failed;
841            plan.detail = Some(format!("the restart did not start: {error:#}"));
842            let _ = save_plan(&self.instance.layout.update_plan(), &plan);
843            *self.current.lock().unwrap_or_else(PoisonError::into_inner) = None;
844            let text = announcement(&plan, env!("CARGO_PKG_VERSION"));
845            self.notifier
846                .deliver(plan.origin.as_ref(), &text, None, &self.cancel)
847                .await;
848            let _ = std::fs::remove_file(self.instance.layout.update_plan());
849        }
850    }
851
852    /// Record the plan as restarting, keep this release's binary, and start
853    /// the watchdog that restarts the unit.
854    async fn hand_over(&self, plan: &mut Plan) -> Result<()> {
855        plan.state = PlanState::Restarting;
856        plan.restart_unix = Some(unix_now());
857        if let Some(components) = self.components.upgrade() {
858            plan.expected = components
859                .lock()
860                .await
861                .status()
862                .components
863                .into_iter()
864                .filter(|health| health.enabled && health.state == ComponentState::Connected)
865                .map(|health| health.id)
866                .collect();
867        }
868        match &self.launcher {
869            Launcher::Systemd => {}
870            #[cfg(test)]
871            Launcher::Record(plans) => {
872                save_plan(&self.instance.layout.update_plan(), plan)?;
873                plans.lock().unwrap().push(plan.clone());
874                return Ok(());
875            }
876        }
877        plan.previous = match keep_previous(&plan.binary) {
878            Ok(path) => Some(path),
879            Err(error) => {
880                tracing::warn!("No rollback copy of this release: {error:#}");
881                None
882            }
883        };
884        let path = self.instance.layout.update_plan();
885        save_plan(&path, plan)?;
886        // The watchdog runs the release known to work: this one.
887        let watchdog = plan.previous.clone().unwrap_or_else(|| plan.binary.clone());
888        let mut command = std::process::Command::new("systemd-run");
889        command.args([
890            "--user",
891            "--quiet",
892            "--collect",
893            &format!("--unit=scv-update-{}", plan.id),
894        ]);
895        // The watchdog selects the same instance and configuration.
896        let layout = &self.instance.layout;
897        let home = (!layout.is_default()).then(|| layout.home());
898        let config = self.instance.overrides.config_file.as_deref();
899        for (variable, value) in [("SCV_HOME", home), ("SCV_CONFIG", config)] {
900            if let Some(value) = value {
901                let mut setting = std::ffi::OsString::from(format!("--setenv={variable}="));
902                setting.push(value);
903                command.arg(setting);
904            }
905        }
906        command
907            .arg(watchdog)
908            .arg("restart-watchdog")
909            .arg("--plan")
910            .arg(&path)
911            .stdin(std::process::Stdio::null());
912        let status = tokio::task::spawn_blocking(move || command.status())
913            .await?
914            .context("run systemd-run")?;
915        if !status.success() {
916            bail!("systemd-run exited with {status}");
917        }
918        tracing::info!(
919            "Handed the restart into v{} to unit scv-update-{}",
920            plan.to_version,
921            plan.id
922        );
923        Ok(())
924    }
925}
926
927/// Copy the running executable (still readable through `/proc/self/exe`
928/// after it was replaced on disk) next to `binary` as `<binary>.prev`.
929fn keep_previous(binary: &Path) -> Result<PathBuf> {
930    let previous = binary.with_file_name(format!(
931        "{}.prev",
932        binary
933            .file_name()
934            .and_then(|name| name.to_str())
935            .unwrap_or("scv")
936    ));
937    install_copy(Path::new("/proc/self/exe"), &previous)?;
938    Ok(previous)
939}
940
941/// Copy `source` to `target` through a temporary file beside it, executable.
942fn install_copy(source: &Path, target: &Path) -> Result<()> {
943    let parent = target
944        .parent()
945        .ok_or_else(|| anyhow!("{} has no parent", target.display()))?;
946    let temporary = tempfile::Builder::new()
947        .prefix(".scv-install")
948        .tempfile_in(parent)?;
949    std::fs::copy(source, temporary.path())
950        .with_context(|| format!("copy {} to {}", source.display(), target.display()))?;
951    #[cfg(unix)]
952    {
953        use std::os::unix::fs::PermissionsExt;
954        std::fs::set_permissions(temporary.path(), std::fs::Permissions::from_mode(0o755))?;
955    }
956    temporary
957        .persist(target)
958        .map_err(|error| error.error)
959        .with_context(|| format!("install {}", target.display()))?;
960    Ok(())
961}
962
963// ---------------------------------------------------------------------------
964// The watchdog, run by `scv restart-watchdog` outside the daemon.
965
966/// Restart the unit, check the new release, and roll back when it fails and
967/// the releases share a config layout. Records the outcome in the plan.
968pub async fn watchdog(layout: &Layout, plan_path: &Path) -> Result<()> {
969    let mut plan = load_plan(plan_path)?.context("no restart plan")?;
970    if plan.state != PlanState::Restarting {
971        bail!("the restart plan is {:?}, not restarting", plan.state);
972    }
973    let socket = layout.socket();
974    eprintln!("Restarting {} into v{}", plan.unit, plan.to_version);
975    systemctl_restart(&plan.unit);
976    let outcome = verify(
977        &socket,
978        &plan.to_version,
979        &plan.expected,
980        plan.verify_seconds,
981    )
982    .await;
983    match outcome {
984        Ok(()) => {
985            eprintln!("v{} is up with its channels", plan.to_version);
986            plan.state = PlanState::Verified;
987        }
988        Err(reason) => {
989            eprintln!("v{} failed: {reason}", plan.to_version);
990            match rollback_refusal(&plan) {
991                None => {
992                    let previous = plan.previous.clone().expect("checked by rollback_refusal");
993                    let detail = match install_copy(&previous, &plan.binary) {
994                        Ok(()) => {
995                            systemctl_restart(&plan.unit);
996                            let seconds = plan.verify_seconds.min(ROLLBACK_SECONDS);
997                            match verify(&socket, &plan.from_version, &[], seconds).await {
998                                Ok(()) => reason,
999                                Err(again) => format!(
1000                                    "{reason}; after the rollback v{} did not come back either ({again})",
1001                                    plan.from_version
1002                                ),
1003                            }
1004                        }
1005                        Err(error) => format!(
1006                            "{reason}; putting v{} back failed: {error:#}",
1007                            plan.from_version
1008                        ),
1009                    };
1010                    plan.state = PlanState::RolledBack;
1011                    plan.detail = Some(detail);
1012                }
1013                Some(refusal) => {
1014                    plan.state = PlanState::Failed;
1015                    plan.detail = Some(format!("{reason}; not rolled back: {refusal}"));
1016                }
1017            }
1018        }
1019    }
1020    save_plan(plan_path, &plan)?;
1021    Ok(())
1022}
1023
1024fn systemctl_restart(unit: &str) {
1025    match std::process::Command::new("systemctl")
1026        .args(["--user", "restart", unit])
1027        .status()
1028    {
1029        Ok(status) if status.success() => {}
1030        Ok(status) => eprintln!("systemctl --user restart {unit} exited with {status}"),
1031        Err(error) => eprintln!("could not run systemctl: {error}"),
1032    }
1033}
1034
1035/// Why the previous binary may not be put back, or `None` when it may.
1036fn rollback_refusal(plan: &Plan) -> Option<String> {
1037    if plan.to_layout != plan.from_layout {
1038        return Some(format!(
1039            "v{} uses config layout {} and v{} uses {}, so the older binary cannot read the \
1040             current configuration",
1041            plan.to_version, plan.to_layout, plan.from_version, plan.from_layout
1042        ));
1043    }
1044    match &plan.previous {
1045        Some(previous) if previous.is_file() => None,
1046        _ => Some(format!("no copy of v{} was kept", plan.from_version)),
1047    }
1048}
1049
1050/// Wait until the daemon reports `version` and every `expected` account is
1051/// connected, or explain what was missing when `seconds` run out.
1052async fn verify(
1053    socket: &Path,
1054    version: &str,
1055    expected: &[String],
1056    seconds: u64,
1057) -> std::result::Result<(), String> {
1058    let deadline = tokio::time::Instant::now() + Duration::from_secs(seconds);
1059    let mut last = format!("v{version} did not start");
1060    loop {
1061        match scv_client::control(socket, DaemonCommand::Status).await {
1062            Ok(status) if status.version == version => {
1063                let missing: Vec<_> = expected
1064                    .iter()
1065                    .filter(|id| {
1066                        !status.components.iter().any(|health| {
1067                            &health.id == *id && health.state == ComponentState::Connected
1068                        })
1069                    })
1070                    .map(String::as_str)
1071                    .collect();
1072                if missing.is_empty() {
1073                    return Ok(());
1074                }
1075                last = format!(
1076                    "v{version} started, but {} did not reconnect",
1077                    missing.join(" and ")
1078                );
1079            }
1080            Ok(status) => last = format!("SCV still reports v{}", status.version),
1081            Err(_) => {}
1082        }
1083        if tokio::time::Instant::now() >= deadline {
1084            return Err(last);
1085        }
1086        tokio::time::sleep(Duration::from_secs(2)).await;
1087    }
1088}
1089
1090// ---------------------------------------------------------------------------
1091// Startup: explain the previous run, then announce.
1092
1093/// What the daemon found at startup about how its predecessor ended.
1094pub(crate) struct Startup {
1095    plan: Option<Plan>,
1096    /// The previous daemon stopped without shutting down: its version and
1097    /// start time.
1098    unclean: Option<(String, u64)>,
1099}
1100
1101#[derive(Serialize, Deserialize)]
1102struct Marker {
1103    pid: u32,
1104    version: String,
1105    started_unix: u64,
1106}
1107
1108/// Read the restart plan and the running marker, tell the hub whether this
1109/// start is a planned restart (before any bridge recovers), and mark this
1110/// daemon running until [`clean_shutdown`].
1111pub(crate) fn startup(layout: &Layout, hub: &Hub) -> Startup {
1112    let plan = load_plan(&layout.update_plan()).unwrap_or_else(|error| {
1113        tracing::warn!("Ignoring an unreadable restart plan: {error:#}");
1114        let _ = std::fs::remove_file(layout.update_plan());
1115        None
1116    });
1117    let planned = plan.as_ref().filter(|plan| {
1118        plan.state != PlanState::Waiting
1119            && plan
1120                .restart_unix
1121                .is_some_and(|at| unix_now().saturating_sub(at) < RESTART_CONTEXT_MAX_AGE)
1122    });
1123    hub.set_restart(planned.map(|plan| Restart {
1124        to_version: plan.to_version.clone(),
1125    }));
1126    let marker = layout.daemon_marker();
1127    let unclean = std::fs::read(&marker)
1128        .ok()
1129        .and_then(|bytes| serde_json::from_slice::<Marker>(&bytes).ok())
1130        .filter(|previous| previous.pid != std::process::id())
1131        .map(|previous| (previous.version, previous.started_unix));
1132    let current = Marker {
1133        pid: std::process::id(),
1134        version: env!("CARGO_PKG_VERSION").into(),
1135        started_unix: unix_now(),
1136    };
1137    if let Err(error) = serde_json::to_vec(&current)
1138        .map_err(anyhow::Error::from)
1139        .and_then(|bytes| write_private(&marker, &bytes))
1140    {
1141        tracing::warn!("Could not record the running daemon: {error:#}");
1142    }
1143    Startup { plan, unclean }
1144}
1145
1146/// The daemon stopped on request: the next one will not report a crash.
1147pub(crate) fn clean_shutdown(layout: &Layout) {
1148    let _ = std::fs::remove_file(layout.daemon_marker());
1149}
1150
1151/// What the next daemon should say about a plan, given its own version.
1152#[derive(Debug, PartialEq, Eq)]
1153enum Decision {
1154    Say(String),
1155    /// The watchdog is still deciding.
1156    Wait,
1157    Drop,
1158}
1159
1160fn decide(plan: &Plan, own: &str, watchdog_overdue: bool) -> Decision {
1161    match plan.state {
1162        PlanState::Waiting => Decision::Say(if own == plan.to_version {
1163            format!(
1164                "SCV is now running {}. It stopped before the planned restart, so work that \
1165                 was running then was stopped.",
1166                plan.label()
1167            )
1168        } else {
1169            format!(
1170                "SCV stopped before it could restart into v{}; it is running v{own}. Deploy \
1171                 again to finish the update.",
1172                plan.to_version
1173            )
1174        }),
1175        PlanState::Restarting if !watchdog_overdue => Decision::Wait,
1176        PlanState::Restarting if own == plan.to_version => Decision::Say(format!(
1177            "SCV is now running {}; the update watchdog did not report back.",
1178            plan.label()
1179        )),
1180        PlanState::Restarting if own == plan.from_version => Decision::Say(format!(
1181            "The update to v{} did not take effect; SCV is still running v{own}.",
1182            plan.to_version
1183        )),
1184        PlanState::Restarting => Decision::Drop,
1185        PlanState::Verified | PlanState::RolledBack | PlanState::Failed => {
1186            Decision::Say(announcement(plan, own))
1187        }
1188    }
1189}
1190
1191/// The announcement of a finished plan.
1192fn announcement(plan: &Plan, own: &str) -> String {
1193    let detail = plan.detail.as_deref().unwrap_or("it did not come up");
1194    let mut text = match plan.state {
1195        PlanState::Verified => format!("SCV updated: now running {}.", plan.label()),
1196        PlanState::RolledBack => format!(
1197            "The update to v{} failed: {detail}. SCV rolled back to v{}.",
1198            plan.to_version, plan.from_version
1199        ),
1200        PlanState::Failed if own == plan.to_version => {
1201            format!("SCV is running {}, but {detail}.", plan.label())
1202        }
1203        _ => format!("The update to v{} failed: {detail}.", plan.to_version),
1204    };
1205    if plan.waited_out {
1206        let minutes = plan
1207            .deadline_unix
1208            .saturating_sub(plan.requested_unix)
1209            .div_ceil(60);
1210        text.push_str(&format!(
1211            " It waited {minutes} minutes for running work, then restarted anyway; work \
1212             still running then was stopped."
1213        ));
1214    }
1215    text
1216}
1217
1218/// Announce how the previous run ended, once the accounts can take it.
1219pub(crate) async fn announce(
1220    layout: Layout,
1221    startup: Startup,
1222    notifier: Notifier,
1223    cancel: CancellationToken,
1224) {
1225    let own = env!("CARGO_PKG_VERSION");
1226    if let Some(mut plan) = startup.plan {
1227        let path = layout.update_plan();
1228        let overdue_at = plan.restart_unix.unwrap_or(plan.requested_unix)
1229            + plan.verify_seconds
1230            + ROLLBACK_SECONDS
1231            + 60;
1232        let text = loop {
1233            match decide(&plan, own, unix_now() >= overdue_at) {
1234                Decision::Say(text) => break Some(text),
1235                Decision::Drop => break None,
1236                Decision::Wait => {}
1237            }
1238            tokio::select! {
1239                () = cancel.cancelled() => return,
1240                () = tokio::time::sleep(Duration::from_secs(2)) => {}
1241            }
1242            match load_plan(&path) {
1243                Ok(Some(reloaded)) if reloaded.id == plan.id => plan = reloaded,
1244                _ => break None,
1245            }
1246        };
1247        if let Some(text) = text {
1248            tracing::info!("{text}");
1249            notifier
1250                .deliver(plan.origin.as_ref(), &text, None, &cancel)
1251                .await;
1252        }
1253        if !cancel.is_cancelled() {
1254            let _ = std::fs::remove_file(&path);
1255        }
1256    } else if let Some((version, started)) = startup.unclean {
1257        let text = format!(
1258            "SCV started again after an unexpected stop (a crash or a host restart); it had run \
1259             v{version} since {}. Work in progress then was stopped.",
1260            format_time(started)
1261        );
1262        tracing::warn!("{text}");
1263        notifier.deliver(None, &text, None, &cancel).await;
1264    }
1265}
1266
1267fn format_time(unix: u64) -> String {
1268    let age = unix_now().saturating_sub(unix);
1269    match age {
1270        0..=119 => "moments before".into(),
1271        120..=7199 => format!("{} minutes before", age / 60),
1272        7200..=172_799 => format!("{} hours before", age / 3600),
1273        _ => format!("{} days before", age / 86_400),
1274    }
1275}
1276
1277// ---------------------------------------------------------------------------
1278// Accounts that stay disconnected.
1279
1280/// Tell the owner, through another account, when an enabled account stays
1281/// disconnected for [`DOWN_NOTICE_AFTER`]; once per outage.
1282pub(crate) async fn monitor(notifier: Notifier, cancel: CancellationToken) {
1283    let mut down: HashMap<String, (tokio::time::Instant, bool)> = HashMap::new();
1284    loop {
1285        tokio::select! {
1286            () = cancel.cancelled() => return,
1287            () = tokio::time::sleep(MONITOR_INTERVAL) => {}
1288        }
1289        let states = notifier.states.get().await;
1290        down.retain(|id, _| {
1291            states
1292                .get(id)
1293                .is_some_and(|state| *state != ComponentState::Connected)
1294        });
1295        for (id, state) in &states {
1296            if *state == ComponentState::Connected {
1297                continue;
1298            }
1299            let (since, told) = down
1300                .entry(id.clone())
1301                .or_insert((tokio::time::Instant::now(), false));
1302            if *told || since.elapsed() < DOWN_NOTICE_AFTER {
1303                continue;
1304            }
1305            *told = true;
1306            let (channel, account) = id.split_once(':').unwrap_or((id, "default"));
1307            let text = format!(
1308                "SCV's {} account {account} has been disconnected for {} minutes; its sign-in \
1309                 may have expired. On the host, check `scv channels status {channel}` and sign \
1310                 in again with `scv channels login {channel}` if needed.",
1311                channel_title(id),
1312                since.elapsed().as_secs() / 60
1313            );
1314            tracing::warn!("{text}");
1315            let notifier = notifier.clone();
1316            let cancel = cancel.clone();
1317            let id = id.clone();
1318            tokio::spawn(async move { notifier.deliver(None, &text, Some(&id), &cancel).await });
1319        }
1320    }
1321}
1322
1323#[cfg(test)]
1324mod tests;