diff --git a/CMakeLists.txt b/CMakeLists.txt index 051597c5..f9f60932 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -76,6 +76,8 @@ find_library(libuci NAMES uci) find_library(libubox NAMES ubox) find_library(libubus NAMES ubus) find_library(libblobmsg_json NAMES blobmsg_json) +find_library(wslay NAMES wslay) +find_path(wslay_include NAMES wslay/wslay.h) find_package(ZLIB) find_library(libmd NAMES libmd.a md) find_library(libffi NAMES ffi) @@ -100,6 +102,10 @@ if(libubox) set(DEFAULT_ULOOP_SUPPORT ON) endif() +if(wslay AND wslay_include AND libubox) + set(DEFAULT_WEBSOCKET_SUPPORT ON) +endif() + if(ZLIB_FOUND) set(DEFAULT_ZLIB_SUPPORT ON) endif() @@ -126,6 +132,7 @@ option(STRUCT_SUPPORT "Struct plugin support" ON) option(ULOOP_SUPPORT "Uloop plugin support" ${DEFAULT_ULOOP_SUPPORT}) option(LOG_SUPPORT "Log plugin support" ON) option(SOCKET_SUPPORT "Socket plugin support" ON) +option(WEBSOCKET_SUPPORT "WebSocket plugin support" ${DEFAULT_WEBSOCKET_SUPPORT}) option(SERIAL_SUPPORT "Serial port plugin support" ON) option(ZLIB_SUPPORT "Zlib plugin support" ${DEFAULT_ZLIB_SUPPORT}) option(DIGEST_SUPPORT "Digest plugin support" ${DEFAULT_DIGEST_SUPPORT}) @@ -421,6 +428,15 @@ if(SOCKET_SUPPORT) target_link_options(socket_lib PRIVATE ${UCODE_MODULE_LINK_OPTIONS}) endif() +if(WEBSOCKET_SUPPORT) + set(LIBRARIES ${LIBRARIES} websocket_lib) + add_library(websocket_lib MODULE lib/websocket.c) + set_target_properties(websocket_lib PROPERTIES OUTPUT_NAME websocket PREFIX "") + target_include_directories(websocket_lib PRIVATE ${wslay_include}) + target_link_options(websocket_lib PRIVATE ${UCODE_MODULE_LINK_OPTIONS}) + target_link_libraries(websocket_lib ${wslay} ${libubox}) +endif() + if(SERIAL_SUPPORT) set(LIBRARIES ${LIBRARIES} serial_lib) add_library(serial_lib MODULE lib/serial.c) diff --git a/WEBSOCKET_PLAN.md b/WEBSOCKET_PLAN.md new file mode 100644 index 00000000..8c14b6a5 --- /dev/null +++ b/WEBSOCKET_PLAN.md @@ -0,0 +1,378 @@ +# ucode `websocket` Module — Design & Tracking Plan + +> Status: **M5 COMPLETE — PR-ready; M6 (TLS) optional** +> Created: 2026-09-01 · Last update: 2026-09-01 +> Target repo path: `lib/websocket.c` (single-file ucode module) + +--- + +## 1. Goals + +- Provide WebSocket (RFC 6455) **client** support to ucode scripts. +- Minimum memory footprint (embedded/OpenWrt targets: mips, arm, 32–128 MB RAM devices). +- Depend only on a small, stable, well-maintained C library (see §3). +- Async, event-driven API integrated with `uloop` (pattern from `lib/uloop.c`). +- Play nice with the existing `socket` module philosophy (resources, errno-style errors). + +### Non-goals (explicit) + +- ❌ WebSocket **server** role (phase 2+, only if a real need appears). +- ❌ `permessage-deflate` compression — rejected: memory/CPU cost defeats the + minimum-footprint goal on embedded targets. +- ❌ Auto-reconnect / reconnect backoff — userland responsibility by design. +- ❌ Thread pool / background threads — single-threaded, uloop-driven only. + +--- + +## 2. Constraints & Budget + +| Metric | Target | +|---|---| +| Code size (stripped .so, mips16) | < 100 KB | +| RSS delta per idle connection | < 60 KB | +| Default recv ring buffer | 4–8 KB (fixed, no per-message malloc) | +| Default max frame size cap | 256 KB (user-tunable, hard floor enforced) | +| External deps | wslay only (vendored), libc, libubox (uloop) | + +--- + +## 3. Library Decision + +| Option | Footprint | Maintenance | Verdict | +|---|---|---|---| +| **wslay** (MIT) | ~30 KB, zero-copy, no threads, no internal buffers (caller supplies buffers + I/O) | Stable; RFC 6455 is a frozen spec (2011) → code does not rot | ✅ **CHOSEN** (framing only) | +| libwebsockets | Large (HTTP/SSL machinery) | Excellent, very active | ❌ Overkill for client-only, too heavy | +| libuwsc | Small | Effectively unmaintained | ❌ Dead project | + +**Architecture:** wslay is used *purely as the RFC 6455 framing engine*. +Transport is plain POSIX non-blocking sockets owned by the module, integrated +with `uloop` for readiness events. TLS (mbedTLS, already packaged in OpenWrt) +is a phase-2 optional add-on behind the same API. + +--- + +## 4. Proposed API + +```javascript +import { connect } from 'websocket'; + +let ws = connect('ws://10.0.0.1:8080/path', { + headers: { Authorization: 'Bearer …' }, + max_frame_size: 262144, + recv_buffer_size: 8192, + timeout: 15000 +}); + +ws.on('open', (ws) => { … }); +ws.on('message', (ws, data, is_text) => { … }); +ws.on('close', (ws, code, reason) => { … }); +ws.on('error', (ws, err) => { … }); + +ws.send('hello'); // string => text frame +ws.send(new Uint8Array(…)); // typed array => binary frame +ws.ping('are you there?'); +ws.close(1000, 'bye'); + +ws.state; // 'connecting' | 'open' | 'closing' | 'closed' +``` + +Error reporting convention: match `lib/socket.c` — store last error, expose +`last_error()` / string form; throw ucode exceptions on programmer errors +(bad arguments), report runtime/network errors via the `error` event. + +--- + +## 5. Architecture Notes + +- Resource type via `ucv_resource_create_ex()` + `ucv_resource_persistent_set()` + (pattern: `lib/uloop.c:88-101`). +- One fixed receive ring buffer per connection; partial frames resume across + uloop wakeups — **no allocation per message**. +- Backpressure: stop registering the read event (`uloop_handle` delete) when the + user's message callback is still running; resume after callback returns. +- Ping/pong and close handshake handled internally and transparently. +- Handshake `Sec-WebSocket-Key`/`Accept`: SHA-1 + base64 implemented locally + (mirror the primitives used by `lib/digest.c`; no new dependency). +- GC finalizer calls `close()` safely if the user forgot (idempotent teardown). + +--- + +## 6. Milestones & Step Tracking + +Legend: `[ ]` pending · `[~]` in progress · `[x]` done · `[!]` blocked (see §9) + +### M1 — Scaffold ✅ (2026-09-01) +- [x] ~~Vendor wslay sources under `lib/wslay/`~~ → superseded by D8: **external lib, zlib pattern** +- [x] CMake: `find_library(wslay)`/`find_path` + auto-gated `WEBSOCKET_SUPPORT` + link (like `zlib_lib`) +- [x] Module registration skeleton + doc header (`@module websocket`) in `lib/websocket.c` + (plus `connect()` stub raising a "not implemented yet" exception) +- [x] Verified `ucode -lwebsocket` loads and registers `connect` (T1 ✅, re-verified after D8 refactor) +- [x] `libwslay` buildroot package (`../openwrt/package/wslay`, v1.1.1, hash pinned) — + **cross-compiled + staged for aarch64** (pulled forward from M5) +- [x] `ucode-mod-websocket` added to `openwrt/ucode/Makefile` (DEPENDS `+libubox +libwslay`) +- Note: `websocket.so` has no wslay DT_NEEDED yet — stub references no symbols and the + linker uses `--as-needed`; real linkage appears in M2 + +### M2 — Core client ✅ (2026-09-01) +- [x] URL parsing (`ws://` supported; `wss://` rejected with clear "TLS not supported yet"; + userinfo rejected; explicit port + query + IPv6 literal `[::1]:8080` handled — + ⚠️ initial parser broke on `[::1]` (first colon read as port separator); found in + pre-M3 review, fixed in M2 commit `f66c9d8` and verified end-to-end against a server on ::1) +- [x] DNS via synchronous `getaddrinfo` (see D12) + non-blocking TCP connect via uloop +- [x] HTTP Upgrade request generation incl. `headers` option (CRLF-injection guarded) +- [x] `Sec-WebSocket-Accept` validation (local SHA-1 + base64, no libmd dep) +- [x] wslay evented send/recv callbacks wired to the non-blocking fd +- [x] `on()` event dispatch (open/message/close/error) — callbacks receive the + connection as **explicit first argument** (ucode arrows have no `this`, see D11) +- [x] `send()` (string → text, array → binary), `ping()`, `close(code, reason)`, + plus `state()` and `fileno()` methods +- [x] Receive path: wslay fixed 4 KiB ibuf + capped message buffer (see D10) +- [x] Smoke-tested end-to-end against stdlib-Python WS echo server: + T2 handshake ✓, T3 text echo ✓, T4 binary echo ✓, T7 ping/pong ✓, + T8 close handshake ✓ (code 1000 round-trip), dead-port → `ERROR Connection + refused` + clean exit, bad URL → type exception +- [x] Committed as `fe57635` on `websocket-module` + +### M3 — Hardening ✅ (2026-09-01) +- [x] Frame-size cap enforcement — **T6 passes**: server frame > `max_frame_size` + → close event code **1009**, completes in ~0.1 s (no timeout linger) +- [x] `EINTR`/`EAGAIN` audit (loops already retry); `EMFILE` path audited + (connect loop closes fd per attempt, error propagates) +- [x] Idempotent `close()` (no-op while closing), abort semantics for `close()` + during connecting/handshake (close 1006 + immediate teardown), GC finalizer + (`ws_free_resource` via `uc_type_declare`) — dropped-resource scenario ASan clean +- [x] uloop fd/timeout deregistration on all error paths (teardown centralized; + `uloop_fd_delete` on unregistered fd verified safe against libubox source) +- [x] Close-phase timeout (5 s) + re-entrancy guard (`dispatching`/`need_flush`) + so `send()/close()` inside event callbacks never call into wslay re-entrantly +- [~] ASan/LSan clean — **zero findings** across happy/IPv6/oversized/abort/ + dead-port/dropped-resource scenarios (gcc `-fsanitize=address` build); + "full cram suite" part lands with M4 +- Bugs found & fixed during M3 (commit `2b30672`): + - **UAF/heap corruption**: connection resource lacked its own `ucv_get()` + reference — C side and script shared one refcount (the uloop.c pattern + exists for a reason) + - Handshake READ polling regression from the M2→M3 refactor (all connects timed out) + - `ws://host:port/` host off-by-one introduced by the IPv6 fix (caught by + re-running the v4 suite — regression tests matter) + +### M4 — Tests & docs ✅ (2026-09-01) +- [x] C fixture WS server: `tests/cram/fixtures/ws-fixture.c` — scripted scenarios + by request path (`/announce /echo /close /big /frag /flood /ping + /reset /badaccept /http200`), independent hand-rolled framing (no wslay + on the server side), dual-stack listener +- [x] Cram tests: `tests/cram/test_websocket.t` — T1, T2, T3+T4, T7, T8, T5, T6, + T11 (flood2000 no-loss), T9, badaccept, http200, T12 — all green +- [x] Test gating: fixture target + test file registered only under + `WEBSOCKET_SUPPORT`; skips gracefully in wslay-less CI builds +- [x] jsdoc: full `connect()` documentation (options, events, example); module + header covers usage patterns +- [x] ASan regression sweep after suite-driven fixes: announce + oversized clean +- Committed as `4543a6c` (history rewritten 2026-09-01: fix commits squashed into their milestone commits — c1bc024→M2, M4's module fixes→M3; each milestone verified to build; final tree byte-identical to pre-rewrite). Suite-driven bug fixes (all in same commit): + un-zeroed resource struct, pipelined handshake bytes lost, close reason + dropped, ECONNRESET unhandled, `on()` callbacks silently dropped (mangled + guard), whole options object serialized as HTTP headers + +### M5 — OpenWrt packaging & cross validation ✅ (2026-09-01) +- [x] Core `package/utils/ucode/Makefile`: `ucode-mod-websocket` added via the + `UcodeModule` macro (`+libubox +libwslay`) — clean upstream-ready change + on branch **`ucode-mod-websocket`** (commit `17829e4c4e`, worktree + `/tmp/opencode/ucode-mod-pr`, based on upstream master) +- [x] Local cross-build: buildroot `package/utils/ucode/Makefile` temporarily + overridden to `PKG_SOURCE_URL:=file:///home/nicolo/openwrt_ucode` pinned + to `4543a6cb` (stock file backed up as `Makefile.stock`) → **aarch64 + `websocket.so` + `ucode-mod-websocket.apk` + `libwslay.apk` all built** +- [x] **Runtime cross-architecture validation**: aarch64/musl ucode executed + under `qemu-aarch64` completed a full WebSocket session (handshake, + pipelined announce message, close handshake) against the x86 fixture — + see §8 for the exact invocation recipe +- [x] `.config`: `CONFIG_PACKAGE_libwslay=m`, `CONFIG_PACKAGE_ucode-mod-websocket=m` +- TLS moved to **M6**; server role remains deferred (non-goal until a use case) + +### M6 — TLS via mbedTLS (future, optional) +- [ ] `wss://` support behind the same API (mbedTLS, non-blocking handshake + integrated into the uloop state machine) +- [ ] Certificate verification policy options (ca_file, verify depth, insecure) +- [ ] Tests with an mbedTLS-based `wss` fixture endpoint + +--- + +## 7. Decisions Log + +| # | Date | Decision | Rationale | +|---|---|---|---| +| D1 | 2026-09-01 | Use wslay for RFC 6455 framing only; own transport on POSIX + uloop | Min memory; frozen spec = eternal; no hidden buffers/threads | +| D2 | 2026-09-01 | No permessage-deflate | Footprint goal on embedded targets | +| D3 | 2026-09-01 | Client role only in phase 1 | YAGNI; halves handshake/state-machine surface | +| D4 | 2026-09-01 | Fixed ring buffer, no per-message malloc | Deterministic memory; predictable RSS | +| D5 | 2026-09-01 | No auto-reconnect | Userland concern; keeps module stateless re: policy | +| D6 | 2026-09-01 | Vendor wslay in-tree, **statically linked into `websocket.so`** (like `ffi_lib` bundles its sources, `CMakeLists.txt:459`); optional `libwslay` feed package only in M5 | OpenWrt does **not** package wslay (verified in openwrt/packages master + snapshot indexes); OpenWrt ucode build compiles from this tree (`openwrt/ucode/Makefile:18`) so in-tree sources need no feed dependency. `ucode-mod-websocket` DEPENDS stays `ucode +libubox` | +| D7 | 2026-09-01 | Pin **wslay v1.1.1**; compile flags need no autotools `config.h` when building externally | Core sources compile warning-free under repo flags; autotools `configure` handles the `HAVE_*` detection when built as a proper library | +| D8 | 2026-09-01 | **Supersedes D6**: no in-tree vendoring — follow the **zlib pattern** instead: `find_library(wslay)` + `find_path(wslay/wslay.h)`, `WEBSOCKET_SUPPORT` auto-gates on discovery, module links external `libwslay`. OpenWrt side: `libwslay` package in the local buildroot (`../openwrt/package/wslay`, staged OK) + `ucode-mod-websocket` DEPENDS `+libwslay`. Local dev: wslay built into `/tmp/opencode/wslay-prefix`, pass `-DCMAKE_PREFIX_PATH` | User preference + matches ucode's zero-vendored-deps convention (`zlib_lib` links `ZLIB::ZLIB`, only the binding is in-tree) | +| D9 | 2026-09-01 | wslay goes to **core OpenWrt** (`package/libs/wslay`), not the packages feed | ucode is a core package (`package/utils/ucode`, source-pinned from jow-/ucode); core cannot depend on feeds; precedent: every existing ucode module dependency (libubox, libnl-tiny, libuci) is core. Feed alternative (standalone feed module, luci's ucode-mod-html pattern) kept as fallback if core review pushes back | +| D10 | 2026-09-01 | Receive memory model: wslay **default buffering** with `max_recv_msg_length` = `max_frame_size` option (default 256 KiB, 1 KiB–16 MiB). No custom ring buffer — wslay already owns fixed 4 KiB ibuf/obuf per context; our layer performs **zero per-message allocations** | Supersedes the "fixed ring buffer" design in §5: wslay's internal buffer is exactly the bounded, reused buffer we planned to write; duplicating it adds code without saving memory. Idle-connection fixed cost ≈ wslay ctx (~8.5 KiB) + handshake buffer (2 KiB) | +| D11 | 2026-09-01 | Event callbacks receive the connection resource as **explicit first argument** (`(ws, data, is_text)`), not via `this` | ucode arrow functions have no `this` binding; explicit args are the only portable form. Unhandled callback exceptions route through the shared `uloop.ex_handler` registry convention (falls back to `uloop_end()`) | +| D12 | 2026-09-01 | DNS resolution is synchronous (`getaddrinfo`) inside `connect()`; TCP connect + handshake + session are fully async | Matches `socket` module behavior (also sync connect); async DNS (uloop process or resolv-based) deferred until a real use case blocks on it | +| D13 | 2026-09-01 | Close completion semantics: when our close frame is **sent** and reads are **disabled** (wslay fatal recv condition), complete immediately with the close code we sent; only wait for the peer reply while reads remain enabled (bounded by 5 s `WS_CLOSE_TIMEOUT`) | wslay disables reads on oversize/protocol errors, so the peer's close reply can never be processed — waiting for it wedged connections for the full timeout | +| D14 | 2026-09-01 | Docs live in jsdoc headers (module + `connect()`), repo README untouched; test docs = cram file itself | ucode README describes the interpreter, not individual modules | +| D15 | 2026-09-01 | Fixture server implements framing independently (no wslay) — protocol bugs cannot hide behind a shared implementation; fixture + test file gated on `WEBSOCKET_SUPPORT` | Independent-implementation testing principle; keeps wslay-less CI green | +| D16 | 2026-09-01 | ucode grammar quirks to remember for tests: `typeof` is a prefix **operator** (`typeof(x)` in argument position misparses — assign to a variable first); heredoc-fed scripts (`ucode - < `max_frame_size`; expect close code 1009 | M3 | +| T7 | ping/pong | server pings; expect pong + no user-visible event | M2 | +| T8 | close-handshake | clean close from both sides; assert code+reason surface in `close` event | M2 | +| T9 | abrupt-reset | server RSTs connection; `error` event + no fd/uloop leak | M3 | +| T10 | gc-teardown | drop reference without `close()`; GC; assert no leak (ASan/LSan run) | M3 | +| T11 | backpressure | slow consumer callback; assert read event paused/resumed (fixture sends unboundedly) | M3 | +| T12 | uri-validation | bad schemes/ports/userinfo → clear exception, no fd created | M2 | + +Infrastructure notes: +- Fixture server: small C program in `tests/cram/fixtures/` (libevent-free, plain poll()) + able to run scripted scenarios: echo, fragment, oversized, reset, ping-flood. +- All tests must run under `ctest` via the existing cram harness (`tests/CMakeLists.txt`). +- Memory assertions: run selected tests under `valgrind --leak-check=full` in CI (x86_64 only). + +--- + +## 11. PR Preparation (2026-09-01) + +### PR 1 — wslay → core OpenWrt (`openwrt/openwrt` master) +- Branch: **`wslay-package`** (in the `../openwrt` fork repo, checked out in the + worktree `/tmp/opencode/wslay-pr`, based on current upstream master `9550b20e42`) +- Commit: `893873a4e8` — `package/libs/wslay: add wslay WebSocket library` (signed off) +- Push + open PR: + ``` + git -C ../openwrt push origin wslay-package + # then PR: hitech95:wslay-package -> openwrt:master + ``` +- Draft description: + > Adds wslay v1.1.1 (MIT), a non-IO WebSocket (RFC 6455) framing library, as + > `package/libs/wslay`. It performs no I/O itself — the caller drives the event + > loop through send/recv callbacks — making it suitable for memory-constrained, + > event-driven use. Needed by the upcoming `ucode-mod-websocket` module + > (companion PR in jow-/ucode, which requires this lib in core since ucode is + > a core package and core cannot depend on feeds). Static + shared libs built; + > InstallDev ships headers + pkg-config. + +### PR 1b — ucode-mod-websocket → core OpenWrt (`openwrt/openwrt` master) +- Branch: **`ucode-mod-websocket`** (worktree `/tmp/opencode/ucode-mod-pr`, + commit `17829e4c4e`): adds the module package via the `UcodeModule` macro + (`+libubox +libwslay`). **Hold** until the ucode-side PR is merged, then + combine with the `PKG_SOURCE_VERSION` bump in the same PR. +- Push: `git -C ../openwrt push origin ucode-mod-websocket` + +### PR 2 — websocket module → upstream ucode (`jow-/ucode` master) +- Branch: **`websocket-module`** (in this repo, `4543a6c` — M1–M4 complete, one clean commit per milestone: + implementation, hardening, cram suite with fixture server, jsdoc) +- M2–M4 are done → **ready to submit** (was: hold until M2–M4) +- Push (needs a jow-/ucode fork first) + open PR: + ``` + git remote add fork git@github.com:/ucode.git + git push -u fork websocket-module + # then PR: :websocket-module -> jow-:master + ``` +- Draft description: + > Adds a `websocket` module providing RFC 6455 client connectivity with an + > event-driven, uloop-integrated API (`connect()`, `on('open'|'message'| + > 'close'|'error')`, `send()`, `close()`). Framing is delegated to wslay + > (linked externally, zlib-style: `WEBSOCKET_SUPPORT` auto-gates on library + > discovery). Fixed receive ring buffer, frame-size caps and uloop-driven + > backpressure keep the memory footprint deterministic. Requires libwslay + > (core PR package/libs/wslay). +- Follow-up after merge: bump `PKG_SOURCE_VERSION` of core `package/utils/ucode` + and add `ucode-mod-websocket` there (sync with `openwrt/ucode/Makefile`). +- `WEBSOCKET_PLAN.md` intentionally untracked (internal tracking doc). + +--- + +## 12. References + +- RFC 6455 — The WebSocket Protocol +- wslay — https://github.com/tatsuhiro-t/wslay (MIT) +- Repo patterns: `lib/socket.c` (resources/errors), `lib/uloop.c:88-101` + (persistent resources + callbacks), `lib/digest.c` (SHA-1 primitives) +- Existing test harness: `tests/cram/`, `tests/CMakeLists.txt` + +### Commit history (post-rewrite 2026-09-01) +- `ce8ae99` M1 skeleton · `f66c9d8` M2 implementation (+IPv6, suite-class fixes) · `2b30672` M3 hardening (+test-uncovered fixes) · `4543a6c` M4 tests+fixture+jsdoc +- Backup of pre-rewrite history: `websocket-module-backup` diff --git a/lib/websocket.c b/lib/websocket.c new file mode 100644 index 00000000..902abfc5 --- /dev/null +++ b/lib/websocket.c @@ -0,0 +1,1745 @@ +/* + * Copyright (C) 2026 ucode contributors + * + * Permission to use, copy, modify, and/or distribute this software for any + * purpose with or without fee is hereby granted, provided that the above + * copyright notice and this permission notice appear in all copies. + * + * THE SOFTWARE IS PROVIDED "AS IS" AND THE AUTHOR DISCLAIMS ALL WARRANTIES + * WITH REGARD TO THIS SOFTWARE INCLUDING ALL IMPLIED WARRANTIES OF + * MERCHANTABILITY AND FITNESS. IN NO EVENT SHALL THE AUTHOR BE LIABLE FOR + * ANY SPECIAL, DIRECT, INDIRECT, OR CONSEQUENTIAL DAMAGES OR ANY DAMAGES + * WHATSOEVER RESULTING FROM LOSS OF USE, DATA OR PROFITS, WHETHER IN AN + * ACTION OF CONTRACT, NEGLIGENCE OR OTHER TORTIOUS ACTION, ARISING OUT OF + * OR IN CONNECTION WITH THE USE OR PERFORMANCE OF THIS SOFTWARE. + */ + +/** + * # WebSocket Module + * + * The `websocket` module provides functions for interacting with WebSocket + * (RFC 6455) servers using an event driven, uloop based API. + * + * Functions can be individually imported and directly accessed using the + * {@link https://developer.mozilla.org/en-US/docs/Web/JavaScript/Reference/Statements/import#named_import named import} + * syntax: + * + * ```javascript + * import { connect } from 'websocket'; + * import * as uloop from 'uloop'; + * + * uloop.init(); + * + * let ws = connect('ws://10.0.0.1:8080/path'); + * + * ws.on('message', (ws, data, is_text) => print(data)); + * + * ws.on('open', (ws) => ws.send('hello')); + * + * ws.on('close', (ws, code, reason) => print(`closed ${code} ${reason}\n`)); + * + * ws.on('error', (ws, error) => print(`error: ${error}\n`)); + * + * uloop.run(); + * ``` + * + * Alternatively, the module namespace can be imported + * using a wildcard import statement: + * + * ```javascript + * import * as websocket from 'websocket'; + * + * let ws = websocket.connect('ws://10.0.0.1:8080/path'); + * ``` + * + * Additionally, the websocket module namespace may also be imported by + * invoking the `ucode` interpreter with the `-lwebsocket` switch. + * + * @module websocket + */ + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#include +#include + +#include "ucode/module.h" + +#define ok_return(expr) do { last_error = 0; return (expr); } while(0) +#define err_return(err) do { last_error = err; return NULL; } while(0) + +#define WS_GUID "258EAFA5-E914-47DA-95CA-C5AB0DC85B11" +#define WS_HS_BUFSIZE 2048 +#define WS_DEFAULT_TIMEOUT 15000 +#define WS_DEFAULT_MAX_MSG (256 * 1024) +#define WS_MIN_TIMEOUT 100 +#define WS_CLOSE_TIMEOUT 5000 + +#ifndef MSG_MORE +#define MSG_MORE 0 +#endif + +#if defined(__APPLE__) +# define SOCK_NONBLOCK (1 << 16) +# define SOCK_CLOEXEC (1 << 17) +#endif + +static int last_error = 0; + +enum { + WS_STATE_CONNECTING, + WS_STATE_HANDSHAKE_WRITE, + WS_STATE_HANDSHAKE_READ, + WS_STATE_OPEN, + WS_STATE_CLOSING, + WS_STATE_CLOSED, +}; + +typedef struct { + uc_vm_t *vm; + uc_value_t *obj; + + struct uloop_fd ufd; + struct uloop_timeout timeout; + + wslay_event_context_ptr ctx; + + int state; + int port; + int io_errno; + bool eof; + int dispatching; + bool need_flush; + bool in_recv; + + char *host; + char *path; + + char *hs_req; + size_t hs_req_len; + size_t hs_sent; + + char hs_buf[WS_HS_BUFSIZE]; + size_t hs_len; + size_t hs_pendoff; + size_t hs_pendlen; + char hs_accept[32]; + + uint64_t max_msg_len; + + int close_code; + char close_reason[124]; +} uc_websocket_t; + +static void ws_flush(uc_websocket_t *ws); +static void ws_readable(uc_websocket_t *ws); + +/* ---------------------------------------------------------------------- */ +/* SHA-1 (RFC 3174) and base64 encoding for the opening handshake */ +/* ---------------------------------------------------------------------- */ + +typedef struct { + uint32_t state[5]; + uint64_t count; + uint8_t buffer[64]; +} sha1_ctx_t; + +#define SHA1_ROTL(x, n) (((x) << (n)) | ((x) >> (32 - (n)))) + +static void +sha1_init(sha1_ctx_t *ctx) +{ + ctx->state[0] = 0x67452301; + ctx->state[1] = 0xEFCDAB89; + ctx->state[2] = 0x98BADCFE; + ctx->state[3] = 0x10325476; + ctx->state[4] = 0xC3D2E1F0; + ctx->count = 0; +} + +static void +sha1_transform(sha1_ctx_t *ctx, const uint8_t *p) +{ + uint32_t w[80], a, b, c, d, e, t; + size_t i; + + for (i = 0; i < 16; i++) + w[i] = ((uint32_t)p[i * 4] << 24) | ((uint32_t)p[i * 4 + 1] << 16) | + ((uint32_t)p[i * 4 + 2] << 8) | (uint32_t)p[i * 4 + 3]; + + for (i = 16; i < 80; i++) + w[i] = SHA1_ROTL(w[i-3] ^ w[i-8] ^ w[i-14] ^ w[i-16], 1); + + a = ctx->state[0]; + b = ctx->state[1]; + c = ctx->state[2]; + d = ctx->state[3]; + e = ctx->state[4]; + + for (i = 0; i < 80; i++) { + if (i < 20) + t = ((b & c) | ((~b) & d)) + 0x5A827999; + else if (i < 40) + t = (b ^ c ^ d) + 0x6ED9EBA1; + else if (i < 60) + t = ((b & c) | (b & d) | (c & d)) + 0x8F1BBCDC; + else + t = (b ^ c ^ d) + 0xCA62C1D6; + + t += SHA1_ROTL(a, 5) + e + w[i]; + e = d; + d = c; + c = SHA1_ROTL(b, 30); + b = a; + a = t; + } + + ctx->state[0] += a; + ctx->state[1] += b; + ctx->state[2] += c; + ctx->state[3] += d; + ctx->state[4] += e; +} + +static void +sha1_update(sha1_ctx_t *ctx, const uint8_t *data, size_t len) +{ + size_t i = 0, n; + + if (ctx->count % 64) { + n = 64 - (ctx->count % 64); + + if (n > len) + n = len; + + memcpy(ctx->buffer + (ctx->count % 64), data, n); + ctx->count += n; + i = n; + + if (ctx->count % 64) + return; + + sha1_transform(ctx, ctx->buffer); + } + + for (; i + 64 <= len; i += 64) { + sha1_transform(ctx, data + i); + ctx->count += 64; + } + + n = len - i; + + if (n) { + memcpy(ctx->buffer, data + i, n); + ctx->count += n; + } +} + +static void +sha1_final(sha1_ctx_t *ctx, uint8_t digest[20]) +{ + static const uint8_t pad[64] = { 0x80 }; + uint64_t bits = ctx->count * 8; + uint8_t tail[8]; + size_t i; + + for (i = 0; i < 8; i++) + tail[i] = (uint8_t)(bits >> (56 - 8 * i)); + + sha1_update(ctx, pad, 1 + ((119 - ctx->count % 64) % 64)); + sha1_update(ctx, tail, 8); + + for (i = 0; i < 5; i++) { + digest[i * 4] = (uint8_t)(ctx->state[i] >> 24); + digest[i * 4 + 1] = (uint8_t)(ctx->state[i] >> 16); + digest[i * 4 + 2] = (uint8_t)(ctx->state[i] >> 8); + digest[i * 4 + 3] = (uint8_t)(ctx->state[i]); + } +} + +static const char b64tab[] = + "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/"; + +static size_t +b64_encode(const uint8_t *in, size_t inlen, char *out, size_t outsize) +{ + size_t i, j = 0; + + for (i = 0; i + 3 <= inlen; i += 3) { + uint32_t v = ((uint32_t)in[i] << 16) | + ((uint32_t)in[i + 1] << 8) | (uint32_t)in[i + 2]; + + out[j++] = b64tab[(v >> 18) & 63]; + out[j++] = b64tab[(v >> 12) & 63]; + out[j++] = b64tab[(v >> 6) & 63]; + out[j++] = b64tab[v & 63]; + } + + if (i < inlen) { + uint32_t v = (uint32_t)in[i] << 16; + + if (i + 1 < inlen) + v |= (uint32_t)in[i + 1] << 8; + + out[j++] = b64tab[(v >> 18) & 63]; + out[j++] = b64tab[(v >> 12) & 63]; + out[j++] = (i + 1 < inlen) ? b64tab[(v >> 6) & 63] : '='; + out[j++] = '='; + i += 3; + } + + (void)i; + out[j] = '\0'; + + return j; +} + +static int +ws_urandom(uint8_t *buf, size_t len) +{ + static int fd = -1; + ssize_t n; + + if (fd == -1) { + fd = open("/dev/urandom", O_RDONLY | O_CLOEXEC); + + if (fd == -1) + return -1; + } + + do { + n = read(fd, buf, len); + } while (n < 0 && errno == EINTR); + + return (n == (ssize_t)len) ? 0 : -1; +} + +/* ---------------------------------------------------------------------- */ +/* URL parsing */ +/* ---------------------------------------------------------------------- */ + +static bool +ws_parse_url(const char *url, char **host, int *port, char **path, bool *tls) +{ + const char *p = url, *host_start, *host_end, *port_start = NULL; + const char *bracket_close = NULL; + char hostbuf[256], *end; + size_t hostlen; + long n; + + if (!strncasecmp(p, "ws://", 5)) { + host_start = p + 5; + *tls = false; + } + else if (!strncasecmp(p, "wss://", 6)) { + host_start = p + 6; + *tls = true; + } + else { + return false; + } + + if (*host_start == '[') { + bracket_close = strchr(host_start, ']'); + + if (!bracket_close || bracket_close == host_start + 1) + return false; + + host_end = bracket_close + 1; + + if (*host_end == ':') { + port_start = host_end + 1; + } + else if (*host_end && *host_end != '/' && *host_end != '?' && *host_end != '#') { + return false; + } + + /* store the IPv6 literal without brackets for getaddrinfo() */ + hostlen = (size_t)(bracket_close - host_start - 1); + + if (hostlen >= sizeof(hostbuf)) + return false; + + memcpy(hostbuf, host_start + 1, hostlen); + hostbuf[hostlen] = '\0'; + } + else { + for (host_end = host_start; *host_end; host_end++) { + if (*host_end == '/' || *host_end == '?' || *host_end == '#') + break; + + if (*host_end == '@') + return false; + + if (*host_end == ':') { + if (port_start) + return false; + + port_start = host_end + 1; + } + } + + hostlen = (size_t)((port_start ? port_start - 1 : host_end) - host_start); + + if (!hostlen || hostlen >= sizeof(hostbuf)) + return false; + + memcpy(hostbuf, host_start, hostlen); + hostbuf[hostlen] = '\0'; + } + + if (port_start) { + n = strtol(port_start, &end, 10); + + if (!*port_start || n < 1 || n > 65535) + return false; + + if (!bracket_close && end != host_end) + return false; + + if (bracket_close) { + if (*end && *end != '/' && *end != '?' && *end != '#') + return false; + + host_end = end; + } + + *port = (int)n; + } + else { + *port = *tls ? 443 : 80; + } + + *host = strdup(hostbuf); + + if (*host_end == '/') + *path = strdup(host_end); + else if (*host_end) + *path = NULL; /* ?query or #fragment: prepend slash below */ + else + *path = strdup("/"); + + if (!*path && *host_end) { + *path = malloc(strlen(host_end) + 2); + + if (*path) + snprintf(*path, strlen(host_end) + 2, "/%s", host_end); + } + + return (*host && *path); +} + +/* ---------------------------------------------------------------------- */ +/* Event dispatching */ +/* ---------------------------------------------------------------------- */ + +static bool +ws_vm_call(uc_vm_t *vm, bool mcall, size_t nargs) +{ + uc_value_t *exh, *val; + + if (uc_vm_call(vm, mcall, nargs) == EXCEPTION_NONE) + return true; + + exh = uc_vm_registry_get(vm, "uloop.ex_handler"); + + if (!ucv_is_callable(exh)) + goto error; + + val = uc_vm_exception_object(vm); + uc_vm_stack_push(vm, ucv_get(exh)); + uc_vm_stack_push(vm, val); + + if (uc_vm_call(vm, false, 1) != EXCEPTION_NONE) + goto error; + + ucv_put(uc_vm_stack_pop(vm)); + + return false; + +error: + uloop_end(); + + return false; +} + +static void +ws_invoke(uc_websocket_t *ws, size_t slot, uc_value_t **args, size_t nargs) +{ + uc_value_t *fn; + size_t i; + + if (!ws->obj || ws->state == WS_STATE_CLOSED) + return; + + fn = ucv_resource_value_get(ws->obj, slot); + + if (!ucv_is_callable(fn)) + return; + + uc_vm_stack_push(ws->vm, ucv_get(ws->obj)); + uc_vm_stack_push(ws->vm, ucv_get(fn)); + uc_vm_stack_push(ws->vm, ucv_get(ws->obj)); + + for (i = 0; i < nargs; i++) + uc_vm_stack_push(ws->vm, ucv_get(args[i])); + + ws->dispatching++; + + if (ws_vm_call(ws->vm, true, nargs + 1)) + ucv_put(uc_vm_stack_pop(ws->vm)); + + ws->dispatching--; + + /* defer wslay_event_send() until we are completely outside of the + * wslay_event_recv() call stack: the epilogue below still runs while + * wslay_event_recv() is active, so an ongoing receive is another + * reason to keep the flush pending */ + if (!ws->dispatching && !ws->in_recv && ws->need_flush) { + ws->need_flush = false; + + if (ws->state == WS_STATE_OPEN || ws->state == WS_STATE_CLOSING) + ws_flush(ws); + } +} + +static void +ws_emit_error(uc_websocket_t *ws, const char *msg) +{ + uc_value_t *args[1] = { ucv_string_new(msg) }; + + ws_invoke(ws, 3, args, 1); + ucv_put(args[0]); +} + +static void +ws_emit_close(uc_websocket_t *ws, int code, const char *reason) +{ + uc_value_t *args[2]; + + args[0] = ucv_int64_new(code); + args[1] = ucv_string_new(reason ? reason : ""); + + ws_invoke(ws, 2, args, 2); + ucv_put(args[0]); + ucv_put(args[1]); +} + +/* ---------------------------------------------------------------------- */ +/* Teardown */ +/* ---------------------------------------------------------------------- */ + +static void +ws_teardown(uc_websocket_t *ws) +{ + uc_value_t *obj; + uc_resource_ext_t *ext; + size_t i; + + if (!ws || ws->state == WS_STATE_CLOSED) + return; + + ws->state = WS_STATE_CLOSED; + + uloop_fd_delete(&ws->ufd); + uloop_timeout_cancel(&ws->timeout); + + if (ws->ctx) { + wslay_event_context_free(ws->ctx); + ws->ctx = NULL; + } + + if (ws->ufd.fd >= 0) { + close(ws->ufd.fd); + ws->ufd.fd = -1; + } + + free(ws->hs_req); + free(ws->host); + free(ws->path); + + obj = ws->obj; + ws->obj = NULL; + + ext = (uc_resource_ext_t *)obj; + + for (i = 0; i < ext->uvcount; i++) + ucv_resource_value_set(obj, i, NULL); + + ucv_resource_persistent_set(obj, false); + ucv_put(obj); +} + +/* ---------------------------------------------------------------------- */ +/* wslay callbacks */ +/* ---------------------------------------------------------------------- */ + +static ssize_t +ws_wslay_recv(wslay_event_context_ptr ctx, uint8_t *buf, size_t len, + int flags, void *user_data) +{ + uc_websocket_t *ws = user_data; + ssize_t n; + size_t chunk; + + (void)flags; + + /* serve bytes that arrived pipelined with the handshake response + * before reading new data from the socket */ + if (ws->hs_pendlen) { + chunk = (ws->hs_pendlen < len) ? ws->hs_pendlen : len; + + memcpy(buf, ws->hs_buf + ws->hs_pendoff, chunk); + + ws->hs_pendoff += chunk; + ws->hs_pendlen -= chunk; + + return (ssize_t)chunk; + } + + do { + n = read(ws->ufd.fd, buf, len); + } while (n < 0 && errno == EINTR); + + if (n > 0) + return n; + + if (n == 0) + ws->eof = true; + else if (errno != EAGAIN && errno != EWOULDBLOCK) + ws->io_errno = errno; + + wslay_event_set_error(ctx, + (n < 0 && (errno == EAGAIN || errno == EWOULDBLOCK)) + ? WSLAY_ERR_WOULDBLOCK : WSLAY_ERR_CALLBACK_FAILURE); + + return -1; +} + +static ssize_t +ws_wslay_send(wslay_event_context_ptr ctx, const uint8_t *buf, size_t len, + int flags, void *user_data) +{ + uc_websocket_t *ws = user_data; + ssize_t n; + + (void)ctx; + + do { + n = send(ws->ufd.fd, buf, len, + (flags & WSLAY_MSG_MORE) ? MSG_MORE : 0); + } while (n < 0 && errno == EINTR); + + if (n >= 0) + return n; + + if (errno != EAGAIN && errno != EWOULDBLOCK) + ws->io_errno = errno; + + wslay_event_set_error(ctx, + (errno == EAGAIN || errno == EWOULDBLOCK) + ? WSLAY_ERR_WOULDBLOCK : WSLAY_ERR_CALLBACK_FAILURE); + + return -1; +} + +static int +ws_wslay_genmask(wslay_event_context_ptr ctx, uint8_t *buf, size_t len, + void *user_data) +{ + (void)ctx; + (void)user_data; + + return ws_urandom(buf, len); +} + +static void +ws_wslay_on_msg(wslay_event_context_ptr ctx, + const struct wslay_event_on_msg_recv_arg *arg, void *user_data) +{ + uc_websocket_t *ws = user_data; + uc_value_t *args[2]; + + (void)ctx; + + switch (arg->opcode) { + case WSLAY_CONNECTION_CLOSE: + /* reply close frame is queued automatically by wslay_event_recv() */ + ws->close_code = (int)arg->status_code; + + if (arg->msg_length > 2) { + size_t n = arg->msg_length - 2; + + if (n > sizeof(ws->close_reason) - 1) + n = sizeof(ws->close_reason) - 1; + + memcpy(ws->close_reason, arg->msg + 2, n); + ws->close_reason[n] = '\0'; + } + + break; + + case WSLAY_TEXT_FRAME: + case WSLAY_BINARY_FRAME: + args[0] = ucv_string_new_length((const char *)arg->msg, arg->msg_length); + args[1] = ucv_boolean_new(arg->opcode == WSLAY_TEXT_FRAME); + + ws_invoke(ws, 1, args, 2); + + ucv_put(args[0]); + ucv_put(args[1]); + break; + + default: + /* ping frames are answered automatically by wslay_event_recv() */ + break; + } +} + +static const struct wslay_event_callbacks ws_wslay_callbacks = { + .recv_callback = ws_wslay_recv, + .send_callback = ws_wslay_send, + .genmask_callback = ws_wslay_genmask, + .on_frame_recv_start_callback = NULL, + .on_frame_recv_chunk_callback = NULL, + .on_frame_recv_end_callback = NULL, + .on_msg_recv_callback = ws_wslay_on_msg, +}; + +/* ---------------------------------------------------------------------- */ +/* Connection lifecycle */ +/* ---------------------------------------------------------------------- */ + +static void +ws_update_poll(uc_websocket_t *ws) +{ + unsigned int flags = 0; + + if (ws->state == WS_STATE_CONNECTING || ws->state == WS_STATE_HANDSHAKE_WRITE) { + flags = ULOOP_WRITE; + } + else if (ws->state == WS_STATE_HANDSHAKE_READ) { + flags = ULOOP_READ; + } + else if (ws->ctx) { + if (wslay_event_get_read_enabled(ws->ctx)) + flags |= ULOOP_READ; + + if (wslay_event_want_write(ws->ctx)) + flags |= ULOOP_WRITE; + } + + if (flags) + uloop_fd_add(&ws->ufd, flags); + else + uloop_fd_delete(&ws->ufd); +} + +static void +ws_check_lifecycle(uc_websocket_t *ws) +{ + if (ws->state != WS_STATE_OPEN && ws->state != WS_STATE_CLOSING) + return; + + /* entering the closing phase: either the peer started the close + * handshake, or wslay disabled reads after queueing an automatic + * close reply (oversized message, protocol error, ...) */ + if (ws->state == WS_STATE_OPEN && ws->ctx && + (!wslay_event_get_read_enabled(ws->ctx) || + wslay_event_get_close_received(ws->ctx))) { + ws->state = WS_STATE_CLOSING; + uloop_timeout_set(&ws->timeout, WS_CLOSE_TIMEOUT); + } + + if ((ws->eof || ws->io_errno) && !wslay_event_get_close_received(ws->ctx)) { + ws_emit_error(ws, ws->io_errno + ? strerror(ws->io_errno) : "connection closed unexpectedly"); + ws_emit_close(ws, 1006, ""); + ws_teardown(ws); + return; + } + + /* once our close frame is sent and no further reads are possible + * (wslay disabled reads after a protocol violation or oversized + * message), the close handshake cannot progress any further */ + if (ws->ctx && wslay_event_get_close_sent(ws->ctx) && + !wslay_event_get_read_enabled(ws->ctx)) { + int code = ws->close_code; + const char *reason = ws->close_reason; + + if (!wslay_event_get_close_received(ws->ctx)) { + code = (int)wslay_event_get_status_code_sent(ws->ctx); + + if (code < 1000 || code > 4999) + code = 1006; + + reason = ""; + } + + ws_emit_close(ws, code, reason); + ws_teardown(ws); + return; + } + + if (wslay_event_get_close_received(ws->ctx) && wslay_event_get_close_sent(ws->ctx)) { + ws_emit_close(ws, ws->close_code, ws->close_reason); + ws_teardown(ws); + } +} + +static void +ws_flush(uc_websocket_t *ws) +{ + int rv; + + if (!ws->ctx) + return; + + rv = wslay_event_send(ws->ctx); + + if (rv < 0) { + ws_emit_error(ws, rv == WSLAY_ERR_NOMEM + ? "out of memory" : "send failure"); + ws_emit_close(ws, 1006, ""); + ws_teardown(ws); + return; + } + + if (ws->state == WS_STATE_OPEN || ws->state == WS_STATE_CLOSING) + ws_update_poll(ws); +} + +static bool +ws_validate_handshake(uc_websocket_t *ws) +{ + const char *eol, *end, *p = ws->hs_buf; + bool have_upgrade = false, have_connection = false, have_accept = false; + char line[256]; + size_t len; + + /* data may be pipelined behind the handshake response, so split + * at the header terminator instead of requiring it at the end */ + end = strstr(ws->hs_buf, "\r\n\r\n"); + + if (!end || ws->hs_len < 4 || + (size_t)(end + 4 - ws->hs_buf) > ws->hs_len) + return false; + + ws->hs_pendoff = (size_t)(end + 4 - ws->hs_buf); + ws->hs_pendlen = ws->hs_len - ws->hs_pendoff; + + if (strncmp(p, "HTTP/1.1 101", 12) && strncmp(p, "HTTP/1.0 101", 12)) + return false; + + while (*p && p < end) { + eol = strstr(p, "\r\n"); + + if (!eol || eol > end) + break; + + len = (size_t)(eol - p); + + if (len >= sizeof(line)) + return false; + + memcpy(line, p, len); + line[len] = '\0'; + + if (!strncasecmp(line, "Upgrade:", 8)) { + if (!strcasestr(line, "websocket")) + return false; + + have_upgrade = true; + } + else if (!strncasecmp(line, "Connection:", 11)) { + if (!strcasestr(line, "upgrade")) + return false; + + have_connection = true; + } + else if (!strncasecmp(line, "Sec-WebSocket-Accept:", 21)) { + p = line + 21; + + while (*p == ' ' || *p == '\t') + p++; + + if (strncmp(p, ws->hs_accept, strlen(ws->hs_accept))) + return false; + + have_accept = true; + } + + p = eol + 2; + } + + return have_upgrade && have_connection && have_accept; +} + +static void +ws_established(uc_websocket_t *ws) +{ + if (wslay_event_context_client_init(&ws->ctx, &ws_wslay_callbacks, ws)) { + ws_emit_error(ws, "out of memory"); + ws_teardown(ws); + return; + } + + wslay_event_config_set_max_recv_msg_length(ws->ctx, ws->max_msg_len); + + ws->state = WS_STATE_OPEN; + + uloop_timeout_cancel(&ws->timeout); + + ws_invoke(ws, 0, NULL, 0); + + if (ws->state == WS_STATE_OPEN || ws->state == WS_STATE_CLOSING) + ws_update_poll(ws); + + /* data pipelined behind the handshake response is already waiting + * in hs_buf; no further read event would ever fire, so drain it + * through wslay right away */ + if (ws->hs_pendlen && ws->state == WS_STATE_OPEN) + ws_readable(ws); +} + +static void +ws_handshake_read(uc_websocket_t *ws) +{ + ssize_t n; + + do { + n = read(ws->ufd.fd, ws->hs_buf + ws->hs_len, + sizeof(ws->hs_buf) - ws->hs_len - 1); + } while (n < 0 && errno == EINTR); + + if (n == 0) { + ws_emit_error(ws, "connection closed during handshake"); + ws_teardown(ws); + return; + } + + if (n < 0) { + if (errno != EAGAIN && errno != EWOULDBLOCK) { + ws_emit_error(ws, strerror(errno)); + ws_teardown(ws); + } + + return; + } + + ws->hs_len += n; + ws->hs_buf[ws->hs_len] = '\0'; + + if (!strstr(ws->hs_buf, "\r\n\r\n")) { + if (ws->hs_len + 1 >= sizeof(ws->hs_buf)) { + ws_emit_error(ws, "handshake response too large"); + ws_teardown(ws); + } + + return; + } + + if (!ws_validate_handshake(ws)) { + ws_emit_error(ws, "invalid handshake response"); + ws_teardown(ws); + return; + } + + ws_established(ws); +} + +static void +ws_handshake_write(uc_websocket_t *ws) +{ + ssize_t n; + + while (ws->hs_sent < ws->hs_req_len) { + do { + n = write(ws->ufd.fd, ws->hs_req + ws->hs_sent, + ws->hs_req_len - ws->hs_sent); + } while (n < 0 && errno == EINTR); + + if (n < 0) { + if (errno != EAGAIN && errno != EWOULDBLOCK) { + ws_emit_error(ws, strerror(errno)); + ws_teardown(ws); + } + + return; + } + + ws->hs_sent += n; + } + + ws->state = WS_STATE_HANDSHAKE_READ; + + ws_update_poll(ws); +} + +static void +ws_connected(uc_websocket_t *ws) +{ + int err = 0; + socklen_t errlen = sizeof(err); + + if (getsockopt(ws->ufd.fd, SOL_SOCKET, SO_ERROR, &err, &errlen) || err) { + ws_emit_error(ws, err ? strerror(err) : "connection failed"); + ws_teardown(ws); + return; + } + + ws->state = WS_STATE_HANDSHAKE_WRITE; + + ws_handshake_write(ws); +} + +static void +ws_readable(uc_websocket_t *ws) +{ + int rv; + + if (ws->state == WS_STATE_HANDSHAKE_READ) { + ws_handshake_read(ws); + return; + } + + ws->in_recv = true; + rv = wslay_event_recv(ws->ctx); + ws->in_recv = false; + + if (rv < 0 && !ws->eof && !ws->io_errno) { + ws_emit_error(ws, rv == WSLAY_ERR_NOMEM + ? "out of memory" : "receive failure"); + ws_teardown(ws); + return; + } + + /* flush data queued from within message callbacks now that the + * receive operation completed; a fatal send error tears the + * connection down safely here */ + if (ws->need_flush && !ws->dispatching) { + ws->need_flush = false; + + if (ws->state == WS_STATE_OPEN || ws->state == WS_STATE_CLOSING) + ws_flush(ws); + } + + ws_check_lifecycle(ws); + + if (ws->state == WS_STATE_OPEN || ws->state == WS_STATE_CLOSING) + ws_update_poll(ws); +} + +static void +ws_ufd_cb(struct uloop_fd *fd, unsigned int events) +{ + uc_websocket_t *ws = container_of(fd, uc_websocket_t, ufd); + + if (ws->state == WS_STATE_CLOSED) + return; + + if (events & ULOOP_WRITE) { + switch (ws->state) { + case WS_STATE_CONNECTING: + ws_connected(ws); + return; + + case WS_STATE_HANDSHAKE_WRITE: + ws_handshake_write(ws); + return; + + default: + ws_flush(ws); + ws_check_lifecycle(ws); + return; + } + } + + if (events & ULOOP_READ) + ws_readable(ws); +} + +static void +ws_timeout_cb(struct uloop_timeout *timeout) +{ + uc_websocket_t *ws = container_of(timeout, uc_websocket_t, timeout); + + if (ws->state == WS_STATE_CLOSING && ws->ctx) { + ws_emit_error(ws, "close handshake timeout"); + ws_emit_close(ws, wslay_event_get_close_received(ws->ctx) + ? ws->close_code : 1006, + wslay_event_get_close_received(ws->ctx) + ? ws->close_reason : ""); + } + else { + ws_emit_error(ws, "connection or handshake timeout"); + ws_emit_close(ws, 1006, ""); + } + + ws_teardown(ws); +} + +/* ---------------------------------------------------------------------- */ +/* Handshake request generation */ +/* ---------------------------------------------------------------------- */ + +static char * +ws_build_request(uc_vm_t *vm, uc_websocket_t *ws, uc_value_t *headers) +{ + char req[WS_HS_BUFSIZE], key[32], accept_src[96], *hvs; + uint8_t nonce[16], digest[20]; + sha1_ctx_t sha1; + size_t len, n; + bool ok = true; + + if (ws_urandom(nonce, sizeof(nonce))) + return NULL; + + b64_encode(nonce, sizeof(nonce), key, sizeof(key)); + + snprintf(accept_src, sizeof(accept_src), "%s%s", key, WS_GUID); + + sha1_init(&sha1); + sha1_update(&sha1, (const uint8_t *)accept_src, strlen(accept_src)); + sha1_final(&sha1, digest); + + b64_encode(digest, sizeof(digest), ws->hs_accept, sizeof(ws->hs_accept)); + + if (strchr(ws->host, ':')) { + /* IPv6 literal: brackets are required in the Host header */ + if (ws->port == 80) + len = snprintf(req, sizeof(req), + "GET %s HTTP/1.1\r\n" + "Host: [%s]\r\n" + "Upgrade: websocket\r\n" + "Connection: Upgrade\r\n" + "Sec-WebSocket-Key: %s\r\n" + "Sec-WebSocket-Version: 13\r\n", + ws->path, ws->host, key); + else + len = snprintf(req, sizeof(req), + "GET %s HTTP/1.1\r\n" + "Host: [%s]:%d\r\n" + "Upgrade: websocket\r\n" + "Connection: Upgrade\r\n" + "Sec-WebSocket-Key: %s\r\n" + "Sec-WebSocket-Version: 13\r\n", + ws->path, ws->host, ws->port, key); + } + else if (ws->port == 80) + len = snprintf(req, sizeof(req), + "GET %s HTTP/1.1\r\n" + "Host: %s\r\n" + "Upgrade: websocket\r\n" + "Connection: Upgrade\r\n" + "Sec-WebSocket-Key: %s\r\n" + "Sec-WebSocket-Version: 13\r\n", + ws->path, ws->host, key); + else + len = snprintf(req, sizeof(req), + "GET %s HTTP/1.1\r\n" + "Host: %s:%d\r\n" + "Upgrade: websocket\r\n" + "Connection: Upgrade\r\n" + "Sec-WebSocket-Key: %s\r\n" + "Sec-WebSocket-Version: 13\r\n", + ws->path, ws->host, ws->port, key); + + if (len >= sizeof(req)) + return NULL; + + if (headers && ucv_type(headers) == UC_OBJECT) { + ucv_object_foreach(headers, hk, hval) { + if (!hk || strchr(hk, '\r') || strchr(hk, '\n') || strchr(hk, ':')) { + ok = false; + break; + } + + hvs = ucv_to_string(vm, hval); + + if (!hvs || strchr(hvs, '\r') || strchr(hvs, '\n')) { + free(hvs); + ok = false; + break; + } + + n = snprintf(req + len, sizeof(req) - len, "%s: %s\r\n", hk, hvs); + + free(hvs); + + if (n >= sizeof(req) - len - 2) { + ok = false; + break; + } + + len += n; + } + } + + if (!ok) + return NULL; + + len += snprintf(req + len, sizeof(req) - len, "\r\n"); + + return strdup(req); +} + +/* ---------------------------------------------------------------------- */ +/* Resource methods */ +/* ---------------------------------------------------------------------- */ + +static int +ws_event_slot(const char *name) +{ + if (!strcmp(name, "open")) return 0; + if (!strcmp(name, "message")) return 1; + if (!strcmp(name, "close")) return 2; + if (!strcmp(name, "error")) return 3; + + return -1; +} + +static uc_value_t * +uc_ws_on(uc_vm_t *vm, size_t nargs) +{ + uc_websocket_t *ws = uc_fn_thisval("websocket.connection"); + uc_value_t *event = uc_fn_arg(0); + uc_value_t *fn = uc_fn_arg(1); + const char *name; + int slot; + + if (!ws) + err_return(EINVAL); + + if (!ws->obj || ws->state == WS_STATE_CLOSED) + err_return(ENOTCONN); + + name = ucv_string_get(event); + + if (!name || (slot = ws_event_slot(name)) < 0) { + uc_vm_raise_exception(vm, EXCEPTION_TYPE, + "Event must be one of 'open', 'message', 'close' or 'error'"); + + return NULL; + } + + if (!ucv_is_callable(fn)) { + uc_vm_raise_exception(vm, EXCEPTION_TYPE, "Callback must be a function"); + + return NULL; + } + + ucv_resource_value_set(ws->obj, slot, ucv_get(fn)); + + ok_return(ucv_boolean_new(true)); +} + +static uc_value_t * +uc_ws_send(uc_vm_t *vm, size_t nargs) +{ + uc_websocket_t *ws = uc_fn_thisval("websocket.connection"); + uc_value_t *data = uc_fn_arg(0); + struct wslay_event_msg msg; + uint8_t *bin = NULL; + const uint8_t *payload; + size_t len, i; + int64_t b; + int rv; + + if (!ws) + err_return(EINVAL); + + if (ws->state != WS_STATE_OPEN) + err_return(ENOTCONN); + + if (ucv_type(data) == UC_STRING) { + payload = (const uint8_t *)ucv_string_get(data); + len = ucv_string_length(data); + msg.opcode = WSLAY_TEXT_FRAME; + } + else if (ucv_type(data) == UC_ARRAY) { + len = ucv_array_length(data); + bin = malloc(len ? len : 1); + + if (!bin) + err_return(ENOMEM); + + for (i = 0; i < len; i++) { + b = ucv_int64_get(ucv_array_get(data, i)); + + if (errno || b < 0 || b > 255) { + free(bin); + err_return(EINVAL); + } + + bin[i] = (uint8_t)b; + } + + payload = bin; + msg.opcode = WSLAY_BINARY_FRAME; + } + else { + uc_vm_raise_exception(vm, EXCEPTION_TYPE, + "Data must be a string or an array of bytes"); + + return NULL; + } + + msg.msg = payload; + msg.msg_length = len; + + rv = wslay_event_queue_msg(ws->ctx, &msg); + + free(bin); + + if (rv < 0) + err_return(rv == WSLAY_ERR_NO_MORE_MSG ? EPIPE : ENOMEM); + + if (ws->dispatching) + ws->need_flush = true; + else + ws_flush(ws); + + ok_return(ucv_boolean_new(true)); +} + +static uc_value_t * +uc_ws_ping(uc_vm_t *vm, size_t nargs) +{ + uc_websocket_t *ws = uc_fn_thisval("websocket.connection"); + uc_value_t *data = uc_fn_arg(0); + struct wslay_event_msg msg; + const char *payload = ""; + size_t len = 0; + int rv; + + if (!ws) + err_return(EINVAL); + + if (ws->state != WS_STATE_OPEN) + err_return(ENOTCONN); + + if (data && ucv_type(data) == UC_STRING) { + payload = ucv_string_get(data); + len = ucv_string_length(data); + + if (len > 125) { + uc_vm_raise_exception(vm, EXCEPTION_TYPE, + "Ping payload must not exceed 125 bytes"); + + return NULL; + } + } + + msg.opcode = WSLAY_PING; + msg.msg = (const uint8_t *)payload; + msg.msg_length = len; + + rv = wslay_event_queue_msg(ws->ctx, &msg); + + if (rv < 0) + err_return(rv == WSLAY_ERR_NO_MORE_MSG ? EPIPE : ENOMEM); + + if (ws->dispatching) + ws->need_flush = true; + else + ws_flush(ws); + + ok_return(ucv_boolean_new(true)); +} + +static uc_value_t * +uc_ws_close(uc_vm_t *vm, size_t nargs) +{ + uc_websocket_t *ws = uc_fn_thisval("websocket.connection"); + uc_value_t *code = uc_fn_arg(0); + uc_value_t *reason = uc_fn_arg(1); + const char *r = ""; + uint16_t c = 1000; + int64_t n = 1000; + int rv; + + if (!ws) + err_return(EINVAL); + + if (ws->state == WS_STATE_CLOSED) + err_return(ENOTCONN); + + /* close() while the connection is still being established aborts it */ + if (ws->state != WS_STATE_OPEN && ws->state != WS_STATE_CLOSING) { + ws_emit_close(ws, 1006, ""); + ws_teardown(ws); + ok_return(ucv_boolean_new(true)); + } + + /* subsequent close() calls while already closing are no-ops */ + if (ws->state == WS_STATE_CLOSING) + ok_return(ucv_boolean_new(true)); + + if (code) { + errno = 0; + n = ucv_int64_get(code); + + if (errno || n < 1000 || n > 4999) { + uc_vm_raise_exception(vm, EXCEPTION_TYPE, + "Close code must be between 1000 and 4999"); + + return NULL; + } + } + + if (reason && ucv_type(reason) == UC_STRING) { + r = ucv_string_get(reason); + + if (strlen(r) > 123 || strchr(r, '\r') || strchr(r, '\n')) { + uc_vm_raise_exception(vm, EXCEPTION_TYPE, + "Close reason must not exceed 123 bytes or contain line breaks"); + + return NULL; + } + } + + c = (uint16_t)n; + + ws->state = WS_STATE_CLOSING; + + uloop_timeout_set(&ws->timeout, WS_CLOSE_TIMEOUT); + + rv = wslay_event_queue_close(ws->ctx, c, (const uint8_t *)r, strlen(r)); + + if (rv < 0) + err_return(EPIPE); + + if (ws->dispatching) + ws->need_flush = true; + else { + ws_flush(ws); + ws_check_lifecycle(ws); + } + + ok_return(ucv_boolean_new(true)); +} + +static uc_value_t * +uc_ws_state(uc_vm_t *vm, size_t nargs) +{ + uc_websocket_t *ws = uc_fn_thisval("websocket.connection"); + const char *state; + + if (!ws) + err_return(EINVAL); + + switch (ws->state) { + case WS_STATE_CONNECTING: + case WS_STATE_HANDSHAKE_WRITE: + case WS_STATE_HANDSHAKE_READ: + state = "connecting"; + break; + + case WS_STATE_OPEN: + state = "open"; + break; + + case WS_STATE_CLOSING: + state = "closing"; + break; + + default: + state = "closed"; + break; + } + + ok_return(ucv_string_new(state)); +} + +static uc_value_t * +uc_ws_fileno(uc_vm_t *vm, size_t nargs) +{ + uc_websocket_t *ws = uc_fn_thisval("websocket.connection"); + + if (!ws) + err_return(EINVAL); + + ok_return(ucv_int64_new(ws->ufd.fd)); +} + +/* ---------------------------------------------------------------------- */ +/* connect() */ +/* ---------------------------------------------------------------------- */ + +/** + * Initiate a WebSocket connection to the given URL. + * + * The URL must use the `ws://` scheme (`wss://` requires TLS support which + * is not implemented yet). IPv6 literals may be given in bracket notation, + * e.g. `ws://[::1]:8080/path`. Userinfo components are rejected. + * + * Name resolution is performed synchronously; the TCP connection, the + * protocol handshake and the subsequent data exchange are fully + * asynchronous and driven by the uloop event loop, so `uloop.run()` must + * be invoked for the connection to progress. + * + * Runtime failures (refused connections, handshake errors, timeouts, + * abrupt resets) are reported through the `error` event, invalid + * arguments raise type exceptions. + * + * @function module:websocket#connect + * + * @param {string} url + * The WebSocket URL to connect to, e.g. `ws://host:port/path`. + * + * @param {Object} [options] + * Additional connection options: + * + * - `timeout` (number): connect and handshake timeout in milliseconds, + * `0` disables the timeout. Default: `15000`. + * - `max_frame_size` (number): maximum accepted message size in bytes; + * larger messages cause a close handshake with status code 1009. + * Clamped to the range 1024 to 16777216. Default: `262144`. + * - `headers` (Object): additional HTTP headers to include in the + * handshake request. Values containing line breaks are rejected. + * + * @returns {?websocket.connection} + * Returns the connection resource or `null` on error, e.g. when name + * resolution fails. Consult `error()` for the cause in that case. + * + * @example + * ```javascript + * import { connect } from 'websocket'; + * import * as uloop from 'uloop'; + * + * uloop.init(); + * + * let ws = connect('ws://10.0.0.1:8080/telemetry', { + * headers: { Authorization: 'Bearer secret' }, + * max_frame_size: 65536 + * }); + * + * ws.on('open', (ws) => ws.send('hello')); + * ws.on('message', (ws, data, is_text) => print(data)); + * ws.on('close', (ws, code, reason) => print(`closed ${code}\n`)); + * ws.on('error', (ws, error) => print(`error: ${error}\n`)); + * + * uloop.run(); + * ``` + */ +static uc_value_t * +uc_ws_connect(uc_vm_t *vm, size_t nargs) +{ + uc_value_t *url = uc_fn_arg(0); + uc_value_t *options = uc_fn_arg(1); + uc_value_t *hv, *hdrs = NULL; + struct addrinfo hints, *res = NULL, *rp; + uc_websocket_t *ws = NULL; + char *host = NULL, *path = NULL; + bool tls = false; + int port = 0, fd = -1, rv, immediate = 0; + int64_t timeout = WS_DEFAULT_TIMEOUT; + uint64_t maxmsg = WS_DEFAULT_MAX_MSG; + + if (ucv_type(url) != UC_STRING) { + uc_vm_raise_exception(vm, EXCEPTION_TYPE, "URL must be a string"); + + return NULL; + } + + if (!ws_parse_url(ucv_string_get(url), &host, &port, &path, &tls)) { + free(host); + free(path); + uc_vm_raise_exception(vm, EXCEPTION_TYPE, "Invalid WebSocket URL"); + + return NULL; + } + + if (tls) { + free(host); + free(path); + uc_vm_raise_exception(vm, EXCEPTION_TYPE, + "TLS (wss://) is not supported yet"); + + return NULL; + } + + if (options && ucv_type(options) == UC_OBJECT) { + hv = ucv_object_get(options, "timeout", NULL); + + if (hv) { + timeout = ucv_int64_get(hv); + + if (timeout < 0) + timeout = 0; + } + + hv = ucv_object_get(options, "max_frame_size", NULL); + + if (hv) { + maxmsg = ucv_uint64_get(hv); + + if (maxmsg < 1024) + maxmsg = 1024; + + if (maxmsg > 16 * 1024 * 1024) + maxmsg = 16 * 1024 * 1024; + } + + hdrs = ucv_object_get(options, "headers", NULL); + } + + if (timeout && timeout < WS_MIN_TIMEOUT) + timeout = WS_MIN_TIMEOUT; + + memset(&hints, 0, sizeof(hints)); + hints.ai_family = AF_UNSPEC; + hints.ai_socktype = SOCK_STREAM; + hints.ai_protocol = IPPROTO_TCP; + + rv = getaddrinfo(host, NULL, &hints, &res); + + if (rv) { + free(host); + free(path); + err_return(EHOSTUNREACH); + } + + for (rp = res; rp; rp = rp->ai_next) { + fd = socket(rp->ai_family, rp->ai_socktype | SOCK_NONBLOCK | SOCK_CLOEXEC, 0); + + if (fd < 0) + continue; + + if (rp->ai_family == AF_INET6) + ((struct sockaddr_in6 *)rp->ai_addr)->sin6_port = htons(port); + else + ((struct sockaddr_in *)rp->ai_addr)->sin_port = htons(port); + + rv = connect(fd, rp->ai_addr, rp->ai_addrlen); + + if (rv == 0) { + immediate = 1; + break; + } + + if (errno == EINPROGRESS) + break; + + close(fd); + fd = -1; + } + + if (fd < 0) { + freeaddrinfo(res); + free(host); + free(path); + err_return(errno ? errno : ECONNREFUSED); + } + + hv = ucv_resource_create_ex(vm, "websocket.connection", + (void **)&ws, 4, sizeof(*ws)); + + if (!hv) { + close(fd); + freeaddrinfo(res); + free(host); + free(path); + err_return(ENOMEM); + } + + memset(ws, 0, sizeof(*ws)); + + ws->vm = vm; + ws->obj = ucv_get(hv); + ws->host = host; + ws->path = path; + ws->port = port; + ws->max_msg_len = maxmsg; + ws->ufd.fd = fd; + ws->ufd.cb = ws_ufd_cb; + ws->timeout.cb = ws_timeout_cb; + ws->state = immediate ? WS_STATE_HANDSHAKE_WRITE : WS_STATE_CONNECTING; + + ucv_resource_persistent_set(ws->obj, true); + + ws->hs_req = ws_build_request(vm, ws, + (hdrs && ucv_type(hdrs) == UC_OBJECT) ? hdrs : NULL); + + if (!ws->hs_req) { + ws_teardown(ws); + ucv_put(hv); + err_return(EINVAL); + } + + ws->hs_req_len = strlen(ws->hs_req); + + freeaddrinfo(res); + + if (uloop_fd_add(&ws->ufd, ULOOP_WRITE) != 0) { + ws_teardown(ws); + ucv_put(hv); + err_return(errno ? errno : EINVAL); + } + + if (timeout) + uloop_timeout_set(&ws->timeout, (int)timeout); + + if (immediate) + ws_handshake_write(ws); + + ok_return(hv); +} + +/* ---------------------------------------------------------------------- */ +/* Registration */ +/* ---------------------------------------------------------------------- */ + +static uc_value_t * +uc_ws_error(uc_vm_t *vm, size_t nargs) +{ + uc_value_t *errmsg; + + if (last_error == 0) + return NULL; + + errmsg = ucv_string_new(strerror(last_error)); + last_error = 0; + + return errmsg; +} + +static const uc_function_list_t ws_methods[] = { + { "on", uc_ws_on }, + { "send", uc_ws_send }, + { "ping", uc_ws_ping }, + { "close", uc_ws_close }, + { "state", uc_ws_state }, + { "fileno", uc_ws_fileno }, +}; + +static const uc_function_list_t ws_functions[] = { + { "connect", uc_ws_connect }, + { "error", uc_ws_error }, +}; + +static void +ws_free_resource(void *ud) +{ + ws_teardown(ud); +} + +void uc_module_init(uc_vm_t *vm, uc_value_t *scope) +{ + uc_function_list_register(scope, ws_functions); + + uc_type_declare(vm, "websocket.connection", ws_methods, ws_free_resource); + + uloop_init(); +} diff --git a/openwrt/ucode/Makefile b/openwrt/ucode/Makefile index 10cd385a..ab9e16b3 100644 --- a/openwrt/ucode/Makefile +++ b/openwrt/ucode/Makefile @@ -160,8 +160,20 @@ define Package/ucode-mod-uloop endef define Package/ucode-mod-uloop/description - The uloop module allows ucode scripts to interact with OpenWrt uloop event - loop implementation. + The uloop module allows ucode scripts to interact with OpenWrt uloop event + loop implementation. +endef + + +define Package/ucode-mod-websocket + $(Package/ucode/default) + TITLE+= (websocket module) + DEPENDS:=ucode +libubox +libwslay +endef + +define Package/ucode-mod-websocket/description + The websocket module allows ucode scripts to interact with WebSocket + (RFC 6455) servers using an event driven, uloop based API. endef @@ -227,6 +239,11 @@ define Package/ucode-mod-uloop/install $(INSTALL_BIN) $(PKG_INSTALL_DIR)/usr/lib/ucode/uloop.so $(1)/usr/lib/ucode/ endef +define Package/ucode-mod-websocket/install + $(INSTALL_DIR) $(1)/usr/lib/ucode + $(INSTALL_BIN) $(PKG_INSTALL_DIR)/usr/lib/ucode/websocket.so $(1)/usr/lib/ucode/ +endef + $(eval $(call BuildPackage,libucode)) $(eval $(call BuildPackage,ucode)) $(eval $(call BuildPackage,ucode-mod-fs)) @@ -238,4 +255,5 @@ $(eval $(call BuildPackage,ucode-mod-struct)) $(eval $(call BuildPackage,ucode-mod-ubus)) $(eval $(call BuildPackage,ucode-mod-uci)) $(eval $(call BuildPackage,ucode-mod-uloop)) +$(eval $(call BuildPackage,ucode-mod-websocket)) $(eval $(call HostBuild)) diff --git a/tests/cram/CMakeLists.txt b/tests/cram/CMakeLists.txt index a93add5a..46954e71 100644 --- a/tests/cram/CMakeLists.txt +++ b/tests/cram/CMakeLists.txt @@ -1,6 +1,12 @@ FIND_PACKAGE(PythonInterp 3 REQUIRED) FILE(GLOB test_cases "test_*.t") +IF(WEBSOCKET_SUPPORT) + ADD_EXECUTABLE(ws-fixture fixtures/ws-fixture.c) +ELSE() + LIST(REMOVE_ITEM test_cases ${CMAKE_CURRENT_SOURCE_DIR}/test_websocket.t) +ENDIF() + SET(PYTHON_VENV_DIR "${CMAKE_CURRENT_BINARY_DIR}/.venv") SET(PYTHON_VENV_PIP "${PYTHON_VENV_DIR}/bin/pip") SET(PYTHON_VENV_CRAM "${PYTHON_VENV_DIR}/bin/cram") @@ -20,6 +26,10 @@ ADD_TEST( SET_PROPERTY(TEST cram APPEND PROPERTY ENVIRONMENT "BUILD_BIN_DIR=$") +IF(WEBSOCKET_SUPPORT) + SET_PROPERTY(TEST cram APPEND PROPERTY ENVIRONMENT "WS_FIXTURE=$") +ENDIF() + IF(CMAKE_C_COMPILER_ID STREQUAL "Clang") SET_PROPERTY(TEST cram APPEND PROPERTY ENVIRONMENT "UCODE_BIN=$") ELSE() diff --git a/tests/cram/fixtures/ws-fixture.c b/tests/cram/fixtures/ws-fixture.c new file mode 100644 index 00000000..1df03981 --- /dev/null +++ b/tests/cram/fixtures/ws-fixture.c @@ -0,0 +1,600 @@ +/* + * Minimal scripted WebSocket server used as test fixture by the + * tests/cram/websocket test cases. + * + * Scenarios are selected through the request path: + * + * /announce send a text message "path= host=" then echo + * /echo echo data, binary and ping frames + * /close initiate close(1001, "server bye") then echo loop + * /big send a single N byte text frame, then echo loop + * /frag send N bytes as one fragmented message in 1 KiB frames + * /flood send N x 1 KiB text messages, then echo loop + * /ping send a ping, then echo loop + * /reset tear the connection down with TCP RST + * /badaccept complete the handshake with a wrong accept hash + * /http200 reply with a plain HTTP 200 instead of 101 + * + * Usage: ws-fixture [connections] + */ + +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +#define WS_GUID "258EAFA5-E914-47DA-95CA-C5AB0DC85B11" + +static uint32_t sha1_state[5] = { + 0x67452301, 0xEFCDAB89, 0x98BADCFE, 0x10325476, 0xC3D2E1F0 +}; + +static uint64_t sha1_count; +static uint8_t sha1_buffer[64]; + +#define ROTL(x, n) (((x) << (n)) | ((x) >> (32 - (n)))) + +static void +sha1_transform(const uint8_t *p) +{ + uint32_t w[80], a, b, c, d, e, t; + size_t i; + + for (i = 0; i < 16; i++) + w[i] = ((uint32_t)p[i * 4] << 24) | ((uint32_t)p[i * 4 + 1] << 16) | + ((uint32_t)p[i * 4 + 2] << 8) | (uint32_t)p[i * 4 + 3]; + + for (i = 16; i < 80; i++) + w[i] = ROTL(w[i-3] ^ w[i-8] ^ w[i-14] ^ w[i-16], 1); + + a = sha1_state[0]; b = sha1_state[1]; c = sha1_state[2]; + d = sha1_state[3]; e = sha1_state[4]; + + for (i = 0; i < 80; i++) { + if (i < 20) + t = ((b & c) | ((~b) & d)) + 0x5A827999; + else if (i < 40) + t = (b ^ c ^ d) + 0x6ED9EBA1; + else if (i < 60) + t = ((b & c) | (b & d) | (c & d)) + 0x8F1BBCDC; + else + t = (b ^ c ^ d) + 0xCA62C1D6; + + t += ROTL(a, 5) + e + w[i]; + e = d; d = c; c = ROTL(b, 30); b = a; a = t; + } + + sha1_state[0] += a; sha1_state[1] += b; sha1_state[2] += c; + sha1_state[3] += d; sha1_state[4] += e; +} + +static void +sha1_reset(void) +{ + sha1_state[0] = 0x67452301; sha1_state[1] = 0xEFCDAB89; + sha1_state[2] = 0x98BADCFE; sha1_state[3] = 0x10325476; + sha1_state[4] = 0xC3D2E1F0; + sha1_count = 0; +} + +static void +sha1_update(const uint8_t *data, size_t len) +{ + size_t i = 0, n; + + while (i < len) { + n = 64 - (sha1_count % 64); + + if (n > len - i) + n = len - i; + + memcpy(sha1_buffer + (sha1_count % 64), data + i, n); + sha1_count += n; + i += n; + + if (sha1_count % 64 == 0) + sha1_transform(sha1_buffer); + } +} + +static void +sha1_final(uint8_t digest[20]) +{ + static const uint8_t pad[64] = { 0x80 }; + uint64_t bits = sha1_count * 8; + uint8_t tail[8]; + size_t i; + + for (i = 0; i < 8; i++) + tail[i] = (uint8_t)(bits >> (56 - 8 * i)); + + sha1_update(pad, 1 + ((119 - sha1_count % 64) % 64)); + sha1_update(tail, 8); + + for (i = 0; i < 5; i++) { + digest[i * 4] = (uint8_t)(sha1_state[i] >> 24); + digest[i * 4 + 1] = (uint8_t)(sha1_state[i] >> 16); + digest[i * 4 + 2] = (uint8_t)(sha1_state[i] >> 8); + digest[i * 4 + 3] = (uint8_t)(sha1_state[i]); + } +} + +static const char b64tab[] = + "ABCDEFGHIJKLMNOPQRSTUVWXYZabcdefghijklmnopqrstuvwxyz0123456789+/"; + +static void +b64_encode(const uint8_t *in, size_t inlen, char *out) +{ + size_t i = 0, j = 0; + + while (i + 3 <= inlen) { + uint32_t v = ((uint32_t)in[i] << 16) | + ((uint32_t)in[i + 1] << 8) | (uint32_t)in[i + 2]; + + out[j++] = b64tab[(v >> 18) & 63]; + out[j++] = b64tab[(v >> 12) & 63]; + out[j++] = b64tab[(v >> 6) & 63]; + out[j++] = b64tab[v & 63]; + i += 3; + } + + if (i < inlen) { + uint32_t v = (uint32_t)in[i] << 16; + + if (i + 1 < inlen) + v |= (uint32_t)in[i + 1] << 8; + + out[j++] = b64tab[(v >> 18) & 63]; + out[j++] = b64tab[(v >> 12) & 63]; + out[j++] = (i + 1 < inlen) ? b64tab[(v >> 6) & 63] : '='; + out[j++] = '='; + } + + out[j] = '\0'; +} + +static int +readn(int fd, void *buf, size_t len) +{ + uint8_t *p = buf; + ssize_t n; + + while (len > 0) { + n = read(fd, p, len); + + if (n == 0) + return -1; + + if (n < 0) { + if (errno == EINTR) + continue; + + return -1; + } + + p += n; + len -= n; + } + + return 0; +} + +static int +writeall(int fd, const void *buf, size_t len) +{ + const uint8_t *p = buf; + ssize_t n; + + while (len > 0) { + n = write(fd, p, len); + + if (n < 0) { + if (errno == EINTR) + continue; + + return -1; + } + + p += n; + len -= n; + } + + return 0; +} + +static int +send_frame(int fd, int opcode, int fin, const uint8_t *payload, size_t len) +{ + uint8_t hdr[10]; + size_t hlen = 2; + + hdr[0] = (fin ? 0x80 : 0x00) | (opcode & 0x0F); + + if (len < 126) { + hdr[1] = (uint8_t)len; + } + else if (len <= 0xFFFF) { + hdr[1] = 126; + hdr[2] = (uint8_t)(len >> 8); + hdr[3] = (uint8_t)len; + hlen = 4; + } + else { + hdr[1] = 127; + for (int i = 0; i < 8; i++) + hdr[2 + i] = (uint8_t)(len >> (56 - 8 * i)); + hlen = 10; + } + + if (writeall(fd, hdr, hlen) < 0) + return -1; + + return writeall(fd, payload, len); +} + +static int +recv_frame(int fd, int *opcode, int *fin, uint8_t **payload, size_t *len) +{ + uint8_t hdr[2], ext[8], mask[4]; + uint64_t plen; + size_t need; + uint8_t *buf; + + if (readn(fd, hdr, 2) < 0) + return -1; + + *fin = (hdr[0] >> 7) & 1; + *opcode = hdr[0] & 0x0F; + + if (!(hdr[1] & 0x80)) + return -1; /* client frames must be masked */ + + plen = hdr[1] & 0x7F; + + if (plen == 126) { + if (readn(fd, ext, 2) < 0) + return -1; + plen = ((uint64_t)ext[0] << 8) | ext[1]; + } + else if (plen == 127) { + if (readn(fd, ext, 8) < 0) + return -1; + plen = 0; + for (int i = 0; i < 8; i++) + plen = (plen << 8) | ext[i]; + } + + if (readn(fd, mask, 4) < 0) + return -1; + + if (plen > 16 * 1024 * 1024) + return -1; + + need = (size_t)plen; + buf = malloc(need ? need : 1); + + if (!buf) + return -1; + + if (readn(fd, buf, need) < 0) { + free(buf); + return -1; + } + + for (size_t i = 0; i < need; i++) + buf[i] ^= mask[i % 4]; + + *payload = buf; + *len = need; + + return 0; +} + +static int +do_handshake(int fd, char *path, size_t pathsize, char *host, size_t hostsize, + char *key, size_t keysize) +{ + char req[4096]; + size_t len = 0; + ssize_t n; + char *line, *eol, *k; + + while (len + 1 < sizeof(req)) { + n = read(fd, req + len, sizeof(req) - len - 1); + + if (n <= 0) + return -1; + + len += n; + req[len] = '\0'; + + if (strstr(req, "\r\n\r\n")) + break; + } + + if (!strstr(req, "\r\n\r\n")) + return -1; + + path[0] = host[0] = key[0] = '\0'; + + for (line = req; (eol = strstr(line, "\r\n")); line = eol + 2) { + if (!strncmp(line, "GET ", 4)) { + k = strchr(line + 4, ' '); + + if (k) { + snprintf(path, pathsize, "%.*s", + (int)(k - (line + 4)), line + 4); + } + } + else if (!strncasecmp(line, "Host:", 5)) { + k = line + 5; + + while (*k == ' ') + k++; + + snprintf(host, hostsize, "%.*s", (int)(eol - k), k); + } + else if (!strncasecmp(line, "Sec-WebSocket-Key:", 18)) { + k = line + 18; + + while (*k == ' ') + k++; + + snprintf(key, keysize, "%.*s", (int)(eol - k), k); + } + } + + return (key[0] && path[0]) ? 0 : -1; +} + +static int +send_upgrade_response(int fd, const char *key, int mode) +{ + char accept_src[128], accept[64], response[512]; + uint8_t digest[20]; + + sha1_reset(); + snprintf(accept_src, sizeof(accept_src), "%s%s", key, WS_GUID); + sha1_update((const uint8_t *)accept_src, strlen(accept_src)); + sha1_final(digest); + b64_encode(digest, 20, accept); + + switch (mode) { + case 1: /* wrong accept hash */ + accept[3] = (accept[3] == 'A') ? 'B' : 'A'; + break; + + case 2: /* plain HTTP reply */ + return writeall(fd, "HTTP/1.1 200 OK\r\nContent-Length: 0\r\n\r\n", 38); + } + + snprintf(response, sizeof(response), + "HTTP/1.1 101 Switching Protocols\r\n" + "Upgrade: websocket\r\n" + "Connection: Upgrade\r\n" + "Sec-WebSocket-Accept: %s\r\n\r\n", accept); + + return writeall(fd, response, strlen(response)); +} + +static int +echo_loop(int fd) +{ + uint8_t *payload; + size_t len; + int opcode, fin; + + while (recv_frame(fd, &opcode, &fin, &payload, &len) == 0) { + if (opcode == 8) { + uint8_t reply[2] = { payload[0], payload[1] }; + + send_frame(fd, 8, 1, reply, len >= 2 ? 2 : 0); + free(payload); + return 0; + } + else if (opcode == 9) { + send_frame(fd, 10, 1, payload, len); + } + else if (opcode == 1 || opcode == 2) { + send_frame(fd, opcode, 1, payload, len); + } + + free(payload); + } + + return -1; +} + +static void +fill_pattern(uint8_t *buf, size_t len) +{ + for (size_t i = 0; i < len; i++) + buf[i] = (uint8_t)('a' + (i % 26)); +} + +static int +handle_connection(int fd) +{ + char path[256], host[256], key[64]; + char *s; + long n; + uint8_t *blob; + int mode = 0; + + if (do_handshake(fd, path, sizeof(path), host, sizeof(host), + key, sizeof(key)) < 0) + return -1; + + if (!strncmp(path, "/badaccept", 10)) + mode = 1; + else if (!strncmp(path, "/http200", 8)) + mode = 2; + + if (send_upgrade_response(fd, key, mode) < 0) + return -1; + + if (!strncmp(path, "/reset", 6)) { + struct linger l = { .l_onoff = 1, .l_linger = 0 }; + + setsockopt(fd, SOL_SOCKET, SO_LINGER, &l, sizeof(l)); + + return -1; + } + + if (!strncmp(path, "/msgreset", 9)) { + struct linger l = { .l_onoff = 1, .l_linger = 0 }; + + send_frame(fd, 1, 1, (const uint8_t *)"bye", 3); + setsockopt(fd, SOL_SOCKET, SO_LINGER, &l, sizeof(l)); + + return -1; + } + + if (!strncmp(path, "/announce", 9)) { + char msg[512]; + + snprintf(msg, sizeof(msg), "path=%.220s host=%.220s", path, host); + send_frame(fd, 1, 1, (const uint8_t *)msg, strlen(msg)); + } + else if (!strncmp(path, "/close", 6)) { + const char *reason = "server bye"; + uint8_t payload[125]; + + payload[0] = 1001 >> 8; + payload[1] = 1001 & 0xFF; + memcpy(payload + 2, reason, strlen(reason)); + + send_frame(fd, 8, 1, payload, 2 + strlen(reason)); + } + else if (!strncmp(path, "/big", 4)) { + n = strtol(path + 4, &s, 10); + + if (n > 0 && n <= 16 * 1024 * 1024) { + blob = malloc(n); + + if (blob) { + fill_pattern(blob, n); + send_frame(fd, 1, 1, blob, n); + free(blob); + } + } + } + else if (!strncmp(path, "/frag", 5)) { + n = strtol(path + 5, &s, 10); + + if (n > 0 && n <= 16 * 1024 * 1024) { + size_t total = (size_t)n; + + blob = malloc(total); + + if (blob) { + size_t off = 0, chunk; + + fill_pattern(blob, total); + + while (off < total) { + chunk = total - off; + + if (chunk > 1024) + chunk = 1024; + + send_frame(fd, (off == 0) ? 1 : 0, + (off + chunk == total) ? 1 : 0, + blob + off, chunk); + off += chunk; + } + + free(blob); + } + } + } + else if (!strncmp(path, "/flood", 6)) { + n = strtol(path + 6, &s, 10); + + if (n > 0 && n <= 100000) { + blob = malloc(1024); + + if (blob) { + fill_pattern(blob, 1024); + + for (long i = 0; i < n; i++) + send_frame(fd, 1, 1, blob, 1024); + + free(blob); + } + } + } + else if (!strncmp(path, "/ping", 5)) { + send_frame(fd, 9, 1, (const uint8_t *)"fixture-ping", 12); + } + + return echo_loop(fd); +} + +int +main(int argc, char **argv) +{ + struct sockaddr_in6 addr; + int srv, fd, one = 1, connections = 1, served = 0; + uint16_t port; + char *s; + + if (argc < 2) { + fprintf(stderr, "Usage: %s [connections]\n", argv[0]); + return 1; + } + + port = (uint16_t)strtoul(argv[1], &s, 10); + + if (argc > 2) + connections = (int)strtoul(argv[2], &s, 10); + + signal(SIGPIPE, SIG_IGN); + + srv = socket(AF_INET6, SOCK_STREAM | SOCK_CLOEXEC, 0); + + if (srv < 0) + return 1; + + setsockopt(srv, SOL_SOCKET, SO_REUSEADDR, &one, sizeof(one)); + setsockopt(srv, IPPROTO_IPV6, IPV6_V6ONLY, &(int){ 0 }, sizeof(int)); + + memset(&addr, 0, sizeof(addr)); + addr.sin6_family = AF_INET6; + addr.sin6_addr = in6addr_any; + addr.sin6_port = htons(port); + + if (bind(srv, (struct sockaddr *)&addr, sizeof(addr)) < 0) { + fprintf(stderr, "bind failed: %s\n", strerror(errno)); + return 1; + } + + if (listen(srv, 8) < 0) + return 1; + + while (served < connections) { + fd = accept4(srv, NULL, NULL, SOCK_CLOEXEC); + + if (fd < 0) { + if (errno == EINTR) + continue; + + return 1; + } + + handle_connection(fd); + + close(fd); + served++; + } + + return 0; +} diff --git a/tests/cram/test_websocket.t b/tests/cram/test_websocket.t new file mode 100644 index 00000000..26f94a97 --- /dev/null +++ b/tests/cram/test_websocket.t @@ -0,0 +1,422 @@ +setup common environment: + + $ [ -n "$BUILD_BIN_DIR" ] && export PATH="$BUILD_BIN_DIR:$PATH" + $ if command -v valgrind >/dev/null 2>&1; then + > alias ucode="$UCODE_BIN" + > else + > alias ucode="$BUILD_BIN_DIR/ucode" + > fi + + $ for m in $BUILD_BIN_DIR/*.so; do + > ln -s "$m" "$(pwd)/$(basename $m)"; \ + > done + +check that the websocket module loads: + + $ ucode -lwebsocket -e 'let f = websocket.connect; print(f)' + function connect(...) { [native code] } (no-eol) + + +test connection establishment and handshake header validation (T2): + + $ PORT=28971 + $ $WS_FIXTURE $PORT 1 & + $ sleep 0.2 + + $ ucode -lwebsocket -luloop - <<'EOF' + > import { connect } from 'websocket'; + > import * as uloop from 'uloop'; + > uloop.init(); + > let ws = connect('ws://127.0.0.1:28971/announce?x=1', { timeout: 5000 }); + > ws.on('message', (w, data, is_text) => { + > print(`MSG ${data}\n`); + > w.close(1000, 'done'); + > }); + > ws.on('close', (w, code) => { print(`CLOSE ${code}\n`); uloop.end(); }); + > ws.on('error', (w, e) => { print(`ERROR ${e}\n`); uloop.end(); }); + > uloop.run(); + > print('DONE\n'); + > EOF + MSG path=/announce?x=1 host=127.0.0.1:28971 + CLOSE 1000 + DONE + + $ wait + + +test text and binary echo (T3, T4): + + $ PORT=28972 + $ $WS_FIXTURE $PORT 1 & + $ sleep 0.2 + + $ ucode -lwebsocket -luloop - <<'EOF' + > import { connect } from 'websocket'; + > import * as uloop from 'uloop'; + > uloop.init(); + > let ws = connect('ws://127.0.0.1:28972/echo', { timeout: 5000 }); + > let count = 0; + > ws.on('open', (w) => { + > w.send('hello'); + > w.send([1, 2, 3, 250]); + > }); + > ws.on('message', (w, data, is_text) => { + > count++; + > print(`MSG ${is_text} ${length(data)}\n`); + > if (count >= 2) + > w.close(1000, ''); + > }); + > ws.on('close', (w, code) => { print(`CLOSE ${code}\n`); uloop.end(); }); + > ws.on('error', (w, e) => { print(`ERROR ${e}\n`); uloop.end(); }); + > uloop.run(); + > print('DONE\n'); + > EOF + MSG true 5 + MSG false 4 + CLOSE 1000 + DONE + + $ wait + + +test ping/pong handling (T7): + + $ PORT=28973 + $ $WS_FIXTURE $PORT 1 & + $ sleep 0.2 + + $ ucode -lwebsocket -luloop - <<'EOF' + > import { connect } from 'websocket'; + > import * as uloop from 'uloop'; + > uloop.init(); + > let ws = connect('ws://127.0.0.1:28973/ping', { timeout: 5000 }); + > ws.on('open', (w) => { + > w.ping('client-ping'); + > uloop.timer(500, () => w.close(1000, '')); + > }); + > ws.on('message', (w, data, is_text) => { + > print(`MSG ${data}\n`); + > w.close(1000, ''); + > }); + > ws.on('close', (w, code) => { print(`CLOSE ${code}\n`); uloop.end(); }); + > ws.on('error', (w, e) => { print(`ERROR ${e}\n`); uloop.end(); }); + > uloop.run(); + > print('DONE\n'); + > EOF + CLOSE 1000 + DONE + + $ wait + + +test server initiated close handshake with code and reason (T8): + + $ PORT=28974 + $ $WS_FIXTURE $PORT 1 & + $ sleep 0.2 + + $ ucode -lwebsocket -luloop - <<'EOF' + > import { connect } from 'websocket'; + > import * as uloop from 'uloop'; + > uloop.init(); + > let ws = connect('ws://127.0.0.1:28974/close', { timeout: 5000 }); + > ws.on('message', (w, data, is_text) => print(`UNEXPECTED MSG\n`)); + > ws.on('close', (w, code, reason) => { + > print(`CLOSE ${code} ${reason}\n`); + > uloop.end(); + > }); + > ws.on('error', (w, e) => { print(`ERROR ${e}\n`); uloop.end(); }); + > uloop.run(); + > print('DONE\n'); + > EOF + CLOSE 1001 server bye + DONE + + $ wait + + +test fragmented message reassembly (T5): + + $ PORT=28975 + $ $WS_FIXTURE $PORT 1 & + $ sleep 0.2 + + $ ucode -lwebsocket -luloop - <<'EOF' + > import { connect } from 'websocket'; + > import * as uloop from 'uloop'; + > uloop.init(); + > let ws = connect('ws://127.0.0.1:28975/frag8192', { timeout: 5000 }); + > ws.on('message', (w, data, is_text) => { + > print(`MSG ${length(data)}\n`); + > w.close(1000, ''); + > }); + > ws.on('close', (w, code) => { print(`CLOSE ${code}\n`); uloop.end(); }); + > ws.on('error', (w, e) => { print(`ERROR ${e}\n`); uloop.end(); }); + > uloop.run(); + > print('DONE\n'); + > EOF + MSG 8192 + CLOSE 1000 + DONE + + $ wait + + +test oversized frame rejection with close code 1009 (T6): + + $ PORT=28976 + $ $WS_FIXTURE $PORT 1 & + $ sleep 0.2 + + $ ucode -lwebsocket -luloop - <<'EOF' + > import { connect } from 'websocket'; + > import * as uloop from 'uloop'; + > uloop.init(); + > let ws = connect('ws://127.0.0.1:28976/big131072', + > { timeout: 5000, max_frame_size: 65536 }); + > ws.on('message', (w, data, is_text) => print('UNEXPECTED MSG\n')); + > ws.on('close', (w, code) => { print(`CLOSE ${code}\n`); uloop.end(); }); + > ws.on('error', (w, e) => { print(`ERROR ${e}\n`); uloop.end(); }); + > uloop.run(); + > print('DONE\n'); + > EOF + CLOSE 1009 + DONE + + $ wait + + +test message flood is processed without loss (T11): + + $ PORT=28977 + $ $WS_FIXTURE $PORT 1 & + $ sleep 0.2 + + $ ucode -lwebsocket -luloop - <<'EOF' + > import { connect } from 'websocket'; + > import * as uloop from 'uloop'; + > uloop.init(); + > let ws = connect('ws://127.0.0.1:28977/flood2000', + > { timeout: 5000, max_frame_size: 65536 }); + > let count = 0; + > ws.on('message', (w, data, is_text) => { + > count++; + > if (count == 2000) { + > print(`MSGS ${count} ${length(data)}\n`); + > w.close(1000, ''); + > } + > }); + > ws.on('close', (w, code) => { print(`CLOSE ${code}\n`); uloop.end(); }); + > ws.on('error', (w, e) => { print(`ERROR ${e}\n`); uloop.end(); }); + > uloop.run(); + > print('DONE\n'); + > EOF + MSGS 2000 1024 + CLOSE 1000 + DONE + + $ wait + + +test abrupt connection reset surfaces an error (T9): + + $ $WS_FIXTURE 28978 1 & + $ sleep 0.2 + + $ ucode -lwebsocket -luloop - <<'EOF' + > import { connect } from 'websocket'; + > import * as uloop from 'uloop'; + > uloop.init(); + > let ws = connect('ws://127.0.0.1:28978/reset', { timeout: 5000 }); + > ws.on('open', (w) => print('OPEN\n')); + > ws.on('close', (w, code) => { print(`CLOSE ${code}\n`); uloop.end(); }); + > ws.on('error', (w, e) => { print(`ERROR\n`); uloop.end(); }); + > uloop.run(); + > print('DONE\n'); + > EOF + OPEN + ERROR + CLOSE 1006 + DONE + + $ wait + + +test handshake failure on wrong Sec-WebSocket-Accept: + + $ $WS_FIXTURE 28979 1 & + $ sleep 0.2 + + $ ucode -lwebsocket -luloop - <<'EOF' + > import { connect } from 'websocket'; + > import * as uloop from 'uloop'; + > uloop.init(); + > let ws = connect('ws://127.0.0.1:28979/badaccept', { timeout: 5000 }); + > ws.on('open', (w) => print('UNEXPECTED OPEN\n')); + > ws.on('close', (w, code) => { print(`CLOSE ${code}\n`); uloop.end(); }); + > ws.on('error', (w, e) => { print('ERROR invalid handshake response\n'); uloop.end(); }); + > uloop.run(); + > print('DONE\n'); + > EOF + ERROR invalid handshake response + DONE + + $ wait + + +test handshake failure on non-101 response: + + $ $WS_FIXTURE 28980 1 & + $ sleep 0.2 + + $ ucode -lwebsocket -luloop - <<'EOF' + > import { connect } from 'websocket'; + > import * as uloop from 'uloop'; + > uloop.init(); + > let ws = connect('ws://127.0.0.1:28980/http200', { timeout: 5000 }); + > ws.on('open', (w) => print('UNEXPECTED OPEN\n')); + > ws.on('close', (w, code) => { print(`CLOSE ${code}\n`); uloop.end(); }); + > ws.on('error', (w, e) => { print('ERROR invalid handshake response\n'); uloop.end(); }); + > uloop.run(); + > print('DONE\n'); + > EOF + ERROR invalid handshake response + DONE + + $ wait + + +test URL validation (T12): + + $ ucode -lwebsocket -e 'import { connect } from "websocket"; connect("http://example.org/");' ; echo "rc=$?" + Type error: Invalid WebSocket URL + In [-e argument], line 1, byte 67: + + `import { connect } from "websocket"; connect("http://example.org/");` + Near here --------------------------------------------------------^ + + + rc=254 + + $ ucode -lwebsocket -e 'import { connect } from "websocket"; connect("ws://user@example.org/");' ; echo "rc=$?" + Type error: Invalid WebSocket URL + In [-e argument], line 1, byte 70: + + `import { connect } from "websocket"; connect("ws://user@example.org/");` + Near here -----------------------------------------------------------^ + + + rc=254 + + $ ucode -lwebsocket -e 'import { connect } from "websocket"; connect("ws://example.org:70000/");' ; echo "rc=$?" + Type error: Invalid WebSocket URL + In [-e argument], line 1, byte 71: + + `import { connect } from "websocket"; connect("ws://example.org:70000/");` + Near here ------------------------------------------------------------^ + + + rc=254 + + $ ucode -lwebsocket -e 'import { connect } from "websocket"; connect("wss://example.org/");' ; echo "rc=$?" + Type error: TLS (wss://) is not supported yet + In [-e argument], line 1, byte 66: + + `import { connect } from "websocket"; connect("wss://example.org/");` + Near here -------------------------------------------------------^ + + + rc=254 + + $ ucode -lwebsocket -e 'import { connect } from "websocket"; connect(123)' ; echo "rc=$?" + Type error: URL must be a string + In [-e argument], line 1, byte 49: + + `import { connect } from "websocket"; connect(123)` + Near here --------------------------------------^ + + + rc=254 + + +test IPv6 literal with port yields a correct request path (review LOW): + + $ $WS_FIXTURE 28981 1 & + $ sleep 0.2 + + $ ucode -lwebsocket -luloop - <<'EOF' + > import { connect } from 'websocket'; + > import * as uloop from 'uloop'; + > uloop.init(); + > let ws = connect('ws://[::1]:28981/announce?v6=1', { timeout: 5000 }); + > ws.on('message', (w, data, is_text) => { + > print(`MSG ${data}\n`); + > w.close(1000, ''); + > }); + > ws.on('close', (w, code) => { print(`CLOSE ${code}\n`); uloop.end(); }); + > ws.on('error', (w, e) => { print(`ERROR ${e}\n`); uloop.end(); }); + > uloop.run(); + > print('DONE\n'); + > EOF + MSG path=/announce?v6=1 host=[::1]:28981 + CLOSE 1000 + DONE + + $ wait + + +test on() after teardown is rejected instead of crashing (review MEDIUM): + + $ $WS_FIXTURE 28982 1 & + $ sleep 0.2 + + $ ucode -lwebsocket -luloop - <<'EOF' + > import { connect } from 'websocket'; + > import * as uloop from 'uloop'; + > uloop.init(); + > let ws = connect('ws://127.0.0.1:28982/echo', { timeout: 5000 }); + > ws.on('open', (w) => w.close(1000, 'bye')); + > ws.on('close', (w, code) => { + > print(`CLOSE ${code}\n`); + > uloop.timer(300, () => { + > let rv = ws.on('message', (w2, d) => print('never\n')); + > print(`ON-After-TEARDOWN ${rv}\n`); + > uloop.end(); + > }); + > }); + > ws.on('error', (w, e) => { print(`ERROR ${e}\n`); uloop.end(); }); + > uloop.run(); + > print('DONE\n'); + > EOF + CLOSE 1000 + ON-After-TEARDOWN null + DONE + + $ wait + + +test send() from message callback surviving peer reset (review HIGH): + + $ $WS_FIXTURE 28983 1 & + $ sleep 0.2 + + $ ucode -lwebsocket -luloop - <<'EOF' + > import { connect } from 'websocket'; + > import * as uloop from 'uloop'; + > uloop.init(); + > let ws = connect('ws://127.0.0.1:28983/msgreset', { timeout: 5000 }); + > ws.on('message', (w, data) => { + > print(`MSG ${data}\n`); + > w.send('ack'); + > }); + > ws.on('close', (w, code) => { print(`CLOSE ${code}\n`); uloop.end(); }); + > ws.on('error', (w, e) => { print('ERROR\n'); uloop.end(); }); + > uloop.run(); + > print('DONE\n'); + > EOF + MSG bye + ERROR + CLOSE 1006 + DONE + + $ wait