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