Skip to main content

mj_controller/session_manager/
spawn.rs

1use 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                            // A watch sender may have published another target while this
68                            // completion was already ready. Reconcile against its newest
69                            // value so an intermediate replacement is never started.
70                            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}