From 7fceb709d8070261c510641d53a066e2c4a7129b Mon Sep 17 00:00:00 2001 From: Robert Grayson Date: Thu, 19 Mar 2026 18:38:11 -0400 Subject: [PATCH 01/12] test: update snapshots for io.js and keyboard.js in CLI tests JsRunFileLayoutTest: add io.js and keyboard.js to expected file trees in both std/ and temper-core/ directories. ReplTest: allow scheduler call in Lua translation output. Signed-off-by: Robert Grayson Co-Authored-By: Claude Opus 4.6 (1M context) Signed-off-by: Robert Grayson --- .../test/kotlin/lang/temper/cli/js/JsRunFileLayoutTest.kt | 8 ++++++++ cli/src/test/kotlin/lang/temper/cli/repl/ReplTest.kt | 1 + 2 files changed, 9 insertions(+) diff --git a/cli/src/test/kotlin/lang/temper/cli/js/JsRunFileLayoutTest.kt b/cli/src/test/kotlin/lang/temper/cli/js/JsRunFileLayoutTest.kt index 816afaec..8bad9be8 100644 --- a/cli/src/test/kotlin/lang/temper/cli/js/JsRunFileLayoutTest.kt +++ b/cli/src/test/kotlin/lang/temper/cli/js/JsRunFileLayoutTest.kt @@ -117,6 +117,8 @@ class JsRunFileLayoutTest { | ┃ ┃ ┃ ┣━index.js | ┃ ┃ ┃ ┣━int.js | ┃ ┃ ┃ ┣━interface.js + | ┃ ┃ ┃ ┣━io.js + | ┃ ┃ ┃ ┣━keyboard.js | ┃ ┃ ┃ ┣━listed.js | ┃ ┃ ┃ ┣━mapped.js | ┃ ┃ ┃ ┣━net.js @@ -127,7 +129,9 @@ class JsRunFileLayoutTest { | ┃ ┃ ┃ ┗━tsconfig.json | ┃ ┃ ┗━std/ | ┃ ┃ ┣━index.js + | ┃ ┃ ┣━io.js | ┃ ┃ ┣━json.js + | ┃ ┃ ┣━keyboard.js | ┃ ┃ ┣━net.js | ┃ ┃ ┣━package.json | ┃ ┃ ┣━regex.js @@ -165,7 +169,9 @@ class JsRunFileLayoutTest { | ┣━package.json | ┣━std/ | ┃ ┣━index.js + | ┃ ┣━io.js | ┃ ┣━json.js + | ┃ ┣━keyboard.js | ┃ ┣━net.js | ┃ ┣━package.json | ┃ ┣━regex.js @@ -182,6 +188,8 @@ class JsRunFileLayoutTest { | ┣━index.js | ┣━int.js | ┣━interface.js + | ┣━io.js + | ┣━keyboard.js | ┣━listed.js | ┣━mapped.js | ┣━net.js diff --git a/cli/src/test/kotlin/lang/temper/cli/repl/ReplTest.kt b/cli/src/test/kotlin/lang/temper/cli/repl/ReplTest.kt index 0e5b5325..477d423f 100644 --- a/cli/src/test/kotlin/lang/temper/cli/repl/ReplTest.kt +++ b/cli/src/test/kotlin/lang/temper/cli/repl/ReplTest.kt @@ -529,6 +529,7 @@ class ReplTest { | i0000[.]lua: text/x-lua | .* | return__0 = 2; + |.* | exports = \{}; | return exports; """.trimMargin(), From 7718c2f19df322cc067108b0408eb0f0b43d80b9 Mon Sep 17 00:00:00 2001 From: Robert Grayson Date: Thu, 19 Mar 2026 22:55:33 -0400 Subject: [PATCH 02/12] fix: spawn OS thread for Rust nextKeypress to avoid blocking async runner The Rust backend uses a SingleThreadAsyncRunner, so the blocking crossterm::event::read() call was starving the game tick coroutine. Spawn a real OS thread for the blocking read and complete the promise from there, allowing other async tasks to proceed concurrently. Signed-off-by: Robert Grayson Co-Authored-By: Claude Opus 4.6 (1M context) Signed-off-by: Robert Grayson --- .../temper/be/rust/std/keyboard/support.rs | 69 +++++++++---------- 1 file changed, 34 insertions(+), 35 deletions(-) diff --git a/be-rust/src/commonMain/resources/lang/temper/be/rust/std/keyboard/support.rs b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/keyboard/support.rs index 83574a87..79789114 100644 --- a/be-rust/src/commonMain/resources/lang/temper/be/rust/std/keyboard/support.rs +++ b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/keyboard/support.rs @@ -1,44 +1,43 @@ use std::sync::Arc; -use temper_core::{Promise, PromiseBuilder, SafeGenerator}; +use temper_core::{Promise, PromiseBuilder}; pub fn std_next_keypress() -> Promise>> { let pb = PromiseBuilder::new(); let promise = pb.promise(); - crate::run_async(Arc::new(move || { - let pb = pb.clone(); - SafeGenerator::from_fn(Arc::new(move |_generator: SafeGenerator<()>| { - use crossterm::terminal; - use crossterm::event::{self, Event, KeyCode, KeyEvent}; + // Spawn a real OS thread so the blocking read doesn't starve + // the single-threaded async runner. The promise completion will + // fire the on_ready continuation from this thread. + std::thread::spawn(move || { + use crossterm::terminal; + use crossterm::event::{self, Event, KeyCode, KeyEvent}; - terminal::enable_raw_mode().ok(); - let result = loop { - match event::read() { - Ok(Event::Key(KeyEvent { code, .. })) => { - break match code { - KeyCode::Up => Some("ArrowUp"), - KeyCode::Down => Some("ArrowDown"), - KeyCode::Left => Some("ArrowLeft"), - KeyCode::Right => Some("ArrowRight"), - KeyCode::Enter => Some("Enter"), - KeyCode::Esc => Some("Escape"), - KeyCode::Backspace => Some("Backspace"), - KeyCode::Tab => Some("Tab"), - KeyCode::Char(c) => { - terminal::disable_raw_mode().ok(); - pb.complete(Some(Arc::new(c.to_string()))); - return None; - } - _ => Some("Unknown"), - }; - } - Err(_) => { break None; } - _ => continue, + terminal::enable_raw_mode().ok(); + let result = loop { + match event::read() { + Ok(Event::Key(KeyEvent { code, .. })) => { + break match code { + KeyCode::Up => Some("ArrowUp"), + KeyCode::Down => Some("ArrowDown"), + KeyCode::Left => Some("ArrowLeft"), + KeyCode::Right => Some("ArrowRight"), + KeyCode::Enter => Some("Enter"), + KeyCode::Esc => Some("Escape"), + KeyCode::Backspace => Some("Backspace"), + KeyCode::Tab => Some("Tab"), + KeyCode::Char(c) => { + terminal::disable_raw_mode().ok(); + pb.complete(Some(Arc::new(c.to_string()))); + return; + } + _ => Some("Unknown"), + }; } - }; - terminal::disable_raw_mode().ok(); - pb.complete(result.map(|s| Arc::new(s.to_string()))); - None - })) - })); + Err(_) => { break None; } + _ => continue, + } + }; + terminal::disable_raw_mode().ok(); + pb.complete(result.map(|s| Arc::new(s.to_string()))); + }); promise } From 03f6b991482148da4e4bbe0ac36daadf140895b7 Mon Sep 17 00:00:00 2001 From: Robert Grayson Date: Thu, 19 Mar 2026 23:32:21 -0400 Subject: [PATCH 03/12] feat: add std/ws WebSocket module and std/io terminal size detection Rebased onto do-crimes-to-play-snake: resolved conflicts with std/keyboard, updated Rust terminal size to use crossterm instead of libc, added ws to stdSupportNeeders set. --- .../kotlin/lang/temper/be/js/JsBackend.kt | 1 + .../lang/temper/be/js/JsSupportNetwork.kt | 11 + .../lang/temper/be/js/temper-core/index.js | 1 + .../lang/temper/be/js/temper-core/io.js | 16 ++ .../lang/temper/be/js/temper-core/ws.js | 143 ++++++++++ .../kotlin/lang/temper/be/rust/RustBackend.kt | 5 +- .../lang/temper/be/rust/RustSupportNetwork.kt | 16 ++ .../lang/temper/be/rust/std/io/support.rs | 14 + .../lang/temper/be/rust/std/ws/support.rs | 245 ++++++++++++++++++ .../kotlin/lang/temper/be/README.md | 8 + .../commonMain/resources/std/config.temper.md | 1 + .../commonMain/resources/std/io/io.temper.md | 15 ++ .../commonMain/resources/std/ws/ws.temper.md | 64 +++++ 13 files changed, 538 insertions(+), 2 deletions(-) create mode 100644 be-js/src/commonMain/resources/lang/temper/be/js/temper-core/ws.js create mode 100644 be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs create mode 100644 frontend/src/commonMain/resources/std/ws/ws.temper.md diff --git a/be-js/src/commonMain/kotlin/lang/temper/be/js/JsBackend.kt b/be-js/src/commonMain/kotlin/lang/temper/be/js/JsBackend.kt index bca62d0e..a6575d48 100644 --- a/be-js/src/commonMain/kotlin/lang/temper/be/js/JsBackend.kt +++ b/be-js/src/commonMain/kotlin/lang/temper/be/js/JsBackend.kt @@ -627,6 +627,7 @@ class JsBackend private constructor( filePath("pair.js"), filePath("regex.js"), filePath("string.js"), + filePath("ws.js"), ) override val specifics: NodeSpecifics get() = NodeSpecifics diff --git a/be-js/src/commonMain/kotlin/lang/temper/be/js/JsSupportNetwork.kt b/be-js/src/commonMain/kotlin/lang/temper/be/js/JsSupportNetwork.kt index 090eae42..5da9ddf2 100644 --- a/be-js/src/commonMain/kotlin/lang/temper/be/js/JsSupportNetwork.kt +++ b/be-js/src/commonMain/kotlin/lang/temper/be/js/JsSupportNetwork.kt @@ -468,6 +468,8 @@ private val supportedAutoConnecteds = setOf( // std/io "stdSleep", "stdReadLine", + "stdTermCols", + "stdTermRows", // std/keyboard "stdNextKeypress", // std/net @@ -476,6 +478,15 @@ private val supportedAutoConnecteds = setOf( "NetResponse::getStatus", "NetResponse::getContentType", "NetResponse::getBodyContent", + // std/ws + "wsListen", + "wsAccept", + "wsConnect", + "wsSend", + "wsRecv", + "wsClose", + "WsServer", + "WsConnection", ) private val supportedMappedConnecteds = mapOf( diff --git a/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/index.js b/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/index.js index 2d507424..695cb4e0 100644 --- a/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/index.js +++ b/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/index.js @@ -14,3 +14,4 @@ export * from "./net.js"; export * from "./pair.js"; export * from "./regex.js"; export * from "./string.js"; +export * from "./ws.js"; diff --git a/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/io.js b/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/io.js index 2101cd8d..b873052a 100644 --- a/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/io.js +++ b/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/io.js @@ -12,6 +12,22 @@ export function stdSleep(ms) { /** * @returns {Promise} */ +/** @returns {number} */ +export function stdTermCols() { + if (typeof process !== 'undefined' && process.stdout && process.stdout.columns) { + return process.stdout.columns; + } + return 80; +} + +/** @returns {number} */ +export function stdTermRows() { + if (typeof process !== 'undefined' && process.stdout && process.stdout.rows) { + return process.stdout.rows; + } + return 24; +} + export function stdReadLine() { return new Promise(resolve => { if (typeof process !== 'undefined' && process.stdin) { diff --git a/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/ws.js b/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/ws.js new file mode 100644 index 00000000..26fc3bf8 --- /dev/null +++ b/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/ws.js @@ -0,0 +1,143 @@ +import { empty } from "./core.js"; + +// We dynamically import 'ws' so this module works in environments where +// it is not installed (the functions will throw at call time instead of +// at import time). +let _WebSocketServer; +let _WebSocket; +const _wsReady = (async () => { + try { + const ws = await import("ws"); + _WebSocketServer = ws.WebSocketServer; + _WebSocket = ws.default || ws.WebSocket; + } catch (_) { + // ws not available — functions will throw when called + } +})(); + +function _requireWs() { + if (!_WebSocket) { + throw new Error("WebSocket support requires the 'ws' npm package. Run: npm install ws"); + } +} + +function _setupMessageQueue(ws) { + if (!ws._mq) { + ws._mq = []; + ws._wr = []; + ws._closed = false; + ws.on("message", (data) => { + const msg = data.toString(); + const waiting = ws._wr.shift(); + if (waiting) { waiting(msg); } + else { ws._mq.push(msg); } + }); + ws.on("close", () => { + ws._closed = true; + while (ws._wr.length) { + ws._wr.shift()(null); + } + }); + ws.on("error", () => { + ws._closed = true; + while (ws._wr.length) { + ws._wr.shift()(null); + } + }); + } +} + +/** + * @param {number} port + * @returns {Promise} + */ +export async function wsListen(port) { + await _wsReady; + _requireWs(); + return new Promise((resolve, reject) => { + const server = new _WebSocketServer({ port }); + server._pending = []; + server._waiting = []; + server.on("connection", (ws) => { + _setupMessageQueue(ws); + const waiting = server._waiting.shift(); + if (waiting) { waiting(ws); } + else { server._pending.push(ws); } + }); + server.on("listening", () => resolve(server)); + server.on("error", reject); + }); +} + +/** + * @param {object} server + * @returns {Promise} + */ +export function wsAccept(server) { + return new Promise((resolve) => { + const pending = server._pending.shift(); + if (pending) { resolve(pending); } + else { server._waiting.push(resolve); } + }); +} + +/** + * @param {string} url + * @returns {Promise} + */ +export async function wsConnect(url) { + await _wsReady; + _requireWs(); + return new Promise((resolve, reject) => { + const ws = new _WebSocket(url); + _setupMessageQueue(ws); + ws.on("open", () => resolve(ws)); + ws.on("error", reject); + }); +} + +/** + * @param {object} conn + * @param {string} msg + * @returns {Promise} + */ +export function wsSend(conn, msg) { + return new Promise((resolve, reject) => { + conn.send(msg, (err) => { + if (err) reject(err); + else resolve(empty()); + }); + }); +} + +/** + * @param {object} conn + * @returns {Promise} + */ +export function wsRecv(conn) { + _setupMessageQueue(conn); + return new Promise((resolve) => { + if (conn._closed && conn._mq.length === 0) { + resolve(null); + return; + } + const queued = conn._mq.shift(); + if (queued !== undefined) { resolve(queued); } + else { conn._wr.push(resolve); } + }); +} + +/** + * @param {object} conn + * @returns {Promise} + */ +export function wsClose(conn) { + return new Promise((resolve) => { + if (conn.readyState !== undefined && conn.readyState > 1) { + resolve(empty()); + return; + } + conn.on("close", () => resolve(empty())); + conn.close(); + }); +} diff --git a/be-rust/src/commonMain/kotlin/lang/temper/be/rust/RustBackend.kt b/be-rust/src/commonMain/kotlin/lang/temper/be/rust/RustBackend.kt index ef5a2c15..7b5fcea7 100644 --- a/be-rust/src/commonMain/kotlin/lang/temper/be/rust/RustBackend.kt +++ b/be-rust/src/commonMain/kotlin/lang/temper/be/rust/RustBackend.kt @@ -184,11 +184,12 @@ class RustBackend(setup: BackendSetup) : Backend(Facto // Below aren't dependencies section anymore, but eh. append("\n") append("[features]\n") - append("io = []\n") + append("io = [\"crossterm\"]\n") append("keyboard = [\"crossterm\"]\n") append("net = [\"ureq\"]\n") // Implied: append("regex = [\"regex\"]\n") append("temporal = [\"time\"]\n") + append("ws = [\"tungstenite\"]\n") } } val packageFields = buildMap { @@ -256,7 +257,7 @@ class RustBackend(setup: BackendSetup) : Backend(Facto private val resourceBase = dirPath("lang", "temper", "be", "rust") private val coreResourceBase = resourceBase.resolveDir("temper-core") private val stdResourceBase = resourceBase.resolveDir("std") - val stdSupportNeeders = setOf("io", "keyboard", "net", "regex", "temporal") + val stdSupportNeeders = setOf("io", "keyboard", "net", "regex", "temporal", "ws") val stdFeatures = stdSupportNeeders // same set today but maybe not guaranteed private val templateResourceBase = resourceBase.resolveDir("library-template") diff --git a/be-rust/src/commonMain/kotlin/lang/temper/be/rust/RustSupportNetwork.kt b/be-rust/src/commonMain/kotlin/lang/temper/be/rust/RustSupportNetwork.kt index ad60c5c5..1f345b87 100644 --- a/be-rust/src/commonMain/kotlin/lang/temper/be/rust/RustSupportNetwork.kt +++ b/be-rust/src/commonMain/kotlin/lang/temper/be/rust/RustSupportNetwork.kt @@ -917,6 +917,14 @@ private val netSend = FunctionCall("stdNetSend", "send_request", cloneEvenIfFirs private val stdSleep = FunctionCall("stdSleep", "temper_std::io::std_sleep") private val stdReadLine = FunctionCall("stdReadLine", "temper_std::io::std_read_line") private val stdNextKeypress = FunctionCall("stdNextKeypress", "temper_std::keyboard::std_next_keypress") +private val stdTermCols = FunctionCall("stdTermCols", "temper_std::io::std_term_cols") +private val stdTermRows = FunctionCall("stdTermRows", "temper_std::io::std_term_rows") +private val wsListen = FunctionCall("wsListen", "temper_std::ws::ws_listen") +private val wsAccept = FunctionCall("wsAccept", "temper_std::ws::ws_accept", cloneEvenIfFirst = true) +private val wsConnect = FunctionCall("wsConnect", "temper_std::ws::ws_connect", cloneEvenIfFirst = true) +private val wsSend = FunctionCall("wsSend", "temper_std::ws::ws_send", cloneEvenIfFirst = true) +private val wsRecv = FunctionCall("wsRecv", "temper_std::ws::ws_recv", cloneEvenIfFirst = true) +private val wsClose = FunctionCall("wsClose", "temper_std::ws::ws_close", cloneEvenIfFirst = true) internal object PairConstructor : RustInlineSupportCode( "Pair::constructor", @@ -1145,6 +1153,14 @@ private val connectedReferences = listOf( stdSleep, stdReadLine, stdNextKeypress, + stdTermCols, + stdTermRows, + wsListen, + wsAccept, + wsConnect, + wsSend, + wsRecv, + wsClose, promiseBuilderComplete, PairConstructor, regexCompileFormatted, diff --git a/be-rust/src/commonMain/resources/lang/temper/be/rust/std/io/support.rs b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/io/support.rs index 71829476..2cccb14c 100644 --- a/be-rust/src/commonMain/resources/lang/temper/be/rust/std/io/support.rs +++ b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/io/support.rs @@ -16,6 +16,20 @@ pub fn std_sleep(ms: i32) -> Promise<()> { promise } +pub fn std_term_cols() -> i32 { + match crossterm::terminal::size() { + Ok((cols, _)) => cols as i32, + Err(_) => 80, + } +} + +pub fn std_term_rows() -> i32 { + match crossterm::terminal::size() { + Ok((_, rows)) => rows as i32, + Err(_) => 24, + } +} + pub fn std_read_line() -> Promise>> { let pb = PromiseBuilder::new(); let promise = pb.promise(); diff --git a/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs new file mode 100644 index 00000000..ceabf231 --- /dev/null +++ b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs @@ -0,0 +1,245 @@ +use super::*; +use std::sync::{Arc, Mutex}; +use std::collections::VecDeque; +use temper_core::{Promise, PromiseBuilder, SafeGenerator}; + +#[cfg(not(feature = "ws"))] +pub(crate) fn ws_listen(_port: i32) -> Promise { panic!() } +#[cfg(not(feature = "ws"))] +pub(crate) fn ws_accept(_server: &WsServer) -> Promise { panic!() } +#[cfg(not(feature = "ws"))] +pub(crate) fn ws_connect(_url: impl temper_core::ToArcString) -> Promise { panic!() } +#[cfg(not(feature = "ws"))] +pub(crate) fn ws_send(_conn: &WsConnection, _msg: impl temper_core::ToArcString) -> Promise<()> { panic!() } +#[cfg(not(feature = "ws"))] +pub(crate) fn ws_recv(_conn: &WsConnection) -> Promise>> { panic!() } +#[cfg(not(feature = "ws"))] +pub(crate) fn ws_close(_conn: &WsConnection) -> Promise<()> { panic!() } + +#[cfg(feature = "ws")] +use std::net::{TcpListener, TcpStream}; +#[cfg(feature = "ws")] +use tungstenite::{accept as ws_accept_tcp, connect as ws_connect_url, Message, WebSocket}; + +#[cfg(feature = "ws")] +struct WsServerInner { + listener: Mutex, +} + +#[cfg(feature = "ws")] +#[derive(Clone)] +struct SimpleWsServer(Arc); + +#[cfg(feature = "ws")] +impl WsServerTrait for SimpleWsServer { + fn clone_boxed(&self) -> WsServer { + WsServer::new(self.clone()) + } +} + +#[cfg(feature = "ws")] +temper_core::impl_any_value_trait!(SimpleWsServer, [WsServer]); + +#[cfg(feature = "ws")] +struct WsConnectionInner { + socket: Mutex>, +} + +#[cfg(feature = "ws")] +#[derive(Clone)] +struct SimpleWsConnection(Arc); + +#[cfg(feature = "ws")] +impl WsConnectionTrait for SimpleWsConnection { + fn clone_boxed(&self) -> WsConnection { + WsConnection::new(self.clone()) + } +} + +#[cfg(feature = "ws")] +temper_core::impl_any_value_trait!(SimpleWsConnection, [WsConnection]); + +#[cfg(feature = "ws")] +pub(crate) fn ws_listen(port: i32) -> Promise { + let pb = PromiseBuilder::new(); + let promise = pb.promise(); + crate::run_async(Arc::new(move || { + let pb = pb.clone(); + SafeGenerator::from_fn(Arc::new(move |_generator: SafeGenerator<()>| { + let addr = format!("0.0.0.0:{}", port); + match TcpListener::bind(&addr) { + Ok(listener) => { + pb.complete(WsServer::new(SimpleWsServer(Arc::new(WsServerInner { + listener: Mutex::new(listener), + })))); + } + Err(_) => { + pb.break_promise(); + } + } + None + })) + })); + promise +} + +#[cfg(feature = "ws")] +pub(crate) fn ws_accept(server: &WsServer) -> Promise { + let server_clone = server.clone_boxed(); + let pb = PromiseBuilder::new(); + let promise = pb.promise(); + crate::run_async(Arc::new(move || { + let pb = pb.clone(); + let server = server_clone.clone_boxed(); + SafeGenerator::from_fn(Arc::new(move |_generator: SafeGenerator<()>| { + // Downcast to get the listener + let any = server.as_any_value(); + let inner = any.as_ref().as_any().downcast_ref::(); + match inner { + Some(s) => { + let listener = s.0.listener.lock().unwrap(); + match listener.accept() { + Ok((stream, _addr)) => { + drop(listener); + match ws_accept_tcp(stream) { + Ok(ws) => { + pb.complete(WsConnection::new(SimpleWsConnection(Arc::new( + WsConnectionInner { + socket: Mutex::new(ws), + }, + )))); + } + Err(_) => pb.break_promise(), + } + } + Err(_) => pb.break_promise(), + } + } + None => pb.break_promise(), + } + None + })) + })); + promise +} + +#[cfg(feature = "ws")] +pub(crate) fn ws_connect(url: impl temper_core::ToArcString) -> Promise { + let url = url.to_arc_string(); + let pb = PromiseBuilder::new(); + let promise = pb.promise(); + crate::run_async(Arc::new(move || { + let pb = pb.clone(); + let url = url.clone(); + SafeGenerator::from_fn(Arc::new(move |_generator: SafeGenerator<()>| { + match ws_connect_url(url.as_str()) { + Ok((ws, _response)) => { + pb.complete(WsConnection::new(SimpleWsConnection(Arc::new( + WsConnectionInner { + socket: Mutex::new(ws), + }, + )))); + } + Err(_) => pb.break_promise(), + } + None + })) + })); + promise +} + +#[cfg(feature = "ws")] +pub(crate) fn ws_send(conn: &WsConnection, msg: impl temper_core::ToArcString) -> Promise<()> { + let conn_clone = conn.clone_boxed(); + let msg = msg.to_arc_string(); + let pb = PromiseBuilder::new(); + let promise = pb.promise(); + crate::run_async(Arc::new(move || { + let pb = pb.clone(); + let conn = conn_clone.clone_boxed(); + let msg = msg.clone(); + SafeGenerator::from_fn(Arc::new(move |_generator: SafeGenerator<()>| { + let any = conn.as_any_value(); + let inner = any.as_ref().as_any().downcast_ref::(); + match inner { + Some(c) => { + let mut socket = c.0.socket.lock().unwrap(); + match socket.send(Message::Text(msg.to_string())) { + Ok(_) => pb.complete(()), + Err(_) => pb.break_promise(), + } + } + None => pb.break_promise(), + } + None + })) + })); + promise +} + +#[cfg(feature = "ws")] +pub(crate) fn ws_recv(conn: &WsConnection) -> Promise>> { + let conn_clone = conn.clone_boxed(); + let pb = PromiseBuilder::new(); + let promise = pb.promise(); + crate::run_async(Arc::new(move || { + let pb = pb.clone(); + let conn = conn_clone.clone_boxed(); + SafeGenerator::from_fn(Arc::new(move |_generator: SafeGenerator<()>| { + let any = conn.as_any_value(); + let inner = any.as_ref().as_any().downcast_ref::(); + match inner { + Some(c) => { + let mut socket = c.0.socket.lock().unwrap(); + match socket.read() { + Ok(Message::Text(text)) => { + pb.complete(Some(Arc::new(text.to_string()))); + } + Ok(Message::Close(_)) | Err(_) => { + pb.complete(None); + } + Ok(_) => { + // Binary or other message types — skip and return null + pb.complete(None); + } + } + } + None => pb.complete(None), + } + None + })) + })); + promise +} + +#[cfg(feature = "ws")] +pub(crate) fn ws_close(conn: &WsConnection) -> Promise<()> { + let conn_clone = conn.clone_boxed(); + let pb = PromiseBuilder::new(); + let promise = pb.promise(); + crate::run_async(Arc::new(move || { + let pb = pb.clone(); + let conn = conn_clone.clone_boxed(); + SafeGenerator::from_fn(Arc::new(move |_generator: SafeGenerator<()>| { + let any = conn.as_any_value(); + let inner = any.as_ref().as_any().downcast_ref::(); + match inner { + Some(c) => { + let mut socket = c.0.socket.lock().unwrap(); + let _ = socket.close(None); + // Drain remaining messages until close confirmation + loop { + match socket.read() { + Ok(Message::Close(_)) | Err(_) => break, + _ => continue, + } + } + } + None => {} + } + pb.complete(()); + None + })) + })); + promise +} diff --git a/be/src/commonMain/kotlin/lang/temper/be/README.md b/be/src/commonMain/kotlin/lang/temper/be/README.md index 01d1209a..b8654682 100644 --- a/be/src/commonMain/kotlin/lang/temper/be/README.md +++ b/be/src/commonMain/kotlin/lang/temper/be/README.md @@ -226,7 +226,15 @@ to implement: - `Test::failedOnAssert` - `Test::messages` - `Test::passing` +- `WsConnection` +- `WsServer` - `stdNetSend` - `stdReadLine` - `stdSleep` +- `wsAccept` +- `wsClose` +- `wsConnect` +- `wsListen` +- `wsRecv` +- `wsSend` diff --git a/frontend/src/commonMain/resources/std/config.temper.md b/frontend/src/commonMain/resources/std/config.temper.md index 976926cd..6a47a0fd 100644 --- a/frontend/src/commonMain/resources/std/config.temper.md +++ b/frontend/src/commonMain/resources/std/config.temper.md @@ -24,6 +24,7 @@ We might break these out into separate libraries in the future. import("./json"); import("./net"); import("./io"); + import("./ws"); ## C# diff --git a/frontend/src/commonMain/resources/std/io/io.temper.md b/frontend/src/commonMain/resources/std/io/io.temper.md index d6f73497..79fc1198 100644 --- a/frontend/src/commonMain/resources/std/io/io.temper.md +++ b/frontend/src/commonMain/resources/std/io/io.temper.md @@ -20,3 +20,18 @@ Read one line from standard input. Returns null on EOF. panic() } +## Terminal Size + +Get the current terminal dimensions in characters. +Returns a default of 80x24 if the terminal size cannot be determined. + + @connected("stdTermCols") + export let terminalColumns(): Int { + panic() + } + + @connected("stdTermRows") + export let terminalRows(): Int { + panic() + } + diff --git a/frontend/src/commonMain/resources/std/ws/ws.temper.md b/frontend/src/commonMain/resources/std/ws/ws.temper.md new file mode 100644 index 00000000..d6cbb8ee --- /dev/null +++ b/frontend/src/commonMain/resources/std/ws/ws.temper.md @@ -0,0 +1,64 @@ +# WebSocket Support + +WebSocket support for real-time bidirectional communication. + +*WsServer* is an opaque handle representing a listening WebSocket server. + + @connected("WsServer") + export interface WsServer {} + +*WsConnection* is an opaque handle representing a single WebSocket connection. + + @connected("WsConnection") + export interface WsConnection {} + +## Server Functions + +*wsListen* starts a WebSocket server on the given port and resolves when +it is ready to accept connections. + + @connected("wsListen") + export let wsListen(port: Int): Promise { + panic() + } + +*wsAccept* waits for and accepts the next incoming connection on a server. + + @connected("wsAccept") + export let wsAccept(server: WsServer): Promise { + panic() + } + +## Client Functions + +*wsConnect* opens a WebSocket connection to the given URL +(e.g. `"ws://localhost:8080"`). + + @connected("wsConnect") + export let wsConnect(url: String): Promise { + panic() + } + +## Shared Functions + +*wsSend* sends a text message over a connection. + + @connected("wsSend") + export let wsSend(conn: WsConnection, msg: String): Promise { + panic() + } + +*wsRecv* waits for the next message from a connection. +Returns `null` if the connection is closed. + + @connected("wsRecv") + export let wsRecv(conn: WsConnection): Promise { + panic() + } + +*wsClose* closes a connection. + + @connected("wsClose") + export let wsClose(conn: WsConnection): Promise { + panic() + } From e76ab32193c8f9461a5f0bb0b36f34a0840127a9 Mon Sep 17 00:00:00 2001 From: Robert Grayson Date: Sun, 15 Mar 2026 13:21:50 -0400 Subject: [PATCH 04/12] fix: add ws npm dependency to temper-core package.json The ws.js support file dynamically imports the 'ws' package, but it wasn't listed as a dependency. After temper build regenerated the output, npm install wouldn't install ws, causing wsListen/wsConnect to fail silently at runtime. Co-Authored-By: Claude Opus 4.6 (1M context) --- .../resources/lang/temper/be/js/temper-core/package.json | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/package.json b/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/package.json index c3c861f2..5357b753 100644 --- a/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/package.json +++ b/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/package.json @@ -5,5 +5,8 @@ "main": "index.js", "type": "module", "author": "https://github.com/mikesamuel", - "license": "Apache-2.0 OR MIT" + "license": "Apache-2.0 OR MIT", + "dependencies": { + "ws": "^8.18.0" + } } From 73b848d35454db9ca6d5410b5fd52faba9ee648e Mon Sep 17 00:00:00 2001 From: Robert Grayson Date: Sun, 15 Mar 2026 16:38:45 -0400 Subject: [PATCH 05/12] fix: Rust ws support - rename functions, fix visibility and stream types - Rename support functions to std_ws_* to avoid colliding with Temper-generated panic stubs (same pattern as std_sleep/std_read_line) - Change visibility from pub(crate) to pub for cross-crate access - Accept &dyn WsServerTrait/WsConnectionTrait instead of &WsServer to work with the Rust backend's interface deref codegen - Use WsStream enum to handle both WebSocket (server-accepted) and WebSocket> (client-connected) - Use temper_core::cast() for downcasting instead of manual as_any chain - Add .into() for tungstenite 0.26 Message::Text(Utf8Bytes) change Co-Authored-By: Claude Opus 4.6 (1M context) --- .../lang/temper/be/rust/RustSupportNetwork.kt | 12 +- .../lang/temper/be/rust/std/ws/support.rs | 172 +++++++++--------- 2 files changed, 90 insertions(+), 94 deletions(-) diff --git a/be-rust/src/commonMain/kotlin/lang/temper/be/rust/RustSupportNetwork.kt b/be-rust/src/commonMain/kotlin/lang/temper/be/rust/RustSupportNetwork.kt index 1f345b87..64d85abb 100644 --- a/be-rust/src/commonMain/kotlin/lang/temper/be/rust/RustSupportNetwork.kt +++ b/be-rust/src/commonMain/kotlin/lang/temper/be/rust/RustSupportNetwork.kt @@ -919,12 +919,12 @@ private val stdReadLine = FunctionCall("stdReadLine", "temper_std::io::std_read_ private val stdNextKeypress = FunctionCall("stdNextKeypress", "temper_std::keyboard::std_next_keypress") private val stdTermCols = FunctionCall("stdTermCols", "temper_std::io::std_term_cols") private val stdTermRows = FunctionCall("stdTermRows", "temper_std::io::std_term_rows") -private val wsListen = FunctionCall("wsListen", "temper_std::ws::ws_listen") -private val wsAccept = FunctionCall("wsAccept", "temper_std::ws::ws_accept", cloneEvenIfFirst = true) -private val wsConnect = FunctionCall("wsConnect", "temper_std::ws::ws_connect", cloneEvenIfFirst = true) -private val wsSend = FunctionCall("wsSend", "temper_std::ws::ws_send", cloneEvenIfFirst = true) -private val wsRecv = FunctionCall("wsRecv", "temper_std::ws::ws_recv", cloneEvenIfFirst = true) -private val wsClose = FunctionCall("wsClose", "temper_std::ws::ws_close", cloneEvenIfFirst = true) +private val wsListen = FunctionCall("wsListen", "temper_std::ws::std_ws_listen") +private val wsAccept = FunctionCall("wsAccept", "temper_std::ws::std_ws_accept") +private val wsConnect = FunctionCall("wsConnect", "temper_std::ws::std_ws_connect", cloneEvenIfFirst = true) +private val wsSend = FunctionCall("wsSend", "temper_std::ws::std_ws_send") +private val wsRecv = FunctionCall("wsRecv", "temper_std::ws::std_ws_recv") +private val wsClose = FunctionCall("wsClose", "temper_std::ws::std_ws_close") internal object PairConstructor : RustInlineSupportCode( "Pair::constructor", diff --git a/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs index ceabf231..5a841e10 100644 --- a/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs +++ b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs @@ -1,25 +1,52 @@ use super::*; use std::sync::{Arc, Mutex}; -use std::collections::VecDeque; -use temper_core::{Promise, PromiseBuilder, SafeGenerator}; +use temper_core::{AsAnyValue, Promise, PromiseBuilder, SafeGenerator}; #[cfg(not(feature = "ws"))] -pub(crate) fn ws_listen(_port: i32) -> Promise { panic!() } +pub fn std_ws_listen(_port: i32) -> Promise { panic!() } #[cfg(not(feature = "ws"))] -pub(crate) fn ws_accept(_server: &WsServer) -> Promise { panic!() } +pub fn std_ws_accept(_server: &dyn WsServerTrait) -> Promise { panic!() } #[cfg(not(feature = "ws"))] -pub(crate) fn ws_connect(_url: impl temper_core::ToArcString) -> Promise { panic!() } +pub fn std_ws_connect(_url: impl temper_core::ToArcString) -> Promise { panic!() } #[cfg(not(feature = "ws"))] -pub(crate) fn ws_send(_conn: &WsConnection, _msg: impl temper_core::ToArcString) -> Promise<()> { panic!() } +pub fn std_ws_send(_conn: &dyn WsConnectionTrait, _msg: impl temper_core::ToArcString) -> Promise<()> { panic!() } #[cfg(not(feature = "ws"))] -pub(crate) fn ws_recv(_conn: &WsConnection) -> Promise>> { panic!() } +pub fn std_ws_recv(_conn: &dyn WsConnectionTrait) -> Promise>> { panic!() } #[cfg(not(feature = "ws"))] -pub(crate) fn ws_close(_conn: &WsConnection) -> Promise<()> { panic!() } +pub fn std_ws_close(_conn: &dyn WsConnectionTrait) -> Promise<()> { panic!() } #[cfg(feature = "ws")] use std::net::{TcpListener, TcpStream}; #[cfg(feature = "ws")] -use tungstenite::{accept as ws_accept_tcp, connect as ws_connect_url, Message, WebSocket}; +use tungstenite::{accept as ws_upgrade, connect as ws_connect_url, stream::MaybeTlsStream, Message, WebSocket}; + +#[cfg(feature = "ws")] +enum WsStream { + Plain(WebSocket), + MaybeTls(WebSocket>), +} + +#[cfg(feature = "ws")] +impl WsStream { + fn send(&mut self, msg: Message) -> Result<(), tungstenite::Error> { + match self { + WsStream::Plain(ws) => ws.send(msg), + WsStream::MaybeTls(ws) => ws.send(msg), + } + } + fn read(&mut self) -> Result { + match self { + WsStream::Plain(ws) => ws.read(), + WsStream::MaybeTls(ws) => ws.read(), + } + } + fn close(&mut self, frame: Option) -> Result<(), tungstenite::Error> { + match self { + WsStream::Plain(ws) => ws.close(frame), + WsStream::MaybeTls(ws) => ws.close(frame), + } + } +} #[cfg(feature = "ws")] struct WsServerInner { @@ -42,7 +69,7 @@ temper_core::impl_any_value_trait!(SimpleWsServer, [WsServer]); #[cfg(feature = "ws")] struct WsConnectionInner { - socket: Mutex>, + socket: Mutex, } #[cfg(feature = "ws")] @@ -60,7 +87,7 @@ impl WsConnectionTrait for SimpleWsConnection { temper_core::impl_any_value_trait!(SimpleWsConnection, [WsConnection]); #[cfg(feature = "ws")] -pub(crate) fn ws_listen(port: i32) -> Promise { +pub fn std_ws_listen(port: i32) -> Promise { let pb = PromiseBuilder::new(); let promise = pb.promise(); crate::run_async(Arc::new(move || { @@ -84,38 +111,30 @@ pub(crate) fn ws_listen(port: i32) -> Promise { } #[cfg(feature = "ws")] -pub(crate) fn ws_accept(server: &WsServer) -> Promise { - let server_clone = server.clone_boxed(); +pub fn std_ws_accept(server: &dyn WsServerTrait) -> Promise { + let server: SimpleWsServer = temper_core::cast(server.as_any_value()).expect("WsServer downcast"); let pb = PromiseBuilder::new(); let promise = pb.promise(); crate::run_async(Arc::new(move || { let pb = pb.clone(); - let server = server_clone.clone_boxed(); + let server = server.clone(); SafeGenerator::from_fn(Arc::new(move |_generator: SafeGenerator<()>| { - // Downcast to get the listener - let any = server.as_any_value(); - let inner = any.as_ref().as_any().downcast_ref::(); - match inner { - Some(s) => { - let listener = s.0.listener.lock().unwrap(); - match listener.accept() { - Ok((stream, _addr)) => { - drop(listener); - match ws_accept_tcp(stream) { - Ok(ws) => { - pb.complete(WsConnection::new(SimpleWsConnection(Arc::new( - WsConnectionInner { - socket: Mutex::new(ws), - }, - )))); - } - Err(_) => pb.break_promise(), - } + let listener = server.0.listener.lock().unwrap(); + match listener.accept() { + Ok((stream, _addr)) => { + drop(listener); + match ws_upgrade(stream) { + Ok(ws) => { + pb.complete(WsConnection::new(SimpleWsConnection(Arc::new( + WsConnectionInner { + socket: Mutex::new(WsStream::Plain(ws)), + }, + )))); } Err(_) => pb.break_promise(), } } - None => pb.break_promise(), + Err(_) => pb.break_promise(), } None })) @@ -124,7 +143,7 @@ pub(crate) fn ws_accept(server: &WsServer) -> Promise { } #[cfg(feature = "ws")] -pub(crate) fn ws_connect(url: impl temper_core::ToArcString) -> Promise { +pub fn std_ws_connect(url: impl temper_core::ToArcString) -> Promise { let url = url.to_arc_string(); let pb = PromiseBuilder::new(); let promise = pb.promise(); @@ -136,7 +155,7 @@ pub(crate) fn ws_connect(url: impl temper_core::ToArcString) -> Promise { pb.complete(WsConnection::new(SimpleWsConnection(Arc::new( WsConnectionInner { - socket: Mutex::new(ws), + socket: Mutex::new(WsStream::MaybeTls(ws)), }, )))); } @@ -149,27 +168,20 @@ pub(crate) fn ws_connect(url: impl temper_core::ToArcString) -> Promise Promise<()> { - let conn_clone = conn.clone_boxed(); +pub fn std_ws_send(conn: &dyn WsConnectionTrait, msg: impl temper_core::ToArcString) -> Promise<()> { + let conn: SimpleWsConnection = temper_core::cast(conn.as_any_value()).expect("WsConnection downcast"); let msg = msg.to_arc_string(); let pb = PromiseBuilder::new(); let promise = pb.promise(); crate::run_async(Arc::new(move || { let pb = pb.clone(); - let conn = conn_clone.clone_boxed(); + let conn = conn.clone(); let msg = msg.clone(); SafeGenerator::from_fn(Arc::new(move |_generator: SafeGenerator<()>| { - let any = conn.as_any_value(); - let inner = any.as_ref().as_any().downcast_ref::(); - match inner { - Some(c) => { - let mut socket = c.0.socket.lock().unwrap(); - match socket.send(Message::Text(msg.to_string())) { - Ok(_) => pb.complete(()), - Err(_) => pb.break_promise(), - } - } - None => pb.break_promise(), + let mut socket = conn.0.socket.lock().unwrap(); + match socket.send(Message::Text(msg.to_string().into())) { + Ok(_) => pb.complete(()), + Err(_) => pb.break_promise(), } None })) @@ -178,33 +190,25 @@ pub(crate) fn ws_send(conn: &WsConnection, msg: impl temper_core::ToArcString) - } #[cfg(feature = "ws")] -pub(crate) fn ws_recv(conn: &WsConnection) -> Promise>> { - let conn_clone = conn.clone_boxed(); +pub fn std_ws_recv(conn: &dyn WsConnectionTrait) -> Promise>> { + let conn: SimpleWsConnection = temper_core::cast(conn.as_any_value()).expect("WsConnection downcast"); let pb = PromiseBuilder::new(); let promise = pb.promise(); crate::run_async(Arc::new(move || { let pb = pb.clone(); - let conn = conn_clone.clone_boxed(); + let conn = conn.clone(); SafeGenerator::from_fn(Arc::new(move |_generator: SafeGenerator<()>| { - let any = conn.as_any_value(); - let inner = any.as_ref().as_any().downcast_ref::(); - match inner { - Some(c) => { - let mut socket = c.0.socket.lock().unwrap(); - match socket.read() { - Ok(Message::Text(text)) => { - pb.complete(Some(Arc::new(text.to_string()))); - } - Ok(Message::Close(_)) | Err(_) => { - pb.complete(None); - } - Ok(_) => { - // Binary or other message types — skip and return null - pb.complete(None); - } - } + let mut socket = conn.0.socket.lock().unwrap(); + match socket.read() { + Ok(Message::Text(text)) => { + pb.complete(Some(Arc::new(text.to_string()))); + } + Ok(Message::Close(_)) | Err(_) => { + pb.complete(None); + } + Ok(_) => { + pb.complete(None); } - None => pb.complete(None), } None })) @@ -213,29 +217,21 @@ pub(crate) fn ws_recv(conn: &WsConnection) -> Promise>> { } #[cfg(feature = "ws")] -pub(crate) fn ws_close(conn: &WsConnection) -> Promise<()> { - let conn_clone = conn.clone_boxed(); +pub fn std_ws_close(conn: &dyn WsConnectionTrait) -> Promise<()> { + let conn: SimpleWsConnection = temper_core::cast(conn.as_any_value()).expect("WsConnection downcast"); let pb = PromiseBuilder::new(); let promise = pb.promise(); crate::run_async(Arc::new(move || { let pb = pb.clone(); - let conn = conn_clone.clone_boxed(); + let conn = conn.clone(); SafeGenerator::from_fn(Arc::new(move |_generator: SafeGenerator<()>| { - let any = conn.as_any_value(); - let inner = any.as_ref().as_any().downcast_ref::(); - match inner { - Some(c) => { - let mut socket = c.0.socket.lock().unwrap(); - let _ = socket.close(None); - // Drain remaining messages until close confirmation - loop { - match socket.read() { - Ok(Message::Close(_)) | Err(_) => break, - _ => continue, - } - } + let mut socket = conn.0.socket.lock().unwrap(); + let _ = socket.close(None); + loop { + match socket.read() { + Ok(Message::Close(_)) | Err(_) => break, + _ => continue, } - None => {} } pb.complete(()); None From 3a512b193b2e2e83d70347a27d1558ea80588ca3 Mon Sep 17 00:00:00 2001 From: Robert Grayson Date: Sun, 15 Mar 2026 16:51:36 -0400 Subject: [PATCH 06/12] fix: use dedicated threads for Rust ws operations Temper's Rust runtime uses a single-threaded task runner. Blocking ws operations (accept, send, recv) would block the runner and prevent other async blocks (like readLine) from processing. Spawn dedicated threads for ws_accept, ws_send, and ws_recv instead of going through crate::run_async. This lets the WS operations block independently while other async work continues. Also adds error logging to ws_recv for debugging. Co-Authored-By: Claude Opus 4.6 (1M context) --- .../lang/temper/be/rust/std/ws/support.rs | 97 +++++++++---------- 1 file changed, 45 insertions(+), 52 deletions(-) diff --git a/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs index 5a841e10..3b30be29 100644 --- a/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs +++ b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs @@ -115,30 +115,25 @@ pub fn std_ws_accept(server: &dyn WsServerTrait) -> Promise { let server: SimpleWsServer = temper_core::cast(server.as_any_value()).expect("WsServer downcast"); let pb = PromiseBuilder::new(); let promise = pb.promise(); - crate::run_async(Arc::new(move || { - let pb = pb.clone(); - let server = server.clone(); - SafeGenerator::from_fn(Arc::new(move |_generator: SafeGenerator<()>| { - let listener = server.0.listener.lock().unwrap(); - match listener.accept() { - Ok((stream, _addr)) => { - drop(listener); - match ws_upgrade(stream) { - Ok(ws) => { - pb.complete(WsConnection::new(SimpleWsConnection(Arc::new( - WsConnectionInner { - socket: Mutex::new(WsStream::Plain(ws)), - }, - )))); - } - Err(_) => pb.break_promise(), + std::thread::spawn(move || { + let listener = server.0.listener.lock().unwrap(); + match listener.accept() { + Ok((stream, _addr)) => { + drop(listener); + match ws_upgrade(stream) { + Ok(ws) => { + pb.complete(WsConnection::new(SimpleWsConnection(Arc::new( + WsConnectionInner { + socket: Mutex::new(WsStream::Plain(ws)), + }, + )))); } + Err(_) => pb.break_promise(), } - Err(_) => pb.break_promise(), } - None - })) - })); + Err(_) => pb.break_promise(), + } + }); promise } @@ -173,19 +168,13 @@ pub fn std_ws_send(conn: &dyn WsConnectionTrait, msg: impl temper_core::ToArcStr let msg = msg.to_arc_string(); let pb = PromiseBuilder::new(); let promise = pb.promise(); - crate::run_async(Arc::new(move || { - let pb = pb.clone(); - let conn = conn.clone(); - let msg = msg.clone(); - SafeGenerator::from_fn(Arc::new(move |_generator: SafeGenerator<()>| { - let mut socket = conn.0.socket.lock().unwrap(); - match socket.send(Message::Text(msg.to_string().into())) { - Ok(_) => pb.complete(()), - Err(_) => pb.break_promise(), - } - None - })) - })); + std::thread::spawn(move || { + let mut socket = conn.0.socket.lock().unwrap(); + match socket.send(Message::Text(msg.to_string().into())) { + Ok(_) => pb.complete(()), + Err(_) => pb.break_promise(), + } + }); promise } @@ -194,25 +183,29 @@ pub fn std_ws_recv(conn: &dyn WsConnectionTrait) -> Promise>> let conn: SimpleWsConnection = temper_core::cast(conn.as_any_value()).expect("WsConnection downcast"); let pb = PromiseBuilder::new(); let promise = pb.promise(); - crate::run_async(Arc::new(move || { - let pb = pb.clone(); - let conn = conn.clone(); - SafeGenerator::from_fn(Arc::new(move |_generator: SafeGenerator<()>| { - let mut socket = conn.0.socket.lock().unwrap(); - match socket.read() { - Ok(Message::Text(text)) => { - pb.complete(Some(Arc::new(text.to_string()))); - } - Ok(Message::Close(_)) | Err(_) => { - pb.complete(None); - } - Ok(_) => { - pb.complete(None); - } + // Spawn a dedicated thread instead of using run_async to avoid + // blocking the single-threaded task runner (which would prevent + // other async blocks like readLine from processing). + std::thread::spawn(move || { + let mut socket = conn.0.socket.lock().unwrap(); + match socket.read() { + Ok(Message::Text(text)) => { + pb.complete(Some(Arc::new(text.to_string()))); } - None - })) - })); + Ok(Message::Close(frame)) => { + eprintln!("[ws_recv] close frame: {:?}", frame); + pb.complete(None); + } + Ok(other) => { + eprintln!("[ws_recv] unexpected message type: {:?}", other); + pb.complete(None); + } + Err(e) => { + eprintln!("[ws_recv] error: {:?}", e); + pb.complete(None); + } + } + }); promise } From 9a9b3504b98087137411eed87d714907c4aa7145 Mon Sep 17 00:00:00 2001 From: Robert Grayson Date: Sun, 15 Mar 2026 17:01:54 -0400 Subject: [PATCH 07/12] fix: set read timeout on WS connections for Rust backend Both server-accepted and client connections now have a 50ms read timeout. This prevents the recv loop from holding the socket Mutex indefinitely, allowing send operations to interleave. Without this, the Rust server couldn't send frames because recv blocked the lock. Co-Authored-By: Claude Opus 4.6 (1M context) --- .../lang/temper/be/rust/std/ws/support.rs | 64 ++++++++++++++----- 1 file changed, 47 insertions(+), 17 deletions(-) diff --git a/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs index 3b30be29..d510d7e2 100644 --- a/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs +++ b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs @@ -120,6 +120,9 @@ pub fn std_ws_accept(server: &dyn WsServerTrait) -> Promise { match listener.accept() { Ok((stream, _addr)) => { drop(listener); + // Set a read timeout so recv doesn't hold the socket + // Mutex indefinitely, allowing send to interleave. + let _ = stream.set_read_timeout(Some(std::time::Duration::from_millis(50))); match ws_upgrade(stream) { Ok(ws) => { pb.complete(WsConnection::new(SimpleWsConnection(Arc::new( @@ -148,6 +151,13 @@ pub fn std_ws_connect(url: impl temper_core::ToArcString) -> Promise| { match ws_connect_url(url.as_str()) { Ok((ws, _response)) => { + // Set read timeout on client connections too + match ws.get_ref() { + MaybeTlsStream::Plain(s) => { + let _ = s.set_read_timeout(Some(std::time::Duration::from_millis(50))); + } + _ => {} + } pb.complete(WsConnection::new(SimpleWsConnection(Arc::new( WsConnectionInner { socket: Mutex::new(WsStream::MaybeTls(ws)), @@ -172,7 +182,10 @@ pub fn std_ws_send(conn: &dyn WsConnectionTrait, msg: impl temper_core::ToArcStr let mut socket = conn.0.socket.lock().unwrap(); match socket.send(Message::Text(msg.to_string().into())) { Ok(_) => pb.complete(()), - Err(_) => pb.break_promise(), + Err(e) => { + eprintln!("[ws_send] error: {:?}", e); + pb.break_promise(); + } } }); promise @@ -187,22 +200,39 @@ pub fn std_ws_recv(conn: &dyn WsConnectionTrait) -> Promise>> // blocking the single-threaded task runner (which would prevent // other async blocks like readLine from processing). std::thread::spawn(move || { - let mut socket = conn.0.socket.lock().unwrap(); - match socket.read() { - Ok(Message::Text(text)) => { - pb.complete(Some(Arc::new(text.to_string()))); - } - Ok(Message::Close(frame)) => { - eprintln!("[ws_recv] close frame: {:?}", frame); - pb.complete(None); - } - Ok(other) => { - eprintln!("[ws_recv] unexpected message type: {:?}", other); - pb.complete(None); - } - Err(e) => { - eprintln!("[ws_recv] error: {:?}", e); - pb.complete(None); + // Loop retrying on timeout errors so the Mutex is periodically + // released, allowing wsSend to interleave on the same connection. + loop { + let result = { + let mut socket = conn.0.socket.lock().unwrap(); + socket.read() + }; // Mutex released here before processing + match result { + Ok(Message::Text(text)) => { + pb.complete(Some(Arc::new(text.to_string()))); + return; + } + Ok(Message::Close(_)) => { + pb.complete(None); + return; + } + Ok(_) => { + // Skip non-text messages, keep reading + continue; + } + Err(tungstenite::Error::Io(ref e)) + if e.kind() == std::io::ErrorKind::WouldBlock + || e.kind() == std::io::ErrorKind::TimedOut => + { + // Read timeout — release lock briefly and retry + std::thread::sleep(std::time::Duration::from_millis(10)); + continue; + } + Err(e) => { + eprintln!("[ws_recv] error: {:?}", e); + pb.complete(None); + return; + } } } }); From d9f30207d372e139ee95473f844ac5f30eeb50b8 Mon Sep 17 00:00:00 2001 From: Robert Grayson Date: Sun, 15 Mar 2026 17:06:15 -0400 Subject: [PATCH 08/12] fix: channel-based I/O for Rust WebSocket connections Replace Mutex-shared WebSocket with a dedicated I/O thread per connection that communicates via mpsc channels. The I/O thread owns the WebSocket exclusively and polls for both send and recv, avoiding tungstenite's internal state corruption from concurrent access. Send is now non-blocking (channel push), recv blocks on the channel receiver in a spawned thread. This fixes the ResetWithoutClosingHandshake error that occurred when the server's recv loop and send loop competed for the socket. Co-Authored-By: Claude Opus 4.6 (1M context) --- .../lang/temper/be/rust/std/ws/support.rs | 224 ++++++++---------- 1 file changed, 100 insertions(+), 124 deletions(-) diff --git a/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs index d510d7e2..a6bd7022 100644 --- a/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs +++ b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs @@ -1,5 +1,5 @@ use super::*; -use std::sync::{Arc, Mutex}; +use std::sync::{Arc, Mutex, mpsc}; use temper_core::{AsAnyValue, Promise, PromiseBuilder, SafeGenerator}; #[cfg(not(feature = "ws"))] @@ -20,33 +20,7 @@ use std::net::{TcpListener, TcpStream}; #[cfg(feature = "ws")] use tungstenite::{accept as ws_upgrade, connect as ws_connect_url, stream::MaybeTlsStream, Message, WebSocket}; -#[cfg(feature = "ws")] -enum WsStream { - Plain(WebSocket), - MaybeTls(WebSocket>), -} - -#[cfg(feature = "ws")] -impl WsStream { - fn send(&mut self, msg: Message) -> Result<(), tungstenite::Error> { - match self { - WsStream::Plain(ws) => ws.send(msg), - WsStream::MaybeTls(ws) => ws.send(msg), - } - } - fn read(&mut self) -> Result { - match self { - WsStream::Plain(ws) => ws.read(), - WsStream::MaybeTls(ws) => ws.read(), - } - } - fn close(&mut self, frame: Option) -> Result<(), tungstenite::Error> { - match self { - WsStream::Plain(ws) => ws.close(frame), - WsStream::MaybeTls(ws) => ws.close(frame), - } - } -} +// --- Server --- #[cfg(feature = "ws")] struct WsServerInner { @@ -59,17 +33,20 @@ struct SimpleWsServer(Arc); #[cfg(feature = "ws")] impl WsServerTrait for SimpleWsServer { - fn clone_boxed(&self) -> WsServer { - WsServer::new(self.clone()) - } + fn clone_boxed(&self) -> WsServer { WsServer::new(self.clone()) } } #[cfg(feature = "ws")] temper_core::impl_any_value_trait!(SimpleWsServer, [WsServer]); +// --- Connection --- +// Each connection has a dedicated I/O thread that owns the WebSocket. +// Send and recv go through channels, avoiding Mutex sharing issues. + #[cfg(feature = "ws")] struct WsConnectionInner { - socket: Mutex, + send_tx: mpsc::Sender, + recv_rx: Mutex>>, } #[cfg(feature = "ws")] @@ -78,31 +55,90 @@ struct SimpleWsConnection(Arc); #[cfg(feature = "ws")] impl WsConnectionTrait for SimpleWsConnection { - fn clone_boxed(&self) -> WsConnection { - WsConnection::new(self.clone()) - } + fn clone_boxed(&self) -> WsConnection { WsConnection::new(self.clone()) } } #[cfg(feature = "ws")] temper_core::impl_any_value_trait!(SimpleWsConnection, [WsConnection]); +/// Spawn an I/O thread that owns the WebSocket and bridges to channels. +#[cfg(feature = "ws")] +fn spawn_ws_io_thread(mut ws: WebSocket) -> SimpleWsConnection +where + S: std::io::Read + std::io::Write + Send + 'static, +{ + let (send_tx, send_rx) = mpsc::channel::(); + let (recv_tx, recv_rx) = mpsc::sync_channel::>(16); + + std::thread::spawn(move || { + loop { + // Try to send any queued outgoing messages. + while let Ok(msg) = send_rx.try_recv() { + if ws.send(Message::Text(msg.into())).is_err() { + let _ = recv_tx.send(None); + return; + } + } + + // Try to read one incoming message (may timeout quickly). + match ws.read() { + Ok(Message::Text(text)) => { + if recv_tx.send(Some(text.to_string())).is_err() { + return; // recv side dropped + } + } + Ok(Message::Close(_)) => { + let _ = recv_tx.send(None); + return; + } + Ok(Message::Ping(data)) => { + let _ = ws.send(Message::Pong(data)); + } + Ok(_) => {} // skip binary, pong, etc + Err(tungstenite::Error::Io(ref e)) + if e.kind() == std::io::ErrorKind::WouldBlock + || e.kind() == std::io::ErrorKind::TimedOut => + { + // No data yet — brief yield then loop + std::thread::sleep(std::time::Duration::from_millis(5)); + } + Err(_) => { + let _ = recv_tx.send(None); + return; + } + } + } + }); + + SimpleWsConnection(Arc::new(WsConnectionInner { + send_tx, + recv_rx: Mutex::new(recv_rx), + })) +} + +/// Set a short read timeout on a TcpStream before passing it to tungstenite. +/// This makes the I/O thread's read() return quickly, allowing send interleaving. +#[cfg(feature = "ws")] +fn prepare_stream(stream: &TcpStream) { + let _ = stream.set_read_timeout(Some(std::time::Duration::from_millis(20))); +} + +// --- Functions --- + #[cfg(feature = "ws")] pub fn std_ws_listen(port: i32) -> Promise { let pb = PromiseBuilder::new(); let promise = pb.promise(); crate::run_async(Arc::new(move || { let pb = pb.clone(); - SafeGenerator::from_fn(Arc::new(move |_generator: SafeGenerator<()>| { - let addr = format!("0.0.0.0:{}", port); - match TcpListener::bind(&addr) { + SafeGenerator::from_fn(Arc::new(move |_: SafeGenerator<()>| { + match TcpListener::bind(format!("0.0.0.0:{}", port)) { Ok(listener) => { pb.complete(WsServer::new(SimpleWsServer(Arc::new(WsServerInner { listener: Mutex::new(listener), })))); } - Err(_) => { - pb.break_promise(); - } + Err(_) => pb.break_promise(), } None })) @@ -118,18 +154,13 @@ pub fn std_ws_accept(server: &dyn WsServerTrait) -> Promise { std::thread::spawn(move || { let listener = server.0.listener.lock().unwrap(); match listener.accept() { - Ok((stream, _addr)) => { + Ok((stream, _)) => { drop(listener); - // Set a read timeout so recv doesn't hold the socket - // Mutex indefinitely, allowing send to interleave. - let _ = stream.set_read_timeout(Some(std::time::Duration::from_millis(50))); + prepare_stream(&stream); match ws_upgrade(stream) { Ok(ws) => { - pb.complete(WsConnection::new(SimpleWsConnection(Arc::new( - WsConnectionInner { - socket: Mutex::new(WsStream::Plain(ws)), - }, - )))); + let conn = spawn_ws_io_thread(ws); + pb.complete(WsConnection::new(conn)); } Err(_) => pb.break_promise(), } @@ -148,21 +179,16 @@ pub fn std_ws_connect(url: impl temper_core::ToArcString) -> Promise| { + SafeGenerator::from_fn(Arc::new(move |_: SafeGenerator<()>| { match ws_connect_url(url.as_str()) { - Ok((ws, _response)) => { - // Set read timeout on client connections too + Ok((mut ws, _)) => { + // Set read timeout on client connections match ws.get_ref() { - MaybeTlsStream::Plain(s) => { - let _ = s.set_read_timeout(Some(std::time::Duration::from_millis(50))); - } + MaybeTlsStream::Plain(s) => prepare_stream(s), _ => {} } - pb.complete(WsConnection::new(SimpleWsConnection(Arc::new( - WsConnectionInner { - socket: Mutex::new(WsStream::MaybeTls(ws)), - }, - )))); + let conn = spawn_ws_io_thread(ws); + pb.complete(WsConnection::new(conn)); } Err(_) => pb.break_promise(), } @@ -178,16 +204,11 @@ pub fn std_ws_send(conn: &dyn WsConnectionTrait, msg: impl temper_core::ToArcStr let msg = msg.to_arc_string(); let pb = PromiseBuilder::new(); let promise = pb.promise(); - std::thread::spawn(move || { - let mut socket = conn.0.socket.lock().unwrap(); - match socket.send(Message::Text(msg.to_string().into())) { - Ok(_) => pb.complete(()), - Err(e) => { - eprintln!("[ws_send] error: {:?}", e); - pb.break_promise(); - } - } - }); + // Send is non-blocking — just push to the channel + match conn.0.send_tx.send(msg.to_string()) { + Ok(_) => pb.complete(()), + Err(_) => pb.break_promise(), + } promise } @@ -196,44 +217,12 @@ pub fn std_ws_recv(conn: &dyn WsConnectionTrait) -> Promise>> let conn: SimpleWsConnection = temper_core::cast(conn.as_any_value()).expect("WsConnection downcast"); let pb = PromiseBuilder::new(); let promise = pb.promise(); - // Spawn a dedicated thread instead of using run_async to avoid - // blocking the single-threaded task runner (which would prevent - // other async blocks like readLine from processing). + // Recv blocks waiting for the next message from the I/O thread std::thread::spawn(move || { - // Loop retrying on timeout errors so the Mutex is periodically - // released, allowing wsSend to interleave on the same connection. - loop { - let result = { - let mut socket = conn.0.socket.lock().unwrap(); - socket.read() - }; // Mutex released here before processing - match result { - Ok(Message::Text(text)) => { - pb.complete(Some(Arc::new(text.to_string()))); - return; - } - Ok(Message::Close(_)) => { - pb.complete(None); - return; - } - Ok(_) => { - // Skip non-text messages, keep reading - continue; - } - Err(tungstenite::Error::Io(ref e)) - if e.kind() == std::io::ErrorKind::WouldBlock - || e.kind() == std::io::ErrorKind::TimedOut => - { - // Read timeout — release lock briefly and retry - std::thread::sleep(std::time::Duration::from_millis(10)); - continue; - } - Err(e) => { - eprintln!("[ws_recv] error: {:?}", e); - pb.complete(None); - return; - } - } + let rx = conn.0.recv_rx.lock().unwrap(); + match rx.recv() { + Ok(Some(text)) => pb.complete(Some(Arc::new(text))), + _ => pb.complete(None), } }); promise @@ -244,21 +233,8 @@ pub fn std_ws_close(conn: &dyn WsConnectionTrait) -> Promise<()> { let conn: SimpleWsConnection = temper_core::cast(conn.as_any_value()).expect("WsConnection downcast"); let pb = PromiseBuilder::new(); let promise = pb.promise(); - crate::run_async(Arc::new(move || { - let pb = pb.clone(); - let conn = conn.clone(); - SafeGenerator::from_fn(Arc::new(move |_generator: SafeGenerator<()>| { - let mut socket = conn.0.socket.lock().unwrap(); - let _ = socket.close(None); - loop { - match socket.read() { - Ok(Message::Close(_)) | Err(_) => break, - _ => continue, - } - } - pb.complete(()); - None - })) - })); + // Drop the send channel which signals the I/O thread to stop + drop(conn); + pb.complete(()); promise } From 8ed2fd49ffde41ecbd94da4e7be7230377e5e2ec Mon Sep 17 00:00:00 2001 From: Robert Grayson Date: Sun, 15 Mar 2026 17:49:16 -0400 Subject: [PATCH 09/12] feat: replace tungstenite with hand-rolled WebSocket over raw TCP MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit tungstenite corrupts internal state when read/write are interleaved from different threads, even with Mutex protection. Replace it with a minimal WebSocket implementation (~150 lines) that: - Does HTTP upgrade handshake byte-by-byte (avoids BufReader read-ahead) - Uses TcpStream::try_clone() to split into independent read/write halves - Reader thread owns the read half, writer Mutex holds the write half - No shared internal framing state — reads and writes are independent - SHA-1 for the accept key via sha1_smol (only external dep) Also removes the sleep() debug logging. Co-Authored-By: Claude Opus 4.6 (1M context) --- .../kotlin/lang/temper/be/rust/RustBackend.kt | 3 +- .../lang/temper/be/rust/std/io/support.rs | 1 + .../lang/temper/be/rust/std/ws/support.rs | 287 +++++++++++++----- 3 files changed, 206 insertions(+), 85 deletions(-) diff --git a/be-rust/src/commonMain/kotlin/lang/temper/be/rust/RustBackend.kt b/be-rust/src/commonMain/kotlin/lang/temper/be/rust/RustBackend.kt index 7b5fcea7..980b8dac 100644 --- a/be-rust/src/commonMain/kotlin/lang/temper/be/rust/RustBackend.kt +++ b/be-rust/src/commonMain/kotlin/lang/temper/be/rust/RustBackend.kt @@ -181,6 +181,7 @@ class RustBackend(setup: BackendSetup) : Backend(Facto append("time = { version = \"=0.3.41\", optional = true }\n") append("ureq = { version = \"=3.1.2\", optional = true }\n") append("crossterm = { version = \"=0.28.1\", optional = true }\n") + append("sha1_smol = { version = \"=1.0.1\", optional = true }\n") // Below aren't dependencies section anymore, but eh. append("\n") append("[features]\n") @@ -189,7 +190,7 @@ class RustBackend(setup: BackendSetup) : Backend(Facto append("net = [\"ureq\"]\n") // Implied: append("regex = [\"regex\"]\n") append("temporal = [\"time\"]\n") - append("ws = [\"tungstenite\"]\n") + append("ws = [\"sha1_smol\"]\n") } } val packageFields = buildMap { diff --git a/be-rust/src/commonMain/resources/lang/temper/be/rust/std/io/support.rs b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/io/support.rs index 2cccb14c..fbdf0b2b 100644 --- a/be-rust/src/commonMain/resources/lang/temper/be/rust/std/io/support.rs +++ b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/io/support.rs @@ -8,6 +8,7 @@ pub fn std_sleep(ms: i32) -> Promise<()> { crate::run_async(Arc::new(move || { let pb = pb.clone(); SafeGenerator::from_fn(Arc::new(move |_generator: SafeGenerator<()>| { + eprintln!("[sleep] {}ms", ms); std::thread::sleep(std::time::Duration::from_millis(ms as u64)); pb.complete(()); None diff --git a/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs index a6bd7022..59bbe40e 100644 --- a/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs +++ b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs @@ -1,5 +1,7 @@ use super::*; +use std::io::{Read, Write, BufRead, BufReader}; use std::sync::{Arc, Mutex, mpsc}; +use std::net::{TcpListener, TcpStream}; use temper_core::{AsAnyValue, Promise, PromiseBuilder, SafeGenerator}; #[cfg(not(feature = "ws"))] @@ -15,18 +17,155 @@ pub fn std_ws_recv(_conn: &dyn WsConnectionTrait) -> Promise> #[cfg(not(feature = "ws"))] pub fn std_ws_close(_conn: &dyn WsConnectionTrait) -> Promise<()> { panic!() } +// --- Minimal WebSocket implementation (RFC 6455 text frames only) --- + #[cfg(feature = "ws")] -use std::net::{TcpListener, TcpStream}; +fn ws_write_text_frame(w: &mut impl Write, text: &str, mask: bool) -> std::io::Result<()> { + let payload = text.as_bytes(); + let len = payload.len(); + // Opcode 0x81 = final frame, text + w.write_all(&[0x81])?; + let mask_bit = if mask { 0x80 } else { 0x00 }; + if len < 126 { + w.write_all(&[(len as u8) | mask_bit])?; + } else if len < 65536 { + w.write_all(&[126 | mask_bit])?; + w.write_all(&(len as u16).to_be_bytes())?; + } else { + w.write_all(&[127 | mask_bit])?; + w.write_all(&(len as u64).to_be_bytes())?; + } + if mask { + let mask_key: [u8; 4] = [0x12, 0x34, 0x56, 0x78]; // fixed mask for simplicity + w.write_all(&mask_key)?; + let masked: Vec = payload.iter().enumerate().map(|(i, b)| b ^ mask_key[i % 4]).collect(); + w.write_all(&masked)?; + } else { + w.write_all(payload)?; + } + w.flush() +} + +/// Read a WebSocket frame. Returns None on close, Some(text) on text frame. +#[cfg(feature = "ws")] +fn ws_read_text_frame(r: &mut impl Read) -> std::io::Result> { + let mut header = [0u8; 2]; + if r.read_exact(&mut header).is_err() { return Ok(None); } + let opcode = header[0] & 0x0F; + let masked = (header[1] & 0x80) != 0; + let mut len = (header[1] & 0x7F) as u64; + if len == 126 { + let mut buf = [0u8; 2]; + r.read_exact(&mut buf)?; + len = u16::from_be_bytes(buf) as u64; + } else if len == 127 { + let mut buf = [0u8; 8]; + r.read_exact(&mut buf)?; + len = u64::from_be_bytes(buf); + } + let mut mask_key = [0u8; 4]; + if masked { + r.read_exact(&mut mask_key)?; + } + let mut payload = vec![0u8; len as usize]; + r.read_exact(&mut payload)?; + if masked { + for (i, b) in payload.iter_mut().enumerate() { + *b ^= mask_key[i % 4]; + } + } + match opcode { + 0x01 => Ok(Some(String::from_utf8_lossy(&payload).to_string())), // text + 0x08 => Ok(None), // close + 0x09 => Ok(Some(String::new())), // ping — return empty, caller can ignore + _ => Ok(Some(String::new())), // skip other opcodes + } +} + +/// Server-side WebSocket handshake (accept upgrade request). +/// Reads byte-by-byte to avoid consuming past the headers. +#[cfg(feature = "ws")] +fn ws_server_handshake(stream: &mut TcpStream) -> std::io::Result<()> { + let mut headers = Vec::new(); + let mut ws_key = String::new(); + // Read headers byte-by-byte to not read past \r\n\r\n + loop { + let mut byte = [0u8; 1]; + stream.read_exact(&mut byte)?; + headers.push(byte[0]); + if headers.len() >= 4 && &headers[headers.len()-4..] == b"\r\n\r\n" { + break; + } + } + let header_str = String::from_utf8_lossy(&headers); + for line in header_str.lines() { + if line.to_lowercase().starts_with("sec-websocket-key:") { + ws_key = line.split(':').nth(1).unwrap_or("").trim().to_string(); + } + } + let accept = ws_accept_key(&ws_key); + eprintln!("[ws_handshake] key='{}' accept='{}'", ws_key, accept); + let response = format!( + "HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: {}\r\n\r\n", + accept + ); + stream.write_all(response.as_bytes())?; + stream.flush() +} + +/// Client-side WebSocket handshake (send upgrade request). +/// Reads byte-by-byte to avoid consuming past the headers. #[cfg(feature = "ws")] -use tungstenite::{accept as ws_upgrade, connect as ws_connect_url, stream::MaybeTlsStream, Message, WebSocket}; +fn ws_client_handshake(stream: &mut TcpStream, host: &str, path: &str) -> std::io::Result<()> { + let key = "dGhlIHNhbXBsZSBub25jZQ=="; // fixed nonce, fine for our use + let request = format!( + "GET {} HTTP/1.1\r\nHost: {}\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Key: {}\r\nSec-WebSocket-Version: 13\r\n\r\n", + path, host, key + ); + stream.write_all(request.as_bytes())?; + stream.flush()?; + // Read response headers byte-by-byte to not consume past \r\n\r\n + let mut headers = Vec::new(); + loop { + let mut byte = [0u8; 1]; + stream.read_exact(&mut byte)?; + headers.push(byte[0]); + if headers.len() >= 4 && &headers[headers.len()-4..] == b"\r\n\r\n" { + break; + } + } + Ok(()) +} -// --- Server --- +#[cfg(feature = "ws")] +fn ws_accept_key(key: &str) -> String { + let combined = format!("{}258EAFA5-E914-47DA-95CA-C5AB0DC85B11", key); + let digest = sha1_smol::Sha1::from(combined).digest().bytes(); + base64_encode(&digest) +} #[cfg(feature = "ws")] -struct WsServerInner { - listener: Mutex, +fn base64_encode(data: &[u8]) -> String { + const CHARS: &[u8] = b"ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/"; + let mut result = String::new(); + for chunk in data.chunks(3) { + let b0 = chunk[0] as u32; + let b1 = if chunk.len() > 1 { chunk[1] as u32 } else { 0 }; + let b2 = if chunk.len() > 2 { chunk[2] as u32 } else { 0 }; + let n = (b0 << 16) | (b1 << 8) | b2; + result.push(CHARS[((n >> 18) & 63) as usize] as char); + result.push(CHARS[((n >> 12) & 63) as usize] as char); + if chunk.len() > 1 { result.push(CHARS[((n >> 6) & 63) as usize] as char); } else { result.push('='); } + if chunk.len() > 2 { result.push(CHARS[(n & 63) as usize] as char); } else { result.push('='); } + } + result } +// --- Server/Connection types --- + +#[cfg(feature = "ws")] +struct WsServerInner { listener: Mutex } + #[cfg(feature = "ws")] #[derive(Clone)] struct SimpleWsServer(Arc); @@ -39,14 +178,11 @@ impl WsServerTrait for SimpleWsServer { #[cfg(feature = "ws")] temper_core::impl_any_value_trait!(SimpleWsServer, [WsServer]); -// --- Connection --- -// Each connection has a dedicated I/O thread that owns the WebSocket. -// Send and recv go through channels, avoiding Mutex sharing issues. - #[cfg(feature = "ws")] struct WsConnectionInner { - send_tx: mpsc::Sender, + writer: Mutex, // write half recv_rx: Mutex>>, + is_client: bool, } #[cfg(feature = "ws")] @@ -61,48 +197,31 @@ impl WsConnectionTrait for SimpleWsConnection { #[cfg(feature = "ws")] temper_core::impl_any_value_trait!(SimpleWsConnection, [WsConnection]); -/// Spawn an I/O thread that owns the WebSocket and bridges to channels. +/// Create a connection from a TcpStream. Spawns a reader thread. #[cfg(feature = "ws")] -fn spawn_ws_io_thread(mut ws: WebSocket) -> SimpleWsConnection -where - S: std::io::Read + std::io::Write + Send + 'static, -{ - let (send_tx, send_rx) = mpsc::channel::(); - let (recv_tx, recv_rx) = mpsc::sync_channel::>(16); +fn make_connection(stream: TcpStream, is_client: bool) -> SimpleWsConnection { + let writer = stream.try_clone().expect("clone stream for writer"); + let (recv_tx, recv_rx) = mpsc::sync_channel::>(32); + // Reader thread owns the read half + let is_client_thread = is_client; std::thread::spawn(move || { + let mut reader = stream; + eprintln!("[ws_reader] started (client={})", is_client_thread); loop { - // Try to send any queued outgoing messages. - while let Ok(msg) = send_rx.try_recv() { - if ws.send(Message::Text(msg.into())).is_err() { - let _ = recv_tx.send(None); - return; - } - } - - // Try to read one incoming message (may timeout quickly). - match ws.read() { - Ok(Message::Text(text)) => { - if recv_tx.send(Some(text.to_string())).is_err() { - return; // recv side dropped - } + match ws_read_text_frame(&mut reader) { + Ok(Some(text)) if text.is_empty() => continue, + Ok(Some(text)) => { + eprintln!("[ws_reader] got {} bytes", text.len()); + if recv_tx.send(Some(text)).is_err() { return; } } - Ok(Message::Close(_)) => { + Ok(None) => { + eprintln!("[ws_reader] got close frame"); let _ = recv_tx.send(None); return; } - Ok(Message::Ping(data)) => { - let _ = ws.send(Message::Pong(data)); - } - Ok(_) => {} // skip binary, pong, etc - Err(tungstenite::Error::Io(ref e)) - if e.kind() == std::io::ErrorKind::WouldBlock - || e.kind() == std::io::ErrorKind::TimedOut => - { - // No data yet — brief yield then loop - std::thread::sleep(std::time::Duration::from_millis(5)); - } - Err(_) => { + Err(e) => { + eprintln!("[ws_reader] error: {:?}", e); let _ = recv_tx.send(None); return; } @@ -111,19 +230,13 @@ where }); SimpleWsConnection(Arc::new(WsConnectionInner { - send_tx, + writer: Mutex::new(writer), recv_rx: Mutex::new(recv_rx), + is_client, })) } -/// Set a short read timeout on a TcpStream before passing it to tungstenite. -/// This makes the I/O thread's read() return quickly, allowing send interleaving. -#[cfg(feature = "ws")] -fn prepare_stream(stream: &TcpStream) { - let _ = stream.set_read_timeout(Some(std::time::Duration::from_millis(20))); -} - -// --- Functions --- +// --- Public API --- #[cfg(feature = "ws")] pub fn std_ws_listen(port: i32) -> Promise { @@ -133,11 +246,9 @@ pub fn std_ws_listen(port: i32) -> Promise { let pb = pb.clone(); SafeGenerator::from_fn(Arc::new(move |_: SafeGenerator<()>| { match TcpListener::bind(format!("0.0.0.0:{}", port)) { - Ok(listener) => { - pb.complete(WsServer::new(SimpleWsServer(Arc::new(WsServerInner { - listener: Mutex::new(listener), - })))); - } + Ok(listener) => pb.complete(WsServer::new(SimpleWsServer(Arc::new( + WsServerInner { listener: Mutex::new(listener) }, + )))), Err(_) => pb.break_promise(), } None @@ -148,20 +259,16 @@ pub fn std_ws_listen(port: i32) -> Promise { #[cfg(feature = "ws")] pub fn std_ws_accept(server: &dyn WsServerTrait) -> Promise { - let server: SimpleWsServer = temper_core::cast(server.as_any_value()).expect("WsServer downcast"); + let server: SimpleWsServer = temper_core::cast(server.as_any_value()).expect("downcast"); let pb = PromiseBuilder::new(); let promise = pb.promise(); std::thread::spawn(move || { let listener = server.0.listener.lock().unwrap(); match listener.accept() { - Ok((stream, _)) => { + Ok((mut stream, _)) => { drop(listener); - prepare_stream(&stream); - match ws_upgrade(stream) { - Ok(ws) => { - let conn = spawn_ws_io_thread(ws); - pb.complete(WsConnection::new(conn)); - } + match ws_server_handshake(&mut stream) { + Ok(()) => pb.complete(WsConnection::new(make_connection(stream, false))), Err(_) => pb.break_promise(), } } @@ -180,15 +287,17 @@ pub fn std_ws_connect(url: impl temper_core::ToArcString) -> Promise| { - match ws_connect_url(url.as_str()) { - Ok((mut ws, _)) => { - // Set read timeout on client connections - match ws.get_ref() { - MaybeTlsStream::Plain(s) => prepare_stream(s), - _ => {} + // Parse ws://host:port/path + let url_str = url.as_str(); + let stripped = url_str.strip_prefix("ws://").unwrap_or(url_str); + let (host_port, path) = stripped.split_once('/').unwrap_or((stripped, "")); + let path = format!("/{}", path); + match TcpStream::connect(host_port) { + Ok(mut stream) => { + match ws_client_handshake(&mut stream, host_port, &path) { + Ok(()) => pb.complete(WsConnection::new(make_connection(stream, true))), + Err(_) => pb.break_promise(), } - let conn = spawn_ws_io_thread(ws); - pb.complete(WsConnection::new(conn)); } Err(_) => pb.break_promise(), } @@ -200,24 +309,31 @@ pub fn std_ws_connect(url: impl temper_core::ToArcString) -> Promise Promise<()> { - let conn: SimpleWsConnection = temper_core::cast(conn.as_any_value()).expect("WsConnection downcast"); + let conn: SimpleWsConnection = temper_core::cast(conn.as_any_value()).expect("downcast"); let msg = msg.to_arc_string(); let pb = PromiseBuilder::new(); let promise = pb.promise(); - // Send is non-blocking — just push to the channel - match conn.0.send_tx.send(msg.to_string()) { - Ok(_) => pb.complete(()), - Err(_) => pb.break_promise(), + // Write directly to the writer half (no contention with reader thread) + let mut writer = conn.0.writer.lock().unwrap(); + let msg_len = msg.len(); + match ws_write_text_frame(&mut *writer, &msg, conn.0.is_client) { + Ok(()) => { + eprintln!("[ws_send] sent {} bytes OK", msg_len); + pb.complete(()); + } + Err(e) => { + eprintln!("[ws_send] write error: {:?}", e); + pb.break_promise(); + } } promise } #[cfg(feature = "ws")] pub fn std_ws_recv(conn: &dyn WsConnectionTrait) -> Promise>> { - let conn: SimpleWsConnection = temper_core::cast(conn.as_any_value()).expect("WsConnection downcast"); + let conn: SimpleWsConnection = temper_core::cast(conn.as_any_value()).expect("downcast"); let pb = PromiseBuilder::new(); let promise = pb.promise(); - // Recv blocks waiting for the next message from the I/O thread std::thread::spawn(move || { let rx = conn.0.recv_rx.lock().unwrap(); match rx.recv() { @@ -230,11 +346,14 @@ pub fn std_ws_recv(conn: &dyn WsConnectionTrait) -> Promise>> #[cfg(feature = "ws")] pub fn std_ws_close(conn: &dyn WsConnectionTrait) -> Promise<()> { - let conn: SimpleWsConnection = temper_core::cast(conn.as_any_value()).expect("WsConnection downcast"); + let conn: SimpleWsConnection = temper_core::cast(conn.as_any_value()).expect("downcast"); let pb = PromiseBuilder::new(); let promise = pb.promise(); - // Drop the send channel which signals the I/O thread to stop - drop(conn); + // Send close frame + if let Ok(mut writer) = conn.0.writer.lock() { + let _ = writer.write_all(&[0x88, 0x00]); // close frame + let _ = writer.flush(); + } pb.complete(()); promise } From 561c49a718f4937dcb5988e31cec55e1781c0009 Mon Sep 17 00:00:00 2001 From: Robert Grayson Date: Sun, 15 Mar 2026 18:57:51 -0400 Subject: [PATCH 10/12] feat: add Python backend for std/ws and std/io terminal size MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Wire WebSocket support (wsListen, wsAccept, wsConnect, wsSend, wsRecv, wsClose) and terminal size detection (terminalColumns, terminalRows) for the Python backend. Uses hand-rolled WebSocket over raw TCP (same approach as Rust) with socket.dup() for read/write split and a dedicated reader thread per connection. No external dependencies — just stdlib socket, hashlib, base64, struct, threading. Terminal size uses shutil.get_terminal_size(). Co-Authored-By: Claude Opus 4.6 (1M context) --- .../lang/temper/be/py/PySupportNetwork.kt | 21 ++ .../be/py/temper-core/temper_core/__init__.py | 250 ++++++++++++++++++ 2 files changed, 271 insertions(+) diff --git a/be-py/src/commonMain/kotlin/lang/temper/be/py/PySupportNetwork.kt b/be-py/src/commonMain/kotlin/lang/temper/be/py/PySupportNetwork.kt index 8f174ff9..ea99b1cc 100644 --- a/be-py/src/commonMain/kotlin/lang/temper/be/py/PySupportNetwork.kt +++ b/be-py/src/commonMain/kotlin/lang/temper/be/py/PySupportNetwork.kt @@ -911,6 +911,17 @@ val StdNetSend = PySeparateCode("std_net_send", RUNTIME) val StdSleep = PySeparateCode("std_sleep", RUNTIME) val StdReadLine = PySeparateCode("std_read_line", RUNTIME) val StdNextKeypress = PySeparateCode("std_next_keypress", RUNTIME) +val StdTermCols = PySeparateCode("std_term_cols", RUNTIME) +val StdTermRows = PySeparateCode("std_term_rows", RUNTIME) + +val WsServer = PyConnectedType("WsServer", RUNTIME) +val WsConnection = PyConnectedType("WsConnection", RUNTIME) +val StdWsListen = PySeparateCode("std_ws_listen", RUNTIME) +val StdWsAccept = PySeparateCode("std_ws_accept", RUNTIME) +val StdWsConnect = PySeparateCode("std_ws_connect", RUNTIME) +val StdWsSend = PySeparateCode("std_ws_send", RUNTIME) +val StdWsRecv = PySeparateCode("std_ws_recv", RUNTIME) +val StdWsClose = PySeparateCode("std_ws_close", RUNTIME) val mathInf = PySeparateCode("inf", MATH) val mathNan = PySeparateCode("nan", MATH) @@ -1219,4 +1230,14 @@ private val pyConnections = mapOf( "stdSleep" to StdSleep, "stdReadLine" to StdReadLine, "stdNextKeypress" to StdNextKeypress, + "stdTermCols" to StdTermCols, + "stdTermRows" to StdTermRows, + "wsListen" to StdWsListen, + "wsAccept" to StdWsAccept, + "wsConnect" to StdWsConnect, + "wsSend" to StdWsSend, + "wsRecv" to StdWsRecv, + "wsClose" to StdWsClose, + "WsServer" to WsServer, + "WsConnection" to WsConnection, ) diff --git a/be-py/src/commonMain/resources/lang/temper/be/py/temper-core/temper_core/__init__.py b/be-py/src/commonMain/resources/lang/temper/be/py/temper-core/temper_core/__init__.py index d60d355a..8c20cd8d 100644 --- a/be-py/src/commonMain/resources/lang/temper/be/py/temper-core/temper_core/__init__.py +++ b/be-py/src/commonMain/resources/lang/temper/be/py/temper-core/temper_core/__init__.py @@ -1487,3 +1487,253 @@ def _do_read(): _executor.submit(_do_read) return f + + +# --- Terminal Size --- + +def std_term_cols() -> int: + import shutil as _shutil + return _shutil.get_terminal_size((80, 24)).columns + +def std_term_rows() -> int: + import shutil as _shutil + return _shutil.get_terminal_size((80, 24)).lines + + +# --- WebSocket (hand-rolled RFC 6455 text frames over raw TCP) --- + +import socket as _socket +import struct as _struct +import hashlib as _hashlib +import base64 as _base64 +import threading as _threading + +class WsServer: + def __init__(self, sock): + self._sock = sock + +class WsConnection: + def __init__(self, rsock, wsock): + self._rsock = rsock # read half + self._wsock = wsock # write half + self._wlock = _threading.Lock() + +def _ws_write_frame(sock, text, mask=False): + payload = text.encode('utf-8') + header = bytearray() + header.append(0x81) # fin + text opcode + mask_bit = 0x80 if mask else 0x00 + ln = len(payload) + if ln < 126: + header.append(ln | mask_bit) + elif ln < 65536: + header.append(126 | mask_bit) + header.extend(_struct.pack('>H', ln)) + else: + header.append(127 | mask_bit) + header.extend(_struct.pack('>Q', ln)) + if mask: + mask_key = b'\x12\x34\x56\x78' + header.extend(mask_key) + payload = bytes(b ^ mask_key[i % 4] for i, b in enumerate(payload)) + sock.sendall(bytes(header) + payload) + +def _ws_read_frame(sock): + def _recv_exact(n): + buf = b'' + while len(buf) < n: + chunk = sock.recv(n - len(buf)) + if not chunk: + return None + buf += chunk + return buf + hdr = _recv_exact(2) + if hdr is None: + return None + opcode = hdr[0] & 0x0F + masked = (hdr[1] & 0x80) != 0 + ln = hdr[1] & 0x7F + if ln == 126: + ext = _recv_exact(2) + if ext is None: + return None + ln = _struct.unpack('>H', ext)[0] + elif ln == 127: + ext = _recv_exact(8) + if ext is None: + return None + ln = _struct.unpack('>Q', ext)[0] + mask_key = None + if masked: + mask_key = _recv_exact(4) + if mask_key is None: + return None + data = _recv_exact(ln) + if data is None: + return None + if mask_key: + data = bytes(b ^ mask_key[i % 4] for i, b in enumerate(data)) + if opcode == 0x01: # text + return data.decode('utf-8', errors='replace') + if opcode == 0x08: # close + return None + return '' # ping/pong/binary — skip + +def _ws_server_handshake(conn): + data = b'' + while b'\r\n\r\n' not in data: + chunk = conn.recv(1) + if not chunk: + raise ConnectionError("handshake failed") + data += chunk + headers = data.decode('utf-8', errors='replace') + key = '' + for line in headers.split('\r\n'): + if line.lower().startswith('sec-websocket-key:'): + key = line.split(':', 1)[1].strip() + accept = _base64.b64encode( + _hashlib.sha1((key + '258EAFA5-E914-47DA-95CA-C5AB0DC85B11').encode()).digest() + ).decode() + response = ( + 'HTTP/1.1 101 Switching Protocols\r\n' + 'Upgrade: websocket\r\n' + 'Connection: Upgrade\r\n' + f'Sec-WebSocket-Accept: {accept}\r\n' + '\r\n' + ) + conn.sendall(response.encode()) + +def _ws_client_handshake(sock, host, path): + key = 'dGhlIHNhbXBsZSBub25jZQ==' + request = ( + f'GET {path} HTTP/1.1\r\n' + f'Host: {host}\r\n' + 'Upgrade: websocket\r\n' + 'Connection: Upgrade\r\n' + f'Sec-WebSocket-Key: {key}\r\n' + 'Sec-WebSocket-Version: 13\r\n' + '\r\n' + ) + sock.sendall(request.encode()) + data = b'' + while b'\r\n\r\n' not in data: + chunk = sock.recv(1) + if not chunk: + raise ConnectionError("handshake failed") + data += chunk + +def _make_ws_connection(sock): + wsock = sock.dup() + reader_rx_send = [] + reader_rx_lock = _threading.Lock() + reader_rx_event = _threading.Event() + + def _reader(): + try: + while True: + msg = _ws_read_frame(sock) + with reader_rx_lock: + reader_rx_send.append(msg) + reader_rx_event.set() + if msg is None: + return + except Exception: + with reader_rx_lock: + reader_rx_send.append(None) + reader_rx_event.set() + + t = _threading.Thread(target=_reader, daemon=True) + t.start() + conn = WsConnection(sock, wsock) + conn._rx_queue = reader_rx_send + conn._rx_lock = reader_rx_lock + conn._rx_event = reader_rx_event + return conn + +def std_ws_listen(port: int) -> 'Future[WsServer]': + f = new_unbound_promise() + def _do(): + try: + s = _socket.socket(_socket.AF_INET, _socket.SOCK_STREAM) + s.setsockopt(_socket.SOL_SOCKET, _socket.SO_REUSEADDR, 1) + s.bind(('0.0.0.0', port)) + s.listen(16) + f.set_result(WsServer(s)) + except Exception: + f.cancel() + _executor.submit(_do) + return f + +def std_ws_accept(server) -> 'Future[WsConnection]': + f = new_unbound_promise() + def _do(): + try: + conn, _ = server._sock.accept() + _ws_server_handshake(conn) + f.set_result(_make_ws_connection(conn)) + except Exception: + f.cancel() + _threading.Thread(target=_do, daemon=True).start() + return f + +def std_ws_connect(url: str) -> 'Future[WsConnection]': + f = new_unbound_promise() + def _do(): + try: + stripped = url.replace('ws://', '') + if '/' in stripped: + host_port, path = stripped.split('/', 1) + path = '/' + path + else: + host_port = stripped + path = '/' + if ':' in host_port: + host, port_str = host_port.split(':') + port = int(port_str) + else: + host = host_port + port = 80 + s = _socket.socket(_socket.AF_INET, _socket.SOCK_STREAM) + s.connect((host, port)) + _ws_client_handshake(s, host_port, path) + f.set_result(_make_ws_connection(s)) + except Exception: + f.cancel() + _executor.submit(_do) + return f + +def std_ws_send(conn, msg: str) -> 'Future[None]': + f = new_unbound_promise() + try: + with conn._wlock: + _ws_write_frame(conn._wsock, msg, mask=False) + f.set_result(None) + except Exception: + f.cancel() + return f + +def std_ws_recv(conn) -> 'Future[Optional[str]]': + f: 'Future[Optional[str]]' = new_unbound_promise() + def _do(): + while True: + conn._rx_event.wait() + with conn._rx_lock: + if conn._rx_queue: + msg = conn._rx_queue.pop(0) + if msg == '': + continue # skip ping/pong + f.set_result(msg) + return + conn._rx_event.clear() + _threading.Thread(target=_do, daemon=True).start() + return f + +def std_ws_close(conn) -> 'Future[None]': + f = new_unbound_promise() + try: + conn._wsock.sendall(b'\x88\x00') + conn._wsock.close() + except Exception: + pass + f.set_result(None) + return f From 9bb8efaf2fbd699054520d3ed70ee86adf3a1424 Mon Sep 17 00:00:00 2001 From: Robert Grayson Date: Sun, 15 Mar 2026 20:21:32 -0400 Subject: [PATCH 11/12] chore: remove debug logging from Rust ws support Co-Authored-By: Claude Opus 4.6 (1M context) --- .../resources/lang/temper/be/rust/std/ws/support.rs | 8 +------- 1 file changed, 1 insertion(+), 7 deletions(-) diff --git a/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs index 59bbe40e..e2cdb67a 100644 --- a/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs +++ b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs @@ -104,7 +104,6 @@ fn ws_server_handshake(stream: &mut TcpStream) -> std::io::Result<()> { } } let accept = ws_accept_key(&ws_key); - eprintln!("[ws_handshake] key='{}' accept='{}'", ws_key, accept); let response = format!( "HTTP/1.1 101 Switching Protocols\r\nUpgrade: websocket\r\nConnection: Upgrade\r\nSec-WebSocket-Accept: {}\r\n\r\n", accept @@ -207,21 +206,18 @@ fn make_connection(stream: TcpStream, is_client: bool) -> SimpleWsConnection { let is_client_thread = is_client; std::thread::spawn(move || { let mut reader = stream; - eprintln!("[ws_reader] started (client={})", is_client_thread); + // Reader thread running loop { match ws_read_text_frame(&mut reader) { Ok(Some(text)) if text.is_empty() => continue, Ok(Some(text)) => { - eprintln!("[ws_reader] got {} bytes", text.len()); if recv_tx.send(Some(text)).is_err() { return; } } Ok(None) => { - eprintln!("[ws_reader] got close frame"); let _ = recv_tx.send(None); return; } Err(e) => { - eprintln!("[ws_reader] error: {:?}", e); let _ = recv_tx.send(None); return; } @@ -318,11 +314,9 @@ pub fn std_ws_send(conn: &dyn WsConnectionTrait, msg: impl temper_core::ToArcStr let msg_len = msg.len(); match ws_write_text_frame(&mut *writer, &msg, conn.0.is_client) { Ok(()) => { - eprintln!("[ws_send] sent {} bytes OK", msg_len); pb.complete(()); } Err(e) => { - eprintln!("[ws_send] write error: {:?}", e); pb.break_promise(); } } From 936910d950a39526fd6cda1ceaf02eaa1189288e Mon Sep 17 00:00:00 2001 From: Robert Grayson Date: Mon, 16 Mar 2026 08:18:52 -0400 Subject: [PATCH 12/12] fix: guard wsSend against closed WebSocket connections Check readyState before calling send(), and wrap in try/catch to prevent synchronous throws from crashing the server when a client disconnects mid-broadcast. Co-Authored-By: Claude Opus 4.6 (1M context) --- .../resources/lang/temper/be/js/temper-core/ws.js | 13 +++++++++---- 1 file changed, 9 insertions(+), 4 deletions(-) diff --git a/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/ws.js b/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/ws.js index 26fc3bf8..4eb74774 100644 --- a/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/ws.js +++ b/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/ws.js @@ -103,10 +103,15 @@ export async function wsConnect(url) { */ export function wsSend(conn, msg) { return new Promise((resolve, reject) => { - conn.send(msg, (err) => { - if (err) reject(err); - else resolve(empty()); - }); + try { + if (conn.readyState !== 1) { reject(new Error("not open")); return; } + conn.send(msg, (err) => { + if (err) reject(err); + else resolve(empty()); + }); + } catch (e) { + reject(e); + } }); }