1use 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
41const NO_PEERS_TIMEOUT: Duration = Duration::from_secs(30);
43
44const NO_PROGRESS_TIMEOUT: Duration = Duration::from_secs(45);
46
47fn 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#[derive(Debug, Clone, Copy, PartialEq, Eq)]
65pub enum Phase {
66 Connecting,
68 Syncing {
70 at: u64,
72 target: u64,
74 },
75 Ready,
78}
79
80#[derive(Debug, Clone, Copy, PartialEq, Eq)]
82pub enum StallReason {
83 NoPeers,
85 NoProgress,
87}
88
89#[derive(Debug, Clone, Copy, PartialEq, Eq)]
91pub enum Health {
92 Ok,
93 Stalled { reason: StallReason },
94}
95
96#[derive(Debug, Clone, Copy, PartialEq, Eq)]
98pub struct LifecycleState {
99 pub phase: Phase,
100 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
115pub(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 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 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 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 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 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
270pub struct LifecycleService {
272 state: Mutex<LifecycleState>,
273 changed: event_listener::Event,
274 num_subscribers: AtomicUsize,
276 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 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 pub async fn current(&self) -> LifecycleState {
307 *self.state.lock().await
308 }
309
310 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 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 self.changed.notify(usize::MAX);
338 self.subscribed.notify(usize::MAX);
339 }
340}
341
342pub struct Subscription {
344 service: Weak<LifecycleService>,
345 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 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 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 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}