Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 4 additions & 4 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down
2 changes: 1 addition & 1 deletion crates/fastled-cli/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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 }"
Expand Down
48 changes: 46 additions & 2 deletions crates/fastled-cli/src/terminal.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<bool> {
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<ClientMessage> {
let Value::ObjectMembers(fields) =
json::parse_members(text.as_bytes()).map_err(io::Error::other)?
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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)?)?
}
Expand Down
85 changes: 73 additions & 12 deletions tests/frontend/test_terminal.py
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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})
Expand Down Expand Up @@ -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 () => {
Expand All @@ -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}"
2 changes: 1 addition & 1 deletion tests/unit/test_kernal_boundary.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Loading