mj_controller/session_manager/
spawn.rs1use super::*;
2
3pub fn spawn_session_manager() -> Result<SessionManagerChannels> {
4 spawn_session_manager_observed(None)
5}
6
7pub(crate) fn spawn_session_manager_observed(
8 delegation: Option<DelegationSender>,
9) -> Result<SessionManagerChannels> {
10 let (targets_tx, mut targets_rx) = watch::channel(Vec::<RelaySessionTarget>::new());
11 let (commands_tx, mut commands_rx) = mpsc::channel(32);
12 let (mut updates_tx, updates_rx) = coalesced_update_channel();
13 updates_tx.delegation = delegation;
14 let session_cpu = updates_tx.cpu.subscribe();
15 let (shutdown_tx, mut shutdown_rx) = oneshot::channel();
16 let task = tokio::spawn(async move {
17 let mut actors = BTreeMap::<String, ActorRegistration>::new();
18 let mut tasks = tokio::task::JoinSet::<String>::new();
19 let mut desired_targets = BTreeMap::<String, RelaySessionTarget>::new();
20 loop {
21 tokio::select! {
22 _ = &mut shutdown_rx => break,
23 changed = targets_rx.changed() => {
24 if changed.is_err() {
25 break;
26 }
27 desired_targets = target_map(&targets_rx.borrow_and_update());
28 reconcile_actors(
29 &desired_targets,
30 &mut actors,
31 &mut tasks,
32 &updates_tx,
33 );
34 }
35 command = commands_rx.recv() => {
36 let Some(ManagerCommand::Session { session_id, reply }) = command else {
37 break;
38 };
39 let handle = actors
40 .get(&session_id)
41 .filter(|actor| !actor.commands.is_closed())
42 .filter(|actor| desired_targets.get(&session_id) == Some(&actor.target))
43 .map(|actor| ManagedSessionHandle {
44 session_id: session_id.clone(),
45 commands: actor.commands.clone(),
46 releases: actor.releases.clone(),
47 view: actor.view.clone(),
48 });
49 if reply.send(handle).is_err() {
50 tracing::debug!(
51 session_id = %session_id,
52 operation = "session_lookup",
53 "session lookup receiver was already closed"
54 );
55 }
56 }
57 joined = tasks.join_next_with_id(), if !tasks.is_empty() => {
58 match joined {
59 Some(Ok((task_id, session_id))) => {
60 let removed = remove_actor_task(&mut actors, task_id);
61 if removed.as_deref().is_some_and(|removed| removed != session_id) {
62 tracing::error!(
63 completed_session_id = session_id,
64 registered_session_id = removed,
65 "session relay actor completed under the wrong registration"
66 );
67 }
68 desired_targets = target_map(&targets_rx.borrow());
72 reconcile_actors(
73 &desired_targets,
74 &mut actors,
75 &mut tasks,
76 &updates_tx,
77 );
78 }
79 Some(Err(error)) if error.is_cancelled() => {
80 let cancelled_task = error.id();
81 let session_id = remove_actor_task(&mut actors, cancelled_task);
82 desired_targets = target_map(&targets_rx.borrow());
83 reconcile_actors(
84 &desired_targets,
85 &mut actors,
86 &mut tasks,
87 &updates_tx,
88 );
89 tracing::warn!(
90 session_id = ?session_id,
91 "cancelled session relay actor was replaced"
92 );
93 }
94 Some(Err(error)) => {
95 let failed_task = error.id();
96 remove_actor_task(&mut actors, failed_task);
97 desired_targets = target_map(&targets_rx.borrow());
98 reconcile_actors(
99 &desired_targets,
100 &mut actors,
101 &mut tasks,
102 &updates_tx,
103 );
104 tracing::error!(%error, "session relay actor failed");
105 }
106 None => {}
107 }
108 }
109 }
110 }
111 shutdown_session_actors(&mut actors, &mut tasks).await;
112 });
113 Ok(SessionManagerChannels {
114 session_cpu: session_cpu.clone(),
115 targets: targets_tx,
116 control: SessionManagerControl {
117 commands: commands_tx,
118 session_cpu,
119 },
120 updates: updates_rx,
121 shutdown: SessionManagerShutdown {
122 signal: Some(shutdown_tx),
123 task: Some(task),
124 },
125 })
126}
127
128pub(super) async fn shutdown_session_actors(
129 actors: &mut BTreeMap<String, ActorRegistration>,
130 tasks: &mut tokio::task::JoinSet<String>,
131) {
132 for actor in actors.values() {
133 actor.retirement.send_replace(true);
134 }
135 actors.clear();
136
137 let graceful = async {
138 while let Some(joined) = tasks.join_next().await {
139 match joined {
140 Ok(_) => {}
141 Err(error) if error.is_cancelled() => {}
142 Err(error) => {
143 tracing::error!(%error, "session relay actor failed during shutdown");
144 }
145 }
146 }
147 };
148 if tokio::time::timeout(SESSION_MANAGER_SHUTDOWN_GRACE, graceful)
149 .await
150 .is_ok()
151 {
152 return;
153 }
154
155 tracing::warn!(
156 timeout_ms = SESSION_MANAGER_SHUTDOWN_GRACE.as_millis(),
157 "session relay actors did not stop before the shutdown deadline; aborting them"
158 );
159 let aborted = tasks.len();
160 let aborting = std::time::Instant::now();
161 tasks.abort_all();
162 while let Some(joined) = tasks.join_next().await {
163 if let Err(error) = joined
164 && !error.is_cancelled()
165 {
166 tracing::error!(%error, "session relay actor failed while being aborted");
167 }
168 }
169 tracing::info!(
172 aborted,
173 duration_ms = aborting.elapsed().as_millis(),
174 "aborted session relay actors stopped"
175 );
176}