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 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                            // A watch sender may have published another target while this
69                            // completion was already ready. Reconcile against its newest
70                            // value so an intermediate replacement is never started.
71                            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    tasks.abort_all();
160    while let Some(joined) = tasks.join_next().await {
161        if let Err(error) = joined
162            && !error.is_cancelled()
163        {
164            tracing::error!(%error, "session relay actor failed while being aborted");
165        }
166    }
167}