Skip to main content

muxr_client/session/
attach.rs

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                // Socket path creation can win the race against listener readiness after spawning the server.
140                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    // The server owns log timestamp generation, and startup failures may happen before attach can return it. Use the
155    // spawned pid as the stable client-known part of the debug hint instead of scanning logs during startup failure.
156    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            // A structured muxr rejection proves the socket is live even if pid metadata is missing or stale.
167            Err(attach_error).attach("socket returned a structured muxr response")
168        }
169        AttachFailure::Unusable(attach_error) => {
170            // Even stale/incompatible servers may still answer Ping; an unusable attach is the compatibility signal.
171            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}