diff --git a/crates/embers-cli/src/interactive.rs b/crates/embers-cli/src/interactive.rs index 8c928173..70eceefb 100644 --- a/crates/embers-cli/src/interactive.rs +++ b/crates/embers-cli/src/interactive.rs @@ -1,3 +1,4 @@ +use std::fs; use std::io::{self, Write}; use std::os::fd::AsRawFd; use std::path::{Path, PathBuf}; @@ -17,6 +18,7 @@ const DEFAULT_SESSION_NAME: &str = "main"; const KEY_SEQUENCE_TIMEOUT: Duration = Duration::from_millis(15); const KEY_SEQUENCE_CONTINUATION_TIMEOUT: Duration = Duration::from_millis(2); const EVENT_POLL_INTERVAL: Duration = Duration::from_millis(20); +const CONFIG_WATCH_POLL_INTERVAL: Duration = Duration::from_millis(250); const BRACKETED_PASTE_END: &[u8] = b"\x1b[201~"; const TERMINAL_ENTER_BASE_SEQUENCE: &str = "\x1b[?1049h\x1b[?1004h\x1b[?2004h\x1b[?25l\x1b[2J\x1b[H"; @@ -37,11 +39,13 @@ pub async fn run( let mut session_id = Some(initial_session_id); let config = ConfigManager::from_process(config_path) .map_err(|error| MuxError::invalid_input(error.to_string()))?; + let watched_config_path = config.active_source().path.clone(); let mut configured = ConfiguredClient::new(client, config); let mut terminal = TerminalGuard::enter(mouse_capture_enabled(&configured))?; let (input_tx, mut input_rx) = mpsc::unbounded_channel(); - let _input_thread = spawn_input_thread(input_tx)?; + let _input_thread = spawn_input_thread(input_tx.clone())?; + let _config_thread = spawn_config_thread(watched_config_path, input_tx)?; let mut terminal_size = terminal.size()?; let mut dirty = true; @@ -117,6 +121,11 @@ pub async fn run( dirty = true; } } + Ok(TerminalEvent::ConfigChanged) => { + let _ = configured.reload_config_if_changed()?; + terminal.write_bytes(&drain_terminal_output(&mut configured))?; + dirty = true; + } Ok(TerminalEvent::InputClosed) => return Ok(()), Ok(TerminalEvent::InputError(message)) => { return Err(MuxError::transport(message)); @@ -390,6 +399,7 @@ enum TerminalEvent { Mouse(MouseEvent), Paste(Vec), Focus(bool), + ConfigChanged, InputClosed, InputError(String), } @@ -426,6 +436,40 @@ fn spawn_input_thread( .map_err(|error| MuxError::internal(format!("failed to spawn input thread: {error}"))) } +fn spawn_config_thread( + config_path: Option, + tx: mpsc::UnboundedSender, +) -> Result>> { + let Some(config_path) = config_path else { + return Ok(None); + }; + + let handle = thread::Builder::new() + .name("embers-config".to_owned()) + .spawn(move || { + let mut last_modified = config_modified(&config_path); + loop { + thread::sleep(CONFIG_WATCH_POLL_INTERVAL); + let next_modified = config_modified(&config_path); + if next_modified != last_modified { + last_modified = next_modified; + if tx.send(TerminalEvent::ConfigChanged).is_err() { + break; + } + } + } + }) + .map_err(|error| MuxError::internal(format!("failed to spawn config thread: {error}")))?; + + Ok(Some(handle)) +} + +fn config_modified(path: &Path) -> Option { + fs::metadata(path) + .and_then(|metadata| metadata.modified()) + .ok() +} + fn read_terminal_event(fd: libc::c_int) -> Result> { let Some(first) = read_byte(fd)? else { return Ok(None); diff --git a/crates/embers-cli/tests/interactive.rs b/crates/embers-cli/tests/interactive.rs index 59450051..4c264e47 100644 --- a/crates/embers-cli/tests/interactive.rs +++ b/crates/embers-cli/tests/interactive.rs @@ -2,13 +2,22 @@ use std::fs; use std::path::{Path, PathBuf}; use std::time::Duration; -use embers_core::PtySize; -use embers_test_support::{ - PtyHarness, TestConnection, TestServer, acquire_test_lock, cargo_bin, cargo_bin_path, +use embers_core::{ActivityState, BufferId, NodeId, PtySize, new_request_id}; +use embers_protocol::{ + BufferRequest, ClientMessage, InputRequest, ServerResponse, SessionSnapshot, + VisibleSnapshotResponse, }; -use tempfile::tempdir; +use embers_test_support::{PtyHarness, TestConnection, TestServer, cargo_bin, cargo_bin_path}; +use tokio::sync::Mutex; -use crate::support::{run_cli, session_snapshot_by_name, stdout}; +use std::sync::OnceLock; + +fn test_lock() -> &'static Mutex<()> { + static LOCK: OnceLock> = OnceLock::new(); + LOCK.get_or_init(|| Mutex::new(())) +} + +use crate::support::{require_pty, run_cli, session_snapshot_by_name, stdout}; const STARTUP_TIMEOUT: Duration = Duration::from_secs(15); const IO_TIMEOUT: Duration = Duration::from_secs(30); @@ -182,10 +191,206 @@ fn first_client_id_finds_attached_row() { assert_eq!(first_client_id(output), 42); } +fn focused_pane_id(snapshot: &SessionSnapshot) -> u64 { + snapshot + .session + .focused_leaf_id + .map(|leaf_id| leaf_id.0) + .expect("session has a focused pane") +} + +fn pane_buffer_id(snapshot: &SessionSnapshot, pane_id: u64) -> BufferId { + snapshot + .nodes + .iter() + .find(|node| node.id == NodeId(pane_id)) + .and_then(|node| node.buffer_view.as_ref()) + .map(|view| view.buffer_id) + .unwrap_or_else(|| panic!("pane {pane_id} buffer view exists")) +} + +fn root_tab_child_id(snapshot: &SessionSnapshot, title: &str) -> NodeId { + snapshot + .nodes + .iter() + .find(|node| node.id == snapshot.session.root_node_id) + .and_then(|node| node.tabs.as_ref()) + .and_then(|tabs| { + tabs.tabs + .iter() + .find(|tab| tab.title == title) + .map(|tab| tab.child_id) + }) + .unwrap_or_else(|| panic!("root tab `{title}` exists")) +} + +fn split_child_order( + snapshot: &SessionSnapshot, + first_pane_id: u64, + second_pane_id: u64, +) -> [u64; 2] { + snapshot + .nodes + .iter() + .find_map(|node| { + let split = node.split.as_ref()?; + let child_ids = split + .child_ids + .iter() + .map(|child| child.0) + .collect::>(); + if child_ids.len() == 2 + && child_ids.contains(&first_pane_id) + && child_ids.contains(&second_pane_id) + { + Some([child_ids[0], child_ids[1]]) + } else { + None + } + }) + .unwrap_or_else(|| { + panic!("split containing panes {first_pane_id} and {second_pane_id} exists") + }) +} + +fn disable_echo_in_pane(server: &TestServer, pane_id: u64) { + run_cli( + server, + [ + "send-keys", + "-t", + &pane_id.to_string(), + "--enter", + "--", + "stty", + "-echo", + ], + ); +} + +async fn send_buffer_input(connection: &mut TestConnection, buffer_id: BufferId, bytes: &[u8]) { + let response = connection + .request(&ClientMessage::Input(InputRequest::Send { + request_id: new_request_id(), + buffer_id, + bytes: bytes.to_vec(), + })) + .await + .expect("send input succeeds"); + assert!( + matches!(response, ServerResponse::Ok(_)), + "expected ok response to input send, got {response:?}" + ); +} + +async fn wait_for_buffer_activity( + connection: &mut TestConnection, + buffer_id: BufferId, + expected: ActivityState, +) { + let deadline = tokio::time::Instant::now() + IO_TIMEOUT; + loop { + let response = connection + .request(&ClientMessage::Buffer(BufferRequest::Get { + request_id: new_request_id(), + buffer_id, + })) + .await + .expect("buffer get succeeds"); + let activity = match response { + ServerResponse::Buffer(response) => response.buffer.activity, + other => panic!("expected buffer response, got {other:?}"), + }; + if activity == expected { + return; + } + + assert!( + tokio::time::Instant::now() < deadline, + "timed out waiting for buffer {buffer_id} activity {expected:?}; last activity {activity:?}" + ); + tokio::time::sleep(Duration::from_millis(50)).await; + } +} + +async fn wait_for_target_pane_buffer( + connection: &mut TestConnection, + session_name: &str, + pane_id: u64, + expected_buffer_id: BufferId, +) -> SessionSnapshot { + let deadline = tokio::time::Instant::now() + IO_TIMEOUT; + loop { + let snapshot = session_snapshot_by_name(connection, session_name).await; + if pane_buffer_id(&snapshot, pane_id) == expected_buffer_id { + return snapshot; + } + + assert!( + tokio::time::Instant::now() < deadline, + "timed out waiting for pane {pane_id} to attach buffer {expected_buffer_id}" + ); + tokio::time::sleep(Duration::from_millis(50)).await; + } +} + +async fn wait_for_split_child_order( + connection: &mut TestConnection, + session_name: &str, + pane_a: u64, + pane_b: u64, + expected: [u64; 2], +) -> SessionSnapshot { + let deadline = tokio::time::Instant::now() + IO_TIMEOUT; + loop { + let snapshot = session_snapshot_by_name(connection, session_name).await; + if split_child_order(&snapshot, pane_a, pane_b) == expected { + return snapshot; + } + + assert!( + tokio::time::Instant::now() < deadline, + "timed out waiting for split order {:?}; last order {:?}", + expected, + split_child_order(&snapshot, pane_a, pane_b) + ); + tokio::time::sleep(Duration::from_millis(50)).await; + } +} + +async fn wait_for_visible_snapshot( + connection: &mut TestConnection, + buffer_id: BufferId, + mut predicate: F, +) -> VisibleSnapshotResponse +where + F: FnMut(&VisibleSnapshotResponse) -> bool, +{ + let deadline = tokio::time::Instant::now() + IO_TIMEOUT; + loop { + let snapshot = connection + .capture_visible_buffer(buffer_id) + .await + .expect("visible capture succeeds"); + if predicate(&snapshot) { + return snapshot; + } + + assert!( + tokio::time::Instant::now() < deadline, + "timed out waiting for visible snapshot predicate; last snapshot: {snapshot:?}" + ); + tokio::time::sleep(Duration::from_millis(50)).await; + } +} + #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn embers_without_subcommand_starts_server_and_client() { - let _guard = acquire_test_lock().await.expect("acquire test lock"); - let tempdir = tempdir().expect("tempdir"); + let _guard = test_lock().lock().await; + if !require_pty() { + return; + } + let tempdir = tempfile::tempdir().expect("tempdir"); let socket_path = tempdir.path().join("embers.sock"); let socket_arg = socket_path.to_string_lossy().into_owned(); let (_spawned, mut harness) = spawn_embers(&["--socket", &socket_arg], socket_path.clone()); @@ -224,7 +429,10 @@ async fn embers_without_subcommand_starts_server_and_client() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn attach_subcommand_connects_to_running_server() { - let _guard = acquire_test_lock().await.expect("acquire test lock"); + let _guard = test_lock().lock().await; + if !require_pty() { + return; + } let server = TestServer::start().await.expect("start server"); let binary = cargo_bin_path("embers"); let binary_dir = binary.parent().expect("binary dir"); @@ -270,7 +478,10 @@ async fn attach_subcommand_connects_to_running_server() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn client_commands_can_switch_and_detach_a_live_attached_client() { - let _guard = acquire_test_lock().await.expect("acquire test lock"); + let _guard = test_lock().lock().await; + if !require_pty() { + return; + } let server = TestServer::start().await.expect("start server"); run_cli(&server, ["new-session", "main"]); @@ -329,7 +540,10 @@ async fn client_commands_can_switch_and_detach_a_live_attached_client() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn buffer_reveal_switches_the_attached_client_to_the_buffer_session() { - let _guard = acquire_test_lock().await.expect("acquire test lock"); + let _guard = test_lock().lock().await; + if !require_pty() { + return; + } let server = TestServer::start().await.expect("start server"); run_cli(&server, ["new-session", "main"]); @@ -398,8 +612,11 @@ async fn buffer_reveal_switches_the_attached_client_to_the_buffer_session() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn page_up_enters_local_scrollback_and_shows_indicator() { - let _guard = acquire_test_lock().await.expect("acquire test lock"); - let tempdir = tempdir().expect("tempdir"); + let _guard = test_lock().lock().await; + if !require_pty() { + return; + } + let tempdir = tempfile::tempdir().expect("tempdir"); let socket_path = tempdir.path().join("embers.sock"); let socket_arg = socket_path.to_string_lossy().into_owned(); let (_spawned, mut harness) = spawn_embers(&["--socket", &socket_arg], socket_path.clone()); @@ -419,8 +636,11 @@ async fn page_up_enters_local_scrollback_and_shows_indicator() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn local_selection_yank_emits_osc52_clipboard_sequence() { - let _guard = acquire_test_lock().await.expect("acquire test lock"); - let tempdir = tempdir().expect("tempdir"); + let _guard = test_lock().lock().await; + if !require_pty() { + return; + } + let tempdir = tempfile::tempdir().expect("tempdir"); let socket_path = tempdir.path().join("embers.sock"); let socket_arg = socket_path.to_string_lossy().into_owned(); let (_spawned, mut harness) = spawn_embers(&["--socket", &socket_arg], socket_path.clone()); @@ -445,3 +665,479 @@ async fn local_selection_yank_emits_osc52_clipboard_sequence() { // spawned.drop() will clean up the orphaned __serve process } + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn scripted_input_bindings_reach_the_live_terminal_in_pty() { + let _guard = test_lock().lock().await; + if !require_pty() { + return; + } + let server = TestServer::start().await.expect("start server"); + + run_cli(&server, ["new-session", "main"]); + run_cli( + &server, + [ + "new-window", + "-t", + "main", + "--title", + "shell", + "--", + "/bin/sh", + ], + ); + + let tempdir = tempfile::tempdir().expect("tempdir"); + let config_path = tempdir.path().join("config.rhai"); + fs::write( + &config_path, + r#"bind("normal", "", action.send_bytes_current("echo scripted-pty\r"))"#, + ) + .expect("write config"); + + let socket_arg = server.socket_path().to_string_lossy().into_owned(); + let socket_path = server.socket_path().to_path_buf(); + let config_arg = config_path.to_string_lossy().into_owned(); + let (_spawned, mut harness) = spawn_embers( + &[ + "attach", + "--socket", + &socket_arg, + "--config", + &config_arg, + "-t", + "main", + ], + socket_path, + ); + harness + .read_until_contains("[main]", STARTUP_TIMEOUT) + .expect("attach client renders"); + harness + .write_all("stty -echo\r") + .expect("disable shell echo in focused pane"); + harness + .wait_for_quiet(QUIET_TIMEOUT, IO_TIMEOUT) + .expect("focused shell settles"); + + let mut connection = TestConnection::connect(server.socket_path()) + .await + .expect("connect protocol client"); + let snapshot = session_snapshot_by_name(&mut connection, "main").await; + let buffer_id = pane_buffer_id(&snapshot, focused_pane_id(&snapshot)); + + harness.write_all("\x07").expect("trigger scripted binding"); + harness + .read_until_contains("scripted-pty", IO_TIMEOUT) + .expect("scripted output renders"); + connection + .wait_for_capture_contains(buffer_id, "scripted-pty", IO_TIMEOUT) + .await + .expect("scripted output reaches focused buffer"); + + harness.write_all("\x11").expect("quit attached client"); + harness.wait().expect("client exits"); + server.shutdown().await.expect("shutdown server"); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn config_reload_updates_live_bindings_without_breaking_terminal_io() { + let _guard = test_lock().lock().await; + if !require_pty() { + return; + } + let server = TestServer::start().await.expect("start server"); + + run_cli(&server, ["new-session", "main"]); + run_cli( + &server, + [ + "new-window", + "-t", + "main", + "--title", + "shell", + "--", + "/bin/sh", + ], + ); + + let tempdir = tempfile::tempdir().expect("tempdir"); + let config_path = tempdir.path().join("config.rhai"); + fs::write( + &config_path, + r#"bind("normal", "", action.send_bytes_current("echo before-reload\r"))"#, + ) + .expect("write initial config"); + + let socket_arg = server.socket_path().to_string_lossy().into_owned(); + let socket_path = server.socket_path().to_path_buf(); + let config_arg = config_path.to_string_lossy().into_owned(); + let (_spawned, mut harness) = spawn_embers( + &[ + "attach", + "--socket", + &socket_arg, + "--config", + &config_arg, + "-t", + "main", + ], + socket_path, + ); + harness + .read_until_contains("[main]", STARTUP_TIMEOUT) + .expect("attach client renders"); + harness + .write_all("stty -echo\r") + .expect("disable shell echo in focused pane"); + harness + .wait_for_quiet(QUIET_TIMEOUT, IO_TIMEOUT) + .expect("focused shell settles"); + + harness.write_all("\x07").expect("trigger initial binding"); + harness + .read_until_contains("before-reload", IO_TIMEOUT) + .expect("initial binding renders"); + + fs::write( + &config_path, + r#"bind("normal", "", action.send_bytes_current("echo after-reload\r"))"#, + ) + .expect("write reloaded config"); + tokio::time::sleep(Duration::from_millis(700)).await; + + harness.write_all("\x07").expect("trigger reloaded binding"); + harness + .read_until_contains("after-reload", IO_TIMEOUT) + .expect("reloaded binding renders"); + + let output = run_pane_command(&mut harness, "echo still-live", "still-live"); + assert!( + output.contains("still-live"), + "regular terminal input must still work after reload:\n{output}" + ); + + harness.write_all("\x11").expect("quit attached client"); + harness.wait().expect("client exits"); + server.shutdown().await.expect("shutdown server"); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn live_pty_client_preserves_buffers_across_layout_and_attachment_changes() { + let _guard = test_lock().lock().await; + if !require_pty() { + return; + } + let server = TestServer::start().await.expect("start server"); + + run_cli(&server, ["new-session", "main"]); + run_cli( + &server, + [ + "new-window", + "-t", + "main", + "--title", + "shell", + "--", + "/bin/sh", + ], + ); + + let socket_arg = server.socket_path().to_string_lossy().into_owned(); + let socket_path = server.socket_path().to_path_buf(); + let (_spawned, mut harness) = spawn_embers( + &["attach", "--socket", &socket_arg, "-t", "main"], + socket_path, + ); + harness + .read_until_contains("[main]", STARTUP_TIMEOUT) + .expect("attach client renders"); + + let split = run_cli(&server, ["split-window", "--", "/bin/sh"]); + let moving_pane_id = stdout(&split) + .trim() + .parse::() + .expect("split-window returns new pane id"); + disable_echo_in_pane(&server, moving_pane_id); + + let mut connection = TestConnection::connect(server.socket_path()) + .await + .expect("connect protocol client"); + let snapshot = session_snapshot_by_name(&mut connection, "main").await; + let anchor_pane_id = snapshot + .nodes + .iter() + .filter(|node| node.buffer_view.is_some()) + .map(|node| node.id.0) + .find(|pane_id| *pane_id != moving_pane_id) + .expect("anchor pane exists"); + let moving_buffer_id = pane_buffer_id(&snapshot, moving_pane_id); + + run_cli( + &server, + [ + "send-keys", + "-t", + &moving_pane_id.to_string(), + "--enter", + "echo", + "split-live", + ], + ); + connection + .wait_for_capture_contains(moving_buffer_id, "split-live", IO_TIMEOUT) + .await + .expect("split pane keeps running"); + harness + .read_until_contains("split-live", IO_TIMEOUT) + .expect("split output renders in attached client"); + + let initial_order = split_child_order(&snapshot, anchor_pane_id, moving_pane_id); + let expected_order = if initial_order == [anchor_pane_id, moving_pane_id] { + run_cli( + &server, + [ + "node", + "move-before", + &moving_pane_id.to_string(), + &anchor_pane_id.to_string(), + ], + ); + [moving_pane_id, anchor_pane_id] + } else { + run_cli( + &server, + [ + "node", + "move-after", + &moving_pane_id.to_string(), + &anchor_pane_id.to_string(), + ], + ); + [anchor_pane_id, moving_pane_id] + }; + let _moved_snapshot = wait_for_split_child_order( + &mut connection, + "main", + anchor_pane_id, + moving_pane_id, + expected_order, + ) + .await; + + run_cli( + &server, + [ + "send-keys", + "-t", + &moving_pane_id.to_string(), + "--enter", + "echo", + "moved-live", + ], + ); + connection + .wait_for_capture_contains(moving_buffer_id, "moved-live", IO_TIMEOUT) + .await + .expect("moved pane keeps running"); + harness + .read_until_contains("moved-live", IO_TIMEOUT) + .expect("moved pane output still renders"); + + let response = connection + .request(&ClientMessage::Buffer(BufferRequest::Detach { + request_id: new_request_id(), + buffer_id: moving_buffer_id, + })) + .await + .expect("detach buffer succeeds"); + assert!( + matches!(response, ServerResponse::Ok(_)), + "expected ok response to buffer detach, got {response:?}" + ); + + send_buffer_input(&mut connection, moving_buffer_id, b"echo detached-live\r").await; + connection + .wait_for_capture_contains(moving_buffer_id, "detached-live", IO_TIMEOUT) + .await + .expect("detached buffer continues receiving output"); + + run_cli( + &server, + [ + "attach-buffer", + &moving_buffer_id.to_string(), + "-t", + &anchor_pane_id.to_string(), + ], + ); + let _reattached_snapshot = + wait_for_target_pane_buffer(&mut connection, "main", anchor_pane_id, moving_buffer_id) + .await; + + send_buffer_input(&mut connection, moving_buffer_id, b"echo reattach-live\r").await; + connection + .wait_for_capture_contains(moving_buffer_id, "reattach-live", IO_TIMEOUT) + .await + .expect("reattached buffer continues receiving output"); + harness + .read_until_contains("reattach-live", IO_TIMEOUT) + .expect("reattached buffer output renders in attached client"); + + harness.write_all("\x11").expect("quit attached client"); + harness.wait().expect("client exits"); + server.shutdown().await.expect("shutdown server"); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn hidden_buffer_bells_surface_in_the_attached_client_and_reveal_buffered_output() { + let _guard = test_lock().lock().await; + if !require_pty() { + return; + } + let server = TestServer::start().await.expect("start server"); + + run_cli(&server, ["new-session", "main"]); + run_cli( + &server, + [ + "new-window", + "-t", + "main", + "--title", + "shell", + "--", + "/bin/sh", + ], + ); + run_cli( + &server, + ["new-window", "-t", "main", "--title", "bg", "--", "/bin/sh"], + ); + run_cli(&server, ["select-window", "-t", "main:0"]); + + let mut connection = TestConnection::connect(server.socket_path()) + .await + .expect("connect protocol client"); + let snapshot = session_snapshot_by_name(&mut connection, "main").await; + let hidden_pane_id = root_tab_child_id(&snapshot, "bg").0; + let hidden_buffer_id = pane_buffer_id(&snapshot, hidden_pane_id); + + let socket_arg = server.socket_path().to_string_lossy().into_owned(); + let socket_path = server.socket_path().to_path_buf(); + let (_spawned, mut harness) = spawn_embers( + &["attach", "--socket", &socket_arg, "-t", "main"], + socket_path, + ); + harness + .read_until_contains("[main]", STARTUP_TIMEOUT) + .expect("attach client renders"); + + send_buffer_input( + &mut connection, + hidden_buffer_id, + b"printf 'hidden-bell\\a\\n'; sleep 0.5\r", + ) + .await; + connection + .wait_for_capture_contains(hidden_buffer_id, "hidden-bell", IO_TIMEOUT) + .await + .expect("hidden buffer output accumulates"); + wait_for_buffer_activity(&mut connection, hidden_buffer_id, ActivityState::Bell).await; + harness + .read_until_contains("!bg", IO_TIMEOUT) + .expect("hidden bell updates tab marker in attached client"); + + run_cli(&server, ["select-window", "-t", "main:bg"]); + harness + .read_until_contains("hidden-bell", IO_TIMEOUT) + .expect("revealed hidden buffer shows accumulated output"); + + harness.write_all("\x11").expect("quit attached client"); + harness.wait().expect("client exits"); + server.shutdown().await.expect("shutdown server"); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn fullscreen_terminal_transitions_render_in_the_live_client_pty() { + let _guard = test_lock().lock().await; + if !require_pty() { + return; + } + let server = TestServer::start().await.expect("start server"); + + run_cli(&server, ["new-session", "main"]); + run_cli( + &server, + [ + "new-window", + "-t", + "main", + "--title", + "shell", + "--", + "/bin/sh", + ], + ); + + let mut connection = TestConnection::connect(server.socket_path()) + .await + .expect("connect protocol client"); + let snapshot = session_snapshot_by_name(&mut connection, "main").await; + let buffer_id = pane_buffer_id(&snapshot, focused_pane_id(&snapshot)); + + let socket_arg = server.socket_path().to_string_lossy().into_owned(); + let socket_path = server.socket_path().to_path_buf(); + let (_spawned, mut harness) = spawn_embers( + &["attach", "--socket", &socket_arg, "-t", "main"], + socket_path, + ); + harness + .read_until_contains("[main]", STARTUP_TIMEOUT) + .expect("attach client renders"); + harness + .write_all("stty -echo\r") + .expect("disable shell echo in focused pane"); + harness + .wait_for_quiet(QUIET_TIMEOUT, IO_TIMEOUT) + .expect("focused shell settles"); + + harness + .write_all( + "printf '\\033[?1049h\\033[2J\\033[HPTY-FULLSCREEN'; sleep 1; printf '\\033[?1049lPTY-RESTORED\\n'\r", + ) + .expect("run fullscreen fixture"); + harness + .read_until_contains("PTY-FULLSCREEN", IO_TIMEOUT) + .expect("fullscreen output renders in live client"); + + let live = wait_for_visible_snapshot(&mut connection, buffer_id, |snapshot| { + snapshot.alternate_screen + && snapshot + .lines + .iter() + .any(|line| line.contains("PTY-FULLSCREEN")) + }) + .await; + assert!(live.alternate_screen); + + harness + .read_until_contains("PTY-RESTORED", IO_TIMEOUT) + .expect("primary screen restoration renders in live client"); + let restored = wait_for_visible_snapshot(&mut connection, buffer_id, |snapshot| { + !snapshot.alternate_screen + && snapshot + .lines + .iter() + .any(|line| line.contains("PTY-RESTORED")) + }) + .await; + assert!(!restored.alternate_screen); + + harness.write_all("\x11").expect("quit attached client"); + harness.wait().expect("client exits"); + server.shutdown().await.expect("shutdown server"); +} diff --git a/crates/embers-client/src/client.rs b/crates/embers-client/src/client.rs index 7552cc2c..06128ca4 100644 --- a/crates/embers-client/src/client.rs +++ b/crates/embers-client/src/client.rs @@ -107,9 +107,8 @@ where } pub async fn process_next_event(&mut self) -> Result { - let event = self.transport.next_event().await?; - self.state.apply_event(&event); - self.resync_for_event(&event).await?; + let event = self.next_event().await?; + self.handle_event(&event).await?; Ok(event) } @@ -125,15 +124,23 @@ where &mut self, timeout: std::time::Duration, ) -> Result> { - let event = match tokio::time::timeout(timeout, self.transport.next_event()).await { + let event = match tokio::time::timeout(timeout, self.next_event()).await { Ok(result) => result?, Err(_) => return Ok(None), }; - self.state.apply_event(&event); - self.resync_for_event(&event).await?; + self.handle_event(&event).await?; Ok(Some(event)) } + pub async fn next_event(&mut self) -> Result { + self.transport.next_event().await + } + + pub async fn handle_event(&mut self, event: &ServerEvent) -> Result<()> { + self.state.apply_event(event); + self.resync_for_event(event).await + } + pub async fn resync_session(&mut self, session_id: SessionId) -> Result<()> { let response = self .transport @@ -174,6 +181,26 @@ where } } + pub async fn refresh_buffer_record(&mut self, buffer_id: BufferId) -> Result<()> { + let response = self + .transport + .request(ClientMessage::Buffer(BufferRequest::Get { + request_id: self.next_request_id(), + buffer_id, + })) + .await?; + + match expect_response(response)? { + ServerResponse::Buffer(response) => { + self.state.apply_buffer_record(response.buffer); + Ok(()) + } + other => Err(MuxError::protocol(format!( + "expected buffer response, got {other:?}" + ))), + } + } + pub async fn capture_buffer(&self, buffer_id: BufferId) -> Result { let response = self .transport @@ -291,10 +318,12 @@ where } Ok(()) } + ServerEvent::RenderInvalidated(event) => { + self.refresh_buffer_record(event.buffer_id).await + } ServerEvent::BufferCreated(_) | ServerEvent::BufferDetached(_) - | ServerEvent::FocusChanged(_) - | ServerEvent::RenderInvalidated(_) => Ok(()), + | ServerEvent::FocusChanged(_) => Ok(()), } } diff --git a/crates/embers-client/src/config/loader.rs b/crates/embers-client/src/config/loader.rs index f6ec9adb..31dca858 100644 --- a/crates/embers-client/src/config/loader.rs +++ b/crates/embers-client/src/config/loader.rs @@ -101,6 +101,21 @@ impl ConfigManager { self.active_script = candidate_script; Ok(()) } + + pub fn reload_if_changed(&mut self) -> Result { + let candidate_source = load_config_source(&self.discovery)?; + if candidate_source == self.active_source { + return Ok(false); + } + + let candidate_script = match candidate_source.origin { + ConfigOrigin::BuiltIn => ScriptEngine::load(&candidate_source)?, + _ => ScriptEngine::load_with_overlay(BUILTIN_CONFIG_SOURCE, &candidate_source)?, + }; + self.active_source = candidate_source; + self.active_script = candidate_script; + Ok(true) + } } pub fn load_config_source(discovery: &ConfigDiscoveryOptions) -> ConfigResult { diff --git a/crates/embers-client/src/configured_client.rs b/crates/embers-client/src/configured_client.rs index 7e3e69a6..dc229001 100644 --- a/crates/embers-client/src/configured_client.rs +++ b/crates/embers-client/src/configured_client.rs @@ -325,8 +325,8 @@ where } pub async fn process_next_event(&mut self) -> Result { - let event = self.client.process_next_event().await?; - self.apply_processed_event(&event).await?; + let event = self.next_event().await?; + self.handle_event(&event).await?; Ok(event) } @@ -337,18 +337,23 @@ where let Some(event) = self.client.process_next_event_timeout(timeout).await? else { return Ok(None); }; - self.apply_processed_event(&event).await?; + self.handle_event(&event).await?; Ok(Some(event)) } - async fn apply_processed_event(&mut self, event: &ServerEvent) -> Result<()> { - if let ServerEvent::RenderInvalidated(event) = &event { + pub async fn next_event(&mut self) -> Result { + self.client.next_event().await + } + + pub async fn handle_event(&mut self, event: &ServerEvent) -> Result<()> { + self.client.handle_event(event).await?; + if let ServerEvent::RenderInvalidated(event) = event { self.client.refresh_buffer_snapshot(event.buffer_id).await?; } let session_id = self.event_session_id(event); let mut event_names = vec![event_name(event).to_owned()]; - if let ServerEvent::RenderInvalidated(render) = &event + if let ServerEvent::RenderInvalidated(render) = event && self .client .state() @@ -418,18 +423,37 @@ where self.config .reload() .map_err(|error| MuxError::invalid_input(error.to_string()))?; + self.finish_config_reload(¤t_mode); + Ok(()) + } + + pub fn reload_config_if_changed(&mut self) -> Result { + match self.config.reload_if_changed() { + Ok(false) => Ok(false), + Ok(true) => { + let current_mode = self.input_state.current_mode().to_owned(); + self.finish_config_reload(¤t_mode); + Ok(true) + } + Err(error) => { + self.record_notification(error.to_string()); + Ok(false) + } + } + } + + fn finish_config_reload(&mut self, current_mode: &str) { if self .config .active_script() .loaded_config() .modes - .contains_key(¤t_mode) + .contains_key(current_mode) { self.input_state.clear_pending(); } else { self.input_state.set_mode(NORMAL_MODE); } - Ok(()) } async fn execute_actions( @@ -1437,6 +1461,9 @@ where presentation: &PresentationModel, actions: &[Action], ) -> bool { + // When the focused pane is acting like a live terminal surface, prefer + // forwarding local search/select/navigation bindings to the program + // instead of stealing keys that fullscreen apps expect to receive. let Some(leaf) = presentation.focused_leaf() else { return false; }; diff --git a/crates/embers-client/src/input/keymap.rs b/crates/embers-client/src/input/keymap.rs index 6cf8f557..5a99b99a 100644 --- a/crates/embers-client/src/input/keymap.rs +++ b/crates/embers-client/src/input/keymap.rs @@ -144,6 +144,52 @@ mod tests { ); } + #[test] + fn prefix_match_keeps_pending_sequence_until_it_resolves() { + let mut state = InputState::default(); + let bindings = bindings(&[("normal", "ab", "target")]); + let modes = builtin_modes(); + + assert_eq!( + resolve_key(&bindings, &modes, &mut state, KeyToken::Char('a')), + InputResolution::PrefixMatch + ); + assert_eq!(state.pending_sequence(), &[KeyToken::Char('a')]); + + assert_eq!( + resolve_key(&bindings, &modes, &mut state, KeyToken::Char('b')), + InputResolution::ExactMatch(super::BindingMatch { + mode: "normal".to_owned(), + sequence: vec![KeyToken::Char('a'), KeyToken::Char('b')], + target: "target".to_owned(), + }) + ); + assert!(state.pending_sequence().is_empty()); + } + + #[test] + fn unmatched_sequence_clears_pending_state_after_prefix_miss() { + let mut state = InputState::default(); + let bindings = bindings(&[("normal", "ab", "target")]); + let modes = builtin_modes(); + + assert_eq!( + resolve_key(&bindings, &modes, &mut state, KeyToken::Char('a')), + InputResolution::PrefixMatch + ); + assert_eq!(state.pending_sequence(), &[KeyToken::Char('a')]); + + assert_eq!( + resolve_key(&bindings, &modes, &mut state, KeyToken::Char('x')), + InputResolution::Unmatched { + mode: "normal".to_owned(), + sequence: vec![KeyToken::Char('a'), KeyToken::Char('x')], + fallback_policy: FallbackPolicy::Passthrough, + } + ); + assert!(state.pending_sequence().is_empty()); + } + fn bindings(entries: &[(&str, &str, &str)]) -> BTreeMap>> { let mut bindings = BTreeMap::>>::new(); for (mode, sequence, target) in entries { diff --git a/crates/embers-client/src/state.rs b/crates/embers-client/src/state.rs index 62a27af4..7ec13342 100644 --- a/crates/embers-client/src/state.rs +++ b/crates/embers-client/src/state.rs @@ -164,6 +164,10 @@ impl ClientState { } } + pub fn apply_buffer_record(&mut self, buffer: BufferRecord) { + self.buffers.insert(buffer.id, buffer); + } + pub fn apply_buffer_snapshot(&mut self, snapshot: VisibleSnapshotResponse) { if let Some(buffer) = self.buffers.get_mut(&snapshot.buffer_id) { buffer.last_snapshot_seq = snapshot.sequence; diff --git a/crates/embers-client/tests/configured_client.rs b/crates/embers-client/tests/configured_client.rs index 02a77478..f5e529e3 100644 --- a/crates/embers-client/tests/configured_client.rs +++ b/crates/embers-client/tests/configured_client.rs @@ -8,16 +8,18 @@ use embers_client::{ }; use embers_core::{ActivityState, BufferId, NodeId, PtySize, RequestId, SessionId, Size}; use embers_protocol::{ - BufferCreatedEvent, BufferRecord, BufferRecordKind, BufferRecordState, BufferViewRecord, - ClientChangedEvent, ClientMessage, ClientRecord, ClientRequest, ClientResponse, - FocusChangedEvent, InputRequest, NodeRecord, NodeRecordKind, NodeRequest, OkResponse, - RenderInvalidatedEvent, ScrollbackSliceResponse, ServerEvent, ServerResponse, SessionRecord, - SessionRequest, SessionSnapshot, SessionSnapshotResponse, SnapshotResponse, + BufferCreatedEvent, BufferRecord, BufferRecordKind, BufferRecordState, BufferResponse, + BufferViewRecord, ClientChangedEvent, ClientMessage, ClientRecord, ClientRequest, + ClientResponse, FocusChangedEvent, InputRequest, NodeRecord, NodeRecordKind, NodeRequest, + OkResponse, RenderInvalidatedEvent, ScrollbackSliceResponse, ServerEvent, ServerResponse, + SessionRecord, SessionRequest, SessionSnapshot, SessionSnapshotResponse, SnapshotResponse, VisibleSnapshotResponse, }; use tempfile::tempdir; -use crate::support::{FOCUSED_LEAF_ID, LEFT_LEAF_ID, SESSION_ID, demo_state, root_focus_state}; +use crate::support::{ + FOCUSED_BUFFER_ID, FOCUSED_LEAF_ID, LEFT_LEAF_ID, SESSION_ID, demo_state, root_focus_state, +}; const SECOND_SESSION_ID: SessionId = SessionId(2); const SECOND_ROOT_ID: NodeId = NodeId(200); @@ -85,6 +87,17 @@ fn visible_snapshot_from_state( snapshot } +fn buffer_response_from_state( + state: &embers_client::ClientState, + buffer_id: BufferId, + request_id: RequestId, +) -> BufferResponse { + BufferResponse { + request_id, + buffer: state.buffers.get(&buffer_id).unwrap().clone(), + } +} + fn scrollback_slice_response( buffer_id: BufferId, request_id: RequestId, @@ -117,6 +130,23 @@ fn snapshot_response( } } +fn push_send_input_refresh_responses( + transport: &FakeTransport, + state: &embers_client::ClientState, + buffer_id: BufferId, +) { + transport.push_response(ServerResponse::Ok(OkResponse { + request_id: RequestId(1), + })); + transport.push_response(ServerResponse::VisibleSnapshot( + visible_snapshot_from_state(state, buffer_id, RequestId(2)), + )); + transport.push_response(ServerResponse::SessionSnapshot(SessionSnapshotResponse { + request_id: RequestId(3), + snapshot: session_snapshot_from_state(state, SESSION_ID), + })); +} + fn second_session_state() -> embers_client::ClientState { let mut state = demo_state(); state.sessions.insert( @@ -253,6 +283,169 @@ async fn configured_keybinding_executes_live_focus_action() { transport.assert_exhausted().unwrap(); } +#[tokio::test] +async fn unmapped_keys_forward_to_the_focused_buffer_in_normal_mode() { + let transport = FakeTransport::default(); + let state = demo_state(); + push_send_input_refresh_responses(&transport, &state, FOCUSED_BUFFER_ID); + + let client = MuxClient::new(transport.clone()); + let (config, _tempdir) = manager_from_source(""); + let mut configured = ConfiguredClient::new(client, config); + *configured.client_mut().state_mut() = state; + + configured + .handle_key( + SESSION_ID, + Size { + width: 80, + height: 20, + }, + KeyEvent::Char('x'), + ) + .await + .unwrap(); + + assert_eq!( + transport.requests()[0], + ClientMessage::Input(InputRequest::Send { + request_id: RequestId(1), + buffer_id: FOCUSED_BUFFER_ID, + bytes: b"x".to_vec(), + }) + ); +} + +#[tokio::test] +async fn leader_prefix_waits_without_forwarding_input() { + let client = MuxClient::new(FakeTransport::default()); + let (config, _tempdir) = manager_from_source( + r#" + fn open_workspace_split(ctx) { action.notify("info", "workspace-split") } + define_action("workspace-split", open_workspace_split); + set_leader(""); + bind("normal", "ws", "workspace-split"); + "#, + ); + let mut configured = ConfiguredClient::new(client, config); + *configured.client_mut().state_mut() = demo_state(); + + configured + .handle_key( + SESSION_ID, + Size { + width: 80, + height: 20, + }, + KeyEvent::Ctrl('a'), + ) + .await + .unwrap(); + configured + .handle_key( + SESSION_ID, + Size { + width: 80, + height: 20, + }, + KeyEvent::Char('w'), + ) + .await + .unwrap(); + + assert!(configured.client().transport().requests().is_empty()); + assert!(configured.notifications().is_empty()); + + configured + .handle_key( + SESSION_ID, + Size { + width: 80, + height: 20, + }, + KeyEvent::Char('s'), + ) + .await + .unwrap(); + + assert!(configured.client().transport().requests().is_empty()); + assert_eq!(configured.notifications(), ["workspace-split"]); +} + +#[tokio::test] +async fn reload_clears_pending_prefix_before_next_unmapped_key() { + let transport = FakeTransport::default(); + let state = demo_state(); + push_send_input_refresh_responses(&transport, &state, FOCUSED_BUFFER_ID); + + let tempdir = tempdir().unwrap(); + let config_path = tempdir.path().join("config.rhai"); + fs::write( + &config_path, + r#" + fn open_workspace_split(ctx) { action.notify("info", "workspace-split") } + define_action("workspace-split", open_workspace_split); + set_leader(""); + bind("normal", "ws", "workspace-split"); + "#, + ) + .unwrap(); + let config = ConfigManager::load( + ConfigDiscoveryOptions::default().with_project_config_dir(tempdir.path()), + ) + .unwrap(); + let client = MuxClient::new(transport.clone()); + let mut configured = ConfiguredClient::new(client, config); + *configured.client_mut().state_mut() = state; + + configured + .handle_key( + SESSION_ID, + Size { + width: 80, + height: 20, + }, + KeyEvent::Ctrl('a'), + ) + .await + .unwrap(); + configured + .handle_key( + SESSION_ID, + Size { + width: 80, + height: 20, + }, + KeyEvent::Char('w'), + ) + .await + .unwrap(); + assert!(transport.requests().is_empty()); + + configured.reload_config().unwrap(); + + configured + .handle_key( + SESSION_ID, + Size { + width: 80, + height: 20, + }, + KeyEvent::Char('x'), + ) + .await + .unwrap(); + + assert_eq!( + transport.requests()[0], + ClientMessage::Input(InputRequest::Send { + request_id: RequestId(1), + buffer_id: FOCUSED_BUFFER_ID, + bytes: b"x".to_vec(), + }) + ); +} + #[tokio::test] async fn configured_render_uses_scripted_tab_bars() { let client = MuxClient::new(FakeTransport::default()); @@ -728,6 +921,47 @@ async fn select_mode_yanks_selection_to_osc52() { ); } +#[tokio::test] +async fn copy_mode_blocks_unmapped_passthrough() { + let transport = FakeTransport::default(); + let client = MuxClient::new(transport.clone()); + let (config, _tempdir) = manager_from_source( + r#" + fn enter_copy(ctx) { action.enter_mode("copy") } + define_action("enter-copy", enter_copy); + unbind("normal", "v"); + bind("normal", "v", "enter-copy"); + "#, + ); + let mut configured = ConfiguredClient::new(client, config); + *configured.client_mut().state_mut() = demo_state(); + + configured + .handle_key( + SESSION_ID, + Size { + width: 80, + height: 20, + }, + KeyEvent::Char('v'), + ) + .await + .unwrap(); + configured + .handle_key( + SESSION_ID, + Size { + width: 80, + height: 20, + }, + KeyEvent::Char('x'), + ) + .await + .unwrap(); + + assert!(transport.requests().is_empty()); +} + #[tokio::test] async fn wheel_mouse_events_scroll_locally_or_forward_to_program() { let mut initial_state = demo_state(); @@ -951,8 +1185,13 @@ async fn render_invalidated_events_use_their_buffer_session_context() { transport.push_event(ServerEvent::RenderInvalidated(RenderInvalidatedEvent { buffer_id: SECOND_BUFFER_ID, })); + transport.push_response(ServerResponse::Buffer(buffer_response_from_state( + &state, + SECOND_BUFFER_ID, + RequestId(1), + ))); transport.push_response(ServerResponse::VisibleSnapshot( - visible_snapshot_from_state(&state, SECOND_BUFFER_ID, RequestId(1)), + visible_snapshot_from_state(&state, SECOND_BUFFER_ID, RequestId(2)), )); let client = MuxClient::new(transport); let (config, _tempdir) = manager_from_source( @@ -980,6 +1219,184 @@ async fn render_invalidated_events_use_their_buffer_session_context() { assert_eq!(configured.notifications(), ["other"]); } +#[tokio::test] +async fn render_invalidated_events_refresh_buffer_activity_before_bell_hooks() { + let mut state = second_session_state(); + state.buffers.get_mut(&SECOND_BUFFER_ID).unwrap().activity = ActivityState::Bell; + + let transport = FakeTransport::default(); + transport.push_event(ServerEvent::RenderInvalidated(RenderInvalidatedEvent { + buffer_id: SECOND_BUFFER_ID, + })); + transport.push_response(ServerResponse::Buffer(buffer_response_from_state( + &state, + SECOND_BUFFER_ID, + RequestId(1), + ))); + transport.push_response(ServerResponse::VisibleSnapshot( + visible_snapshot_from_state(&state, SECOND_BUFFER_ID, RequestId(2)), + )); + let client = MuxClient::new(transport.clone()); + let (config, _tempdir) = manager_from_source( + r#" + fn on_bell(ctx) { action.notify("info", ctx.current_session().name()) } + on("buffer_bell", on_bell); + "#, + ); + let mut configured = ConfiguredClient::new(client, config); + *configured.client_mut().state_mut() = second_session_state(); + + let event = configured.process_next_event().await.unwrap(); + + assert!(matches!(event, ServerEvent::RenderInvalidated(_))); + assert_eq!(configured.notifications(), ["other"]); + assert_eq!( + transport.requests(), + vec![ + ClientMessage::Buffer(embers_protocol::BufferRequest::Get { + request_id: RequestId(1), + buffer_id: SECOND_BUFFER_ID, + }), + ClientMessage::Buffer(embers_protocol::BufferRequest::CaptureVisible { + request_id: RequestId(2), + buffer_id: SECOND_BUFFER_ID, + }), + ] + ); +} + +#[tokio::test] +async fn render_session_refreshes_invalidated_snapshot_before_rendering_title_and_content() { + let transport = FakeTransport::default(); + let mut stale_state = demo_state(); + stale_state.apply_event(&ServerEvent::RenderInvalidated(RenderInvalidatedEvent { + buffer_id: FOCUSED_BUFFER_ID, + })); + + let mut refreshed_state = demo_state(); + let snapshot = refreshed_state + .snapshots + .get_mut(&FOCUSED_BUFFER_ID) + .unwrap(); + snapshot.lines = vec!["fresh render line".to_owned()]; + snapshot.title = Some("fresh-title".to_owned()); + + transport.push_response(ServerResponse::VisibleSnapshot( + visible_snapshot_from_state(&refreshed_state, FOCUSED_BUFFER_ID, RequestId(1)), + )); + + let client = MuxClient::new(transport.clone()); + let (config, _tempdir) = manager_from_source(""); + let mut configured = ConfiguredClient::new(client, config); + *configured.client_mut().state_mut() = stale_state; + + let grid = configured + .render_session( + SESSION_ID, + Size { + width: 80, + height: 20, + }, + ) + .await + .unwrap(); + let rendered = grid.render(); + let presentation = PresentationModel::project( + configured.client().state(), + SESSION_ID, + Size { + width: 80, + height: 20, + }, + ) + .expect("projection succeeds"); + + assert!(rendered.contains("fresh render line")); + assert!(!rendered.contains("logs visible")); + assert_eq!( + configured + .client() + .state() + .buffers + .get(&FOCUSED_BUFFER_ID) + .expect("focused buffer") + .title, + "fresh-title" + ); + assert_eq!( + presentation.focused_leaf().expect("focused leaf").title, + "fresh-title" + ); + assert!(configured.client().state().invalidated_buffers.is_empty()); + assert_eq!( + transport.requests(), + vec![ClientMessage::Buffer( + embers_protocol::BufferRequest::CaptureVisible { + request_id: RequestId(1), + buffer_id: FOCUSED_BUFFER_ID, + } + )] + ); +} + +#[tokio::test] +async fn render_session_replaces_stale_scrolled_cache_when_snapshot_switches_to_alternate_screen() { + let transport = FakeTransport::default(); + let mut stale_state = demo_state(); + let view = stale_state + .view_state_mut(FOCUSED_LEAF_ID) + .expect("focused view state"); + view.follow_output = false; + view.scroll_top_line = 12; + view.total_line_count = 60; + view.visible_lines = vec!["stale scrolled line".to_owned()]; + stale_state.apply_event(&ServerEvent::RenderInvalidated(RenderInvalidatedEvent { + buffer_id: FOCUSED_BUFFER_ID, + })); + + let mut refreshed_state = demo_state(); + let snapshot = refreshed_state + .snapshots + .get_mut(&FOCUSED_BUFFER_ID) + .unwrap(); + snapshot.lines = vec!["alternate screen live".to_owned()]; + snapshot.alternate_screen = true; + snapshot.viewport_top_line = 0; + snapshot.total_lines = 24; + + transport.push_response(ServerResponse::VisibleSnapshot( + visible_snapshot_from_state(&refreshed_state, FOCUSED_BUFFER_ID, RequestId(1)), + )); + + let client = MuxClient::new(transport); + let (config, _tempdir) = manager_from_source(""); + let mut configured = ConfiguredClient::new(client, config); + *configured.client_mut().state_mut() = stale_state; + + let grid = configured + .render_session( + SESSION_ID, + Size { + width: 80, + height: 20, + }, + ) + .await + .unwrap(); + let rendered = grid.render(); + + assert!(rendered.contains("alternate screen live")); + assert!(!rendered.contains("stale scrolled line")); + assert!(!rendered.contains("13/60")); + let view = configured + .client() + .state() + .view_state(FOCUSED_LEAF_ID) + .expect("focused view state"); + assert!(view.alternate_screen); + assert_eq!(view.visible_lines, vec!["alternate screen live".to_owned()]); +} + #[tokio::test] async fn detached_buffer_events_do_not_fall_back_to_the_active_session() { let transport = FakeTransport::default(); @@ -1138,6 +1555,85 @@ async fn event_hook_executes_real_actions() { transport.assert_exhausted().unwrap(); } +#[tokio::test] +async fn scripted_send_keys_current_forwards_to_the_focused_buffer() { + let transport = FakeTransport::default(); + let state = demo_state(); + push_send_input_refresh_responses(&transport, &state, FOCUSED_BUFFER_ID); + + let client = MuxClient::new(transport.clone()); + let (config, _tempdir) = manager_from_source( + r#" + fn send_current(ctx) { action.send_keys_current("abc") } + define_action("send-current", send_current); + bind("normal", "", "send-current"); + "#, + ); + let mut configured = ConfiguredClient::new(client, config); + *configured.client_mut().state_mut() = state; + + configured + .handle_key( + SESSION_ID, + Size { + width: 80, + height: 20, + }, + KeyEvent::Ctrl('g'), + ) + .await + .unwrap(); + + assert_eq!( + transport.requests()[0], + ClientMessage::Input(InputRequest::Send { + request_id: RequestId(1), + buffer_id: FOCUSED_BUFFER_ID, + bytes: b"abc".to_vec(), + }) + ); +} + +#[tokio::test] +async fn scripted_send_bytes_can_target_a_specific_buffer() { + let transport = FakeTransport::default(); + let state = demo_state(); + let target_buffer_id = BufferId(5); + push_send_input_refresh_responses(&transport, &state, target_buffer_id); + + let client = MuxClient::new(transport.clone()); + let (config, _tempdir) = manager_from_source( + r#" + fn send_popup(ctx) { action.send_bytes(5, "popup") } + define_action("send-popup", send_popup); + bind("normal", "", "send-popup"); + "#, + ); + let mut configured = ConfiguredClient::new(client, config); + *configured.client_mut().state_mut() = state; + + configured + .handle_key( + SESSION_ID, + Size { + width: 80, + height: 20, + }, + KeyEvent::Ctrl('p'), + ) + .await + .unwrap(); + + assert_eq!( + transport.requests()[0], + ClientMessage::Input(InputRequest::Send { + request_id: RequestId(1), + buffer_id: target_buffer_id, + bytes: b"popup".to_vec(), + }) + ); +} + #[tokio::test] async fn keybinding_runtime_errors_become_notifications() { let client = MuxClient::new(FakeTransport::default()); diff --git a/crates/embers-client/tests/e2e.rs b/crates/embers-client/tests/e2e.rs index 48c227ce..38cfc518 100644 --- a/crates/embers-client/tests/e2e.rs +++ b/crates/embers-client/tests/e2e.rs @@ -6,8 +6,8 @@ use embers_core::{ ActivityState, BufferId, FloatGeometry, NodeId, SessionId, Size, SplitDirection, new_request_id, }; use embers_protocol::{ - BufferRequest, BufferResponse, BuffersResponse, ClientMessage, FloatingRequest, NodeRequest, - ServerResponse, SessionRequest, SessionSnapshot, + BufferRecord, BufferRequest, BufferResponse, BuffersResponse, ClientMessage, FloatingRequest, + NodeRequest, ServerResponse, SessionRequest, SessionSnapshot, VisibleSnapshotResponse, }; use embers_test_support::{TestConnection, TestServer, cargo_bin}; @@ -48,12 +48,20 @@ async fn create_session(connection: &mut TestConnection, name: &str) -> SessionS async fn create_buffer( connection: &mut TestConnection, title: &str, +) -> embers_protocol::BufferRecord { + create_buffer_with_command(connection, title, vec!["/bin/sh".to_owned()]).await +} + +async fn create_buffer_with_command( + connection: &mut TestConnection, + title: &str, + command: Vec, ) -> embers_protocol::BufferRecord { let response = connection .request(&ClientMessage::Buffer(BufferRequest::Create { request_id: new_request_id(), title: Some(title.to_owned()), - command: vec!["/bin/sh".to_owned()], + command, cwd: None, env: Default::default(), })) @@ -144,6 +152,164 @@ async fn render_session( Renderer.render(client.state(), &model).render() } +async fn wait_for_visible_snapshot( + connection: &mut TestConnection, + buffer_id: BufferId, + timeout: Duration, + mut predicate: F, +) -> VisibleSnapshotResponse +where + F: FnMut(&VisibleSnapshotResponse) -> bool, +{ + let deadline = tokio::time::Instant::now() + timeout; + + loop { + let snapshot = connection + .capture_visible_buffer(buffer_id) + .await + .expect("visible capture succeeds"); + if predicate(&snapshot) { + return snapshot; + } + + if tokio::time::Instant::now() >= deadline { + panic!("timed out waiting for visible snapshot; last snapshot: {snapshot:?}"); + } + + tokio::time::sleep(Duration::from_millis(50)).await; + } +} + +async fn buffer_record(connection: &mut TestConnection, buffer_id: BufferId) -> BufferRecord { + match connection + .request(&ClientMessage::Buffer(BufferRequest::Get { + request_id: new_request_id(), + buffer_id, + })) + .await + .expect("get buffer succeeds") + { + ServerResponse::Buffer(BufferResponse { buffer, .. }) => buffer, + other => panic!("expected buffer response, got {other:?}"), + } +} + +async fn wait_for_buffer_activity( + connection: &mut TestConnection, + buffer_id: BufferId, + expected: ActivityState, + timeout: Duration, +) -> BufferRecord { + let deadline = tokio::time::Instant::now() + timeout; + + loop { + let buffer = buffer_record(connection, buffer_id).await; + if buffer.activity == expected { + return buffer; + } + + if tokio::time::Instant::now() >= deadline { + panic!( + "timed out waiting for buffer {buffer_id} activity {expected:?}; last activity: {:?}", + buffer.activity + ); + } + + tokio::time::sleep(Duration::from_millis(50)).await; + } +} + +struct HiddenTabFixture { + nested_tabs_id: NodeId, + hidden_buffer: BufferRecord, +} + +async fn create_hidden_tab_fixture(connection: &mut TestConnection) -> HiddenTabFixture { + let hidden_buffer = create_buffer(connection, "hidden").await; + create_hidden_tab_fixture_with_buffer(connection, hidden_buffer).await +} + +async fn create_hidden_tab_fixture_with_buffer( + connection: &mut TestConnection, + hidden_buffer: BufferRecord, +) -> HiddenTabFixture { + let session = create_session(connection, "alpha").await; + let buffer_a = create_buffer(connection, "main").await; + let session = match connection + .request(&ClientMessage::Session(SessionRequest::AddRootTab { + request_id: new_request_id(), + session_id: session.session.id, + title: "main".to_owned(), + buffer_id: Some(buffer_a.id), + child_node_id: None, + })) + .await + .expect("add root tab succeeds") + { + ServerResponse::SessionSnapshot(response) => response.snapshot, + other => panic!("expected session snapshot response, got {other:?}"), + }; + let main_leaf = session.session.focused_leaf_id.expect("main leaf exists"); + + let session = match connection + .request(&ClientMessage::Node(NodeRequest::WrapInTabs { + request_id: new_request_id(), + node_id: main_leaf, + title: "main".to_owned(), + })) + .await + .expect("wrap main leaf in tabs succeeds") + { + ServerResponse::SessionSnapshot(response) => response.snapshot, + other => panic!("expected session snapshot response, got {other:?}"), + }; + let nested_tabs_id = node(&session, main_leaf) + .parent_id + .expect("wrapped main leaf has tabs parent"); + + let _ = connection + .request(&ClientMessage::Node(NodeRequest::AddTab { + request_id: new_request_id(), + tabs_node_id: nested_tabs_id, + title: "bg".to_owned(), + buffer_id: Some(hidden_buffer.id), + child_node_id: None, + index: 1, + })) + .await + .expect("add hidden tab succeeds"); + let _ = connection + .request(&ClientMessage::Node(NodeRequest::SelectTab { + request_id: new_request_id(), + tabs_node_id: nested_tabs_id, + index: 0, + })) + .await + .expect("select visible tab succeeds"); + + HiddenTabFixture { + nested_tabs_id, + hidden_buffer, + } +} + +fn fullscreen_fixture_command( + live_title: &str, + restored_title: &str, + sleep_secs: &str, +) -> Vec { + vec![ + "/bin/sh".to_owned(), + "-lc".to_owned(), + format!( + "printf 'main-before\\n'; \ + printf '\\033]0;{live_title}\\007\\033[?1049h\\033[2J\\033[Hfullscreen-live\\033[3;10Hcursor-target'; \ + sleep {sleep_secs}; \ + printf '\\033]0;{restored_title}\\007\\033[?1049lrestored-after\\n'" + ), + ] +} + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn basic_cli_workflow_renders_split_output() { let server = TestServer::start().await.expect("server starts"); @@ -463,6 +629,52 @@ async fn move_and_detach_workflows_preserve_running_buffers() { })) .await .expect("detach succeeds"); + assert!( + buffer_record(&mut connection, buffer_a.id) + .await + .attachment_node_id + .is_none() + ); + + let _ = connection + .request(&ClientMessage::Input(embers_protocol::InputRequest::Send { + request_id: new_request_id(), + buffer_id: buffer_a.id, + bytes: b"printf detached-output\\n\r".to_vec(), + })) + .await + .expect("send detached output succeeds"); + connection + .wait_for_capture_contains(buffer_a.id, "detached-output", Duration::from_secs(3)) + .await + .expect("detached buffer captures output"); + wait_for_buffer_activity( + &mut connection, + buffer_a.id, + ActivityState::Activity, + Duration::from_secs(3), + ) + .await; + + let _ = connection + .request(&ClientMessage::Input(embers_protocol::InputRequest::Send { + request_id: new_request_id(), + buffer_id: buffer_a.id, + bytes: b"printf 'detached-bell\\a\\n'; sleep 0.5\r".to_vec(), + })) + .await + .expect("send detached bell succeeds"); + connection + .wait_for_capture_contains(buffer_a.id, "detached-bell", Duration::from_secs(3)) + .await + .expect("detached buffer captures bell marker"); + wait_for_buffer_activity( + &mut connection, + buffer_a.id, + ActivityState::Bell, + Duration::from_secs(3), + ) + .await; let popup = match connection .request(&ClientMessage::Floating(FloatingRequest::Create { @@ -481,6 +693,10 @@ async fn move_and_detach_workflows_preserve_running_buffers() { ServerResponse::Floating(response) => response.floating, other => panic!("expected floating response, got {other:?}"), }; + assert_eq!( + buffer_record(&mut connection, buffer_a.id).await.activity, + ActivityState::Idle + ); let _ = connection .request(&ClientMessage::Input(embers_protocol::InputRequest::Send { request_id: new_request_id(), @@ -566,73 +782,31 @@ async fn hidden_activity_is_visible_and_reconnect_rehydrates_state() { .await .expect("protocol connection"); - let session = create_session(&mut connection, "alpha").await; - let buffer_a = create_buffer(&mut connection, "main").await; - let session = match connection - .request(&ClientMessage::Session(SessionRequest::AddRootTab { - request_id: new_request_id(), - session_id: session.session.id, - title: "main".to_owned(), - buffer_id: Some(buffer_a.id), - child_node_id: None, - })) - .await - .expect("add root tab succeeds") - { - ServerResponse::SessionSnapshot(response) => response.snapshot, - other => panic!("expected session snapshot response, got {other:?}"), - }; - let main_leaf = session.session.focused_leaf_id.expect("main leaf exists"); - - let session = match connection - .request(&ClientMessage::Node(NodeRequest::WrapInTabs { - request_id: new_request_id(), - node_id: main_leaf, - title: "main".to_owned(), - })) - .await - .expect("wrap main leaf in tabs succeeds") - { - ServerResponse::SessionSnapshot(response) => response.snapshot, - other => panic!("expected session snapshot response, got {other:?}"), - }; - let nested_tabs_id = node(&session, main_leaf) - .parent_id - .expect("wrapped main leaf has tabs parent"); - - let buffer_b = create_buffer(&mut connection, "hidden").await; - let _ = connection - .request(&ClientMessage::Node(NodeRequest::AddTab { - request_id: new_request_id(), - tabs_node_id: nested_tabs_id, - title: "bg".to_owned(), - buffer_id: Some(buffer_b.id), - child_node_id: None, - index: 1, - })) - .await - .expect("add hidden tab succeeds"); - let _ = connection - .request(&ClientMessage::Node(NodeRequest::SelectTab { - request_id: new_request_id(), - tabs_node_id: nested_tabs_id, - index: 0, - })) - .await - .expect("select visible tab succeeds"); + let fixture = create_hidden_tab_fixture(&mut connection).await; let _ = connection .request(&ClientMessage::Input(embers_protocol::InputRequest::Send { request_id: new_request_id(), - buffer_id: buffer_b.id, + buffer_id: fixture.hidden_buffer.id, bytes: b"printf hidden-activity\\n\r".to_vec(), })) .await .expect("send to hidden buffer succeeds"); connection - .wait_for_capture_contains(buffer_b.id, "hidden-activity", Duration::from_secs(3)) + .wait_for_capture_contains( + fixture.hidden_buffer.id, + "hidden-activity", + Duration::from_secs(3), + ) .await .expect("hidden buffer captures output"); + wait_for_buffer_activity( + &mut connection, + fixture.hidden_buffer.id, + ActivityState::Activity, + Duration::from_secs(3), + ) + .await; let mut first_client = MuxClient::connect(server.socket_path()) .await @@ -657,12 +831,12 @@ async fn hidden_activity_is_visible_and_reconnect_rehydrates_state() { let tabs = model .tab_bars .iter() - .find(|tabs| tabs.node_id == nested_tabs_id) + .find(|tabs| tabs.node_id == fixture.nested_tabs_id) .expect("nested tabs frame exists"); if tabs .tabs .iter() - .any(|tab| tab.title == "bg" && tab.activity != ActivityState::Idle) + .any(|tab| tab.title == "bg" && tab.activity == ActivityState::Activity) { saw_hidden_activity = true; break; @@ -679,11 +853,18 @@ async fn hidden_activity_is_visible_and_reconnect_rehydrates_state() { let _ = connection .request(&ClientMessage::Node(NodeRequest::SelectTab { request_id: new_request_id(), - tabs_node_id: nested_tabs_id, + tabs_node_id: fixture.nested_tabs_id, index: 1, })) .await .expect("select hidden tab succeeds"); + wait_for_buffer_activity( + &mut connection, + fixture.hidden_buffer.id, + ActivityState::Idle, + Duration::from_secs(3), + ) + .await; let mut second_client = MuxClient::connect(server.socket_path()) .await @@ -693,3 +874,303 @@ async fn hidden_activity_is_visible_and_reconnect_rehydrates_state() { server.shutdown().await.expect("server shuts down"); } + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn hidden_bell_is_visible_to_clients_until_revealed() { + let server = TestServer::start().await.expect("server starts"); + let mut connection = TestConnection::connect(server.socket_path()) + .await + .expect("protocol connection"); + + let fixture = create_hidden_tab_fixture(&mut connection).await; + + let _ = connection + .request(&ClientMessage::Input(embers_protocol::InputRequest::Send { + request_id: new_request_id(), + buffer_id: fixture.hidden_buffer.id, + bytes: b"printf 'hidden-bell\\a\\n'; sleep 0.5\r".to_vec(), + })) + .await + .expect("send hidden bell succeeds"); + connection + .wait_for_capture_contains( + fixture.hidden_buffer.id, + "hidden-bell", + Duration::from_secs(3), + ) + .await + .expect("hidden bell marker appears"); + wait_for_buffer_activity( + &mut connection, + fixture.hidden_buffer.id, + ActivityState::Bell, + Duration::from_secs(3), + ) + .await; + + let mut client = MuxClient::connect(server.socket_path()) + .await + .expect("client connects"); + client + .resync_all_sessions() + .await + .expect("client resyncs sessions"); + refresh_all_snapshots(&mut client).await; + let session_id = session_id_by_name(&client, "alpha"); + let model = PresentationModel::project( + client.state(), + session_id, + Size { + width: 80, + height: 24, + }, + ) + .expect("projection succeeds"); + let tabs = model + .tab_bars + .iter() + .find(|tabs| tabs.node_id == fixture.nested_tabs_id) + .expect("nested tabs frame exists"); + assert!( + tabs.tabs + .iter() + .any(|tab| tab.title == "bg" && tab.activity == ActivityState::Bell) + ); + + let _ = connection + .request(&ClientMessage::Node(NodeRequest::SelectTab { + request_id: new_request_id(), + tabs_node_id: fixture.nested_tabs_id, + index: 1, + })) + .await + .expect("select hidden tab succeeds"); + wait_for_buffer_activity( + &mut connection, + fixture.hidden_buffer.id, + ActivityState::Idle, + Duration::from_secs(3), + ) + .await; + + server.shutdown().await.expect("server shuts down"); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn fullscreen_fixture_enters_alternate_screen_and_restores_primary_screen() { + let server = TestServer::start().await.expect("server starts"); + let mut connection = TestConnection::connect(server.socket_path()) + .await + .expect("protocol connection"); + + let session = create_session(&mut connection, "alpha").await; + let buffer = create_buffer_with_command( + &mut connection, + "fullscreen", + fullscreen_fixture_command("fullscreen-live-title", "primary-restored-title", "1.0"), + ) + .await; + let _ = connection + .request(&ClientMessage::Session(SessionRequest::AddRootTab { + request_id: new_request_id(), + session_id: session.session.id, + title: "fullscreen".to_owned(), + buffer_id: Some(buffer.id), + child_node_id: None, + })) + .await + .expect("add fullscreen tab succeeds"); + + let live = wait_for_visible_snapshot( + &mut connection, + buffer.id, + Duration::from_secs(3), + |snapshot| { + let text = snapshot.lines.join("\n"); + snapshot.alternate_screen + && snapshot.title.as_deref() == Some("fullscreen-live-title") + && text.contains("fullscreen-live") + && text.contains("cursor-target") + }, + ) + .await; + let live_text = live.lines.join("\n"); + assert!(!live_text.contains("main-before")); + + let mut client = MuxClient::connect(server.socket_path()) + .await + .expect("client connects"); + let render = render_session(&mut client, "alpha").await; + assert!(render.contains("fullscreen-live")); + assert!(render.contains("cursor-target")); + assert!(!render.contains("main-before")); + + let restored = wait_for_visible_snapshot( + &mut connection, + buffer.id, + Duration::from_secs(4), + |snapshot| { + let text = snapshot.lines.join("\n"); + !snapshot.alternate_screen + && snapshot.title.as_deref() == Some("primary-restored-title") + && text.contains("main-before") + && text.contains("restored-after") + }, + ) + .await; + let restored_text = restored.lines.join("\n"); + assert!(!restored_text.contains("fullscreen-live")); + + server.shutdown().await.expect("server shuts down"); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn hidden_fullscreen_buffer_reveals_live_alternate_screen_coherently() { + let server = TestServer::start().await.expect("server starts"); + let mut connection = TestConnection::connect(server.socket_path()) + .await + .expect("protocol connection"); + + let hidden_buffer = create_buffer_with_command( + &mut connection, + "fullscreen-hidden", + fullscreen_fixture_command( + "fullscreen-hidden-live", + "fullscreen-hidden-restored", + "1.2", + ), + ) + .await; + let fixture = create_hidden_tab_fixture_with_buffer(&mut connection, hidden_buffer).await; + + wait_for_visible_snapshot( + &mut connection, + fixture.hidden_buffer.id, + Duration::from_secs(3), + |snapshot| { + snapshot.alternate_screen + && snapshot.title.as_deref() == Some("fullscreen-hidden-live") + && snapshot.lines.join("\n").contains("fullscreen-live") + }, + ) + .await; + + let _ = connection + .request(&ClientMessage::Node(NodeRequest::SelectTab { + request_id: new_request_id(), + tabs_node_id: fixture.nested_tabs_id, + index: 1, + })) + .await + .expect("select fullscreen tab succeeds"); + + let mut client = MuxClient::connect(server.socket_path()) + .await + .expect("client connects"); + let live_render = render_session(&mut client, "alpha").await; + assert!(live_render.contains("fullscreen-live")); + assert!(live_render.contains("cursor-target")); + assert!(!live_render.contains("main-before")); + + let restored = wait_for_visible_snapshot( + &mut connection, + fixture.hidden_buffer.id, + Duration::from_secs(4), + |snapshot| { + let text = snapshot.lines.join("\n"); + !snapshot.alternate_screen + && snapshot.title.as_deref() == Some("fullscreen-hidden-restored") + && text.contains("main-before") + && text.contains("restored-after") + }, + ) + .await; + assert!(!restored.lines.join("\n").contains("fullscreen-live")); + + let restored_render = render_session(&mut client, "alpha").await; + assert!(restored_render.contains("main-before")); + assert!(restored_render.contains("restored-after")); + assert!(!restored_render.contains("fullscreen-live")); + + server.shutdown().await.expect("server shuts down"); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 4)] +async fn rapid_terminal_output_renders_latest_visible_snapshot() { + let server = TestServer::start().await.expect("server starts"); + let mut connection = TestConnection::connect(server.socket_path()) + .await + .expect("protocol connection"); + + let session = create_session(&mut connection, "alpha").await; + let buffer = create_buffer_with_command( + &mut connection, + "burst", + vec![ + "/bin/sh".to_owned(), + "-lc".to_owned(), + "i=1; while [ $i -le 80 ]; do printf 'burst-%02d\\n' \"$i\"; i=$((i+1)); done" + .to_owned(), + ], + ) + .await; + let _ = connection + .request(&ClientMessage::Session(SessionRequest::AddRootTab { + request_id: new_request_id(), + session_id: session.session.id, + title: "burst".to_owned(), + buffer_id: Some(buffer.id), + child_node_id: None, + })) + .await + .expect("add burst tab succeeds"); + + connection + .wait_for_capture_contains(buffer.id, "burst-80", Duration::from_secs(3)) + .await + .expect("rapid output finishes"); + wait_for_visible_snapshot( + &mut connection, + buffer.id, + Duration::from_secs(3), + |snapshot| snapshot.total_lines >= 80 && snapshot.lines.join("\n").contains("burst-80"), + ) + .await; + + let mut client = MuxClient::connect(server.socket_path()) + .await + .expect("client connects"); + let render = render_session(&mut client, "alpha").await; + let session_id = session_id_by_name(&client, "alpha"); + let presentation = PresentationModel::project( + client.state(), + session_id, + Size { + width: 80, + height: 24, + }, + ) + .expect("projection succeeds"); + let visible_rows = presentation + .focused_leaf() + .expect("focused leaf") + .rect + .size + .height + .saturating_sub(1) as usize; + let latest_rendered_line = client + .state() + .snapshots + .get(&buffer.id) + .expect("burst snapshot") + .lines + .iter() + .take(visible_rows) + .rev() + .find(|line| line.starts_with("burst-")) + .expect("latest rendered burst line"); + assert!(render.contains(latest_rendered_line)); + assert!(!render.contains("burst-01")); + + server.shutdown().await.expect("server shuts down"); +} diff --git a/crates/embers-server/src/buffer_runtime.rs b/crates/embers-server/src/buffer_runtime.rs index b165ebda..ab33dee8 100644 --- a/crates/embers-server/src/buffer_runtime.rs +++ b/crates/embers-server/src/buffer_runtime.rs @@ -367,19 +367,16 @@ impl BufferRuntimeHandle { impl BufferRuntimeInner { fn join_threads_blocking(&self) { self.stop.store(true, Ordering::Relaxed); - let mut threads = match self.threads.lock() { - Ok(threads) => threads, + let poller = match self.threads.lock() { + Ok(mut threads) => threads.poller.take(), Err(poisoned) => { error!( %self.buffer_id, "buffer runtime thread registry lock poisoned during shutdown" ); - poisoned.into_inner() + poisoned.into_inner().poller.take() } }; - let poller = threads.poller.take(); - drop(threads); - if let Some(poller) = poller && poller.thread().id() != thread::current().id() { @@ -713,6 +710,18 @@ fn handle_keeper_request( } impl KeeperRuntime { + fn ensure_running(&self) -> Result<()> { + if self + .exit_code + .lock() + .map_err(|_| MuxError::internal("runtime keeper exit lock poisoned"))? + .is_some() + { + return Err(MuxError::conflict("buffer runtime has already exited")); + } + Ok(()) + } + fn status(&self) -> Result { let exit_code = *self .exit_code @@ -739,14 +748,7 @@ impl KeeperRuntime { } fn write(&self, bytes: Vec) -> Result<()> { - if self - .exit_code - .lock() - .map_err(|_| MuxError::internal("runtime keeper exit lock poisoned"))? - .is_some() - { - return Err(MuxError::conflict("buffer runtime has already exited")); - } + self.ensure_running()?; let mut writer = self .writer .lock() @@ -757,6 +759,7 @@ impl KeeperRuntime { } fn resize(&self, size: PtySize) -> Result<()> { + self.ensure_running()?; let master = self .master .lock() @@ -802,6 +805,7 @@ impl KeeperRuntime { } fn kill(&self) -> Result<()> { + self.ensure_running()?; let mut killer = self .killer .lock() diff --git a/crates/embers-server/src/state.rs b/crates/embers-server/src/state.rs index 1a3dc947..7fd48eac 100644 --- a/crates/embers-server/src/state.rs +++ b/crates/embers-server/src/state.rs @@ -1455,6 +1455,7 @@ impl ServerState { pub fn focus_leaf(&mut self, session_id: SessionId, leaf_id: NodeId) -> Result<()> { self.ensure_leaf_belongs_to(leaf_id, session_id)?; self.ensure_leaf_is_focusable(session_id, leaf_id)?; + let buffer_id = self.buffer_view_buffer_id(leaf_id)?; self.clear_session_focus(session_id)?; self.set_leaf_focus(leaf_id, true)?; @@ -1492,6 +1493,8 @@ impl ServerState { child = parent; } + self.set_buffer_activity(buffer_id, ActivityState::Idle)?; + Ok(()) } diff --git a/crates/embers-server/src/terminal_backend.rs b/crates/embers-server/src/terminal_backend.rs index df25749f..796569da 100644 --- a/crates/embers-server/src/terminal_backend.rs +++ b/crates/embers-server/src/terminal_backend.rs @@ -38,6 +38,11 @@ pub enum BackendDamage { Partial(Vec), } +/// Terminal emulation boundary used by the runtime keeper. +/// +/// Raw PTY bytes are routed through `RawByteRouter` and then ingested here. The backend owns +/// terminal parsing, alternate-screen state, scrollback, snapshots, cursor metadata, and render +/// damage tracking. pub trait TerminalBackend: Send { fn ingest_bytes(&mut self, bytes: &[u8]); fn resize(&mut self, size: PtySize); @@ -58,10 +63,18 @@ pub trait TerminalBackend: Send { pub struct RawByteRouter; impl RawByteRouter { + /// Route client-originated bytes before they reach the PTY. + /// + /// The current implementation is intentionally passthrough, but the method is the explicit + /// seam for future prefix/passthrough-aware interception. pub fn route_input(&self, bytes: Vec) -> Vec { bytes } + /// Route PTY output bytes before terminal emulation. + /// + /// Today this forwards output directly into the backend, making the raw-routing seam explicit + /// without introducing policy beyond passthrough. pub fn route_output(&mut self, backend: &mut dyn TerminalBackend, bytes: &[u8]) { backend.ingest_bytes(bytes); } @@ -346,8 +359,73 @@ impl TerminalBackend for AlacrittyTerminalBackend { #[cfg(test)] mod tests { - use super::{AlacrittyTerminalBackend, BackendDamage, TerminalBackend}; - use embers_core::{ActivityState, CursorShape, PtySize}; + use std::path::PathBuf; + + use super::{ + AlacrittyTerminalBackend, BackendDamage, BackendMetadata, BackendScrollbackSlice, + RawByteRouter, TerminalBackend, + }; + use embers_core::{ActivityState, CursorShape, PtySize, TerminalSnapshot}; + + #[derive(Default)] + struct StubBackend { + ingested: Vec, + } + + impl TerminalBackend for StubBackend { + fn ingest_bytes(&mut self, bytes: &[u8]) { + self.ingested.extend_from_slice(bytes); + } + + fn resize(&mut self, _size: PtySize) {} + + fn visible_snapshot( + &self, + sequence: u64, + size: PtySize, + cwd: Option, + ) -> embers_core::TerminalSnapshot { + let mut snapshot = embers_core::TerminalSnapshot::from_lines( + sequence, + size, + [String::from_utf8_lossy(&self.ingested).into_owned()], + ); + snapshot.cwd = cwd; + snapshot + } + + fn capture_scrollback(&self) -> Vec { + vec![String::from_utf8_lossy(&self.ingested).into_owned()] + } + + fn capture_scrollback_slice( + &self, + start_line: u64, + _line_count: u32, + ) -> BackendScrollbackSlice { + BackendScrollbackSlice { + start_line, + total_lines: 1, + lines: vec![String::from_utf8_lossy(&self.ingested).into_owned()], + } + } + + fn metadata(&self) -> BackendMetadata { + BackendMetadata::default() + } + + fn take_activity(&mut self) -> ActivityState { + ActivityState::Activity + } + + fn take_damage(&mut self) -> BackendDamage { + BackendDamage::None + } + } + + fn snapshot_lines(snapshot: TerminalSnapshot) -> Vec { + snapshot.lines.into_iter().map(|line| line.text).collect() + } #[test] fn visible_snapshot_extracts_plain_text_lines() { @@ -357,7 +435,7 @@ mod tests { backend.ingest_bytes(b"hello\r\nworld"); let snapshot = backend.visible_snapshot(3, PtySize::new(8, 3), None); - let lines: Vec<_> = snapshot.lines.into_iter().map(|line| line.text).collect(); + let lines = snapshot_lines(snapshot.clone()); assert_eq!(lines, vec!["hello", "world", ""]); assert_eq!(snapshot.total_lines, 3); assert_eq!(snapshot.viewport_top_line, 0); @@ -367,6 +445,50 @@ mod tests { )); } + #[test] + fn carriage_return_overwrites_cells_without_advancing_the_row() { + let mut backend = AlacrittyTerminalBackend::new(PtySize::new(8, 2)); + let _ = backend.take_damage(); + + backend.ingest_bytes(b"hello\rHEY"); + + let lines = snapshot_lines(backend.visible_snapshot(1, PtySize::new(8, 2), None)); + assert_eq!(lines, vec!["HEYlo", ""]); + } + + #[test] + fn automatic_wrap_moves_following_bytes_to_the_next_row() { + let mut backend = AlacrittyTerminalBackend::new(PtySize::new(4, 2)); + let _ = backend.take_damage(); + + backend.ingest_bytes(b"abcdX"); + + let lines = snapshot_lines(backend.visible_snapshot(1, PtySize::new(4, 2), None)); + assert_eq!(lines, vec!["abcd", "X"]); + } + + #[test] + fn erase_in_line_clears_trailing_cells_from_the_cursor() { + let mut backend = AlacrittyTerminalBackend::new(PtySize::new(6, 1)); + let _ = backend.take_damage(); + + backend.ingest_bytes(b"abcdef\rabc\x1b[K"); + + let lines = snapshot_lines(backend.visible_snapshot(1, PtySize::new(6, 1), None)); + assert_eq!(lines, vec!["abc"]); + } + + #[test] + fn clear_screen_resets_visible_cells_before_new_output() { + let mut backend = AlacrittyTerminalBackend::new(PtySize::new(6, 2)); + let _ = backend.take_damage(); + + backend.ingest_bytes(b"one\r\ntwo\x1b[2J\x1b[Hdone"); + + let lines = snapshot_lines(backend.visible_snapshot(1, PtySize::new(6, 2), None)); + assert_eq!(lines, vec!["done", ""]); + } + #[test] fn scrollback_capture_preserves_history_beyond_viewport() { let mut backend = AlacrittyTerminalBackend::new(PtySize::new(6, 2)); @@ -425,6 +547,26 @@ mod tests { assert!(metadata.bracketed_paste); } + #[test] + fn metadata_mode_flags_clear_when_disable_sequences_arrive() { + let mut backend = AlacrittyTerminalBackend::new(PtySize::new(10, 2)); + let _ = backend.take_damage(); + + backend.ingest_bytes(b"\x1b[?1049h\x1b[?1000h\x1b[?1004h\x1b[?2004h"); + let enabled = backend.metadata(); + assert!(enabled.alternate_screen); + assert!(enabled.mouse_reporting); + assert!(enabled.focus_reporting); + assert!(enabled.bracketed_paste); + + backend.ingest_bytes(b"\x1b[?1049l\x1b[?1000l\x1b[?1004l\x1b[?2004l"); + let disabled = backend.metadata(); + assert!(!disabled.alternate_screen); + assert!(!disabled.mouse_reporting); + assert!(!disabled.focus_reporting); + assert!(!disabled.bracketed_paste); + } + #[test] fn bell_activity_is_consumed_separately_from_metadata() { let mut backend = AlacrittyTerminalBackend::new(PtySize::new(10, 2)); @@ -440,4 +582,59 @@ mod tests { assert_eq!(metadata.title.as_deref(), Some("embers")); assert_eq!(backend.take_activity(), ActivityState::Activity); } + + #[test] + fn raw_byte_router_is_explicit_passthrough_for_input_and_output() { + let mut router = RawByteRouter; + let mut backend = StubBackend::default(); + let input = b"\x1b[200~paste\x1b[201~".to_vec(); + + assert_eq!(router.route_input(input.clone()), input); + + router.route_output(&mut backend, b"hello"); + router.route_output(&mut backend, b" world"); + + assert_eq!(backend.ingested, b"hello world"); + } + + #[test] + fn alternate_screen_visible_snapshot_tracks_active_screen_and_restores_primary_screen() { + let mut backend = AlacrittyTerminalBackend::new(PtySize::new(20, 4)); + let _ = backend.take_damage(); + + backend.ingest_bytes(b"main-one\r\nmain-two"); + backend.ingest_bytes(b"\x1b[?1049h\x1b[Halt-screen"); + + let alternate = backend.visible_snapshot(2, PtySize::new(20, 4), None); + let alternate_lines: Vec<_> = alternate + .lines + .iter() + .map(|line| line.text.as_str()) + .collect(); + assert!(alternate.modes.alternate_screen); + assert!( + alternate_lines + .iter() + .any(|line| line.contains("alt-screen")), + "alternate visible lines: {alternate_lines:?}" + ); + + backend.ingest_bytes(b"\x1b[?1049l"); + + let restored = backend.visible_snapshot(3, PtySize::new(20, 4), None); + let restored_lines: Vec<_> = restored + .lines + .iter() + .map(|line| line.text.as_str()) + .collect(); + assert!(!restored.modes.alternate_screen); + assert!( + restored_lines.iter().any(|line| line.contains("main-one")), + "restored visible lines: {restored_lines:?}" + ); + assert!( + restored_lines.iter().any(|line| line.contains("main-two")), + "restored visible lines: {restored_lines:?}" + ); + } } diff --git a/crates/embers-server/tests/buffer_lifecycle.rs b/crates/embers-server/tests/buffer_lifecycle.rs index 5da155f7..2d35d34b 100644 --- a/crates/embers-server/tests/buffer_lifecycle.rs +++ b/crates/embers-server/tests/buffer_lifecycle.rs @@ -106,6 +106,43 @@ fn closing_a_view_detaches_but_preserves_running_buffer() { )); } +#[test] +fn focusing_a_leaf_clears_recorded_activity() { + let mut state = ServerState::new(); + let session_id = state.create_session("main"); + let first_buffer = state.create_buffer("first", vec!["/bin/sh".to_owned()], None); + let first_view = state + .create_buffer_view(session_id, first_buffer) + .expect("create first view"); + state + .add_root_tab(session_id, "first", first_view) + .expect("attach first view"); + + let second_buffer = state.create_buffer("second", vec!["/bin/sh".to_owned()], None); + let second_view = state + .create_buffer_view(session_id, second_buffer) + .expect("create second view"); + state + .add_root_tab(session_id, "second", second_view) + .expect("attach second view"); + + state + .set_buffer_activity(first_buffer, ActivityState::Bell) + .expect("mark first buffer active"); + + state + .focus_leaf(session_id, first_view) + .expect("focus hidden first leaf"); + + assert_eq!( + state + .buffer(first_buffer) + .expect("first buffer exists") + .activity, + ActivityState::Idle + ); +} + #[test] fn resize_updates_buffer_size_for_attached_and_detached_buffers() { let mut state = ServerState::new(); @@ -135,3 +172,28 @@ fn resize_updates_buffer_size_for_attached_and_detached_buffers() { PtySize::new(90, 20) ); } + +#[test] +fn exited_detached_buffers_can_be_removed_cleanly() { + let mut state = ServerState::new(); + let buffer_id = state.create_buffer("shell", vec!["/bin/sh".to_owned()], None); + + state + .mark_buffer_running(buffer_id, Some(88)) + .expect("mark running"); + state + .mark_buffer_exited(buffer_id, Some(0)) + .expect("mark exited"); + + let removed = state + .remove_buffer(buffer_id) + .expect("remove detached exited buffer"); + assert!(matches!( + removed.state, + BufferState::Exited(ref exited) if exited.exit_code == Some(0) + )); + assert!(matches!( + state.buffer(buffer_id), + Err(MuxError::NotFound(_)) + )); +} diff --git a/crates/embers-test-support/src/server.rs b/crates/embers-test-support/src/server.rs index c9ae6fb5..e7de5380 100644 --- a/crates/embers-test-support/src/server.rs +++ b/crates/embers-test-support/src/server.rs @@ -18,6 +18,7 @@ pub struct TestServer { impl TestServer { pub async fn start() -> Result { init_test_tracing(); + reap_stale_helper_processes(); let tempdir = tempfile::tempdir()?; let socket_path = tempdir.path().join("mux.sock"); @@ -36,13 +37,9 @@ impl TestServer { &self.socket_path } - /// Shuts down the server and kills any orphaned __serve processes - /// that were spawned for this socket during the test. + /// Shuts down the server and kills any orphaned embers helper processes that + /// were spawned for this socket during the test. pub async fn shutdown(mut self) -> Result<()> { - // First, kill any orphaned __serve processes for our socket - self.kill_orphaned_servers(); - - // Then shutdown our own server with a timeout if let Some(handle) = self.handle.take() { match tokio::time::timeout(SHUTDOWN_TIMEOUT, handle.shutdown()).await { Ok(Ok(())) => {} @@ -54,48 +51,99 @@ impl TestServer { } } } + self.kill_orphaned_processes(); Ok(()) } - /// Kill any orphaned embers __serve processes that were spawned - /// for this server's socket but are no longer needed. - fn kill_orphaned_servers(&self) { + /// Kill any orphaned embers helper processes that were spawned for this + /// server's socket but are no longer needed. + fn kill_orphaned_processes(&self) { let socket_path_str = self.socket_path.to_string_lossy(); + let runtime_dir = self.socket_path.with_extension("runtimes"); + let runtime_dir_str = runtime_dir.to_string_lossy(); let pid_path = self.socket_path.with_extension("pid"); - // First try to kill via PID file (for __serve processes) if let Ok(pid_str) = std::fs::read_to_string(&pid_path) && let Ok(pid) = pid_str.trim().parse::() { let _ = Command::new("kill").arg(pid.to_string()).output(); } - // Also try to find and kill any __serve processes referencing our socket - // This handles cases where the PID file wasn't cleaned up or we need - // to find the process by socket path if let Ok(output) = Command::new("ps").args(["-eo", "pid,args"]).output() { for line in String::from_utf8_lossy(&output.stdout).lines() { let line = line.trim(); - // Look for __serve processes with our socket path - if line.contains("__serve") - && line.contains(&*socket_path_str) + let is_server = line.contains("__serve") && line.contains(&*socket_path_str); + let is_runtime_keeper = + line.contains("__runtime-keeper") && line.contains(&*runtime_dir_str); + if (is_server || is_runtime_keeper) && let Some(pid_str) = line.split_whitespace().next() && let Ok(pid) = pid_str.parse::() { let _ = Command::new("kill").arg("-9").arg(pid.to_string()).output(); - tracing::debug!(pid, "killed orphaned __serve process"); + if is_server { + tracing::debug!(pid, "killed orphaned __serve process"); + } else { + tracing::debug!(pid, "killed orphaned __runtime-keeper process"); + } } } } - // Clean up the pid file let _ = std::fs::remove_file(&pid_path); + let _ = std::fs::remove_dir_all(&runtime_dir); } } impl Drop for TestServer { fn drop(&mut self) { - // Ensure any spawned servers are killed even if shutdown wasn't called - self.kill_orphaned_servers(); + self.kill_orphaned_processes(); + } +} + +fn reap_stale_helper_processes() { + if let Ok(output) = Command::new("ps").args(["-eo", "pid,args"]).output() { + for line in String::from_utf8_lossy(&output.stdout).lines() { + let Some((pid, socket_path, helper_kind)) = parse_helper_process(line) else { + continue; + }; + let Some(parent) = socket_path.parent() else { + continue; + }; + if parent.exists() { + continue; + } + let _ = Command::new("kill").arg("-9").arg(pid.to_string()).output(); + tracing::debug!( + pid, + helper = helper_kind, + socket = %socket_path.display(), + "killed stale helper process" + ); + } } } + +fn parse_helper_process(line: &str) -> Option<(i32, std::path::PathBuf, &'static str)> { + let line = line.trim(); + let is_server = line.contains("__serve"); + let is_runtime_keeper = line.contains("__runtime-keeper"); + if !is_server && !is_runtime_keeper { + return None; + } + + let mut fields = line.split_whitespace(); + let pid = fields.next()?.parse::().ok()?; + let args = fields.collect::>(); + let socket_path = args.windows(2).find_map(|window| { + (window[0] == "--socket").then_some(std::path::PathBuf::from(window[1])) + })?; + Some(( + pid, + socket_path, + if is_server { + "__serve" + } else { + "__runtime-keeper" + }, + )) +} diff --git a/crates/embers-test-support/tests/buffer_runtime.rs b/crates/embers-test-support/tests/buffer_runtime.rs index 67005419..621cc71e 100644 --- a/crates/embers-test-support/tests/buffer_runtime.rs +++ b/crates/embers-test-support/tests/buffer_runtime.rs @@ -347,3 +347,145 @@ async fn scrollback_slice_returns_history_while_full_capture_stays_available() { server.shutdown().await.expect("shutdown server"); } + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn repeated_snapshot_reads_stay_stable_without_new_output() { + let _guard = acquire_test_lock().await.expect("acquire test lock"); + let server = TestServer::start().await.expect("start server"); + let mut connection = TestConnection::connect(server.socket_path()) + .await + .expect("connect protocol client"); + + let buffer = create_buffer( + &mut connection, + &[ + "/bin/sh", + "-lc", + "i=1; while [ $i -le 8 ]; do printf 'line-%02d\\n' \"$i\"; i=$((i+1)); done", + ], + ) + .await; + + wait_for_capture_contains(&mut connection, buffer.id, "line-08").await; + wait_for_exit(&mut connection, buffer.id).await; + + let first_full = capture_buffer(&mut connection, buffer.id).await; + let second_full = capture_buffer(&mut connection, buffer.id).await; + assert_eq!(first_full.sequence, second_full.sequence); + assert_eq!(first_full.size, second_full.size); + assert_eq!(first_full.lines, second_full.lines); + assert_eq!(first_full.title, second_full.title); + assert_eq!(first_full.cwd, second_full.cwd); + + let first_visible = connection + .capture_visible_buffer(buffer.id) + .await + .expect("first visible capture succeeds"); + let second_visible = connection + .capture_visible_buffer(buffer.id) + .await + .expect("second visible capture succeeds"); + assert_eq!(first_visible.sequence, second_visible.sequence); + assert_eq!(first_visible.size, second_visible.size); + assert_eq!(first_visible.lines, second_visible.lines); + assert_eq!(first_visible.title, second_visible.title); + assert_eq!(first_visible.cwd, second_visible.cwd); + assert_eq!( + first_visible.viewport_top_line, + second_visible.viewport_top_line + ); + assert_eq!(first_visible.total_lines, second_visible.total_lines); + assert_eq!( + first_visible.alternate_screen, + second_visible.alternate_screen + ); + assert_eq!( + first_visible.mouse_reporting, + second_visible.mouse_reporting + ); + assert_eq!( + first_visible.focus_reporting, + second_visible.focus_reporting + ); + assert_eq!( + first_visible.bracketed_paste, + second_visible.bracketed_paste + ); + assert_eq!(first_visible.cursor, second_visible.cursor); + + let first_slice = connection + .capture_scrollback_slice(buffer.id, 2, 3) + .await + .expect("first scrollback slice succeeds"); + let second_slice = connection + .capture_scrollback_slice(buffer.id, 2, 3) + .await + .expect("second scrollback slice succeeds"); + assert_eq!(first_slice.start_line, second_slice.start_line); + assert_eq!(first_slice.total_lines, second_slice.total_lines); + assert_eq!(first_slice.lines, second_slice.lines); + + server.shutdown().await.expect("shutdown server"); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn detached_visible_capture_tracks_latest_size_and_output() { + let _guard = acquire_test_lock().await.expect("acquire test lock"); + let server = TestServer::start().await.expect("start server"); + let mut connection = TestConnection::connect(server.socket_path()) + .await + .expect("connect protocol client"); + + let buffer = create_buffer( + &mut connection, + &[ + "/bin/sh", + "-lc", + "printf '\\033]0;detached-preview\\007ready\\n'; while IFS= read -r line; do printf 'seen:%s\\n' \"$line\"; done", + ], + ) + .await; + assert_eq!(buffer.attachment_node_id, None); + + wait_for_capture_contains(&mut connection, buffer.id, "ready").await; + let initial_visible = connection + .capture_visible_buffer(buffer.id) + .await + .expect("initial visible capture succeeds"); + assert_eq!(initial_visible.size, PtySize::new(80, 24)); + assert_eq!(initial_visible.title.as_deref(), Some("detached-preview")); + assert!(initial_visible.lines.join("\n").contains("ready")); + + resize_buffer(&mut connection, buffer.id, 96, 18).await; + let resized_visible = connection + .capture_visible_buffer(buffer.id) + .await + .expect("resized visible capture succeeds"); + assert_eq!(resized_visible.size, PtySize::new(96, 18)); + assert!(resized_visible.total_lines >= 1); + + send_input(&mut connection, buffer.id, "after-resize\n").await; + wait_for_capture_contains(&mut connection, buffer.id, "seen:after-resize").await; + + let visible = connection + .capture_visible_buffer(buffer.id) + .await + .expect("final visible capture succeeds"); + assert_eq!(visible.size, PtySize::new(96, 18)); + assert_eq!(visible.title.as_deref(), Some("detached-preview")); + assert!(visible.lines.join("\n").contains("seen:after-resize")); + + let captured = capture_buffer(&mut connection, buffer.id).await; + assert_eq!(captured.size, PtySize::new(96, 18)); + assert!(captured.lines.join("\n").contains("ready")); + assert!(captured.lines.join("\n").contains("seen:after-resize")); + + let slice = connection + .capture_scrollback_slice(buffer.id, 0, 4) + .await + .expect("detached scrollback slice succeeds"); + assert!(slice.total_lines >= 2); + assert!(slice.lines.join("\n").contains("ready")); + + server.shutdown().await.expect("shutdown server"); +} diff --git a/crates/embers-test-support/tests/detach_move.rs b/crates/embers-test-support/tests/detach_move.rs index b67fea2b..3a189e27 100644 --- a/crates/embers-test-support/tests/detach_move.rs +++ b/crates/embers-test-support/tests/detach_move.rs @@ -1,5 +1,6 @@ use std::time::{Duration, Instant}; +use crate::support::integration_test_lock; use embers_core::{SplitDirection, new_request_id}; use embers_protocol::{ BufferRecord, BufferRequest, BuffersResponse, ClientMessage, InputRequest, NodeRequest, @@ -140,6 +141,25 @@ async fn send_input( assert!(matches!(response, ServerResponse::Ok(_))); } +async fn resize_buffer( + connection: &mut TestConnection, + buffer_id: embers_core::BufferId, + cols: u16, + rows: u16, +) { + let response = connection + .request(&ClientMessage::Input(InputRequest::Resize { + request_id: new_request_id(), + buffer_id, + cols, + rows, + })) + .await + .expect("resize request succeeds"); + + assert!(matches!(response, ServerResponse::Ok(_))); +} + async fn wait_for_capture_contains( connection: &mut TestConnection, buffer_id: embers_core::BufferId, @@ -161,6 +181,7 @@ async fn wait_for_capture_contains( #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn detach_list_capture_and_reattach_buffer_via_socket() { + let _guard = integration_test_lock().lock().await; let server = TestServer::start().await.expect("start server"); let mut connection = TestConnection::connect(server.socket_path()) .await @@ -234,6 +255,7 @@ async fn detach_list_capture_and_reattach_buffer_via_socket() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn move_request_replaces_target_leaf_without_killing_buffer() { + let _guard = integration_test_lock().lock().await; let server = TestServer::start().await.expect("start server"); let mut connection = TestConnection::connect(server.socket_path()) .await @@ -301,3 +323,77 @@ async fn move_request_replaces_target_leaf_without_killing_buffer() { server.shutdown().await.expect("shutdown server"); } + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn detach_and_reattach_preserve_runtime_identity_size_and_capture() { + let _guard = integration_test_lock().lock().await; + let server = TestServer::start().await.expect("start server"); + let mut connection = TestConnection::connect(server.socket_path()) + .await + .expect("connect protocol client"); + + let session = create_session(&mut connection, "main").await; + let session_id = session.snapshot.session.id; + + let primary = create_echo_buffer(&mut connection, "primary").await; + let attached = add_root_tab(&mut connection, session_id, "primary", primary.id).await; + let original_leaf = attached + .session + .focused_leaf_id + .expect("primary tab focuses leaf"); + wait_for_capture_contains(&mut connection, primary.id, "ready").await; + + let before_detach = get_buffer(&mut connection, primary.id).await; + resize_buffer(&mut connection, primary.id, 120, 33).await; + let resized = get_buffer(&mut connection, primary.id).await; + assert_eq!(resized.pty_size, embers_core::PtySize::new(120, 33)); + + let detach = connection + .request(&ClientMessage::Buffer(BufferRequest::Detach { + request_id: new_request_id(), + buffer_id: primary.id, + })) + .await + .expect("detach request succeeds"); + assert!(matches!(detach, ServerResponse::Ok(_))); + + let detached = get_buffer(&mut connection, primary.id).await; + assert_eq!(detached.attachment_node_id, None); + assert_eq!(detached.pid, before_detach.pid); + assert_eq!(detached.pty_size, embers_core::PtySize::new(120, 33)); + assert_eq!(detached.state, before_detach.state); + + send_input(&mut connection, primary.id, "detached-still-live\n").await; + let detached_capture = + wait_for_capture_contains(&mut connection, primary.id, "seen:detached-still-live").await; + assert!(detached_capture.lines.join("\n").contains("ready")); + + let replacement = create_echo_buffer(&mut connection, "replacement").await; + let replacement_snapshot = + add_root_tab(&mut connection, session_id, "replacement", replacement.id).await; + let target_leaf = replacement_snapshot + .session + .focused_leaf_id + .expect("replacement tab focuses target leaf"); + assert_ne!(target_leaf, original_leaf); + + let moved = connection + .request(&ClientMessage::Node(NodeRequest::MoveBufferToNode { + request_id: new_request_id(), + buffer_id: primary.id, + target_leaf_node_id: target_leaf, + })) + .await + .expect("move request succeeds"); + assert!(matches!(moved, ServerResponse::SessionSnapshot(_))); + + let reattached = get_buffer(&mut connection, primary.id).await; + assert_eq!(reattached.attachment_node_id, Some(target_leaf)); + assert_eq!(reattached.pid, before_detach.pid); + assert_eq!(reattached.pty_size, embers_core::PtySize::new(120, 33)); + + send_input(&mut connection, primary.id, "reattached\n").await; + wait_for_capture_contains(&mut connection, primary.id, "seen:reattached").await; + + server.shutdown().await.expect("shutdown server"); +} diff --git a/crates/embers-test-support/tests/floating_windows.rs b/crates/embers-test-support/tests/floating_windows.rs index 00d08632..e7df94ca 100644 --- a/crates/embers-test-support/tests/floating_windows.rs +++ b/crates/embers-test-support/tests/floating_windows.rs @@ -1,3 +1,4 @@ +use crate::support::integration_test_lock; use embers_core::{FloatGeometry, new_request_id}; use embers_protocol::{ BufferRecord, BufferRequest, ClientMessage, FloatingRequest, NodeRequest, ServerResponse, @@ -99,6 +100,7 @@ async fn get_buffer( #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn create_focus_move_and_close_floating_window_via_socket() { + let _guard = integration_test_lock().lock().await; let server = TestServer::start().await.expect("start server"); let mut connection = TestConnection::connect(server.socket_path()) .await @@ -203,6 +205,7 @@ async fn create_focus_move_and_close_floating_window_via_socket() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn closing_last_tab_in_floating_tabs_removes_popup() { + let _guard = integration_test_lock().lock().await; let server = TestServer::start().await.expect("start server"); let mut connection = TestConnection::connect(server.socket_path()) .await diff --git a/crates/embers-test-support/tests/integration.rs b/crates/embers-test-support/tests/integration.rs index c3559161..c63ad901 100644 --- a/crates/embers-test-support/tests/integration.rs +++ b/crates/embers-test-support/tests/integration.rs @@ -7,3 +7,4 @@ mod pty_smoke; mod server_harness; mod session_root_tabs; mod split_layout; +mod support; diff --git a/crates/embers-test-support/tests/nested_tabs.rs b/crates/embers-test-support/tests/nested_tabs.rs index 10ceeb2e..82c5aa19 100644 --- a/crates/embers-test-support/tests/nested_tabs.rs +++ b/crates/embers-test-support/tests/nested_tabs.rs @@ -1,3 +1,4 @@ +use crate::support::integration_test_lock; use embers_core::{SplitDirection, new_request_id}; use embers_protocol::{ BufferRecord, BufferRequest, ClientMessage, NodeRequest, ServerResponse, SessionRequest, @@ -108,6 +109,7 @@ fn split_record(snapshot: &SessionSnapshot, node_id: embers_core::NodeId) -> Spl #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn nested_tab_mutations_round_trip_through_socket() { + let _guard = integration_test_lock().lock().await; let server = TestServer::start().await.expect("start server"); let mut connection = TestConnection::connect(server.socket_path()) .await @@ -251,6 +253,7 @@ async fn nested_tab_mutations_round_trip_through_socket() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn get_tree_returns_nested_tab_structure() { + let _guard = integration_test_lock().lock().await; let server = TestServer::start().await.expect("start server"); let mut connection = TestConnection::connect(server.socket_path()) .await diff --git a/crates/embers-test-support/tests/protocol_server.rs b/crates/embers-test-support/tests/protocol_server.rs index 2bbf8361..913dd861 100644 --- a/crates/embers-test-support/tests/protocol_server.rs +++ b/crates/embers-test-support/tests/protocol_server.rs @@ -1,5 +1,6 @@ use std::time::Duration; +use crate::support::integration_test_lock; use embers_core::{ErrorCode, MuxError, RequestId, SessionId, new_request_id}; use embers_protocol::{ ClientMessage, FrameType, NodeRequest, PingRequest, RawFrame, ServerEnvelope, ServerEvent, @@ -55,6 +56,7 @@ fn encode_frame_bytes(frame: &RawFrame) -> Vec { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn subscriptions_fan_out_to_multiple_clients_with_session_filters() { + let _guard = integration_test_lock().lock().await; let server = TestServer::start().await.expect("start server"); let mut actor = TestConnection::connect(server.socket_path()) .await @@ -139,6 +141,7 @@ async fn subscriptions_fan_out_to_multiple_clients_with_session_filters() { #[tokio::test] async fn fragmented_request_frames_round_trip_and_preserve_correlation_id() { + let _guard = integration_test_lock().lock().await; let server = TestServer::start().await.expect("start server"); let mut stream = UnixStream::connect(server.socket_path()) .await @@ -177,6 +180,7 @@ async fn fragmented_request_frames_round_trip_and_preserve_correlation_id() { #[tokio::test] async fn malformed_payloads_return_protocol_violation_errors() { + let _guard = integration_test_lock().lock().await; let server = TestServer::start().await.expect("start server"); let mut stream = UnixStream::connect(server.socket_path()) .await @@ -208,6 +212,7 @@ async fn malformed_payloads_return_protocol_violation_errors() { #[tokio::test] async fn typed_errors_cover_invalid_ids_and_impossible_mutations() { + let _guard = integration_test_lock().lock().await; let server = TestServer::start().await.expect("start server"); let mut connection = TestConnection::connect(server.socket_path()) .await @@ -264,6 +269,7 @@ async fn typed_errors_cover_invalid_ids_and_impossible_mutations() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn disconnected_subscribers_are_cleaned_up_without_breaking_remaining_clients() { + let _guard = integration_test_lock().lock().await; let server = TestServer::start().await.expect("start server"); let mut actor = TestConnection::connect(server.socket_path()) .await diff --git a/crates/embers-test-support/tests/pty_smoke.rs b/crates/embers-test-support/tests/pty_smoke.rs index 7004a981..b0232bf5 100644 --- a/crates/embers-test-support/tests/pty_smoke.rs +++ b/crates/embers-test-support/tests/pty_smoke.rs @@ -1,11 +1,13 @@ use std::time::Duration; +use crate::support::integration_test_lock; use embers_core::PtySize; use embers_test_support::PtyHarness; #[test] #[ignore = "exercises the PTY smoke harness in CI and later end-to-end runs"] fn pty_round_trips_input() { + let _guard = integration_test_lock().blocking_lock(); let mut harness = PtyHarness::spawn( "sh", &["-lc", "read line; printf '%s' \"$line\""], diff --git a/crates/embers-test-support/tests/server_harness.rs b/crates/embers-test-support/tests/server_harness.rs index e12cb79b..a3ffa68a 100644 --- a/crates/embers-test-support/tests/server_harness.rs +++ b/crates/embers-test-support/tests/server_harness.rs @@ -1,7 +1,9 @@ +use crate::support::integration_test_lock; use embers_test_support::{TestConnection, TestServer}; #[tokio::test] async fn harness_starts_server_and_pings_it() { + let _guard = integration_test_lock().lock().await; let server = TestServer::start().await.expect("start server"); let mut connection = TestConnection::connect(server.socket_path()) .await diff --git a/crates/embers-test-support/tests/session_root_tabs.rs b/crates/embers-test-support/tests/session_root_tabs.rs index 73b633a6..48058896 100644 --- a/crates/embers-test-support/tests/session_root_tabs.rs +++ b/crates/embers-test-support/tests/session_root_tabs.rs @@ -1,3 +1,4 @@ +use crate::support::integration_test_lock; use embers_core::{ErrorCode, new_request_id}; use embers_protocol::{ BufferRecord, BufferRequest, ClientMessage, ServerResponse, SessionRequest, @@ -83,6 +84,7 @@ fn root_tabs(snapshot: &embers_protocol::SessionSnapshot) -> Option #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn create_list_get_and_close_sessions_via_socket() { + let _guard = integration_test_lock().lock().await; let server = TestServer::start().await.expect("start server"); let mut connection = TestConnection::connect(server.socket_path()) .await @@ -147,6 +149,7 @@ async fn create_list_get_and_close_sessions_via_socket() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn create_select_rename_and_close_root_tabs_via_socket() { + let _guard = integration_test_lock().lock().await; let server = TestServer::start().await.expect("start server"); let mut connection = TestConnection::connect(server.socket_path()) .await diff --git a/crates/embers-test-support/tests/split_layout.rs b/crates/embers-test-support/tests/split_layout.rs index 82685721..06749666 100644 --- a/crates/embers-test-support/tests/split_layout.rs +++ b/crates/embers-test-support/tests/split_layout.rs @@ -1,3 +1,4 @@ +use crate::support::integration_test_lock; use embers_core::{SplitDirection, new_request_id}; use embers_protocol::{ BufferRecord, BufferRequest, ClientMessage, NodeRequest, ServerResponse, SessionRequest, @@ -107,6 +108,7 @@ fn split_record(snapshot: &SessionSnapshot, node_id: embers_core::NodeId) -> Spl #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn split_and_resize_requests_build_nested_layouts_via_socket() { + let _guard = integration_test_lock().lock().await; let server = TestServer::start().await.expect("start server"); let mut connection = TestConnection::connect(server.socket_path()) .await @@ -188,6 +190,7 @@ async fn split_and_resize_requests_build_nested_layouts_via_socket() { #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn focus_and_close_requests_normalize_layout_and_detach_buffers() { + let _guard = integration_test_lock().lock().await; let server = TestServer::start().await.expect("start server"); let mut connection = TestConnection::connect(server.socket_path()) .await diff --git a/crates/embers-test-support/tests/support.rs b/crates/embers-test-support/tests/support.rs new file mode 100644 index 00000000..b52e5e2b --- /dev/null +++ b/crates/embers-test-support/tests/support.rs @@ -0,0 +1,20 @@ +use std::sync::OnceLock; + +use tokio::sync::{Mutex, MutexGuard}; + +pub struct IntegrationTestLock(Mutex<()>); + +impl IntegrationTestLock { + pub async fn lock(&self) -> MutexGuard<'_, ()> { + self.0.lock().await + } + + pub fn blocking_lock(&self) -> MutexGuard<'_, ()> { + self.0.blocking_lock() + } +} + +pub fn integration_test_lock() -> &'static IntegrationTestLock { + static LOCK: OnceLock = OnceLock::new(); + LOCK.get_or_init(|| IntegrationTestLock(Mutex::new(()))) +} diff --git a/docs/activity-bell-policy.md b/docs/activity-bell-policy.md new file mode 100644 index 00000000..4fc1d27f --- /dev/null +++ b/docs/activity-bell-policy.md @@ -0,0 +1,21 @@ +## Activity and bell policy + +Phase 6 locks down how background terminal updates are reflected in buffer metadata. + +### What counts as activity + +Any PTY output that advances a buffer snapshot marks that buffer as `Activity`. This is stored on the durable `Buffer` record, not on the current view, so hidden and detached buffers keep their activity state even while they are off-screen. + +### What counts as a bell + +If the terminal backend observes a bell during a given output update, that update is recorded as `Bell` instead of plain `Activity`. Bell wins over ordinary activity for that update, which lets client and automation layers distinguish attention-worthy output from normal background chatter. + +### Hidden and detached buffers + +Hidden tabs, inactive panes, and detached buffers continue to ingest PTY output, update captures, and overwrite their stored activity state as new output arrives. Reconnecting clients recover that state from the server snapshot, and automation can observe bell updates from the normal render-invalidated path. + +### Reset on reveal or focus + +When a buffer becomes the focused leaf, the server clears its stored activity back to `Idle`. This acknowledges the previously hidden background signal without destroying terminal state or capture history. + +That reset only applies to the activity that had accumulated before focus. If the focused program produces more output afterward, later runtime updates can mark the buffer active again. diff --git a/docs/config-api/action.md b/docs/config-api/action.md index 42cab5c7..71771054 100644 --- a/docs/config-api/action.md +++ b/docs/config-api/action.md @@ -209,7 +209,7 @@ fn commit_search(_: ActionApi) -> Action
- diff --git a/docs/config-api/registration-action.md b/docs/config-api/registration-action.md index c5a44f6b..f0594de3 100644 --- a/docs/config-api/registration-action.md +++ b/docs/config-api/registration-action.md @@ -209,7 +209,7 @@ fn commit_search(_: ActionApi) -> Action
- diff --git a/docs/input-routing-policy.md b/docs/input-routing-policy.md new file mode 100644 index 00000000..4e186cd5 --- /dev/null +++ b/docs/input-routing-policy.md @@ -0,0 +1,83 @@ +# Input routing policy + +This note defines how Embers decides whether input is handled locally or forwarded to the focused terminal buffer. + +## Routing ownership + +Input-routing policy currently lives in the client layer. + +- `ConfiguredClient` owns key, paste, mouse, and focus routing decisions +- the server receives terminal-bound bytes as `InputRequest::Send` +- the server/runtime path writes those bytes directly to the focused buffer runtime + +Today, routing is decided before the protocol request is sent. The server does not reinterpret keybindings after it receives `InputRequest::Send`. + +## Key routing decision tree + +For normal key events, `ConfiguredClient::handle_key` follows this order: + +1. If the client is in search mode, search-prompt handling wins. +2. Otherwise, the key is resolved against the current mode's bindings. +3. Exact matches execute configured actions. +4. Prefix matches stay pending and do not forward bytes yet. +5. Unmatched sequences follow the current mode's fallback policy: + - `Passthrough`: send the key sequence to the focused buffer + - `Ignore`: consume it locally without terminal output + +That means partial leader or prefix sequences must never leak into the terminal while they are still unresolved. + +## Modes and fallback + +Built-in fallback policy is: + +- `normal`: passthrough +- `copy`: ignore +- `search`: ignore +- `select`: ignore + +Changing modes clears any pending key sequence. Reloading config also clears pending input if the current mode still exists, or resets the client back to `normal` if it no longer does. + +## Prefix and leader behavior + +Leader bindings are expanded into ordinary key sequences during config compilation. At runtime there is no separate "leader state" object; the pending sequence in `InputState` is the source of truth. + +While a sequence is still a prefix: + +- no bytes are sent to the buffer +- no local fallback runs yet +- the client waits for the next key to decide whether the sequence resolves or falls back + +## Local actions vs terminal passthrough + +Most exact matches execute locally, but Embers intentionally avoids stealing some keys from terminal apps. + +Configured bindings are forwarded to the terminal instead of executing locally when: + +- the focused view is in alternate screen and the bound actions are all local search/select/scroll actions, or +- the client is in normal mode with `follow_output` enabled, has no active search or selection state, and the bound actions require local search/select context + +This keeps fullscreen terminal apps from losing keys that would otherwise be meaningful inside the application. + +## Script-generated terminal input + +Scripted actions fit into the same pipeline after binding resolution: + +- `Action::SendKeys` converts notation into bytes +- `Action::SendBytes` uses the provided bytes directly +- both resolve a target buffer (`current` or explicit) +- both send `InputRequest::Send` + +So scripted terminal input uses the same protocol/runtime path as unmapped passthrough keys. + +## Special input paths + +- paste uses `handle_paste`; if the focused buffer has bracketed paste enabled, Embers wraps the payload in `ESC [ 200~` / `ESC [ 201~` +- focus events use `handle_focus_event`; they forward `ESC [ I` / `ESC [ O` only when the program requested focus reporting +- mouse events use `handle_mouse`; they either drive local scroll/focus behavior or send encoded mouse bytes when mouse reporting is enabled + +## Code anchors + +- Input resolution: `crates/embers-client/src/input/keymap.rs` +- Mode state and fallback policy: `crates/embers-client/src/input/modes.rs` +- Live routing and scripted action delivery: `crates/embers-client/src/configured_client.rs` +- Protocol input dispatch: `crates/embers-server/src/server.rs` diff --git a/docs/input-routing.md b/docs/input-routing.md new file mode 100644 index 00000000..7fdce3e4 --- /dev/null +++ b/docs/input-routing.md @@ -0,0 +1,34 @@ +# Input routing + +This is the short-form map of how Embers routes terminal input today. + +The detailed contract and rationale live in `docs/input-routing-policy.md`. This file exists as the +plan-facing entry point for the same behavior. + +## Summary + +- `ConfiguredClient` decides whether a key is handled locally or passed through to the focused + buffer. +- Leader and prefix sequences stay client-local until they resolve or are cleared. +- Copy/select/search modes suppress passthrough for the bindings they own. +- Scripted `send_keys*` and `send_bytes*` actions bypass local rendering logic and write bytes to + the target buffer runtime. +- Hidden or detached buffers do not receive arbitrary typed input unless a script or explicit + command targets them directly. + +## Current boundaries + +- Terminal-facing input ultimately lands in the buffer runtime, not in layout state. +- The `RawByteRouter` remains the explicit seam for raw-byte policy, but normal live PTY input is + currently decided earlier by the configured client and keymap. +- Config reload clears pending prefixes before the next key is interpreted under the new bindings. + +## Primary regression coverage + +- `crates/embers-client/tests/configured_client.rs` + - prefix no-leak behavior + - copy-mode passthrough suppression + - scripted `send_keys` / `send_bytes` + - reload clearing pending prefixes + +For the full policy text, see `docs/input-routing-policy.md`. diff --git a/docs/render-source-contract.md b/docs/render-source-contract.md new file mode 100644 index 00000000..02cc1362 --- /dev/null +++ b/docs/render-source-contract.md @@ -0,0 +1,29 @@ +## Render-source contract + +Phase 8 locks down which server surfaces the client uses for terminal rendering. + +### Authoritative sources + +The server remains authoritative for both layout and terminal state. + +- `SessionSnapshot` provides layout topology plus durable buffer metadata such as title, activity, attachment, and PTY size. +- `VisibleSnapshotResponse` provides the current visible terminal surface for one buffer, including visible lines, cursor state, viewport position, alternate-screen mode, and other terminal-mode flags. +- Full capture and scrollback slices stay on-demand APIs and are not part of the normal render loop. + +The client does not consume terminal diffs. It renders from full visible snapshots, with `RenderInvalidated` acting as a hint that a buffer should be refreshed before the next user-visible render. + +### Freshness expectations + +`RenderInvalidated` means the visible snapshot for that buffer may be stale. The configured client refreshes invalidated visible leaves before rendering and re-projects the presentation after those refreshes so updated titles, alternate-screen flags, and visible lines are used together. + +For event handling, the client also refreshes the affected `BufferRecord` before dispatching render-invalidated hooks. That keeps metadata-only consumers such as bell automation aligned with the latest server state. + +### Metadata synchronization + +Visible snapshots may carry title and mode changes that affect UI immediately. Buffer metadata still lives on the durable `BufferRecord`, so activity, bell state, and detached-buffer discovery remain queryable even when a buffer is hidden. + +Hidden buffers do not eagerly fetch fresh visible snapshots just because they were invalidated. Their visible state is refreshed when they become visible or when a caller explicitly requests capture. Their metadata, however, continues to flow through buffer/session refresh paths. + +### Detached buffers + +Detached buffers are discovered through `BufferRequest::List` / `Get`, and their visible surface is queried explicitly through the same capture endpoints as attached buffers. That keeps detached previews and background metadata within the same contract as attached terminal runtimes. diff --git a/docs/terminal-backend-boundary.md b/docs/terminal-backend-boundary.md new file mode 100644 index 00000000..c6835b25 --- /dev/null +++ b/docs/terminal-backend-boundary.md @@ -0,0 +1,83 @@ +# Terminal backend boundary + +This note defines the PTY-to-render pipeline boundary Embers relies on today. + +## Pipeline + +The active PTY pipeline is: + +```text +PTY bytes + -> RawByteRouter + -> TerminalBackend + -> snapshot / metadata / damage + -> protocol responses and render invalidation +``` + +For live PTY buffers, that pipeline runs inside the runtime keeper (`KeeperSurface` in `crates/embers-server/src/buffer_runtime.rs`). The server process talks to the keeper over the runtime socket and does not own a second terminal parser. + +## Raw routing seam + +`RawByteRouter` is the only layer allowed to inspect or rewrite raw terminal bytes before they hit the backend. + +Today it is intentionally minimal: + +- input routing is passthrough +- output routing forwards bytes directly to the backend + +That seam exists so future work can add: + +- protocol-aware passthrough decisions +- special-case interception +- metadata extraction that must happen before backend ingestion + +The router should not own terminal screen state, scrollback, or view/layout concerns. + +## Backend ownership + +`TerminalBackend` owns terminal emulation state: + +- ANSI/terminal parsing +- primary/alternate screen state +- cursor state +- scrollback capture +- visible snapshot generation +- damage tracking +- terminal mode reporting (alternate screen, mouse, focus, bracketed paste) + +Backends must be able to: + +- ingest output bytes +- resize +- produce a visible snapshot +- produce full capture / scrollback slices +- surface metadata +- surface one-shot activity and damage signals + +## Metadata outside the backend + +The server still owns buffer-level metadata records: + +- buffer title field mirrored onto the `Buffer` +- activity/bell state mirrored onto the `Buffer` +- last snapshot sequence on the `Buffer` +- render invalidation events + +The backend reports the current terminal metadata; the server persists the last observed values onto the durable buffer record. + +## Alternate-screen policy + +Alternate-screen state is owned by the backend. + +- `visible_snapshot` always reflects the currently active screen +- the `alternate_screen` mode bit tells clients whether the visible snapshot is from the alternate buffer +- when alternate screen exits, the visible snapshot returns to the primary screen state managed by the backend +- full capture and scrollback are backend-defined state, not layout-defined state + +In practice, Embers currently inherits alternate-screen capture semantics from the alacritty backend implementation, and the tests lock that behavior down at the backend boundary. + +## Code anchors + +- Raw routing + backend trait: `crates/embers-server/src/terminal_backend.rs` +- Keeper-owned surface: `crates/embers-server/src/buffer_runtime.rs` +- Server/runtime wiring: `crates/embers-server/src/server.rs` diff --git a/docs/terminal-capture-model.md b/docs/terminal-capture-model.md new file mode 100644 index 00000000..3370ef7a --- /dev/null +++ b/docs/terminal-capture-model.md @@ -0,0 +1,92 @@ +# Terminal capture model + +This note defines the terminal snapshot and capture semantics Embers exposes from PTY-backed buffers. + +## Capture surfaces + +Embers exposes three related but distinct capture surfaces for PTY buffers: + +- `capture_snapshot`: the full backend-defined capture for the buffer +- `capture_visible_snapshot`: the renderer-facing view of the currently visible screen +- `capture_scrollback_slice`: a paged read over the backend-defined scrollback history + +All three are sourced from the durable buffer runtime (`BufferRuntimeHandle` -> runtime keeper -> `TerminalBackend`), not from layout state. + +## Full snapshot semantics + +`capture_snapshot` is the "capture pane/buffer" source of truth for PTY buffers. + +It returns: + +- the current snapshot sequence +- the buffer's current PTY size +- the backend's full captured lines +- the terminal title if the backend has one +- the buffer cwd tracked by the server + +For PTY buffers, this is a runtime capture, not a view capture. Moving, detaching, or reattaching the buffer does not change which terminal state is returned. + +## Visible snapshot semantics + +`capture_visible_snapshot` is the authoritative render input for a PTY buffer. + +It returns: + +- the current visible lines from the active screen +- viewport position and total line count +- terminal mode bits such as alternate screen, mouse reporting, focus reporting, and bracketed paste +- cursor metadata +- size, sequence, title, and cwd + +The visible snapshot always reflects the screen the backend considers active at the moment of capture. + +## Scrollback semantics + +`capture_scrollback_slice` pages through backend-defined history without changing visible state. + +It returns: + +- `start_line`: the effective start of the returned slice +- `total_lines`: the full scrollback length at capture time +- `lines`: the requested window into that history + +Repeated reads without new output should be stable: the same buffer state should yield the same full snapshot, visible snapshot, and scrollback slice. + +## Detached and exited buffers + +Detached PTY buffers remain capturable. + +While detached: + +- full capture remains available +- visible snapshot remains available +- scrollback slices remain available +- title and other backend metadata remain available +- the most recent PTY size remains the size reported by capture APIs until another resize arrives + +Exited PTY buffers also remain capturable as long as the buffer record still exists. Writes, resizes, and kills must fail after exit, but capture reads continue to work. + +## Alternate-screen policy + +Alternate-screen ownership stays in the backend. + +- `capture_visible_snapshot` reflects the active screen and reports whether alternate screen is active +- full capture and scrollback semantics follow the backend's own history model +- leaving alternate screen returns visible capture to the primary screen state the backend preserved + +Embers locks down the observable behavior with tests rather than layering extra server-side alternate-screen state on top of the backend. + +## Resize behavior + +Resize updates the PTY and the keeper-owned backend surface. Future full and visible captures report the latest size and remain coherent across: + +- direct resize requests +- detach while preserving the retained size +- reattach into differently sized views once a new resize is applied + +## Code anchors + +- Runtime capture implementation: `crates/embers-server/src/server.rs` +- Runtime keeper capture surface: `crates/embers-server/src/buffer_runtime.rs` +- Backend-visible snapshot behavior: `crates/embers-server/src/terminal_backend.rs` +- PTY integration coverage: `crates/embers-test-support/tests/buffer_runtime.rs` diff --git a/docs/terminal-runtime.md b/docs/terminal-runtime.md new file mode 100644 index 00000000..947e69ae --- /dev/null +++ b/docs/terminal-runtime.md @@ -0,0 +1,88 @@ +# Terminal runtime contract + +This note locks down the runtime contract Embers uses for PTY-backed buffers. + +## Runtime ownership + +`Buffer` is the durable terminal runtime record. A PTY-backed `Buffer` owns: + +- the PTY/process command, cwd, and environment hints +- runtime identity and lifecycle state (`Created`, `Running`, `Interrupted`, `Exited`) +- the runtime keeper socket path used to reconnect after restore +- attachment state (`Attached(NodeId)` or `Detached`) +- the authoritative PTY size policy +- terminal-facing metadata surfaced outside layout code: + - title + - activity/bell state + - last snapshot sequence +- the snapshot source of truth via the buffer runtime/keeper backend + +`BufferView` is a renderer/layout attachment only. It owns: + +- session/layout placement (`NodeId`, parentage, split/tab membership) +- focus/zoom/follow-output view flags +- the last render size used by that view + +`BufferView` does not own PTY handles, process lifetime, terminal parsing state, or scrollback. + +## Closing a view vs killing a buffer + +Closing a view removes the `BufferView` node and transitions the buffer to `Detached`. + +- the PTY/process keeps running +- output continues accumulating +- snapshots and scrollback remain queryable +- title/activity metadata remains on the buffer record + +Killing a buffer targets the runtime, not the view tree. + +- the PTY child is terminated +- the buffer transitions to `Exited` +- capture remains available from the terminal backend snapshot state +- the buffer record may later be cleaned up once it is detached + +## Attachment and move semantics + +`Attached -> Detached` happens when the owning view is closed or explicitly detached. + +`Detached -> Attached` happens when the buffer is attached to a new leaf. + +`Attached -> Moved` is modeled as: + +1. detach/close the old view +2. keep the buffer runtime alive +3. attach the same buffer to the target leaf + +Moving tabs, leaves, or floating roots must never reset PTY state because the runtime stays with the `Buffer`, not the `BufferView`. + +## Detached buffer policy + +Detached PTY buffers keep the last assigned `pty_size` until another explicit resize arrives. + +While detached: + +- output continues to accumulate +- full and visible snapshots remain available +- scrollback remains available +- activity/bell state continues to update on the buffer record + +Reattaching a detached buffer does not recreate the process or reset terminal state. It only gives the durable runtime a new view attachment. + +## Runtime transitions + +The intended transition graph is: + +- `Created -> Running` when the runtime keeper is spawned or restored +- `Running -> Exited` when the child process terminates normally or is explicitly killed +- `Running/Created -> Interrupted` when a restore cannot reconnect to a keeper +- `Attached -> Detached` when a view closes +- `Detached -> Attached` when the buffer is attached to a leaf +- `Attached -> Moved` when the old view is removed and the same buffer is attached elsewhere +- `Exited -> cleaned up` when a detached exited buffer is removed from server state + +## Code anchors + +- Runtime model: `crates/embers-server/src/model.rs` +- State transitions: `crates/embers-server/src/state.rs` +- Runtime keeper and PTY lifecycle: `crates/embers-server/src/buffer_runtime.rs` +- Server/runtime wiring: `crates/embers-server/src/server.rs` diff --git a/docs/terminal-test-matrix.md b/docs/terminal-test-matrix.md new file mode 100644 index 00000000..1965976e --- /dev/null +++ b/docs/terminal-test-matrix.md @@ -0,0 +1,43 @@ +# Terminal test matrix + +This matrix maps the terminal-validation plan to the concrete regression suites that now guard the +behavior. + +## Runtime and backend contracts + +| Concern | Primary tests | Notes | +| --- | --- | --- | +| Buffer runtime ownership and PTY lifecycle | `crates/embers-test-support/tests/buffer_runtime.rs` | Locks down runtime state transitions and detached-buffer policy. | +| Backend boundary and capture semantics | `crates/embers-server/tests/backend.rs` and `crates/embers-server/tests/buffer_lifecycle.rs` | Verifies backend ownership, activity bookkeeping, and server-side lifecycle rules. | +| Byte-stream features and alternate-screen parsing | `crates/terminal-backend/tests/*` | Covers escape-sequence handling, visible snapshots, and parser/backend invariants. | + +## Client and render-source contracts + +| Concern | Primary tests | Notes | +| --- | --- | --- | +| Input routing and scripted actions | `crates/embers-client/tests/configured_client.rs` | Prefix handling, passthrough rules, scripted `send_keys` / `send_bytes`, reload behavior. | +| Activity, bell, and hidden-buffer metadata | `crates/embers-client/tests/e2e.rs` and `crates/embers-server/tests/buffer_lifecycle.rs` | Hidden and detached buffers keep activity, bell, and continuity metadata coherent. | +| Render invalidation and snapshot freshness | `crates/embers-client/tests/configured_client.rs` and `crates/embers-client/tests/e2e.rs` | Confirms the client renders from refreshed authoritative snapshots, not stale cache. | +| Full-screen and alternate-screen behavior | `crates/embers-client/tests/e2e.rs` | Verifies enter/exit semantics, hidden fullscreen buffers, and primary-screen restoration. | + +## Real PTY end-to-end workflows + +| Workflow | Primary tests | What it proves | +| --- | --- | --- | +| Spawn real client in PTY and run shell I/O | `crates/embers-cli/tests/interactive.rs::embers_without_subcommand_starts_server_and_client` | Embers can host a real shell and stay reachable through both the live client and CLI. | +| Attach, switch, detach, and reveal live clients | `crates/embers-cli/tests/interactive.rs` | Live PTY clients stay synchronized with server session targeting. | +| Local scrollback and OSC52 selection | `crates/embers-cli/tests/interactive.rs` | Client-local terminal UX features work under a real PTY. | +| Scripted input path in a live PTY client | `crates/embers-cli/tests/interactive.rs::scripted_input_bindings_reach_the_live_terminal_in_pty` | Scripted bindings still reach the focused terminal runtime end to end. | +| Config reload during terminal interaction | `crates/embers-cli/tests/interactive.rs::config_reload_updates_live_bindings_without_breaking_terminal_io` | Reloaded bindings take effect without breaking ongoing terminal input/output. | +| Split, move, detach, and reattach continuity | `crates/embers-cli/tests/interactive.rs::live_pty_client_preserves_buffers_across_layout_and_attachment_changes` | Layout and attachment changes do not reset the underlying shell process. | +| Hidden-buffer bell visibility and reveal continuity | `crates/embers-cli/tests/interactive.rs::hidden_buffer_bells_surface_in_the_attached_client_and_reveal_buffered_output` | Hidden buffers keep accumulating output, surface bell state, and reveal coherent content later. | +| Full-screen app entry and exit in the live client | `crates/embers-cli/tests/interactive.rs::fullscreen_terminal_transitions_render_in_the_live_client_pty` | The attached PTY client renders alternate-screen transitions the same way the backend models them. | + +## Harness notes + +- `crates/embers-test-support/src/pty.rs` provides the reusable PTY harness used by the CLI + integration tests. +- PTY read helpers include the recent output tail in timeout errors so failures are debuggable + without rerunning under a separate recorder. +- `crates/embers-cli/tests/interactive.rs` serializes PTY-heavy tests with a shared lock to reduce + flaky PTY pressure in CI.