Skip to main content

mj_controller/session_manager/
spawn.rs

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