Skip to content

Commit 39fa562

Browse files
authored
Merge pull request #258 from zackees/fix/terminal-bounded-write
fix(terminal): release a session slot when its client disconnects mid-write
2 parents 008d464 + c6f45bd commit 39fa562

6 files changed

Lines changed: 126 additions & 21 deletions

File tree

‎Cargo.lock‎

Lines changed: 4 additions & 4 deletions
Some generated files are not rendered by default. Learn more about customizing how changed files appear on GitHub.

‎Cargo.toml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -13,7 +13,7 @@ repository = "https://github.com/zackees/fastled-wasm"
1313
homepage = "https://github.com/zackees/fastled-wasm"
1414

1515
[workspace.dependencies]
16-
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"] }
16+
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"] }
1717

1818
[profile.release]
1919
debug = "line-tables-only"

‎crates/fastled-cli/Cargo.toml‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,7 @@ kernal-api = { workspace = true }
3030
# same-named build-dependency too and compile Tauri, GTK and glib for the
3131
# host build script.
3232
[build-dependencies]
33-
kernal-api-build = { git = "https://github.com/zackees/kernal-api.git", tag = "v0.1.10" }
33+
kernal-api-build = { git = "https://github.com/zackees/kernal-api.git", tag = "v0.1.11" }
3434

3535
[package.metadata.binstall]
3636
pkg-url = "{ repo }/releases/download/v{ version }/fastled-{ target }{ archive-suffix }"

‎crates/fastled-cli/src/terminal.rs‎

Lines changed: 46 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -93,6 +93,39 @@ fn exit_frame(status: impl std::fmt::Display) -> String {
9393
format!(r#"{{"exit":{status}}}"#)
9494
}
9595

96+
/// How long one slice of client input may wait for room in the terminal queue.
97+
const INPUT_SLICE_TIMEOUT: Duration = Duration::from_millis(50);
98+
99+
/// Write client input to the terminal, giving up as soon as the client is gone.
100+
///
101+
/// Returns `false` when the client has disconnected; the caller stops and lets
102+
/// the session drop.
103+
///
104+
/// A plain blocking write parks in the kernel once the terminal input queue
105+
/// fills, and nothing releases it: not cancelling the thread, and not ending
106+
/// the process reading the other end — killing the child leaves the writer
107+
/// exactly where it was. This worker owns both the session and its semaphore
108+
/// permit, so parking here holds a terminal slot until the foreground program
109+
/// decides to read again, which for a program that never reads is forever.
110+
///
111+
/// Writing in slices bounded by [`INPUT_SLICE_TIMEOUT`] is what keeps the
112+
/// disconnect observable. A zero-length result means the queue stayed full for
113+
/// that slice, so the loop re-checks the flag and waits again — the wait is the
114+
/// backpressure, not a spin.
115+
fn write_input(session: &mut PtySession, bytes: &[u8], closed: &AtomicBool) -> io::Result<bool> {
116+
let mut written = 0;
117+
while written < bytes.len() {
118+
if closed.load(Ordering::Relaxed) {
119+
return Ok(false);
120+
}
121+
match session.write_available(&bytes[written..], INPUT_SLICE_TIMEOUT)? {
122+
0 => {}
123+
accepted => written += accepted,
124+
}
125+
}
126+
Ok(true)
127+
}
128+
96129
fn parse_text(text: &str) -> io::Result<ClientMessage> {
97130
let Value::ObjectMembers(fields) =
98131
json::parse_members(text.as_bytes()).map_err(io::Error::other)?
@@ -181,6 +214,9 @@ pub(crate) async fn connect(
181214
let input_closed = Arc::new(AtomicBool::new(false));
182215
let input_pending_writes = Arc::clone(&pending_writes);
183216
let input_closed_task = Arc::clone(&input_closed);
217+
// The worker watches this so a disconnect ends a write it would otherwise
218+
// be parked inside.
219+
let input_closed_worker = Arc::clone(&input_closed);
184220
let input_task = async_engine::launch(async move {
185221
while let Ok(Some(message)) = reader.receive().await {
186222
let input = match message {
@@ -232,8 +268,16 @@ pub(crate) async fn connect(
232268
})?;
233269
loop {
234270
match input_rx.try_recv() {
235-
Ok(ClientMessage::Input(data)) => session.write(data.as_bytes())?,
236-
Ok(ClientMessage::Binary(data)) => session.write(&data)?,
271+
Ok(ClientMessage::Input(data)) => {
272+
if !write_input(&mut session, data.as_bytes(), &input_closed_worker)? {
273+
break;
274+
}
275+
}
276+
Ok(ClientMessage::Binary(data)) => {
277+
if !write_input(&mut session, &data, &input_closed_worker)? {
278+
break;
279+
}
280+
}
237281
Ok(ClientMessage::Resize { cols, rows }) => {
238282
session.resize(size(cols, rows)?)?
239283
}

‎tests/frontend/test_terminal.py‎

Lines changed: 73 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -104,7 +104,7 @@ def read_output() -> None:
104104
assert process.poll() is None, "".join(lines)
105105
time.sleep(0.05)
106106
assert url, "server did not announce its URL: " + "".join(lines)
107-
yield url, served
107+
yield url, served, process.pid
108108
finally:
109109
process.terminate()
110110
process.wait(timeout=10)
@@ -115,7 +115,7 @@ def test_terminal_240_interactive_browser(
115115
terminal_server: Any, browser_name: str
116116
) -> None:
117117
playwright = pytest.importorskip("playwright.sync_api")
118-
url, expected_cwd = terminal_server
118+
url, expected_cwd, _ = terminal_server
119119
with playwright.sync_playwright() as manager:
120120
browser = getattr(manager, browser_name).launch()
121121
page = browser.new_page(viewport={"width": 1200, "height": 900})
@@ -245,13 +245,49 @@ def wait_output(text: str) -> None:
245245
browser.close()
246246

247247

248+
def _descendants(pid: int) -> list[int]:
249+
found: list[int] = []
250+
pending = [pid]
251+
while pending:
252+
parent = pending.pop()
253+
for task in Path(f"/proc/{parent}/task").glob("*"):
254+
try:
255+
children = (task / "children").read_text().split()
256+
except OSError:
257+
continue
258+
for child in map(int, children):
259+
found.append(child)
260+
pending.append(child)
261+
return found
262+
263+
264+
def _pty_holders(pid: int) -> list[int]:
265+
holders = []
266+
for candidate in [pid, *_descendants(pid)]:
267+
try:
268+
fds = list(Path(f"/proc/{candidate}/fd").iterdir())
269+
except OSError:
270+
continue
271+
for fd in fds:
272+
try:
273+
if os.readlink(fd) == "/dev/ptmx":
274+
holders.append(candidate)
275+
break
276+
except OSError:
277+
continue
278+
return holders
279+
280+
248281
@pytest.mark.skipif(os.name == "nt", reason="Unix stty regression")
249-
def test_terminal_240_disconnect_with_blocked_stdin(terminal_server: Any) -> None:
282+
@pytest.mark.parametrize("browser_name", ["chromium", "webkit"])
283+
def test_terminal_240_disconnect_with_blocked_stdin(
284+
terminal_server: Any, browser_name: str
285+
) -> None:
250286
"""RED: four blocked writers leaked every slot; fifth upgrade was HTTP 429."""
251287
playwright = pytest.importorskip("playwright.sync_api")
252-
url, _ = terminal_server
288+
url, _, server_pid = terminal_server
253289
with playwright.sync_playwright() as manager:
254-
browser = manager.chromium.launch()
290+
browser = getattr(manager, browser_name).launch()
255291
page = browser.new_page()
256292
page.goto(url)
257293
result = page.evaluate("""async () => {
@@ -277,13 +313,38 @@ def test_terminal_240_disconnect_with_blocked_stdin(terminal_server: Any) -> Non
277313
}
278314
};
279315
});
280-
await new Promise(resolve => setTimeout(resolve, 400));
281316
}
282-
return await new Promise(resolve => {
283-
const ws = new WebSocket(url);
284-
ws.onopen = () => { ws.close(); resolve(true); };
285-
ws.onerror = () => resolve(false);
286-
});
317+
// Every slot must come back promptly, not just one of them, and long
318+
// before the blocked `sleep 30` foreground programs would exit.
319+
const started = performance.now();
320+
const deadline = started + 5000;
321+
while (performance.now() < deadline) {
322+
const sockets = [];
323+
const opened = await Promise.all([0, 1, 2, 3].map(() => new Promise(resolve => {
324+
const ws = new WebSocket(url);
325+
sockets.push(ws);
326+
ws.onopen = () => resolve(true);
327+
ws.onerror = () => resolve(false);
328+
})));
329+
sockets.forEach(ws => ws.close());
330+
if (opened.every(Boolean)) return performance.now() - started;
331+
await new Promise(resolve => setTimeout(resolve, 100));
332+
}
333+
return null;
287334
}""")
288-
assert result, "disconnected blocked writers leaked all terminal slots"
335+
assert result is not None, "disconnected blocked writers leaked terminal slots"
289336
browser.close()
337+
deadline = time.monotonic() + 10
338+
holders = _pty_holders(server_pid)
339+
while holders and time.monotonic() < deadline:
340+
time.sleep(0.1)
341+
holders = _pty_holders(server_pid)
342+
assert not holders, f"sessions outlived their clients: {holders}"
343+
sleepers = []
344+
for pid in _descendants(server_pid):
345+
try:
346+
if b"sleep" in Path(f"/proc/{pid}/cmdline").read_bytes():
347+
sleepers.append(pid)
348+
except OSError:
349+
continue
350+
assert not sleepers, f"blocked foreground programs survived: {sleepers}"

‎tests/unit/test_kernal_boundary.py‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -195,7 +195,7 @@ def test_kernal_api_is_the_only_rust_dependency():
195195
assert not section.get("build-dependencies"), section
196196
assert set(package["dependencies"]) == {"kernal-api"}
197197
kernal = workspace["workspace"]["dependencies"]["kernal-api"]
198-
assert kernal["tag"] == "v0.1.10"
198+
assert kernal["tag"] == "v0.1.11"
199199
assert "rev" not in kernal and "path" not in kernal
200200
assert "hash-sha256" in kernal["features"]
201201
assert "patch" not in workspace

0 commit comments

Comments
 (0)