1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
//! Lossy event-bus relay, alive tracking, and subscriber listener lifecycle.
use std::sync::Arc;
use tokio::sync::broadcast;
use super::SupervisorCore;
use crate::{
core::alive::AliveTracker, events::Event, identity::TaskId, subscribers::SubscriberSet,
};
impl SupervisorCore {
/// Returns a best-effort sorted list of task names currently marked alive.
///
/// This is an event-derived activity view, not registry membership.
/// A registered actor can be marked not alive while it waits to retry.
pub(crate) async fn snapshot(&self) -> Vec<Arc<str>> {
self.alive.snapshot().await
}
/// Returns true if any task with this name is currently marked alive.
///
/// This is a best-effort activity query from the alive tracker, not a registry membership check.
pub(crate) async fn is_alive(&self, name: &str) -> bool {
self.alive.is_alive(name).await
}
/// Applies one event to alive tracking and subscriber fan-out.
async fn distribute(alive: &AliveTracker, set: &SubscriberSet, ev: Arc<Event>) {
alive.update(&ev).await;
set.emit_arc(ev);
}
/// Drains retained events from a bus receiver.
///
/// Used when the subscriber listener is shutting down.
/// Broadcast lag gaps are skipped so the retained tail can still be delivered to alive tracking and subscribers.
pub(super) async fn drain_pending(
rx: &mut broadcast::Receiver<Arc<Event>>,
alive: &AliveTracker,
set: &SubscriberSet,
) {
loop {
match rx.try_recv() {
Ok(ev) => Self::distribute(alive, set, ev).await,
Err(broadcast::error::TryRecvError::Lagged(_)) => continue,
Err(_) => break,
}
}
}
/// Starts the subscriber listener task.
///
/// The listener relays bus events to the alive tracker and subscriber queues.
pub(super) fn subscriber_listener(&self) {
let mut rx = self.bus.subscribe();
let set = Arc::clone(&self.subs);
let alive = Arc::clone(&self.alive);
let registry = Arc::clone(&self.registry);
let rt = self.runtime_token.clone();
let handle = tokio::spawn(async move {
loop {
tokio::select! {
biased;
msg = rx.recv() => match msg {
Ok(arc_ev) => {
if let Err(panic) = crate::core::panic_guard::guarded(
Self::distribute(&alive, &set, arc_ev),
)
.await
{
set.emit_arc(Arc::new(Event::runtime_failure(
"subscriber_listener",
format!("listener panic: {panic}"),
)));
}
}
Err(broadcast::error::RecvError::Lagged(skipped)) => {
let arc_e = Arc::new(Event::subscriber_overflow(
"subscriber_listener",
format!("lagged({skipped})"),
));
alive.update(&arc_e).await;
set.emit_arc(arc_e);
let live: std::collections::HashSet<TaskId> = registry
.list()
.await
.into_iter()
.map(|(id, _)| id)
.collect();
alive.reconcile(&live).await;
}
Err(broadcast::error::RecvError::Closed) => break,
},
_ = rt.cancelled() => {
Self::drain_pending(&mut rx, &alive, &set).await;
break;
}
}
}
});
*self.subscriber_handle.lock().unwrap() = Some(handle);
}
/// Awaits the subscriber listener.
///
/// Returns `false` when Tokio reports that the listener did not join cleanly.
pub(super) async fn join_subscriber_listener(&self) -> bool {
let handle = self
.subscriber_handle
.lock()
.unwrap_or_else(|error| error.into_inner())
.take();
let Some(handle) = handle else {
return true;
};
match handle.await {
Ok(()) => true,
Err(error) => {
self.subs.emit_arc(Arc::new(Event::runtime_failure(
"subscriber_listener",
format!("listener join failed: {error}"),
)));
false
}
}
}
/// Aborts the subscriber listener so shutdown join-failure handling can be tested.
#[cfg(test)]
pub(super) fn abort_subscriber_listener_for_test(&self) {
if let Some(handle) = self
.subscriber_handle
.lock()
.unwrap_or_else(|error| error.into_inner())
.as_ref()
{
handle.abort();
}
}
}