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