1use std::fs;
2use std::io;
3use std::path::Path;
4use std::time::Duration;
5use std::time::Instant;
6
7use muxr_core::AttachRequest;
8use muxr_core::ClientRequest;
9use muxr_core::LayoutSnapshot;
10use muxr_core::PaneRegionsSnapshot;
11use muxr_core::ServerEvent;
12use muxr_core::SessionName;
13use muxr_core::SessionPaths;
14use muxr_core::TerminalSize;
15use muxr_transport::ClientConnection;
16use muxr_transport::ClientEventReader;
17use muxr_transport::ClientRequestWriter;
18use rootcause::prelude::ResultExt;
19use rootcause::report;
20
21use crate::session::start;
22
23const ATTACH_TIMEOUT: Duration = Duration::from_secs(2);
24const SERVER_READY_TIMEOUT: Duration = Duration::from_secs(2);
25
26pub struct AttachedSession {
27 pub layout: LayoutSnapshot,
28 pub pane_regions: PaneRegionsSnapshot,
29 pub reader: ClientEventReader,
30 pub writer: ClientRequestWriter,
31}
32
33enum AttachFailure {
34 Rejected(rootcause::Report),
35 Unusable(rootcause::Report),
36}
37
38pub async fn open_session(
39 session: &SessionName,
40 terminal_size: TerminalSize,
41 server_executable: &Path,
42 external_layout: Option<&Path>,
43) -> rootcause::Result<AttachedSession> {
44 let paths = SessionPaths::from_home(session)?;
45 self::open_session_with_paths(session, &paths, terminal_size, server_executable, external_layout).await
46}
47
48async fn open_session_with_paths(
49 session: &SessionName,
50 paths: &SessionPaths,
51 terminal_size: TerminalSize,
52 server_executable: &Path,
53 external_layout: Option<&Path>,
54) -> rootcause::Result<AttachedSession> {
55 if let Some(external_layout) = external_layout {
56 self::guard_external_start_seed(paths, session, external_layout).await?;
57 }
58
59 match self::attach(session, paths, terminal_size.clone()).await {
60 Ok(attached_session) => return Ok(attached_session),
61 Err(attach_failure) => {
62 handle_attach_failure(attach_failure)?;
63 start::cleanup_stale_session_files(paths)?;
64 }
65 }
66
67 let spawned_server = start::spawn_server_process(session, paths, server_executable, external_layout)?;
68 self::attach_started_server(session, paths, terminal_size, &spawned_server.log_locator).await
69}
70
71async fn attach(
72 session: &SessionName,
73 paths: &SessionPaths,
74 terminal_size: TerminalSize,
75) -> Result<AttachedSession, AttachFailure> {
76 let mut connection = self::connect_with_timeout(paths).await?;
77
78 tokio::time::timeout(
79 ATTACH_TIMEOUT,
80 connection.send_request(&ClientRequest::Attach(AttachRequest {
81 session: session.clone(),
82 terminal_size,
83 })),
84 )
85 .await
86 .map_err(|_| AttachFailure::Unusable(report!("timed out writing muxr attach request")))?
87 .map_err(AttachFailure::Unusable)?;
88
89 let (layout, pane_regions) = match tokio::time::timeout(ATTACH_TIMEOUT, connection.recv_event())
90 .await
91 .map_err(|_| AttachFailure::Unusable(report!("timed out waiting for muxr attach response")))?
92 .map_err(AttachFailure::Unusable)?
93 {
94 Some(ServerEvent::Attached(attached)) => (attached.layout, attached.pane_regions),
95 Some(ServerEvent::Error(error)) => {
96 return Err(AttachFailure::Rejected(
97 report!("muxr server rejected attach")
98 .attach(format!("code={}", error.code()))
99 .attach(format!("msg={}", error.msg())),
100 ));
101 }
102 Some(event) => {
103 return Err(AttachFailure::Unusable(
104 report!("unexpected muxr server attach event").attach(format!("{event:?}")),
105 ));
106 }
107 None => return Err(AttachFailure::Unusable(report!("muxr server closed before attach"))),
108 };
109
110 let (reader, writer) = connection.split();
111 Ok(AttachedSession {
112 layout,
113 pane_regions,
114 reader,
115 writer,
116 })
117}
118
119async fn connect_with_timeout(paths: &SessionPaths) -> Result<ClientConnection, AttachFailure> {
120 tokio::time::timeout(ATTACH_TIMEOUT, ClientConnection::connect(&paths.socket))
121 .await
122 .map_err(|_| AttachFailure::Unusable(report!("timed out connecting muxr session socket")))?
123 .map_err(AttachFailure::Unusable)
124}
125
126async fn attach_started_server(
127 session: &SessionName,
128 paths: &SessionPaths,
129 terminal_size: TerminalSize,
130 server_log_locator: &start::ServerLogLocator,
131) -> rootcause::Result<AttachedSession> {
132 let started_at = Instant::now();
133
134 loop {
135 match self::attach(session, paths, terminal_size.clone()).await {
136 Ok(attached_session) => return Ok(attached_session),
137 Err(AttachFailure::Rejected(error)) => return Err(error),
138 Err(AttachFailure::Unusable(error)) => {
139 if started_at.elapsed() > SERVER_READY_TIMEOUT {
141 return Err(self::server_startup_failure_error(error, server_log_locator));
142 }
143 }
144 }
145
146 tokio::time::sleep(Duration::from_millis(10)).await;
147 }
148}
149
150fn server_startup_failure_error(
151 error: rootcause::Report,
152 server_log_locator: &start::ServerLogLocator,
153) -> rootcause::Report {
154 error
157 .attach("muxr server did not become attachable after start")
158 .attach(format!("server_pid={}", server_log_locator.pid))
159 .attach(format!("logs_dir={}", server_log_locator.logs_dir.display()))
160 .attach(format!("log_pattern={}", server_log_locator.file_pattern))
161}
162
163fn handle_attach_failure(attach_failure: AttachFailure) -> rootcause::Result<()> {
164 match attach_failure {
165 AttachFailure::Rejected(attach_error) => {
166 Err(attach_error).attach("socket returned a structured muxr response")
168 }
169 AttachFailure::Unusable(attach_error) => {
170 drop(attach_error);
172 Ok(())
173 }
174 }
175}
176
177async fn guard_external_start_seed(
178 paths: &SessionPaths,
179 session: &SessionName,
180 external_layout: &Path,
181) -> rootcause::Result<()> {
182 match fs::read(&paths.layout) {
183 Ok(_) => {
184 return Err(self::external_layout_existing_session_error(
185 session,
186 external_layout,
187 "persisted-layout",
188 ));
189 }
190 Err(error) if error.kind() == io::ErrorKind::NotFound => {}
191 Err(error) => return Err(error).context("failed to read muxr session layout metadata")?,
192 }
193
194 if crate::session::list::session_state_async(paths).await? == crate::session::list::SessionState::Live {
195 return Err(self::external_layout_existing_session_error(
196 session,
197 external_layout,
198 "live",
199 ));
200 }
201
202 Ok(())
203}
204
205fn external_layout_existing_session_error(
206 session: &SessionName,
207 external_layout: &Path,
208 state: &str,
209) -> rootcause::Report {
210 report!("muxr external layout can only seed a new session")
211 .attach(format!("session={session}"))
212 .attach(format!("layout={}", external_layout.display()))
213 .attach(format!("state={state}"))
214}
215
216#[cfg(test)]
217mod tests {
218 use std::fs;
219 use std::os::unix::fs::PermissionsExt;
220 use std::path::Path;
221
222 use muxr_core::AttachAccepted;
223 use muxr_core::LayoutSnapshot;
224 use muxr_core::PaneId;
225 use muxr_core::PaneMouseMode;
226 use muxr_core::PaneRegionSnapshot;
227 use muxr_core::PaneRegionsSnapshot;
228 use muxr_core::PaneSnapshot;
229 use muxr_core::ServerError;
230 use muxr_core::SessionPaths;
231 use muxr_core::TabId;
232 use muxr_core::TabSnapshot;
233 use muxr_core::TrackedProcessState;
234 use muxr_transport::ServerListener;
235 use test_that::prelude::*;
236
237 use super::*;
238
239 #[test]
240 fn test_guard_external_start_seed_when_layout_metadata_exists_returns_error() -> rootcause::Result<()> {
241 let tempdir = tempfile::tempdir()?;
242 let (session, paths) = self::session_paths(tempdir.path(), "work")?;
243 let layout = Path::new("../.config/muxr/layouts/work.json");
244 fs::create_dir_all(&paths.root)?;
245 fs::write(&paths.layout, b"not necessarily valid json")?;
246
247 let error = self::runtime()?
248 .block_on(guard_external_start_seed(&paths, &session, layout))
249 .expect_err("expected persisted layout to block external layout seed");
250
251 assert_that!(
252 error.to_string(),
253 contains_substring("muxr external layout can only seed a new session")
254 );
255 assert_that!(error.to_string(), contains_substring("state=persisted-layout"));
256 Ok(())
257 }
258
259 #[test]
260 fn test_guard_external_start_seed_when_socket_is_live_returns_error() -> rootcause::Result<()> {
261 let tempdir = tempfile::tempdir()?;
262 let (session, paths) = self::session_paths(tempdir.path(), "work")?;
263 let layout = Path::new("../.config/muxr/layouts/work.json");
264 fs::create_dir_all(&paths.root)?;
265 let runtime = self::runtime()?;
266 let error = runtime.block_on(async {
267 let listener = ServerListener::bind(&paths.socket)?;
268 let handle = tokio::spawn(async move {
269 let mut connection = listener.accept().await?;
270 assert_that!(connection.recv_request().await?, eq(Some(ClientRequest::Ping)));
271 connection.send_event(&ServerEvent::Pong).await?;
272 Ok::<(), rootcause::Report>(())
273 });
274
275 let error = guard_external_start_seed(&paths, &session, layout)
276 .await
277 .expect_err("expected live session to block external layout seed");
278 handle
279 .await
280 .map_err(|error| report!("muxr live layout guard test task panicked").attach(format!("{error}")))??;
281 Ok::<_, rootcause::Report>(error)
282 })?;
283
284 assert_that!(
285 error.to_string(),
286 contains_substring("muxr external layout can only seed a new session")
287 );
288 assert_that!(error.to_string(), contains_substring("state=live"));
289 Ok(())
290 }
291
292 #[test]
293 fn test_handle_attach_failure_when_server_rejects_and_pid_is_missing_returns_error() -> rootcause::Result<()> {
294 self::runtime()?.block_on(async {
295 let tempdir = tempfile::tempdir()?;
296 let (_, paths) = self::session_paths(tempdir.path(), "work")?;
297 fs::create_dir_all(&paths.root)?;
298 let _listener = ServerListener::bind(&paths.socket)?;
299
300 assert_that!(
301 handle_attach_failure(AttachFailure::Rejected(report!("already attached"))),
302 err(anything())
303 );
304 assert_that!(paths.socket.exists(), eq(true));
305 Ok(())
306 })
307 }
308
309 #[test]
310 fn test_attach_when_server_rejects_returns_rejected_error() -> rootcause::Result<()> {
311 self::runtime()?.block_on(async {
312 let tempdir = tempfile::tempdir()?;
313 let (session, paths) = self::session_paths(tempdir.path(), "work")?;
314 fs::create_dir_all(&paths.root)?;
315 let listener = ServerListener::bind(&paths.socket)?;
316 let handle = tokio::spawn(async move {
317 let mut connection = listener.accept().await?;
318 assert_that!(
319 connection.recv_request().await?,
320 some(matches_pattern!(ClientRequest::Attach(anything())))
321 );
322 connection
323 .send_event(&ServerEvent::Error(ServerError::ClientAlreadyAttached))
324 .await?;
325 Ok::<(), rootcause::Report>(())
326 });
327
328 let attach_error = attach(&session, &paths, TerminalSize::new(80, 24)?).await.map_or_else(
329 |failure| match failure {
330 AttachFailure::Rejected(error) | AttachFailure::Unusable(error) => error,
331 },
332 |_| report!("expected rejected attach"),
333 );
334
335 assert_that!(
336 attach_error.to_string(),
337 contains_substring("muxr server rejected attach")
338 );
339 handle
340 .await
341 .map_err(|error| report!("muxr rejected attach test task panicked").attach(format!("{error}")))??;
342 Ok(())
343 })
344 }
345
346 #[test]
347 fn test_open_session_when_live_session_exists_with_missing_runner_attaches_without_spawning()
348 -> rootcause::Result<()> {
349 self::runtime()?.block_on(async {
350 let tempdir = tempfile::tempdir()?;
351 let (session, paths) = self::session_paths(tempdir.path(), "work")?;
352 fs::create_dir_all(&paths.root)?;
353 let listener = ServerListener::bind(&paths.socket)?;
354 let handle = tokio::spawn(async move {
355 let mut connection = listener.accept().await?;
356 assert_that!(
357 connection.recv_request().await?,
358 some(matches_pattern!(ClientRequest::Attach(anything())))
359 );
360 connection.send_event(&self::attached_event()?).await?;
361 Ok::<(), rootcause::Report>(())
362 });
363 let missing_runner = tempdir.path().join("missing-muxr-server");
364
365 let attached_session =
366 open_session_with_paths(&session, &paths, TerminalSize::new(80, 24)?, &missing_runner, None).await?;
367
368 assert_that!(attached_session.layout.active_tab(), eq(&TabId::new(1)?));
369 handle
370 .await
371 .map_err(|error| report!("muxr live attach test task panicked").attach(format!("{error}")))??;
372 Ok(())
373 })
374 }
375
376 #[test]
377 fn test_open_session_when_started_server_exits_before_attach_returns_log_locator() -> rootcause::Result<()> {
378 self::runtime()?.block_on(async {
379 let tempdir = tempfile::tempdir()?;
380 let (session, paths) = self::session_paths(tempdir.path(), "work")?;
381 let runner = tempdir.path().join("muxr-server");
382 fs::write(&runner, "#!/bin/sh\nexit 17\n").context("failed to write fake muxr server")?;
383 fs::set_permissions(&runner, fs::Permissions::from_mode(0o755))
384 .context("failed to make fake muxr server executable")?;
385
386 let open_result =
387 open_session_with_paths(&session, &paths, TerminalSize::new(80, 24)?, &runner, None).await;
388 assert_that!(
389 open_result.as_ref().map(|_| ()).map_err(ToString::to_string),
390 err(contains_substring("muxr server did not become attachable after start"))
391 );
392 let error = open_result.map_or_else(|error| error.to_string(), |_| String::new());
393 assert_that!(error, contains_substring("server_pid="));
394 assert_that!(
395 error,
396 contains_substring(format!("logs_dir={}", paths.logs_root()?.display()))
397 );
398 assert_that!(error, contains_substring("log_pattern=work-*-"));
399 assert_that!(error, not(contains_substring("server_log=")));
400 assert_that!(error, contains_substring(".log"));
401 Ok(())
402 })
403 }
404
405 fn attached_event() -> rootcause::Result<ServerEvent> {
406 let pane_id = PaneId::new(1)?;
407 let tab_id = TabId::new(1)?;
408 let pane = PaneSnapshot {
409 cmd_label: None,
410 cwd: "/tmp".to_string(),
411 focus_seq: 0,
412 id: pane_id,
413 title: "shell".to_string(),
414 tracked_process_state: TrackedProcessState::None,
415 };
416 let tab = TabSnapshot::new(tab_id, "default", pane_id, vec![pane])?;
417 let layout = LayoutSnapshot::new(tab_id, vec![tab])?;
418 let region = PaneRegionSnapshot::new(pane_id, 0, 0, 80, 24, PaneMouseMode::None, 0)?;
419 let pane_regions = PaneRegionsSnapshot::new(vec![region])?;
420 Ok(ServerEvent::Attached(AttachAccepted { layout, pane_regions }))
421 }
422
423 fn session_paths(base: &Path, raw: &str) -> rootcause::Result<(SessionName, SessionPaths)> {
424 let session = raw.parse()?;
425 let root = base.join("sessions").join(raw);
426
427 Ok((
428 session,
429 SessionPaths {
430 socket: root.join("server.sock"),
431 pid: root.join("server.pid"),
432 layout: root.join("layout.json"),
433 panes: root.join("panes"),
434 root,
435 },
436 ))
437 }
438
439 fn runtime() -> rootcause::Result<tokio::runtime::Runtime> {
440 Ok(tokio::runtime::Runtime::new().context("failed to build muxr test runtime")?)
441 }
442}