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