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/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" + } } 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..4eb74774 --- /dev/null +++ b/be-js/src/commonMain/resources/lang/temper/be/js/temper-core/ws.js @@ -0,0 +1,148 @@ +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) => { + 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); + } + }); +} + +/** + * @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-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 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..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,14 +181,16 @@ 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") - 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 = [\"sha1_smol\"]\n") } } val packageFields = buildMap { @@ -256,7 +258,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..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 @@ -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::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", @@ -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..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 @@ -16,6 +17,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/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 } 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..e2cdb67a --- /dev/null +++ b/be-rust/src/commonMain/resources/lang/temper/be/rust/std/ws/support.rs @@ -0,0 +1,353 @@ +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"))] +pub fn std_ws_listen(_port: i32) -> Promise { panic!() } +#[cfg(not(feature = "ws"))] +pub fn std_ws_accept(_server: &dyn WsServerTrait) -> Promise { panic!() } +#[cfg(not(feature = "ws"))] +pub fn std_ws_connect(_url: impl temper_core::ToArcString) -> Promise { panic!() } +#[cfg(not(feature = "ws"))] +pub fn std_ws_send(_conn: &dyn WsConnectionTrait, _msg: impl temper_core::ToArcString) -> Promise<()> { panic!() } +#[cfg(not(feature = "ws"))] +pub fn std_ws_recv(_conn: &dyn WsConnectionTrait) -> Promise>> { panic!() } +#[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")] +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); + 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")] +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(()) +} + +#[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")] +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); + +#[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 { + writer: Mutex, // write half + recv_rx: Mutex>>, + is_client: bool, +} + +#[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]); + +/// Create a connection from a TcpStream. Spawns a reader thread. +#[cfg(feature = "ws")] +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; + // Reader thread running + loop { + match ws_read_text_frame(&mut reader) { + Ok(Some(text)) if text.is_empty() => continue, + Ok(Some(text)) => { + if recv_tx.send(Some(text)).is_err() { return; } + } + Ok(None) => { + let _ = recv_tx.send(None); + return; + } + Err(e) => { + let _ = recv_tx.send(None); + return; + } + } + } + }); + + SimpleWsConnection(Arc::new(WsConnectionInner { + writer: Mutex::new(writer), + recv_rx: Mutex::new(recv_rx), + is_client, + })) +} + +// --- Public API --- + +#[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 |_: 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(), + } + None + })) + })); + promise +} + +#[cfg(feature = "ws")] +pub fn std_ws_accept(server: &dyn WsServerTrait) -> Promise { + 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((mut stream, _)) => { + drop(listener); + match ws_server_handshake(&mut stream) { + Ok(()) => pb.complete(WsConnection::new(make_connection(stream, false))), + Err(_) => pb.break_promise(), + } + } + Err(_) => pb.break_promise(), + } + }); + promise +} + +#[cfg(feature = "ws")] +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(); + crate::run_async(Arc::new(move || { + let pb = pb.clone(); + let url = url.clone(); + SafeGenerator::from_fn(Arc::new(move |_: SafeGenerator<()>| { + // 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(), + } + } + Err(_) => pb.break_promise(), + } + None + })) + })); + promise +} + +#[cfg(feature = "ws")] +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("downcast"); + let msg = msg.to_arc_string(); + let pb = PromiseBuilder::new(); + let promise = pb.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(()) => { + pb.complete(()); + } + Err(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("downcast"); + let pb = PromiseBuilder::new(); + let promise = pb.promise(); + std::thread::spawn(move || { + let rx = conn.0.recv_rx.lock().unwrap(); + match rx.recv() { + Ok(Some(text)) => pb.complete(Some(Arc::new(text))), + _ => pb.complete(None), + } + }); + promise +} + +#[cfg(feature = "ws")] +pub fn std_ws_close(conn: &dyn WsConnectionTrait) -> Promise<()> { + let conn: SimpleWsConnection = temper_core::cast(conn.as_any_value()).expect("downcast"); + let pb = PromiseBuilder::new(); + let promise = pb.promise(); + // 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 +} 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/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(), 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() + }