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;
60const MAIL_DRAIN: Duration = Duration::from_secs(60);
62
63#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
65pub struct BuildInfo {
66 pub(crate) version: String,
67 pub(crate) config_layout: u32,
68}
69
70pub 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,
83 Restarting,
85 Verified,
87 RolledBack,
89 Failed,
92}
93
94#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
96pub(crate) struct Requester {
97 pub(crate) handle: String,
98 pub(crate) session: String,
99}
100
101#[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 pub(crate) binary: PathBuf,
116 #[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 #[serde(default, skip_serializing_if = "Option::is_none")]
123 pub(crate) origin: Option<Origin>,
124 #[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 #[serde(default, skip_serializing_if = "std::ops::Not::not")]
134 pub(crate) waited_out: bool,
135 #[serde(default, skip_serializing_if = "Option::is_none")]
137 pub(crate) detail: Option<String>,
138 #[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
197fn 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
215fn 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
222async 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
240struct 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
252pub(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 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
291fn 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#[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 Wait,
329 Nothing,
330}
331
332fn 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), ®istered, 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
374fn 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#[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#[derive(Clone)]
416pub(crate) struct Notifier {
417 hub: Arc<Hub>,
418 states: States,
419 grace: Duration,
421 give_up: Duration,
423 poll: Duration,
424 instance: Instance,
426 #[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 #[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 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 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 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 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
597enum Launcher {
602 Systemd,
604 #[cfg(test)]
606 Record(Arc<SyncMutex<Vec<Plan>>>),
607}
608
609pub(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 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 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 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 self.arm(plan)
734 }
735
736 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 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 fn waiting_for(&self, plan: &Plan) -> Option<String> {
768 if let Some(requester) = &plan.requester {
769 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 () = 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 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 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 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 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
927fn 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
941fn 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
963pub 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
1035fn 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
1050async 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
1090pub(crate) struct Startup {
1095 plan: Option<Plan>,
1096 unclean: Option<(String, u64)>,
1099}
1100
1101#[derive(Serialize, Deserialize)]
1102struct Marker {
1103 pid: u32,
1104 version: String,
1105 started_unix: u64,
1106}
1107
1108pub(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(¤t)
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
1146pub(crate) fn clean_shutdown(layout: &Layout) {
1148 let _ = std::fs::remove_file(layout.daemon_marker());
1149}
1150
1151#[derive(Debug, PartialEq, Eq)]
1153enum Decision {
1154 Say(String),
1155 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
1191fn 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
1218pub(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
1277pub(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;