From c6f45bd056dc8298cab16277758a088babb46148 Mon Sep 17 00:00:00 2001 From: Zach Vorhies Date: Wed, 16 Sep 2026 19:26:17 -0700 Subject: [PATCH] fix(terminal): release a session slot when its client disconnects mid-write A client input write parked on a full PTY input queue held the worker, the session and its semaphore permit until the foreground program exited; killing the process tree does not release a blocked write. Follow kernal-api to v0.1.11 and write client input through PtySession::write_available in 50 ms slices, checking the disconnect flag between them. The blocked-stdin regression now runs on Chromium and WebKit, requires all four slots back within 5 s (the blocked program is `sleep 30`), and asserts no descendant of the CLI still holds a PTY master or runs the blocked program. Fixes #256 Co-Authored-By: Claude Opus 5 (1M context) --- Cargo.lock | 8 +-- Cargo.toml | 2 +- crates/fastled-cli/Cargo.toml | 2 +- crates/fastled-cli/src/terminal.rs | 48 ++++++++++++++++- tests/frontend/test_terminal.py | 85 +++++++++++++++++++++++++----- tests/unit/test_kernal_boundary.py | 2 +- 6 files changed, 126 insertions(+), 21 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index 8e62a22f..17ac9d20 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2024,8 +2024,8 @@ dependencies = [ [[package]] name = "kernal-api" -version = "0.1.10" -source = "git+https://github.com/zackees/kernal-api.git?tag=v0.1.10#f2fd5f5c75c60d1b445ca924fdc8059dce9c52e5" +version = "0.1.11" +source = "git+https://github.com/zackees/kernal-api.git?tag=v0.1.11#da40ae5c20e86ab40b89ad3b5a665614f40df3cc" dependencies = [ "blake3", "bytes", @@ -2080,8 +2080,8 @@ dependencies = [ [[package]] name = "kernal-api-build" -version = "0.1.10" -source = "git+https://github.com/zackees/kernal-api.git?tag=v0.1.10#f2fd5f5c75c60d1b445ca924fdc8059dce9c52e5" +version = "0.1.11" +source = "git+https://github.com/zackees/kernal-api.git?tag=v0.1.11#da40ae5c20e86ab40b89ad3b5a665614f40df3cc" dependencies = [ "embed-resource", ] diff --git a/Cargo.toml b/Cargo.toml index ac98f51c..21370d80 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -13,7 +13,7 @@ repository = "https://github.com/zackees/fastled-wasm" homepage = "https://github.com/zackees/fastled-wasm" [workspace.dependencies] -kernal-api = { git = "https://github.com/zackees/kernal-api.git", tag = "v0.1.10", features = ["fs", "fs-watch", "hash-sha256", "archive", "http-client", "http-server", "websocket", "event-stream", "secure-random", "text-similarity", "pty", "terminal-input", "terminal-style", "command-arguments", "command-schema", "config-toml", "source-cpp", "json", "error-context"] } +kernal-api = { git = "https://github.com/zackees/kernal-api.git", tag = "v0.1.11", features = ["fs", "fs-watch", "hash-sha256", "archive", "http-client", "http-server", "websocket", "event-stream", "secure-random", "text-similarity", "pty", "terminal-input", "terminal-style", "command-arguments", "command-schema", "config-toml", "source-cpp", "json", "error-context"] } [profile.release] debug = "line-tables-only" diff --git a/crates/fastled-cli/Cargo.toml b/crates/fastled-cli/Cargo.toml index 1c003995..1229aa4e 100644 --- a/crates/fastled-cli/Cargo.toml +++ b/crates/fastled-cli/Cargo.toml @@ -30,7 +30,7 @@ kernal-api = { workspace = true } # same-named build-dependency too and compile Tauri, GTK and glib for the # host build script. [build-dependencies] -kernal-api-build = { git = "https://github.com/zackees/kernal-api.git", tag = "v0.1.10" } +kernal-api-build = { git = "https://github.com/zackees/kernal-api.git", tag = "v0.1.11" } [package.metadata.binstall] pkg-url = "{ repo }/releases/download/v{ version }/fastled-{ target }{ archive-suffix }" diff --git a/crates/fastled-cli/src/terminal.rs b/crates/fastled-cli/src/terminal.rs index ee913485..ad7419fa 100644 --- a/crates/fastled-cli/src/terminal.rs +++ b/crates/fastled-cli/src/terminal.rs @@ -93,6 +93,39 @@ fn exit_frame(status: impl std::fmt::Display) -> String { format!(r#"{{"exit":{status}}}"#) } +/// How long one slice of client input may wait for room in the terminal queue. +const INPUT_SLICE_TIMEOUT: Duration = Duration::from_millis(50); + +/// Write client input to the terminal, giving up as soon as the client is gone. +/// +/// Returns `false` when the client has disconnected; the caller stops and lets +/// the session drop. +/// +/// A plain blocking write parks in the kernel once the terminal input queue +/// fills, and nothing releases it: not cancelling the thread, and not ending +/// the process reading the other end — killing the child leaves the writer +/// exactly where it was. This worker owns both the session and its semaphore +/// permit, so parking here holds a terminal slot until the foreground program +/// decides to read again, which for a program that never reads is forever. +/// +/// Writing in slices bounded by [`INPUT_SLICE_TIMEOUT`] is what keeps the +/// disconnect observable. A zero-length result means the queue stayed full for +/// that slice, so the loop re-checks the flag and waits again — the wait is the +/// backpressure, not a spin. +fn write_input(session: &mut PtySession, bytes: &[u8], closed: &AtomicBool) -> io::Result { + let mut written = 0; + while written < bytes.len() { + if closed.load(Ordering::Relaxed) { + return Ok(false); + } + match session.write_available(&bytes[written..], INPUT_SLICE_TIMEOUT)? { + 0 => {} + accepted => written += accepted, + } + } + Ok(true) +} + fn parse_text(text: &str) -> io::Result { let Value::ObjectMembers(fields) = json::parse_members(text.as_bytes()).map_err(io::Error::other)? @@ -181,6 +214,9 @@ pub(crate) async fn connect( let input_closed = Arc::new(AtomicBool::new(false)); let input_pending_writes = Arc::clone(&pending_writes); let input_closed_task = Arc::clone(&input_closed); + // The worker watches this so a disconnect ends a write it would otherwise + // be parked inside. + let input_closed_worker = Arc::clone(&input_closed); let input_task = async_engine::launch(async move { while let Ok(Some(message)) = reader.receive().await { let input = match message { @@ -232,8 +268,16 @@ pub(crate) async fn connect( })?; loop { match input_rx.try_recv() { - Ok(ClientMessage::Input(data)) => session.write(data.as_bytes())?, - Ok(ClientMessage::Binary(data)) => session.write(&data)?, + Ok(ClientMessage::Input(data)) => { + if !write_input(&mut session, data.as_bytes(), &input_closed_worker)? { + break; + } + } + Ok(ClientMessage::Binary(data)) => { + if !write_input(&mut session, &data, &input_closed_worker)? { + break; + } + } Ok(ClientMessage::Resize { cols, rows }) => { session.resize(size(cols, rows)?)? } diff --git a/tests/frontend/test_terminal.py b/tests/frontend/test_terminal.py index f5e259a8..4963446d 100644 --- a/tests/frontend/test_terminal.py +++ b/tests/frontend/test_terminal.py @@ -104,7 +104,7 @@ def read_output() -> None: assert process.poll() is None, "".join(lines) time.sleep(0.05) assert url, "server did not announce its URL: " + "".join(lines) - yield url, served + yield url, served, process.pid finally: process.terminate() process.wait(timeout=10) @@ -115,7 +115,7 @@ def test_terminal_240_interactive_browser( terminal_server: Any, browser_name: str ) -> None: playwright = pytest.importorskip("playwright.sync_api") - url, expected_cwd = terminal_server + url, expected_cwd, _ = terminal_server with playwright.sync_playwright() as manager: browser = getattr(manager, browser_name).launch() page = browser.new_page(viewport={"width": 1200, "height": 900}) @@ -245,13 +245,49 @@ def wait_output(text: str) -> None: browser.close() +def _descendants(pid: int) -> list[int]: + found: list[int] = [] + pending = [pid] + while pending: + parent = pending.pop() + for task in Path(f"/proc/{parent}/task").glob("*"): + try: + children = (task / "children").read_text().split() + except OSError: + continue + for child in map(int, children): + found.append(child) + pending.append(child) + return found + + +def _pty_holders(pid: int) -> list[int]: + holders = [] + for candidate in [pid, *_descendants(pid)]: + try: + fds = list(Path(f"/proc/{candidate}/fd").iterdir()) + except OSError: + continue + for fd in fds: + try: + if os.readlink(fd) == "/dev/ptmx": + holders.append(candidate) + break + except OSError: + continue + return holders + + @pytest.mark.skipif(os.name == "nt", reason="Unix stty regression") -def test_terminal_240_disconnect_with_blocked_stdin(terminal_server: Any) -> None: +@pytest.mark.parametrize("browser_name", ["chromium", "webkit"]) +def test_terminal_240_disconnect_with_blocked_stdin( + terminal_server: Any, browser_name: str +) -> None: """RED: four blocked writers leaked every slot; fifth upgrade was HTTP 429.""" playwright = pytest.importorskip("playwright.sync_api") - url, _ = terminal_server + url, _, server_pid = terminal_server with playwright.sync_playwright() as manager: - browser = manager.chromium.launch() + browser = getattr(manager, browser_name).launch() page = browser.new_page() page.goto(url) result = page.evaluate("""async () => { @@ -277,13 +313,38 @@ def test_terminal_240_disconnect_with_blocked_stdin(terminal_server: Any) -> Non } }; }); - await new Promise(resolve => setTimeout(resolve, 400)); } - return await new Promise(resolve => { - const ws = new WebSocket(url); - ws.onopen = () => { ws.close(); resolve(true); }; - ws.onerror = () => resolve(false); - }); + // Every slot must come back promptly, not just one of them, and long + // before the blocked `sleep 30` foreground programs would exit. + const started = performance.now(); + const deadline = started + 5000; + while (performance.now() < deadline) { + const sockets = []; + const opened = await Promise.all([0, 1, 2, 3].map(() => new Promise(resolve => { + const ws = new WebSocket(url); + sockets.push(ws); + ws.onopen = () => resolve(true); + ws.onerror = () => resolve(false); + }))); + sockets.forEach(ws => ws.close()); + if (opened.every(Boolean)) return performance.now() - started; + await new Promise(resolve => setTimeout(resolve, 100)); + } + return null; }""") - assert result, "disconnected blocked writers leaked all terminal slots" + assert result is not None, "disconnected blocked writers leaked terminal slots" browser.close() + deadline = time.monotonic() + 10 + holders = _pty_holders(server_pid) + while holders and time.monotonic() < deadline: + time.sleep(0.1) + holders = _pty_holders(server_pid) + assert not holders, f"sessions outlived their clients: {holders}" + sleepers = [] + for pid in _descendants(server_pid): + try: + if b"sleep" in Path(f"/proc/{pid}/cmdline").read_bytes(): + sleepers.append(pid) + except OSError: + continue + assert not sleepers, f"blocked foreground programs survived: {sleepers}" diff --git a/tests/unit/test_kernal_boundary.py b/tests/unit/test_kernal_boundary.py index 8bc1658f..0e93dbe2 100644 --- a/tests/unit/test_kernal_boundary.py +++ b/tests/unit/test_kernal_boundary.py @@ -195,7 +195,7 @@ def test_kernal_api_is_the_only_rust_dependency(): assert not section.get("build-dependencies"), section assert set(package["dependencies"]) == {"kernal-api"} kernal = workspace["workspace"]["dependencies"]["kernal-api"] - assert kernal["tag"] == "v0.1.10" + assert kernal["tag"] == "v0.1.11" assert "rev" not in kernal and "path" not in kernal assert "hash-sha256" in kernal["features"] assert "patch" not in workspace