smoldot_light/
lifecycle_service.rs

1// Smoldot
2// Copyright (C) 2019-2022  Parity Technologies (UK) Ltd.
3// SPDX-License-Identifier: GPL-3.0-or-later WITH Classpath-exception-2.0
4
5// This program is free software: you can redistribute it and/or modify
6// it under the terms of the GNU General Public License as published by
7// the Free Software Foundation, either version 3 of the License, or
8// (at your option) any later version.
9
10// This program is distributed in the hope that it will be useful,
11// but WITHOUT ANY WARRANTY; without even the implied warranty of
12// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE.  See the
13// GNU General Public License for more details.
14
15// You should have received a copy of the GNU General Public License
16// along with this program.  If not, see <http://www.gnu.org/licenses/>.
17
18//! Per-chain lifecycle state.
19//!
20//! Holds a small [`LifecycleState`] value (bootstrap phase, peer presence, stall verdict) and
21//! lets consumers subscribe to changes, so that an embedder can show what the light client is
22//! doing without parsing log output. See issue #3301.
23//!
24//! This is a "latest value" broadcast, not an event log. A subscriber receives the current
25//! state when it subscribes and then the newest state after every change. A subscriber that
26//! reads slowly simply skips intermediate states. Nothing is buffered, so a slow subscriber
27//! can never fall behind or slow down syncing.
28//!
29//! `start` creates the service of a chain together with the two tasks that keep it up to
30//! date: one maps the sync service's status to the phase, the other polls the network service
31//! for the peer count and derives the stall verdict.
32//!
33//! The schema is unstable.
34
35use crate::{network_service, platform::PlatformRef, sync_service};
36use alloc::sync::{Arc, Weak};
37use async_lock::Mutex;
38use core::sync::atomic::{AtomicUsize, Ordering};
39use core::time::Duration;
40
41/// Time without any connected peer after which the chain is reported as stalled.
42const NO_PEERS_TIMEOUT: Duration = Duration::from_secs(30);
43
44/// Time without warp sync progress after which the chain is reported as stalled.
45const NO_PROGRESS_TIMEOUT: Duration = Duration::from_secs(45);
46
47/// Decides the [`Health`] from how long the chain has had no peer and, if a warp sync is in
48/// progress, how long it has not advanced.
49fn health_verdict(no_peers_for: Duration, no_progress_for: Option<Duration>) -> Health {
50    if no_peers_for >= NO_PEERS_TIMEOUT {
51        Health::Stalled {
52            reason: StallReason::NoPeers,
53        }
54    } else if no_progress_for.is_some_and(|d| d >= NO_PROGRESS_TIMEOUT) {
55        Health::Stalled {
56            reason: StallReason::NoProgress,
57        }
58    } else {
59        Health::Ok
60    }
61}
62
63/// Bootstrap progress of the chain.
64#[derive(Debug, Clone, Copy, PartialEq, Eq)]
65pub enum Phase {
66    /// The chain has been added and no block is being streamed yet.
67    Connecting,
68    /// A GrandPa warp sync is in progress.
69    Syncing {
70        /// Highest block proven finalized by the warp sync fragments verified so far.
71        at: u64,
72        /// Highest best block advertised by a connected peer. Never below `at`.
73        target: u64,
74    },
75    /// The sync service is streaming new blocks. Not terminal: a later warp sync moves the
76    /// chain back to [`Phase::Syncing`], then to `Ready` again.
77    Ready,
78}
79
80/// Why the watchdog considers the chain stalled.
81#[derive(Debug, Clone, Copy, PartialEq, Eq)]
82pub enum StallReason {
83    /// No peer has been connected for a while.
84    NoPeers,
85    /// A warp sync is in progress but hasn't advanced for a while.
86    NoProgress,
87}
88
89/// Verdict of the stall watchdog.
90#[derive(Debug, Clone, Copy, PartialEq, Eq)]
91pub enum Health {
92    Ok,
93    Stalled { reason: StallReason },
94}
95
96/// Lifecycle state of a chain.
97#[derive(Debug, Clone, Copy, PartialEq, Eq)]
98pub struct LifecycleState {
99    pub phase: Phase,
100    /// Number of peers currently connected on this chain.
101    pub num_peers: u32,
102    pub health: Health,
103}
104
105impl Default for LifecycleState {
106    fn default() -> Self {
107        LifecycleState {
108            phase: Phase::Connecting,
109            num_peers: 0,
110            health: Health::Ok,
111        }
112    }
113}
114
115/// Creates the [`LifecycleService`] of a chain and spawns the two tasks that keep it up to date.
116/// Both tasks hold only weak references so that they stop, rather than keep the chain alive,
117/// once the chain is removed.
118pub(crate) fn start<TPlat: PlatformRef>(
119    platform: &TPlat,
120    sync_service: &Arc<sync_service::SyncService<TPlat>>,
121    network_service_chain: &Arc<network_service::NetworkServiceChain<TPlat>>,
122) -> Arc<LifecycleService> {
123    let lifecycle_service = LifecycleService::new();
124
125    // Drives `LifecycleState::phase` from the sync service's own status: `Syncing` while a warp
126    // sync is in progress, `Ready` once the sync service serves the chain. Ends when the sync
127    // service is gone.
128    platform.spawn_task("lifecycle-phase".into(), {
129        let lifecycle_service = Arc::downgrade(&lifecycle_service);
130        let sync_service = Arc::downgrade(sync_service);
131        let platform = platform.clone();
132        async move {
133            // Warp sync fragments can verify at dozens per second. Every status is applied
134            // as soon as it is received, but after a progress update the task pauses for
135            // this interval and then applies only the newest status that arrived meanwhile,
136            // so that consumers see at most a couple of progress updates per second while
137            // still seeing the start of a warp sync and its end without delay.
138            const PROGRESS_BATCH_INTERVAL: Duration = Duration::from_millis(500);
139
140            let sync_status = {
141                let Some(sync_service) = sync_service.upgrade() else {
142                    return;
143                };
144                sync_service.subscribe_sync_status().await
145            };
146
147            let apply = |status: sync_service::SyncStatus| {
148                let lifecycle_service = lifecycle_service.clone();
149                async move {
150                    let lifecycle_service = lifecycle_service.upgrade()?;
151                    let phase = match status {
152                        sync_service::SyncStatus::WarpSyncing { at, target } => {
153                            Phase::Syncing { at, target }
154                        }
155                        sync_service::SyncStatus::Ready => Phase::Ready,
156                    };
157                    lifecycle_service.update(|s| s.phase = phase).await;
158                    Some(())
159                }
160            };
161
162            while let Ok(status) = sync_status.recv().await {
163                if apply(status).await.is_none() {
164                    return;
165                }
166                if matches!(status, sync_service::SyncStatus::WarpSyncing { .. }) {
167                    platform.sleep(PROGRESS_BATCH_INTERVAL).await;
168                    let mut newest = None;
169                    while let Ok(newer) = sync_status.try_recv() {
170                        newest = Some(newer);
171                    }
172                    if let Some(newest) = newest
173                        && apply(newest).await.is_none()
174                    {
175                        return;
176                    }
177                }
178            }
179        }
180    });
181
182    // Drives `LifecycleState::num_peers` and `LifecycleState::health` by polling the network
183    // service. Polling (rather than subscribing to network events) keeps this task from ever
184    // slowing down the networking. The poll is frequent during the first minutes after a
185    // subscriber appears, where an embedder is most likely to display the state, and relaxed
186    // afterwards or once the chain is running with peers. Nothing is polled while the state
187    // has no subscriber.
188    platform.spawn_task("lifecycle-watchdog".into(), {
189        let lifecycle_service = Arc::downgrade(&lifecycle_service);
190        let network_service_chain = Arc::downgrade(network_service_chain);
191        let platform = platform.clone();
192        async move {
193            const FAST_POLL_WINDOW: Duration = Duration::from_secs(120);
194
195            let mut started = platform.now();
196            let mut last_peer_seen = started.clone();
197            // Warp sync height last observed, and when it was first observed.
198            let mut last_progress: Option<(u64, TPlat::Instant)> = None;
199
200            loop {
201                let wait = {
202                    let Some(lifecycle_service) = lifecycle_service.upgrade() else {
203                        return;
204                    };
205                    lifecycle_service.wait_for_subscriber()
206                };
207                if let Some(wait) = wait {
208                    wait.await;
209                    // The time spent without a subscriber was not observed, so the stall
210                    // clocks restart.
211                    started = platform.now();
212                    last_peer_seen = started.clone();
213                    last_progress = None;
214                    continue;
215                }
216
217                let num_peers = {
218                    let Some(network_service_chain) = network_service_chain.upgrade() else {
219                        return;
220                    };
221                    u32::try_from(network_service_chain.peers_list().await.count())
222                        .unwrap_or(u32::MAX)
223                };
224                let has_peers = num_peers > 0;
225                let Some(lifecycle_service) = lifecycle_service.upgrade() else {
226                    return;
227                };
228
229                let now = platform.now();
230                if has_peers {
231                    last_peer_seen = now.clone();
232                }
233
234                let state = lifecycle_service.current().await;
235                let no_progress_for = match (state.phase, &last_progress) {
236                    (Phase::Syncing { at, .. }, Some((seen_at, since))) if *seen_at == at => {
237                        Some(now.clone() - since.clone())
238                    }
239                    (Phase::Syncing { at, .. }, _) => {
240                        last_progress = Some((at, now.clone()));
241                        Some(Duration::ZERO)
242                    }
243                    _ => {
244                        last_progress = None;
245                        None
246                    }
247                };
248                let health = health_verdict(now.clone() - last_peer_seen.clone(), no_progress_for);
249
250                lifecycle_service
251                    .update(|s| {
252                        s.num_peers = num_peers;
253                        s.health = health;
254                    })
255                    .await;
256                drop(lifecycle_service);
257
258                let settled = has_peers && matches!(state.phase, Phase::Ready);
259                let fast = !settled && now - started.clone() < FAST_POLL_WINDOW;
260                platform
261                    .sleep(Duration::from_secs(if fast { 1 } else { 5 }))
262                    .await;
263            }
264        }
265    });
266
267    lifecycle_service
268}
269
270/// Holder of the [`LifecycleState`] of one chain.
271pub struct LifecycleService {
272    state: Mutex<LifecycleState>,
273    changed: event_listener::Event,
274    /// Number of live [`Subscription`]s.
275    num_subscribers: AtomicUsize,
276    /// Notified when [`LifecycleService::num_subscribers`] goes from zero to one, and when the
277    /// service is dropped.
278    subscribed: event_listener::Event,
279}
280
281impl LifecycleService {
282    pub fn new() -> Arc<Self> {
283        Arc::new(LifecycleService {
284            state: Mutex::new(LifecycleState::default()),
285            changed: event_listener::Event::new(),
286            num_subscribers: AtomicUsize::new(0),
287            subscribed: event_listener::Event::new(),
288        })
289    }
290
291    /// If no [`Subscription`] exists, returns a listener that resolves once one is created or
292    /// the service is dropped. Returns `None` if a subscription already exists and there is
293    /// nothing to wait for.
294    ///
295    /// Lets the tasks that maintain the state stay idle while nobody is watching.
296    pub fn wait_for_subscriber(&self) -> Option<event_listener::EventListener> {
297        let listener = self.subscribed.listen();
298        if self.num_subscribers.load(Ordering::Acquire) > 0 {
299            None
300        } else {
301            Some(listener)
302        }
303    }
304
305    /// Returns the current state.
306    pub async fn current(&self) -> LifecycleState {
307        *self.state.lock().await
308    }
309
310    /// Modifies the state in place. Subscribers are woken up only if the state actually changed.
311    pub async fn update(&self, f: impl FnOnce(&mut LifecycleState)) {
312        let mut state = self.state.lock().await;
313        let before = *state;
314        f(&mut state);
315        if *state != before {
316            self.changed.notify(usize::MAX);
317        }
318    }
319
320    /// Subscribes to state changes. The subscription holds only a weak reference, so it never
321    /// keeps the chain alive.
322    pub fn subscribe(self: &Arc<Self>) -> Subscription {
323        if self.num_subscribers.fetch_add(1, Ordering::AcqRel) == 0 {
324            self.subscribed.notify(usize::MAX);
325        }
326        Subscription {
327            service: Arc::downgrade(self),
328            last_seen: None,
329        }
330    }
331}
332
333impl Drop for LifecycleService {
334    fn drop(&mut self) {
335        // Wake up subscribers waiting in `Subscription::next` and tasks waiting in
336        // `wait_for_subscriber` so that they observe the end.
337        self.changed.notify(usize::MAX);
338        self.subscribed.notify(usize::MAX);
339    }
340}
341
342/// Handle returned by [`LifecycleService::subscribe`].
343pub struct Subscription {
344    service: Weak<LifecycleService>,
345    /// Last state returned by [`Subscription::next`]. `None` before the first call.
346    last_seen: Option<LifecycleState>,
347}
348
349impl Drop for Subscription {
350    fn drop(&mut self) {
351        if let Some(service) = self.service.upgrade() {
352            service.num_subscribers.fetch_sub(1, Ordering::AcqRel);
353        }
354    }
355}
356
357impl Subscription {
358    /// Returns the current state on the first call, then the newest state after each change.
359    /// Returns `None` once the [`LifecycleService`] has been dropped, which happens when the
360    /// chain is removed.
361    pub async fn next(&mut self) -> Option<LifecycleState> {
362        loop {
363            let service = self.service.upgrade()?;
364            let listener = {
365                let state = service.state.lock().await;
366                if self.last_seen != Some(*state) {
367                    self.last_seen = Some(*state);
368                    return Some(*state);
369                }
370                // The listener is created while the lock is held, so a change that happens after
371                // the comparison above is guaranteed to wake it up.
372                service.changed.listen()
373            };
374            drop(service);
375            listener.await;
376        }
377    }
378}
379
380#[cfg(test)]
381mod tests {
382    use super::*;
383    use futures_lite::future::{block_on, poll_once};
384
385    #[test]
386    fn first_next_returns_current_state() {
387        block_on(async {
388            let svc = LifecycleService::new();
389            svc.update(|s| s.num_peers = 3).await;
390
391            let mut sub = svc.subscribe();
392            let state = sub.next().await.unwrap();
393
394            assert_eq!(state.num_peers, 3);
395            assert_eq!(state.phase, Phase::Connecting);
396            assert!(poll_once(sub.next()).await.is_none());
397        });
398    }
399
400    #[test]
401    fn updates_are_coalesced_to_the_latest_state() {
402        block_on(async {
403            let svc = LifecycleService::new();
404            let mut sub = svc.subscribe();
405            assert_eq!(sub.next().await.unwrap(), LifecycleState::default());
406
407            svc.update(|s| s.phase = Phase::Syncing { at: 1, target: 10 })
408                .await;
409            svc.update(|s| s.phase = Phase::Syncing { at: 2, target: 10 })
410                .await;
411            svc.update(|s| s.phase = Phase::Ready).await;
412
413            assert_eq!(sub.next().await.unwrap().phase, Phase::Ready);
414            assert!(poll_once(sub.next()).await.is_none());
415        });
416    }
417
418    #[test]
419    fn unchanged_update_does_not_wake_subscribers() {
420        block_on(async {
421            let svc = LifecycleService::new();
422            let mut sub = svc.subscribe();
423            sub.next().await.unwrap();
424
425            svc.update(|s| s.num_peers = 0).await;
426
427            assert!(poll_once(sub.next()).await.is_none());
428        });
429    }
430
431    #[test]
432    fn subscribers_are_independent() {
433        block_on(async {
434            let svc = LifecycleService::new();
435            let mut fast = svc.subscribe();
436            let mut slow = svc.subscribe();
437            fast.next().await.unwrap();
438            slow.next().await.unwrap();
439
440            svc.update(|s| s.num_peers = 3).await;
441            assert_eq!(fast.next().await.unwrap().num_peers, 3);
442            svc.update(|s| s.phase = Phase::Ready).await;
443            assert_eq!(fast.next().await.unwrap().phase, Phase::Ready);
444
445            let seen_by_slow = slow.next().await.unwrap();
446            assert_eq!(seen_by_slow.num_peers, 3);
447            assert_eq!(seen_by_slow.phase, Phase::Ready);
448            assert!(poll_once(slow.next()).await.is_none());
449        });
450    }
451
452    #[test]
453    fn wait_for_subscriber_tracks_live_subscriptions() {
454        block_on(async {
455            let svc = LifecycleService::new();
456            let listener = svc.wait_for_subscriber().unwrap();
457            assert!(poll_once(listener).await.is_none());
458
459            let listener = svc.wait_for_subscriber().unwrap();
460            let sub = svc.subscribe();
461            assert!(poll_once(listener).await.is_some());
462            assert!(svc.wait_for_subscriber().is_none());
463
464            drop(sub);
465            assert!(svc.wait_for_subscriber().is_some());
466        });
467    }
468
469    #[test]
470    fn wait_for_subscriber_wakes_when_service_is_dropped() {
471        block_on(async {
472            let svc = LifecycleService::new();
473            let listener = svc.wait_for_subscriber().unwrap();
474            drop(svc);
475            assert!(poll_once(listener).await.is_some());
476        });
477    }
478
479    #[test]
480    fn next_returns_none_after_service_is_dropped() {
481        block_on(async {
482            let svc = LifecycleService::new();
483            let mut sub = svc.subscribe();
484            sub.next().await.unwrap();
485
486            let pending = poll_once(sub.next()).await;
487            assert!(pending.is_none());
488
489            drop(svc);
490            assert!(sub.next().await.is_none());
491        });
492    }
493
494    #[test]
495    fn health_verdict_thresholds() {
496        let ok = Health::Ok;
497        let no_peers = Health::Stalled {
498            reason: StallReason::NoPeers,
499        };
500        let no_progress = Health::Stalled {
501            reason: StallReason::NoProgress,
502        };
503        let just_under = |d: Duration| d - Duration::from_millis(1);
504
505        assert_eq!(health_verdict(Duration::ZERO, None), ok);
506        assert_eq!(health_verdict(just_under(NO_PEERS_TIMEOUT), None), ok);
507        assert_eq!(health_verdict(NO_PEERS_TIMEOUT, None), no_peers);
508        assert_eq!(
509            health_verdict(Duration::ZERO, Some(just_under(NO_PROGRESS_TIMEOUT))),
510            ok
511        );
512        assert_eq!(
513            health_verdict(Duration::ZERO, Some(NO_PROGRESS_TIMEOUT)),
514            no_progress
515        );
516        // No peers explains the missing progress, so it wins.
517        assert_eq!(
518            health_verdict(NO_PEERS_TIMEOUT, Some(NO_PROGRESS_TIMEOUT)),
519            no_peers
520        );
521    }
522
523    #[test]
524    fn late_subscriber_sees_only_the_latest_state() {
525        block_on(async {
526            let svc = LifecycleService::new();
527            svc.update(|s| s.phase = Phase::Syncing { at: 5, target: 9 })
528                .await;
529            svc.update(|s| {
530                s.health = Health::Stalled {
531                    reason: StallReason::NoProgress,
532                }
533            })
534            .await;
535
536            let mut sub = svc.subscribe();
537            let state = sub.next().await.unwrap();
538
539            assert_eq!(state.phase, Phase::Syncing { at: 5, target: 9 });
540            assert_eq!(
541                state.health,
542                Health::Stalled {
543                    reason: StallReason::NoProgress
544                }
545            );
546        });
547    }
548}